diff --git a/.cirrus.yml b/.cirrus.yml index a7ce0d9d4..abf07ca48 100644 --- a/.cirrus.yml +++ b/.cirrus.yml @@ -1,7 +1,7 @@ only_if: $CIRRUS_TAG == '' && ($CIRRUS_PR != '' || $CIRRUS_BRANCH == 'master' || $CIRRUS_BRANCH =~ 'tokio-.*') auto_cancellation: $CIRRUS_BRANCH != 'master' && $CIRRUS_BRANCH !=~ 'tokio-.*' freebsd_instance: - image_family: freebsd-14-1 + image_family: freebsd-14-2 env: RUST_STABLE: stable RUST_NIGHTLY: nightly-2024-05-05 diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index 9514a2613..9db061611 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -475,10 +475,18 @@ jobs: runs-on: ubuntu-latest steps: - uses: actions/checkout@v4 - - name: Check semver + - name: Check `tokio` semver uses: obi1kenobi/cargo-semver-checks-action@v2 with: rust-toolchain: ${{ env.rust_stable }} + package: tokio + release-type: minor + - name: Check semver for rest of the workspace + if: ${{ !startsWith(github.event.pull_request.base.ref, 'tokio-1.') }} + uses: obi1kenobi/cargo-semver-checks-action@v2 + with: + rust-toolchain: ${{ env.rust_stable }} + exclude: tokio release-type: minor cross-check: @@ -710,7 +718,14 @@ jobs: toolchain: ${{ env.rust_min }} - uses: Swatinem/rust-cache@v2 - name: "check --workspace --all-features" - run: cargo check --workspace --all-features + run: | + if [[ "${{ github.event.pull_request.base.ref }}" =~ ^tokio-1\..* ]]; then + # Only check `tokio` crate as the PR is backporting to an earlier tokio release. + cargo check -p tokio --all-features + else + # Check all crates in the workspace + cargo check --workspace --all-features + fi env: RUSTFLAGS: "" # remove -Dwarnings @@ -1006,10 +1021,10 @@ jobs: targets: ${{ matrix.target }} # Install dependencies - - name: Install cargo-hack, wasmtime, and cargo-wasi + - name: Install cargo-hack, wasmtime uses: taiki-e/install-action@v2 with: - tool: cargo-hack,wasmtime,cargo-wasi + tool: cargo-hack,wasmtime - uses: Swatinem/rust-cache@v2 - name: WASI test tokio full @@ -1035,9 +1050,12 @@ jobs: - name: test tests-integration --features wasi-rt # TODO: this should become: `cargo hack wasi test --each-feature` - run: cargo wasi test --test rt_yield --features wasi-rt + run: cargo test --target ${{ matrix.target }} --test rt_yield --features wasi-rt if: matrix.target == 'wasm32-wasip1' working-directory: tests-integration + env: + CARGO_TARGET_WASM32_WASIP1_RUNNER: "wasmtime run --" + RUSTFLAGS: -Dwarnings -C target-feature=+atomics,+bulk-memory -C link-args=--max-memory=67108864 - name: test tests-integration --features wasi-threads-rt run: cargo test --target ${{ matrix.target }} --features wasi-threads-rt diff --git a/Cargo.toml b/Cargo.toml index 2238deac7..618b310e3 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -17,3 +17,16 @@ members = [ [workspace.metadata.spellcheck] config = "spellcheck.toml" + +[workspace.lints.rust] +unexpected_cfgs = { level = "warn", check-cfg = [ + 'cfg(fuzzing)', + 'cfg(loom)', + 'cfg(mio_unsupported_force_poll_poll)', + 'cfg(tokio_allow_from_blocking_fd)', + 'cfg(tokio_internal_mt_counters)', + 'cfg(tokio_no_parking_lot)', + 'cfg(tokio_no_tuning_tests)', + 'cfg(tokio_taskdump)', + 'cfg(tokio_unstable)', +] } diff --git a/examples/Cargo.toml b/examples/Cargo.toml index 54f2ecb8a..84112c08d 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -95,3 +95,6 @@ path = "named-pipe-multi-client.rs" [[example]] name = "dump" path = "dump.rs" + +[lints] +workspace = true diff --git a/tokio/CHANGELOG.md b/tokio/CHANGELOG.md index d2cb6a204..2ba7e8a0c 100644 --- a/tokio/CHANGELOG.md +++ b/tokio/CHANGELOG.md @@ -303,6 +303,20 @@ Yanked. Please use 1.39.1 instead. [#6709]: https://github.com/tokio-rs/tokio/pull/6709 [#6710]: https://github.com/tokio-rs/tokio/pull/6710 +# 1.38.2 (April 2nd, 2025) + +This release fixes a soundness issue in the broadcast channel. The channel +accepts values that are `Send` but `!Sync`. Previously, the channel called +`clone()` on these values without synchronizing. This release fixes the channel +by synchronizing calls to `.clone()` (Thanks Austin Bonander for finding and +reporting the issue). + +### Fixed + +- sync: synchronize `clone()` call in broadcast channel ([#7232]) + +[#7232]: https://github.com/tokio-rs/tokio/pull/7232 + # 1.38.1 (July 16th, 2024) This release fixes the bug identified as ([#6682]), which caused timers not diff --git a/tokio/Cargo.toml b/tokio/Cargo.toml index 860178716..2b0c1127a 100644 --- a/tokio/Cargo.toml +++ b/tokio/Cargo.toml @@ -173,3 +173,6 @@ allowed_external_types = [ "bytes::buf::buf_mut::BufMut", "tokio_macros::*", ] + +[lints] +workspace = true diff --git a/tokio/src/runtime/tests/queue.rs b/tokio/src/runtime/tests/queue.rs index 8a57ae428..9047f4ad7 100644 --- a/tokio/src/runtime/tests/queue.rs +++ b/tokio/src/runtime/tests/queue.rs @@ -1,5 +1,4 @@ use crate::runtime::scheduler::multi_thread::{queue, Stats}; -use crate::runtime::task::{self, Schedule, Task, TaskHarnessScheduleHooks}; use std::cell::RefCell; use std::thread; @@ -272,22 +271,3 @@ fn stress2() { assert_eq!(num_pop, NUM_TASKS); } } - -#[allow(dead_code)] -struct Runtime; - -impl Schedule for Runtime { - fn release(&self, _task: &Task) -> Option> { - None - } - - fn schedule(&self, _task: task::Notified) { - unreachable!(); - } - - fn hooks(&self) -> TaskHarnessScheduleHooks { - TaskHarnessScheduleHooks { - task_terminate_callback: None, - } - } -} diff --git a/tokio/src/sync/broadcast.rs b/tokio/src/sync/broadcast.rs index 3c3ca98f8..b48493be2 100644 --- a/tokio/src/sync/broadcast.rs +++ b/tokio/src/sync/broadcast.rs @@ -118,7 +118,7 @@ use crate::loom::cell::UnsafeCell; use crate::loom::sync::atomic::{AtomicBool, AtomicUsize}; -use crate::loom::sync::{Arc, Mutex, MutexGuard, RwLock, RwLockReadGuard}; +use crate::loom::sync::{Arc, Mutex, MutexGuard}; use crate::runtime::coop::cooperative; use crate::util::linked_list::{self, GuardedLinkedList, LinkedList}; use crate::util::WakeList; @@ -304,7 +304,7 @@ use self::error::{RecvError, SendError, TryRecvError}; /// Data shared between senders and receivers. struct Shared { /// slots in the channel. - buffer: Box<[RwLock>]>, + buffer: Box<[Mutex>]>, /// Mask a position -> index. mask: usize, @@ -348,7 +348,7 @@ struct Slot { /// /// The value is set by `send` when the write lock is held. When a reader /// drops, `rem` is decremented. When it hits zero, the value is dropped. - val: UnsafeCell>, + val: Option, } /// An entry in the wait queue. @@ -386,7 +386,7 @@ generate_addr_of_methods! { } struct RecvGuard<'a, T> { - slot: RwLockReadGuard<'a, Slot>, + slot: MutexGuard<'a, Slot>, } /// Receive a value future. @@ -395,11 +395,15 @@ struct Recv<'a, T> { receiver: &'a mut Receiver, /// Entry in the waiter `LinkedList`. - waiter: UnsafeCell, + waiter: WaiterCell, } -unsafe impl<'a, T: Send> Send for Recv<'a, T> {} -unsafe impl<'a, T: Send> Sync for Recv<'a, T> {} +// The wrapper around `UnsafeCell` isolates the unsafe impl `Send` and `Sync` +// from `Recv`. +struct WaiterCell(UnsafeCell); + +unsafe impl Send for WaiterCell {} +unsafe impl Sync for WaiterCell {} /// Max number of receivers. Reserve space to lock. const MAX_RECEIVERS: usize = usize::MAX >> 2; @@ -467,12 +471,6 @@ pub fn channel(capacity: usize) -> (Sender, Receiver) { (tx, rx) } -unsafe impl Send for Sender {} -unsafe impl Sync for Sender {} - -unsafe impl Send for Receiver {} -unsafe impl Sync for Receiver {} - impl Sender { /// Creates the sending-half of the [`broadcast`] channel. /// @@ -511,10 +509,10 @@ impl Sender { let mut buffer = Vec::with_capacity(capacity); for i in 0..capacity { - buffer.push(RwLock::new(Slot { + buffer.push(Mutex::new(Slot { rem: AtomicUsize::new(0), pos: (i as u64).wrapping_sub(capacity as u64), - val: UnsafeCell::new(None), + val: None, })); } @@ -600,7 +598,7 @@ impl Sender { tail.pos = tail.pos.wrapping_add(1); // Get the slot - let mut slot = self.shared.buffer[idx].write(); + let mut slot = self.shared.buffer[idx].lock(); // Track the position slot.pos = pos; @@ -609,7 +607,7 @@ impl Sender { slot.rem.with_mut(|v| *v = rem); // Write the value - slot.val = UnsafeCell::new(Some(value)); + slot.val = Some(value); // Release the slot lock before notifying the receivers. drop(slot); @@ -696,7 +694,7 @@ impl Sender { while low < high { let mid = low + (high - low) / 2; let idx = base_idx.wrapping_add(mid) & self.shared.mask; - if self.shared.buffer[idx].read().rem.load(SeqCst) == 0 { + if self.shared.buffer[idx].lock().rem.load(SeqCst) == 0 { low = mid + 1; } else { high = mid; @@ -738,7 +736,7 @@ impl Sender { let tail = self.shared.tail.lock(); let idx = (tail.pos.wrapping_sub(1) & self.shared.mask as u64) as usize; - self.shared.buffer[idx].read().rem.load(SeqCst) == 0 + self.shared.buffer[idx].lock().rem.load(SeqCst) == 0 } /// Returns the number of active receivers. @@ -1058,7 +1056,7 @@ impl Receiver { let idx = (self.next & self.shared.mask as u64) as usize; // The slot holding the next value to read - let mut slot = self.shared.buffer[idx].read(); + let mut slot = self.shared.buffer[idx].lock(); if slot.pos != self.next { // Release the `slot` lock before attempting to acquire the `tail` @@ -1075,7 +1073,7 @@ impl Receiver { let mut tail = self.shared.tail.lock(); // Acquire slot lock again - slot = self.shared.buffer[idx].read(); + slot = self.shared.buffer[idx].lock(); // Make sure the position did not change. This could happen in the // unlikely event that the buffer is wrapped between dropping the @@ -1367,12 +1365,12 @@ impl<'a, T> Recv<'a, T> { fn new(receiver: &'a mut Receiver) -> Recv<'a, T> { Recv { receiver, - waiter: UnsafeCell::new(Waiter { + waiter: WaiterCell(UnsafeCell::new(Waiter { queued: AtomicBool::new(false), waker: None, pointers: linked_list::Pointers::new(), _p: PhantomPinned, - }), + })), } } @@ -1384,7 +1382,7 @@ impl<'a, T> Recv<'a, T> { is_unpin::<&mut Receiver>(); let me = self.get_unchecked_mut(); - (me.receiver, &me.waiter) + (me.receiver, &me.waiter.0) } } } @@ -1418,6 +1416,7 @@ impl<'a, T> Drop for Recv<'a, T> { // `Shared::notify_rx` before we drop the object. let queued = self .waiter + .0 .with(|ptr| unsafe { (*ptr).queued.load(Acquire) }); // If the waiter is queued, we need to unlink it from the waiters list. @@ -1432,6 +1431,7 @@ impl<'a, T> Drop for Recv<'a, T> { // `Relaxed` order suffices because we hold the tail lock. let queued = self .waiter + .0 .with_mut(|ptr| unsafe { (*ptr).queued.load(Relaxed) }); if queued { @@ -1440,7 +1440,7 @@ impl<'a, T> Drop for Recv<'a, T> { // safety: tail lock is held and the wait node is verified to be in // the list. unsafe { - self.waiter.with_mut(|ptr| { + self.waiter.0.with_mut(|ptr| { tail.waiters.remove((&mut *ptr).into()); }); } @@ -1486,7 +1486,7 @@ impl<'a, T> RecvGuard<'a, T> { where T: Clone, { - self.slot.val.with(|ptr| unsafe { (*ptr).clone() }) + self.slot.val.clone() } } @@ -1494,8 +1494,7 @@ impl<'a, T> Drop for RecvGuard<'a, T> { fn drop(&mut self) { // Decrement the remaining counter if 1 == self.slot.rem.fetch_sub(1, SeqCst) { - // Safety: Last receiver, drop the value - self.slot.val.with_mut(|ptr| unsafe { *ptr = None }); + self.slot.val = None; } } }