From d4fb4bb9e6e78f740ce64e0b1729fa25fba6ab12 Mon Sep 17 00:00:00 2001 From: soreavis <263610811+soreavis@users.noreply.github.com> Date: Sun, 9 Aug 2026 18:06:29 +0200 Subject: [PATCH] sync: use acquire/release orderings in `Notify` (#8325) Every atomic operation on `Notify::state` used `SeqCst`. Nothing needs the global total order: every lock-avoidance decision is made by an RMW, which always reads the latest value in that atomic's modification order, and every stale load is re-validated either by a following RMW or by a re-load under the `waiters` mutex. What the `state` orderings do have to carry is the happens-before for data published before `notify_one`, consumed through the permit compare-exchange, and for data published before `notify_waiters`, consumed through the counter check. Acquire/release on `state` provides both. The waiter list is ordered by the mutex, and `AtomicNotification` by its own release/acquire pair, so neither depends on these orderings. Loads become `Acquire`, stores `Release`, compare-exchange `(AcqRel, Acquire)`, and the `notify_waiters` counter increment `AcqRel`. All nineteen sites are converted, so no `SeqCst` is left alongside weaker orderings. Fixes: #6266 --- tokio/src/sync/notify.rs | 46 ++++++++++++++++++++-------------------- 1 file changed, 23 insertions(+), 23 deletions(-) diff --git a/tokio/src/sync/notify.rs b/tokio/src/sync/notify.rs index c8270f17d..d850aa3b6 100644 --- a/tokio/src/sync/notify.rs +++ b/tokio/src/sync/notify.rs @@ -16,7 +16,7 @@ use std::marker::PhantomPinned; use std::panic::{RefUnwindSafe, UnwindSafe}; use std::pin::Pin; use std::ptr::NonNull; -use std::sync::atomic::Ordering::{self, Acquire, Relaxed, Release, SeqCst}; +use std::sync::atomic::Ordering::{self, AcqRel, Acquire, Relaxed, Release}; use std::sync::Arc; use std::task::{Context, Poll, Waker}; @@ -464,7 +464,7 @@ fn inc_num_notify_waiters_calls(data: usize) -> usize { } fn atomic_inc_num_notify_waiters_calls(data: &AtomicUsize) { - data.fetch_add(1 << NOTIFY_WAITERS_SHIFT, SeqCst); + data.fetch_add(1 << NOTIFY_WAITERS_SHIFT, AcqRel); } impl Notify { @@ -562,7 +562,7 @@ impl Notify { pub fn notified(&self) -> Notified<'_> { // we load the number of times notify_waiters // was called and store that in the future. - let state = self.state.load(SeqCst); + let state = self.state.load(Acquire); Notified { notify: self, state: State::Init, @@ -610,7 +610,7 @@ impl Notify { pub fn notified_owned(self: Arc) -> OwnedNotified { // we load the number of times notify_waiters // was called and store that in the future. - let state = self.state.load(SeqCst); + let state = self.state.load(Acquire); OwnedNotified { notify: self, state: State::Init, @@ -673,7 +673,7 @@ impl Notify { fn notify_with_strategy(&self, strategy: NotifyOneStrategy) { // Load the current state - let mut curr = self.state.load(SeqCst); + let mut curr = self.state.load(Acquire); // If the state is `EMPTY`, transition to `NOTIFIED` and return. while let EMPTY | NOTIFIED = get_state(curr) { @@ -681,7 +681,7 @@ impl Notify { // happens-before synchronization must happen between this atomic // operation and a task calling `notified().await`. let new = set_state(curr, NOTIFIED); - let res = self.state.compare_exchange(curr, new, SeqCst, SeqCst); + let res = self.state.compare_exchange(curr, new, AcqRel, Acquire); match res { // No waiters, no further work to do @@ -697,7 +697,7 @@ impl Notify { // The state must be reloaded while the lock is held. The state may only // transition out of WAITING while the lock is held. - curr = self.state.load(SeqCst); + curr = self.state.load(Acquire); if let Some(waker) = notify_locked(&mut waiters, &self.state, curr, strategy) { drop(waiters); @@ -756,7 +756,7 @@ impl Notify { // Increment the number of times this method was called // and transition to empty. let new_state = set_state(inc_num_notify_waiters_calls(curr), EMPTY); - self.state.store(new_state, SeqCst); + self.state.store(new_state, Release); // It is critical for `GuardedLinkedList` safety that the guard node is // pinned in memory and is not dropped until the guarded list is dropped. @@ -819,7 +819,7 @@ impl Notify { // The state must be loaded while the lock is held. The state may only // transition out of WAITING while the lock is held. - let current_state = self.state.load(SeqCst); + let current_state = self.state.load(Acquire); NotifyGuard { guarded_notify: self, @@ -846,14 +846,14 @@ fn notify_locked( ) -> Option { match get_state(curr) { EMPTY | NOTIFIED => { - let res = state.compare_exchange(curr, set_state(curr, NOTIFIED), SeqCst, SeqCst); + let res = state.compare_exchange(curr, set_state(curr, NOTIFIED), AcqRel, Acquire); match res { Ok(_) => None, Err(actual) => { let actual_state = get_state(actual); assert!(actual_state == EMPTY || actual_state == NOTIFIED); - state.store(set_state(actual, NOTIFIED), SeqCst); + state.store(set_state(actual, NOTIFIED), Release); None } } @@ -885,7 +885,7 @@ fn notify_locked( // must be transitioned to `EMPTY`. As transitioning // **from** `WAITING` requires the lock to be held, a // `store` is sufficient. - state.store(set_state(curr, EMPTY), SeqCst); + state.store(set_state(curr, EMPTY), Release); } waker } @@ -1113,7 +1113,7 @@ impl NotifiedProject<'_> { 'outer_loop: loop { match *state { State::Init => { - let curr = notify.state.load(SeqCst); + let curr = notify.state.load(Acquire); // Check if `notify_waiters` was called before attempting to acquire // the `NOTIFIED` state. If a broadcast occurred, we will be woken by it, @@ -1127,8 +1127,8 @@ impl NotifiedProject<'_> { let res = notify.state.compare_exchange( set_state(curr, NOTIFIED), set_state(curr, EMPTY), - SeqCst, - SeqCst, + AcqRel, + Acquire, ); if res.is_ok() { @@ -1146,7 +1146,7 @@ impl NotifiedProject<'_> { let mut waiters = notify.waiters.lock(); // Reload the state with the lock held - let mut curr = notify.state.load(SeqCst); + let mut curr = notify.state.load(Acquire); // if notify_waiters has been called after the future // was created, then we are done @@ -1163,8 +1163,8 @@ impl NotifiedProject<'_> { let res = notify.state.compare_exchange( set_state(curr, EMPTY), set_state(curr, WAITING), - SeqCst, - SeqCst, + AcqRel, + Acquire, ); if let Err(actual) = res { @@ -1180,8 +1180,8 @@ impl NotifiedProject<'_> { let res = notify.state.compare_exchange( set_state(curr, NOTIFIED), set_state(curr, EMPTY), - SeqCst, - SeqCst, + AcqRel, + Acquire, ); match res { @@ -1264,7 +1264,7 @@ impl NotifiedProject<'_> { } // Load the state with the lock held. - let curr = notify.state.load(SeqCst); + let curr = notify.state.load(Acquire); if get_num_notify_waiters_calls(curr) != *notify_waiters_calls { // Before we add a waiter to the list we check if these numbers are @@ -1339,7 +1339,7 @@ impl NotifiedProject<'_> { // longer stored in the linked list. if matches!(*state, State::Waiting) { let mut waiters = notify.waiters.lock(); - let mut notify_state = notify.state.load(SeqCst); + let mut notify_state = notify.state.load(Acquire); // We hold the lock, so this field is not concurrently accessed by // `notify_*` functions and we can use the relaxed ordering. @@ -1354,7 +1354,7 @@ impl NotifiedProject<'_> { if waiters.is_empty() && get_state(notify_state) == WAITING { notify_state = set_state(notify_state, EMPTY); - notify.state.store(notify_state, SeqCst); + notify.state.store(notify_state, Release); } // See if the node was notified but not received. In this case, if