From 91c5b151ece66b2815ff84dfd6b8ee1c9b6d549b Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 21:35:19 +0800 Subject: [PATCH 01/12] chore: fine tune waker wake Signed-off-by: tison --- asyncband/src/barrier/mod.rs | 6 +- asyncband/src/broadcast/mpmc/bounded/mod.rs | 42 ++++------ asyncband/src/broadcast/mpmc/unbounded/mod.rs | 14 ++-- asyncband/src/condvar/mod.rs | 2 +- asyncband/src/event/manual_reset.rs | 3 +- asyncband/src/internal/arena.rs | 17 +--- asyncband/src/internal/semaphore.rs | 2 +- asyncband/src/internal/wakerset.rs | 15 +--- asyncband/src/phaser/mod.rs | 83 +++++++++---------- asyncband/src/watch/mod.rs | 54 ++++++------ 10 files changed, 92 insertions(+), 146 deletions(-) diff --git a/asyncband/src/barrier/mod.rs b/asyncband/src/barrier/mod.rs index b3a0637e..c5760c95 100644 --- a/asyncband/src/barrier/mod.rs +++ b/asyncband/src/barrier/mod.rs @@ -55,7 +55,6 @@ use std::task::Poll; use crate::internal::mutex::Mutex; use crate::internal::wake_all; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -189,10 +188,9 @@ impl Barrier { if state.arrived == self.n { state.arrived = 0; state.generation += 1; - let mut wakers = WakerBatch::new(); - state.waiters.drain_into(&mut wakers); + let wakers = state.waiters.take_all(); drop(state); - wake_all(&mut wakers); + wake_all(wakers); return BarrierWaitResult(true); } diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index 5cc64edd..a1875515 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -120,7 +120,6 @@ use crate::internal::mutex::Mutex; use crate::internal::semaphore::Acquire; use crate::internal::semaphore::Semaphore; use crate::internal::wake_all; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerToken; #[cfg(test)] @@ -212,9 +211,9 @@ impl Shared { // backlog, which would otherwise leave nearly `capacity` permits behind. // // Capping cannot lose a wake-up, by the same argument that lets this read the count at all: - // a producer this load observes is one the release covers, and one it misses incremented - // after the load, which it does before taking the channel lock to recheck — so its recheck - // runs after the reclaim and finds the capacity itself. + // a producer this load observes is the one the release covers, and one it misses + // incremented after the load, which it does before taking the channel lock to recheck — so + // its recheck runs after the reclaim and finds the capacity itself. let waiting = self.waiting_senders(); if freed > 0 && waiting > 0 { self.tx_permits.release_if_nonempty(freed.min(waiting)); @@ -429,7 +428,7 @@ impl BoundedSender { self.publish(msg, |msg| msg) } - /// The publish step both send paths share. + /// The publishing step both send paths share. /// /// `into_msg` is called only once this decides the message will actually be retained, which is /// what lets `try_send` defer its allocation past the capacity check while `try_publish` hands @@ -439,26 +438,21 @@ impl BoundedSender { /// observe an empty buffer and park after this message became visible. fn publish

(&self, payload: P, into_msg: impl FnOnce(P) -> Arc) -> Result<(), P> { let mut discarded = None; - let mut wakers = WakerBatch::new(); - { - let mut inner = self.shared.inner.lock(); - - if !inner.log.has_receivers() { - // Nothing can read this message. The payload leaves the critical section with us - // and is dropped below, so `T::drop` never runs under the lock. - inner.log.publish_discarded(); - discarded = Some(payload); - } else if inner.log.retained() == self.shared.capacity { - // Nothing was published, so there is no wait set to drain. - return Err(payload); - } else { - inner.log.publish_retained(into_msg(payload)); - } - - inner.waiters.drain_into(&mut wakers); + let mut inner = self.shared.inner.lock(); + if !inner.log.has_receivers() { + // Nothing can read this message. The payload leaves the critical section with us + // and is dropped below, so `T::drop` never runs under the lock. + inner.log.publish_discarded(); + discarded = Some(payload); + } else if inner.log.retained() == self.shared.capacity { + // Nothing was published, so there is no wait set to drain. + return Err(payload); + } else { + inner.log.publish_retained(into_msg(payload)); } - - wake_all(&mut wakers); + let wakers = inner.waiters.take_all(); + drop(inner); + wake_all(wakers); drop(discarded); Ok(()) } diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 663e5a6b..19422b9b 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -71,7 +71,6 @@ use super::error::TryRecvError; use crate::internal::arena::SlotId; use crate::internal::mutex::Mutex; use crate::internal::wake_all; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerToken; #[cfg(test)] @@ -177,17 +176,14 @@ impl UnboundedSender { // Publishing and draining the wait set share one critical section, so a receiver can never // observe an empty buffer and park after this message became visible. - let mut wakers = WakerBatch::new(); - let unretained = { - let mut inner = self.shared.inner.lock(); - let unretained = inner.log.publish(msg); - inner.waiters.drain_into(&mut wakers); - unretained - }; + let mut inner = self.shared.inner.lock(); + let unretained = inner.log.publish(msg); + let wakers = inner.waiters.take_all(); + drop(inner); // Notify all waiting receivers. An unsent message is dropped here too, once the lock is // released. - wake_all(&mut wakers); + wake_all(wakers); drop(unretained); } diff --git a/asyncband/src/condvar/mod.rs b/asyncband/src/condvar/mod.rs index f2b06524..47567d33 100644 --- a/asyncband/src/condvar/mod.rs +++ b/asyncband/src/condvar/mod.rs @@ -175,7 +175,7 @@ impl Condvar { {} } - wake_all(&mut wakers); + wake_all(wakers); } /// Waits for a notification, atomically releasing and then reacquiring the mutex. diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index 4e0e2ce9..d2b0bd47 100644 --- a/asyncband/src/event/manual_reset.rs +++ b/asyncband/src/event/manual_reset.rs @@ -117,8 +117,7 @@ impl ManualResetEvent { } } } - - wake_all(&mut wakers); + wake_all(wakers); } /// Clears the set state. diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index ce0bc7bd..fcd11de6 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -166,21 +166,6 @@ impl Arena { value } - /// Drains every occupied value in slot order while retaining the allocation for reuse. - /// - /// After a non-empty drain, every previously issued slot ID becomes invalid, including IDs for - /// slots that were already vacant. Consumers that retain IDs across this operation must supply - /// their own epoch check. - #[inline] - pub fn drain(&mut self) -> impl Iterator + '_ { - self.vacant_head = None; - self.len = 0; - self.slots.drain(..).filter_map(|slot| match slot { - Slot::Occupied(value) => Some(value), - Slot::Vacant { .. } => None, - }) - } - /// Takes every occupied value and the backing allocation in slot order. #[inline] pub fn take_all(&mut self) -> impl Iterator + use { @@ -227,7 +212,7 @@ mod tests { let capacity = arena.slots.capacity(); arena.remove(second); - assert_eq!(arena.drain().collect::>(), vec![1, 3]); + assert_eq!(arena.take_all().collect::>(), vec![1, 3]); assert_eq!(arena.len(), 0); assert_eq!(arena.slots.capacity(), capacity); diff --git a/asyncband/src/internal/semaphore.rs b/asyncband/src/internal/semaphore.rs index 51c7baed..065d2f07 100644 --- a/asyncband/src/internal/semaphore.rs +++ b/asyncband/src/internal/semaphore.rs @@ -173,7 +173,7 @@ impl Semaphore { } } drop(waiters); - wake_all(&mut wakers); + wake_all(wakers); } fn insert_permits_with_lock<'a>( diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index c74a89b5..e6cc16f6 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -26,13 +26,11 @@ use std::task::Waker; use crate::internal::arena::Arena; use crate::internal::arena::SlotId; -use crate::internal::waker_batch::WakerBatch; /// An exclusive handle to one waker slot in a [`WakerSet`]. /// /// This token deliberately does not implement `Clone` or `Copy`. Its owner must not pass it back -/// to the set after the registration has been detached by [`WakerSet::drain_into`] or -/// [`WakerSet::take_all`]. +/// to the set after the registration has been detached by [`WakerSet::take_all`]. #[derive(Debug)] pub struct WakerToken(SlotId); @@ -57,17 +55,6 @@ impl WakerSet { } } - /// Drains all registered wakers into `batch` while retaining slot capacity. - /// - /// The batch is filled in place because its inline storage is too large to move for free: - /// returning it by value costs every publish about 6ns even when nothing is registered. The - /// caller must invalidate every outstanding token and consume or drop the batch after - /// releasing the lock that protects this set. - #[inline] - pub fn drain_into(&mut self, batch: &mut WakerBatch) { - batch.extend(self.wakers.drain()); - } - /// Takes all registered wakers together with the set's backing allocation. /// /// The caller must invalidate every outstanding token and consume or drop the iterator after diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 9b00e7c5..78f2be73 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -131,7 +131,6 @@ use std::task::Poll; use crate::internal::mutex::Mutex; use crate::internal::wake_all; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -172,14 +171,13 @@ struct State { } impl State { - /// Completes the phase once every participant has arrived, moving its waiters into `wakers`. - fn advance_if_ready(&mut self, wakers: &mut WakerBatch) { + /// Completes the phase once every participant has arrived. + fn advance_if_ready(&mut self) { if self.closed || self.unarrived != 0 { return; } self.phase = self.phase.wrapping_add(1); self.unarrived = self.registered; - self.waiters.drain_into(wakers); } fn completion(&self, observed: u64) -> Poll> { @@ -385,16 +383,15 @@ impl Drop for PhaserParticipants { if self.remaining == 0 { return; } - let mut wakers = WakerBatch::new(); - { - let mut state = self.phaser.state.lock(); - // Unyielded participants have never arrived and prevent their phase from advancing. - state.registered -= self.remaining; - state.unarrived -= self.remaining; - self.remaining = 0; - state.advance_if_ready(&mut wakers); - } - wake_all(&mut wakers); + let mut state = self.phaser.state.lock(); + // Unyielded participants have never arrived and prevent their phase from advancing. + state.registered -= self.remaining; + state.unarrived -= self.remaining; + self.remaining = 0; + state.advance_if_ready(); + let wakers = state.waiters.take_all(); + drop(state); + wake_all(wakers); } } @@ -433,21 +430,19 @@ impl PhaserParticipant { /// Repeated calls within one phase count only once. After advancement, an explicit new call /// arrives in the new phase and replaces any previous pending observation. pub fn arrive(&mut self) -> Result { - let mut wakers = WakerBatch::new(); - let phase = { - let mut state = self.phaser.state.lock(); - if state.closed { - return Err(Closed(())); - } - let phase = state.phase; - if self.pending != Some(phase) { - state.unarrived -= 1; - } - self.pending = Some(phase); - state.advance_if_ready(&mut wakers); - phase - }; - wake_all(&mut wakers); + let mut state = self.phaser.state.lock(); + if state.closed { + return Err(Closed(())); + } + let phase = state.phase; + if self.pending != Some(phase) { + state.unarrived -= 1; + } + self.pending = Some(phase); + state.advance_if_ready(); + let wakers = state.waiters.take_all(); + drop(state); + wake_all(wakers); Ok(phase) } @@ -482,23 +477,21 @@ impl PhaserParticipant { } fn do_deregister(&mut self) -> Result { - let mut wakers = WakerBatch::new(); - let result = { - let mut state = self.phaser.state.lock(); - self.registered = false; - state.registered -= 1; - if self.pending != Some(state.phase) { - state.unarrived -= 1; - } - let result = if state.closed { - Err(Closed(())) - } else { - Ok(state.phase) - }; - state.advance_if_ready(&mut wakers); - result + let mut state = self.phaser.state.lock(); + self.registered = false; + state.registered -= 1; + if self.pending != Some(state.phase) { + state.unarrived -= 1; + } + let result = if state.closed { + Err(Closed(())) + } else { + Ok(state.phase) }; - wake_all(&mut wakers); + state.advance_if_ready(); + let wakers = state.waiters.take_all(); + drop(state); + wake_all(wakers); result } } diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index 90b71923..48d35138 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -75,7 +75,6 @@ pub use self::error::RecvError; pub use self::error::SendError; use crate::internal::mutex::Mutex; use crate::internal::wake_all; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -169,23 +168,21 @@ impl Sender { /// /// Panics if the channel has already published `u64::MAX` updates. pub fn send(&self, value: T) -> Result<(), SendError> { - let mut wakers = WakerBatch::new(); - let replaced = { - let mut state = self.shared.state.lock(); - if state.receivers == 0 { - return Err(SendError::new(value)); - } - let version = state - .version - .checked_add(1) - .expect("watch channel version counter overflowed"); - let replaced = mem::replace(&mut state.value, value); - state.version = version; - state.waiters.drain_into(&mut wakers); - replaced - }; + let mut state = self.shared.state.lock(); + if state.receivers == 0 { + return Err(SendError::new(value)); + } + let version = state + .version + .checked_add(1) + .expect("watch channel version counter overflowed"); + let replaced = mem::replace(&mut state.value, value); + state.version = version; + let wakers = state.waiters.take_all(); + drop(state); + // Waker callbacks and the replaced value's destructor may reenter this channel. - wake_all(&mut wakers); + wake_all(wakers); drop(replaced); Ok(()) } @@ -199,19 +196,16 @@ impl Sender { /// /// Panics if the channel has already published `u64::MAX` updates. pub fn send_replace(&self, value: T) -> T { - let mut wakers = WakerBatch::new(); - let replaced = { - let mut state = self.shared.state.lock(); - let version = state - .version - .checked_add(1) - .expect("watch channel version counter overflowed"); - let replaced = mem::replace(&mut state.value, value); - state.version = version; - state.waiters.drain_into(&mut wakers); - replaced - }; - wake_all(&mut wakers); + let mut state = self.shared.state.lock(); + let version = state + .version + .checked_add(1) + .expect("watch channel version counter overflowed"); + let replaced = mem::replace(&mut state.value, value); + state.version = version; + let wakers = state.waiters.take_all(); + drop(state); + wake_all(wakers); replaced } From 6e1722d35e37edee6f1edb581c29590b928e4dcc Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 21:43:10 +0800 Subject: [PATCH 02/12] fixup Signed-off-by: tison --- asyncband/src/internal/arena.rs | 17 ----------------- 1 file changed, 17 deletions(-) diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index fcd11de6..911b3a53 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -203,23 +203,6 @@ mod tests { assert_eq!(arena.get(second), Some(&"second")); } - #[test] - fn drain_restarts_slot_id_allocation() { - let mut arena = Arena::with_capacity(3); - let first = arena.insert(1); - let second = arena.insert(2); - let third = arena.insert(3); - let capacity = arena.slots.capacity(); - arena.remove(second); - - assert_eq!(arena.take_all().collect::>(), vec![1, 3]); - assert_eq!(arena.len(), 0); - assert_eq!(arena.slots.capacity(), capacity); - - let slot_ids = [arena.insert(4), arena.insert(5), arena.insert(6)]; - assert_eq!(slot_ids, [first, second, third]); - } - #[test] fn take_all_releases_the_backing_allocation() { let mut arena = Arena::new(); From 1dcbd46d6faf5fb91240e09f308439d10a9daa31 Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 21:49:35 +0800 Subject: [PATCH 03/12] fix(phaser): retain waiters until phase advancement Signed-off-by: tison --- asyncband/src/phaser/mod.rs | 27 ++++++++------- tests-integration/tests/phaser_test.rs | 47 ++++++++++++++++++++++++++ 2 files changed, 61 insertions(+), 13 deletions(-) diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 78f2be73..b0b7d12c 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -128,6 +128,7 @@ use std::pin::Pin; use std::sync::Arc; use std::task::Context; use std::task::Poll; +use std::task::Waker; use crate::internal::mutex::Mutex; use crate::internal::wake_all; @@ -171,13 +172,16 @@ struct State { } impl State { - /// Completes the phase once every participant has arrived. - fn advance_if_ready(&mut self) { - if self.closed || self.unarrived != 0 { - return; - } - self.phase = self.phase.wrapping_add(1); - self.unarrived = self.registered; + /// Advances a completed phase and returns its waiters for waking outside the lock. + fn advance_if_ready(&mut self) -> impl Iterator + 'static { + let wakers = if self.closed || self.unarrived != 0 { + None + } else { + self.phase = self.phase.wrapping_add(1); + self.unarrived = self.registered; + Some(self.waiters.take_all()) + }; + wakers.into_iter().flatten() } fn completion(&self, observed: u64) -> Poll> { @@ -388,8 +392,7 @@ impl Drop for PhaserParticipants { state.registered -= self.remaining; state.unarrived -= self.remaining; self.remaining = 0; - state.advance_if_ready(); - let wakers = state.waiters.take_all(); + let wakers = state.advance_if_ready(); drop(state); wake_all(wakers); } @@ -439,8 +442,7 @@ impl PhaserParticipant { state.unarrived -= 1; } self.pending = Some(phase); - state.advance_if_ready(); - let wakers = state.waiters.take_all(); + let wakers = state.advance_if_ready(); drop(state); wake_all(wakers); Ok(phase) @@ -488,8 +490,7 @@ impl PhaserParticipant { } else { Ok(state.phase) }; - state.advance_if_ready(); - let wakers = state.waiters.take_all(); + let wakers = state.advance_if_ready(); drop(state); wake_all(wakers); result diff --git a/tests-integration/tests/phaser_test.rs b/tests-integration/tests/phaser_test.rs index 28c225a1..af81580d 100644 --- a/tests-integration/tests/phaser_test.rs +++ b/tests-integration/tests/phaser_test.rs @@ -343,6 +343,53 @@ fn advancing_a_phase_wakes_every_registered_waiter_once() { )); } +#[test] +fn partial_arrivals_and_withdrawals_preserve_pending_waits() { + let phaser = Phaser::new(); + let observed = phaser.phase(); + let mut participants = phaser.register(4).unwrap(); + let mut first = participants.next().unwrap(); + let withdrawing = participants.next().unwrap(); + let mut last = participants.next().unwrap(); + let (waker, counter) = WakeCounter::new(); + let mut wait = Box::pin(phaser.wait(observed)); + assert!(poll_with(wait.as_mut(), &waker).is_pending()); + + first.arrive().unwrap(); + first.arrive().unwrap(); + withdrawing.deregister().unwrap(); + drop(participants); + let early_wakes = counter.count(); + + last.arrive().unwrap(); + assert_eq!(early_wakes, 0); + assert_eq!(counter.count(), 1); + assert_eq!(poll_once(wait.as_mut()), Poll::Ready(Ok(phaser.phase()))); +} + +#[test] +fn cancelling_after_partial_arrival_preserves_other_waiters() { + let phaser = Phaser::new(); + let observed = phaser.phase(); + let mut first = phaser.register_one().unwrap(); + let mut second = phaser.register_one().unwrap(); + let mut cancelled = Box::pin(phaser.wait(observed)); + assert!(poll_once(cancelled.as_mut()).is_pending()); + + first.arrive().unwrap(); + let (waker, counter) = WakeCounter::new(); + let mut remaining = Box::pin(phaser.wait(observed)); + assert!(poll_with(remaining.as_mut(), &waker).is_pending()); + drop(cancelled); + + second.arrive().unwrap(); + assert_eq!(counter.count(), 1); + assert_eq!( + poll_once(remaining.as_mut()), + Poll::Ready(Ok(phaser.phase())) + ); +} + #[test] fn repolling_updates_the_task_that_will_be_notified() { let phaser = Phaser::new(); From b780e36d4a5c4c66fa5ba9e5ccf82bd379df902f Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 21:59:03 +0800 Subject: [PATCH 04/12] perf(barrier): reuse waiter storage across generations Signed-off-by: tison --- asyncband/src/barrier/mod.rs | 6 +++++- asyncband/src/internal/arena.rs | 31 ++++++++++++++++++++++++++++++ asyncband/src/internal/wakerset.rs | 12 +++++++++++- 3 files changed, 47 insertions(+), 2 deletions(-) diff --git a/asyncband/src/barrier/mod.rs b/asyncband/src/barrier/mod.rs index c5760c95..cb426468 100644 --- a/asyncband/src/barrier/mod.rs +++ b/asyncband/src/barrier/mod.rs @@ -55,6 +55,7 @@ use std::task::Poll; use crate::internal::mutex::Mutex; use crate::internal::wake_all; +use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -188,7 +189,10 @@ impl Barrier { if state.arrived == self.n { state.arrived = 0; state.generation += 1; - let wakers = state.waiters.take_all(); + let mut wakers = WakerBatch::new(); + for waker in state.waiters.drain() { + wakers.push(waker); + } drop(state); wake_all(wakers); return BarrierWaitResult(true); diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index 911b3a53..fd9a0bd0 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -166,6 +166,19 @@ impl Arena { value } + /// Drains occupied values in slot order, retaining the allocation for reuse. + /// + /// All previous slot IDs become invalid and may be reused after the drain. + #[inline] + pub fn drain(&mut self) -> impl Iterator + '_ { + self.vacant_head = None; + self.len = 0; + self.slots.drain(..).filter_map(|slot| match slot { + Slot::Occupied(value) => Some(value), + Slot::Vacant { .. } => None, + }) + } + /// Takes every occupied value and the backing allocation in slot order. #[inline] pub fn take_all(&mut self) -> impl Iterator + use { @@ -203,6 +216,24 @@ mod tests { assert_eq!(arena.get(second), Some(&"second")); } + #[test] + fn drain_retains_capacity_and_restarts_slot_ids() { + let mut arena = Arena::with_capacity(3); + let first = arena.insert(1); + let second = arena.insert(2); + let third = arena.insert(3); + let capacity = arena.slots.capacity(); + arena.remove(second); + + assert_eq!(arena.drain().collect::>(), vec![1, 3]); + assert_eq!(arena.len(), 0); + assert_eq!(arena.slots.capacity(), capacity); + assert_eq!( + [arena.insert(4), arena.insert(5), arena.insert(6)], + [first, second, third] + ); + } + #[test] fn take_all_releases_the_backing_allocation() { let mut arena = Arena::new(); diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index e6cc16f6..658dbd28 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -30,7 +30,8 @@ use crate::internal::arena::SlotId; /// An exclusive handle to one waker slot in a [`WakerSet`]. /// /// This token deliberately does not implement `Clone` or `Copy`. Its owner must not pass it back -/// to the set after the registration has been detached by [`WakerSet::take_all`]. +/// to the set after the registration has been detached by [`WakerSet::drain`] or +/// [`WakerSet::take_all`]. #[derive(Debug)] pub struct WakerToken(SlotId); @@ -55,6 +56,15 @@ impl WakerSet { } } + /// Drains registered wakers while retaining slot capacity for reuse. + /// + /// The caller must invalidate outstanding tokens and collect the wakers under the set's lock, + /// then wake or drop them after releasing it. + #[inline] + pub fn drain(&mut self) -> impl Iterator + '_ { + self.wakers.drain() + } + /// Takes all registered wakers together with the set's backing allocation. /// /// The caller must invalidate every outstanding token and consume or drop the iterator after From abb9c2a6d9ee65d5753172216284687105b75a5b Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 22:00:13 +0800 Subject: [PATCH 05/12] refactor(condvar): simplify the notify_all loop Signed-off-by: tison --- asyncband/src/condvar/mod.rs | 20 ++++++++------------ 1 file changed, 8 insertions(+), 12 deletions(-) diff --git a/asyncband/src/condvar/mod.rs b/asyncband/src/condvar/mod.rs index 47567d33..b3821239 100644 --- a/asyncband/src/condvar/mod.rs +++ b/asyncband/src/condvar/mod.rs @@ -161,18 +161,14 @@ impl Condvar { { let mut waiters = self.waiters.lock(); - while waiters - .unlink_first_waiter(|node| { - let WaitState::Waiting(waker) = - mem::replace(&mut node.state, WaitState::NotifiedAll) - else { - unreachable!("only waiting tasks remain linked") - }; - wakers.push(waker); - true - }) - .is_some() - {} + while let Some((_, node)) = waiters.unlink_first_waiter(|_| true) { + let WaitState::Waiting(waker) = + mem::replace(&mut node.state, WaitState::NotifiedAll) + else { + unreachable!("only waiting tasks remain linked") + }; + wakers.push(waker); + } } wake_all(wakers); From c1a36ff3bd5a2e04ae0ccc7f5b447798bc8c0059 Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 22:01:32 +0800 Subject: [PATCH 06/12] refactor: trim unused waker batch interfaces Signed-off-by: tison --- asyncband/src/internal/mod.rs | 6 ------ asyncband/src/internal/waker_batch.rs | 29 +++++++++------------------ 2 files changed, 10 insertions(+), 25 deletions(-) diff --git a/asyncband/src/internal/mod.rs b/asyncband/src/internal/mod.rs index 76d3ce80..602d7e66 100644 --- a/asyncband/src/internal/mod.rs +++ b/asyncband/src/internal/mod.rs @@ -124,15 +124,9 @@ pub(crate) mod waitlist; feature = "barrier", feature = "broadcast", feature = "event", - feature = "completion", - feature = "latch", feature = "mutex", - feature = "once", - feature = "phaser", feature = "rwlock", feature = "semaphore", - feature = "waitgroup", - feature = "watch", ))] // Only the semaphore refills a batch and asks whether it will spill, so other feature subsets // leave that method unused. diff --git a/asyncband/src/internal/waker_batch.rs b/asyncband/src/internal/waker_batch.rs index a41c0f37..2a41105f 100644 --- a/asyncband/src/internal/waker_batch.rs +++ b/asyncband/src/internal/waker_batch.rs @@ -22,11 +22,8 @@ use std::task::Waker; /// An owning FIFO of wakers that stores the first [`Self::STACK_SIZE`] entries without allocating. /// -/// The batch is filled through [`WakerBatch::push`] or [`Extend`] and consumed as its own -/// iterator. Entries are written only as they are pushed, so constructing an empty or small batch -/// touches nothing beyond the two indices. Once every inline entry has been yielded the batch -/// reuses that storage, so a caller that alternates between filling and draining, as the -/// semaphore does, keeps running on the stack. +/// Only pushed entries are initialized. Consuming all inline entries makes their storage +/// available for the semaphore's next notification batch. pub struct WakerBatch { /// The initialized entries are exactly `start..end`. inline: [MaybeUninit; Self::STACK_SIZE], @@ -75,14 +72,6 @@ impl WakerBatch { } } -impl Extend for WakerBatch { - fn extend>(&mut self, iter: I) { - for waker in iter { - self.push(waker); - } - } -} - impl Iterator for WakerBatch { type Item = Waker; @@ -153,8 +142,10 @@ mod tests { })) } - fn wakers(log: &Arc, count: usize) -> impl Iterator + '_ { - (0..count).map(move |id| waker(log, id)) + fn push_wakers(batch: &mut WakerBatch, log: &Arc, count: usize) { + for id in 0..count { + batch.push(waker(log, id)); + } } fn alive(log: &Arc) -> usize { @@ -170,7 +161,7 @@ mod tests { let log = log(); let count = STACK_SIZE + 8; let mut batch = WakerBatch::new(); - batch.extend(wakers(&log, count)); + push_wakers(&mut batch, &log, count); assert!(batch.will_spill()); for waker in &mut batch { @@ -188,7 +179,7 @@ mod tests { let count = STACK_SIZE + 8; for consumed in [0, 5, STACK_SIZE, STACK_SIZE + 3, count] { let mut batch = WakerBatch::new(); - batch.extend(wakers(&log, count)); + push_wakers(&mut batch, &log, count); for _ in 0..consumed { drop(batch.next().unwrap()); } @@ -204,7 +195,7 @@ mod tests { let log = log(); let mut batch = WakerBatch::new(); for _ in 0..3 { - batch.extend(wakers(&log, STACK_SIZE)); + push_wakers(&mut batch, &log, STACK_SIZE); assert!(batch.will_spill()); assert_eq!(batch.by_ref().count(), STACK_SIZE); assert!(!batch.will_spill()); @@ -216,7 +207,7 @@ mod tests { fn keeps_push_order_while_spilled() { let log = log(); let mut batch = WakerBatch::new(); - batch.extend(wakers(&log, STACK_SIZE + 1)); + push_wakers(&mut batch, &log, STACK_SIZE + 1); // Free inline room; the spilled entry must still come out before anything pushed now. for _ in 0..4 { batch.next().unwrap().wake(); From 043d77f639558fd47c5a833859fd0a84d2df50f4 Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 22:03:49 +0800 Subject: [PATCH 07/12] docs(broadcast): consolidate sender notification invariants Signed-off-by: tison --- asyncband/src/broadcast/mpmc/bounded/mod.rs | 51 ++++++--------------- 1 file changed, 13 insertions(+), 38 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index a1875515..7d1c8dba 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -180,64 +180,39 @@ struct Shared { /// and parks again if another producer took the space first. The semaphore starts empty and /// only ever grows when a reclaim finds someone waiting, so an idle channel accumulates none. tx_permits: Semaphore, - /// How many producers are somewhere inside the waiting path of [`BoundedSender::send`]. - /// - /// An upper bound on the number of parked producers, and the only thing either release path - /// consults. It answers both questions a reclaim has — whether to wake anyone, and how many - /// permits are worth handing out — without taking the semaphore's lock. Reclaiming is far more - /// frequent than blocking — under fan-out every message is reclaimed, while a channel with - /// headroom never blocks at all — so paying an atomic load there instead of a lock acquisition - /// is what keeps an uncontended receive off the semaphore entirely. + /// Number of producers in the waiting path of [`BoundedSender::send`]. + /// + /// This upper bound lets receives skip the semaphore lock when no producer is waiting. blocked_senders: AtomicUsize, } impl Shared { - /// Hands `freed` released slots back to producers parked in `send`. - /// - /// Capacity is `retained()`, which is `buffer.len()`. The buffer grows only in - /// `Backlog::publish_retained` and shrinks only in `Backlog::reclaim_vacated`, which is - /// reachable from exactly two places: a receive that vacates the last cursor at the backlog - /// head, and removing a subscription. Those are the only callers of this method, so no path - /// can free capacity without waking a producer. Subscribing cannot: a new cursor starts at the - /// tail and never lowers `retained()`. + /// Notifies blocked producers after a receive or subscription removal frees capacity. /// - /// Callers must invoke this with the channel unlocked, and — on the receive path — before - /// touching the payload, since `common::take_msg` runs user code that may panic. + /// Call with the channel unlocked, before payload cloning or destruction can panic. fn release_reclaimed(&self, freed: usize) { - // Release no more permits than there are producers to wake. A permit the semaphore cannot - // hand to a waiter is kept as slack, and the next producer to block has to burn it off one - // futile publish attempt — a channel lock apiece — at a time before it can park. Freeing a - // large prefix at once is not exotic: dropping a lagging subscription reclaims the whole - // backlog, which would otherwise leave nearly `capacity` permits behind. - // - // Capping cannot lose a wake-up, by the same argument that lets this read the count at all: - // a producer this load observes is the one the release covers, and one it misses - // incremented after the load, which it does before taking the channel lock to recheck — so - // its recheck runs after the reclaim and finds the capacity itself. + // Cap permits at the number of blocked producers. Surplus permits make later sends retry + // a full channel instead of parking, especially after dropping a lagging subscription. let waiting = self.waiting_senders(); if freed > 0 && waiting > 0 { self.tx_permits.release_if_nonempty(freed.min(waiting)); } } - /// Wakes every parked producer, however many slots came back. + /// Wakes every blocked producer when the last subscription leaves. /// - /// The last subscription leaving is not a reclaim of some number of slots — it removes the - /// limit itself, because a channel with no receivers discards instead of retaining. Releasing - /// only as many permits as that final reclaim freed would strand every producer beyond that - /// count, so this is the one release that must be unbounded. + /// Sends now discard payloads without consuming capacity, so every producer can proceed. fn release_all(&self) { if self.waiting_senders() > 0 { self.tx_permits.notify_all(); } } - /// How many producers might be waiting, answered without touching the semaphore's lock. + /// Returns an upper bound on parked producers without locking the semaphore. /// - /// This cannot miss a wake-up. A producer increments the count before it ever takes the - /// channel lock to recheck capacity, and every caller here loads it after releasing that same - /// lock, so the mutex orders the two: either this load observes the producer, or the - /// producer's recheck runs after the change and finds the capacity itself. + /// Producers increment before rechecking capacity under the channel lock; reclaim paths load + /// after releasing it. A producer is either covered by this count or rechecks capacity after + /// the reclaim, so skipping or capping notifications cannot strand it. fn waiting_senders(&self) -> usize { self.blocked_senders.load(Ordering::Acquire) } From 048f0601b729e9421fe3bd71e45fe08adacdeaf1 Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 22:33:31 +0800 Subject: [PATCH 08/12] refactor: remove the generic wake_all wrapper Signed-off-by: tison --- asyncband/src/barrier/mod.rs | 4 ++-- asyncband/src/broadcast/mpmc/bounded/mod.rs | 4 ++-- asyncband/src/broadcast/mpmc/common.rs | 4 ++-- asyncband/src/broadcast/mpmc/unbounded/mod.rs | 4 ++-- asyncband/src/completion/mod.rs | 6 +++--- asyncband/src/condvar/mod.rs | 3 +-- asyncband/src/event/manual_reset.rs | 3 +-- asyncband/src/internal/countdown.rs | 4 ++-- asyncband/src/internal/mod.rs | 10 ---------- asyncband/src/internal/semaphore.rs | 5 ++--- asyncband/src/phaser/mod.rs | 9 ++++----- asyncband/src/waitgroup/mod.rs | 4 ++-- asyncband/src/watch/mod.rs | 8 ++++---- 13 files changed, 27 insertions(+), 41 deletions(-) diff --git a/asyncband/src/barrier/mod.rs b/asyncband/src/barrier/mod.rs index cb426468..20e5d364 100644 --- a/asyncband/src/barrier/mod.rs +++ b/asyncband/src/barrier/mod.rs @@ -52,9 +52,9 @@ use std::future::Future; use std::pin::Pin; use std::task::Context; use std::task::Poll; +use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -194,7 +194,7 @@ impl Barrier { wakers.push(waker); } drop(state); - wake_all(wakers); + wakers.by_ref().for_each(Waker::wake); return BarrierWaitResult(true); } diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index 7d1c8dba..6ca514c0 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -108,6 +108,7 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; +use std::task::Waker; use super::common; use super::common::Backlog; @@ -119,7 +120,6 @@ use crate::internal::arena::SlotId; use crate::internal::mutex::Mutex; use crate::internal::semaphore::Acquire; use crate::internal::semaphore::Semaphore; -use crate::internal::wake_all; use crate::internal::wakerset::WakerToken; #[cfg(test)] @@ -427,7 +427,7 @@ impl BoundedSender { } let wakers = inner.waiters.take_all(); drop(inner); - wake_all(wakers); + wakers.for_each(Waker::wake); drop(discarded); Ok(()) } diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index 4b2d321b..db09605b 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -31,13 +31,13 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; +use std::task::Waker; use super::error::RecvError; use super::error::TryRecvError; use crate::internal::arena::Arena; use crate::internal::arena::SlotId; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -436,7 +436,7 @@ pub fn disconnect(inner: &Mutex>) { let mut inner = inner.lock(); inner.waiters.take_all() }; - wake_all(wakers); + wakers.for_each(Waker::wake); } /// Releases a cancelled receive's waker registration, dropping the waker unlocked. diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 19422b9b..6725f69c 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -62,6 +62,7 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; +use std::task::Waker; use super::common; use super::common::Backlog; @@ -70,7 +71,6 @@ use super::error::RecvError; use super::error::TryRecvError; use crate::internal::arena::SlotId; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerToken; #[cfg(test)] @@ -183,7 +183,7 @@ impl UnboundedSender { // Notify all waiting receivers. An unsent message is dropped here too, once the lock is // released. - wake_all(wakers); + wakers.for_each(Waker::wake); drop(unretained); } diff --git a/asyncband/src/completion/mod.rs b/asyncband/src/completion/mod.rs index 2055ec1c..cba4dfd0 100644 --- a/asyncband/src/completion/mod.rs +++ b/asyncband/src/completion/mod.rs @@ -56,9 +56,9 @@ use std::sync::OnceLock; use std::sync::Weak; use std::task::Context; use std::task::Poll; +use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -135,7 +135,7 @@ impl Completer { // `complete` consumes the only completer. Disarm its destructor before invoking arbitrary // wake callbacks; the completed state no longer needs abandonment handling. self.shared = Weak::new(); - wake_all(wakers); + wakers.for_each(Waker::wake); Ok(()) } } @@ -151,7 +151,7 @@ impl Drop for Completer { assert!(shared.result.set(None).is_ok()); wakers }; - wake_all(wakers); + wakers.for_each(Waker::wake); } } diff --git a/asyncband/src/condvar/mod.rs b/asyncband/src/condvar/mod.rs index b3821239..9cc2e8db 100644 --- a/asyncband/src/condvar/mod.rs +++ b/asyncband/src/condvar/mod.rs @@ -68,7 +68,6 @@ use std::task::Waker; use crate::internal::mutex::Mutex; use crate::internal::waitlist::WaitList; use crate::internal::waitlist::WaiterId; -use crate::internal::wake_all; use crate::internal::waker_batch::WakerBatch; use crate::mutex; use crate::mutex::MutexGuard; @@ -171,7 +170,7 @@ impl Condvar { } } - wake_all(wakers); + wakers.by_ref().for_each(Waker::wake); } /// Waits for a notification, atomically releasing and then reacquiring the mutex. diff --git a/asyncband/src/event/manual_reset.rs b/asyncband/src/event/manual_reset.rs index d2b0bd47..39d5bcf9 100644 --- a/asyncband/src/event/manual_reset.rs +++ b/asyncband/src/event/manual_reset.rs @@ -27,7 +27,6 @@ use crate::internal::mutex::Mutex; use crate::internal::register_waker; use crate::internal::waitlist::WaitList; use crate::internal::waitlist::WaiterId; -use crate::internal::wake_all; use crate::internal::waker_batch::WakerBatch; /// A reusable signal that releases all waiters and remains set until explicitly reset. @@ -117,7 +116,7 @@ impl ManualResetEvent { } } } - wake_all(wakers); + wakers.by_ref().for_each(Waker::wake); } /// Clears the set state. diff --git a/asyncband/src/internal/countdown.rs b/asyncband/src/internal/countdown.rs index 53ff5e35..92ec06b1 100644 --- a/asyncband/src/internal/countdown.rs +++ b/asyncband/src/internal/countdown.rs @@ -19,9 +19,9 @@ use std::sync::atomic::AtomicU32; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; +use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -60,7 +60,7 @@ impl CountdownState { waiters.take_all() }; - wake_all(wakers); + wakers.for_each(Waker::wake); } /// Polls for zero, registering the current waker if the countdown is still active. diff --git a/asyncband/src/internal/mod.rs b/asyncband/src/internal/mod.rs index 602d7e66..da1423e7 100644 --- a/asyncband/src/internal/mod.rs +++ b/asyncband/src/internal/mod.rs @@ -33,16 +33,6 @@ pub fn register_waker(slot: &mut Option, waker: &Waker) -> Option } } -/// Wakes every waker. -#[inline] -// A no-feature or blocking-only build has no primitive that fans notifications out. -#[allow(dead_code)] -pub(crate) fn wake_all(wakers: impl Iterator) { - for waker in wakers { - waker.wake(); - } -} - #[cfg(any( feature = "barrier", feature = "broadcast", diff --git a/asyncband/src/internal/semaphore.rs b/asyncband/src/internal/semaphore.rs index 065d2f07..4ce4ad71 100644 --- a/asyncband/src/internal/semaphore.rs +++ b/asyncband/src/internal/semaphore.rs @@ -37,7 +37,6 @@ use crate::internal::mutex::Mutex; use crate::internal::register_waker; use crate::internal::waitlist::WaitList; use crate::internal::waitlist::WaiterId; -use crate::internal::wake_all; use crate::internal::waker_batch::WakerBatch; /// The internal semaphore that provides low-level async primitives. @@ -173,7 +172,7 @@ impl Semaphore { } } drop(waiters); - wake_all(wakers); + wakers.by_ref().for_each(Waker::wake); } fn insert_permits_with_lock<'a>( @@ -212,7 +211,7 @@ impl Semaphore { } drop(waiters); - wake_all(&mut batch); + batch.by_ref().for_each(Waker::wake); if rem == 0 { return; } diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index b0b7d12c..1be948a8 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -131,7 +131,6 @@ use std::task::Poll; use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -251,7 +250,7 @@ impl Phaser { state.closed = true; state.waiters.take_all() }; - wake_all(wakers); + wakers.for_each(Waker::wake); } /// Returns an instantaneous count of registered participants, including those already arrived. @@ -394,7 +393,7 @@ impl Drop for PhaserParticipants { self.remaining = 0; let wakers = state.advance_if_ready(); drop(state); - wake_all(wakers); + wakers.for_each(Waker::wake); } } @@ -444,7 +443,7 @@ impl PhaserParticipant { self.pending = Some(phase); let wakers = state.advance_if_ready(); drop(state); - wake_all(wakers); + wakers.for_each(Waker::wake); Ok(phase) } @@ -492,7 +491,7 @@ impl PhaserParticipant { }; let wakers = state.advance_if_ready(); drop(state); - wake_all(wakers); + wakers.for_each(Waker::wake); result } } diff --git a/asyncband/src/waitgroup/mod.rs b/asyncband/src/waitgroup/mod.rs index b7e445a4..79238f30 100644 --- a/asyncband/src/waitgroup/mod.rs +++ b/asyncband/src/waitgroup/mod.rs @@ -65,9 +65,9 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; +use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -106,7 +106,7 @@ impl State { let mut waiters = self.waiters.lock(); waiters.take_all() }; - wake_all(wakers); + wakers.for_each(Waker::wake); } fn poll_wait(&self, token: &mut Option, cx: &mut Context<'_>) -> Poll<()> { diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index 48d35138..989e0ea7 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -70,11 +70,11 @@ use std::pin::Pin; use std::sync::Arc; use std::task::Context; use std::task::Poll; +use std::task::Waker; pub use self::error::RecvError; pub use self::error::SendError; use crate::internal::mutex::Mutex; -use crate::internal::wake_all; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -155,7 +155,7 @@ impl Drop for Sender { } state.waiters.take_all() }; - wake_all(wakers); + wakers.for_each(Waker::wake); } } @@ -182,7 +182,7 @@ impl Sender { drop(state); // Waker callbacks and the replaced value's destructor may reenter this channel. - wake_all(wakers); + wakers.for_each(Waker::wake); drop(replaced); Ok(()) } @@ -205,7 +205,7 @@ impl Sender { state.version = version; let wakers = state.waiters.take_all(); drop(state); - wake_all(wakers); + wakers.for_each(Waker::wake); replaced } From a83bc93338c2509133f52425d9629b2e81a3731c Mon Sep 17 00:00:00 2001 From: tison Date: Fri, 2 Oct 2026 22:52:21 +0800 Subject: [PATCH 09/12] refactor: distinguish reusable and terminal waker collection Signed-off-by: tison --- asyncband/src/barrier/mod.rs | 6 +-- asyncband/src/broadcast/mpmc/bounded/mod.rs | 4 +- asyncband/src/broadcast/mpmc/common.rs | 2 +- asyncband/src/broadcast/mpmc/unbounded/mod.rs | 4 +- asyncband/src/completion/mod.rs | 4 +- asyncband/src/internal/arena.rs | 41 +++++++++---------- asyncband/src/internal/countdown.rs | 2 +- asyncband/src/internal/mod.rs | 6 +++ asyncband/src/internal/waker_batch.rs | 22 ++++++++-- asyncband/src/internal/wakerset.rs | 28 +++++++------ asyncband/src/phaser/mod.rs | 31 +++++++------- asyncband/src/waitgroup/mod.rs | 2 +- asyncband/src/watch/mod.rs | 10 ++--- 13 files changed, 89 insertions(+), 73 deletions(-) diff --git a/asyncband/src/barrier/mod.rs b/asyncband/src/barrier/mod.rs index 20e5d364..f864beda 100644 --- a/asyncband/src/barrier/mod.rs +++ b/asyncband/src/barrier/mod.rs @@ -55,7 +55,6 @@ use std::task::Poll; use std::task::Waker; use crate::internal::mutex::Mutex; -use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -189,10 +188,7 @@ impl Barrier { if state.arrived == self.n { state.arrived = 0; state.generation += 1; - let mut wakers = WakerBatch::new(); - for waker in state.waiters.drain() { - wakers.push(waker); - } + let mut wakers = state.waiters.take_all(); drop(state); wakers.by_ref().for_each(Waker::wake); return BarrierWaitResult(true); diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index 6ca514c0..e3590897 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -425,9 +425,9 @@ impl BoundedSender { } else { inner.log.publish_retained(into_msg(payload)); } - let wakers = inner.waiters.take_all(); + let mut wakers = inner.waiters.take_all(); drop(inner); - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); drop(discarded); Ok(()) } diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index db09605b..64285fde 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -434,7 +434,7 @@ impl Inner { pub fn disconnect(inner: &Mutex>) { let wakers = { let mut inner = inner.lock(); - inner.waiters.take_all() + inner.waiters.take_all_and_release() }; wakers.for_each(Waker::wake); } diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 6725f69c..91866ac1 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -178,12 +178,12 @@ impl UnboundedSender { // observe an empty buffer and park after this message became visible. let mut inner = self.shared.inner.lock(); let unretained = inner.log.publish(msg); - let wakers = inner.waiters.take_all(); + let mut wakers = inner.waiters.take_all(); drop(inner); // Notify all waiting receivers. An unsent message is dropped here too, once the lock is // released. - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); drop(unretained); } diff --git a/asyncband/src/completion/mod.rs b/asyncband/src/completion/mod.rs index cba4dfd0..b91fe87e 100644 --- a/asyncband/src/completion/mod.rs +++ b/asyncband/src/completion/mod.rs @@ -127,7 +127,7 @@ impl Completer { }; let wakers = { let mut waiters = shared.waiters.lock(); - let wakers = waiters.take_all(); + let wakers = waiters.take_all_and_release(); // The single completer publishes only after every waiter token has been invalidated. assert!(shared.result.set(Some(value)).is_ok()); wakers @@ -147,7 +147,7 @@ impl Drop for Completer { }; let wakers = { let mut waiters = shared.waiters.lock(); - let wakers = waiters.take_all(); + let wakers = waiters.take_all_and_release(); assert!(shared.result.set(None).is_ok()); wakers }; diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index fd9a0bd0..d4c41c6d 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -166,30 +166,29 @@ impl Arena { value } - /// Drains occupied values in slot order, retaining the allocation for reuse. + /// Collects occupied values in slot order, retaining the allocation for reuse. /// - /// All previous slot IDs become invalid and may be reused after the drain. + /// All previous slot IDs become invalid and may be reused after collection. #[inline] - pub fn drain(&mut self) -> impl Iterator + '_ { + pub fn take_all>(&mut self) -> C { self.vacant_head = None; self.len = 0; - self.slots.drain(..).filter_map(|slot| match slot { - Slot::Occupied(value) => Some(value), - Slot::Vacant { .. } => None, - }) - } - - /// Takes every occupied value and the backing allocation in slot order. - #[inline] - pub fn take_all(&mut self) -> impl Iterator + use { - self.vacant_head = None; - self.len = 0; - mem::take(&mut self.slots) - .into_iter() + self.slots + .drain(..) .filter_map(|slot| match slot { Slot::Occupied(value) => Some(value), Slot::Vacant { .. } => None, }) + .collect() + } + + /// Consumes the arena, yielding occupied values in slot order and releasing its allocation. + #[inline] + pub fn into_values(self) -> impl Iterator { + self.slots.into_iter().filter_map(|slot| match slot { + Slot::Occupied(value) => Some(value), + Slot::Vacant { .. } => None, + }) } } @@ -217,7 +216,7 @@ mod tests { } #[test] - fn drain_retains_capacity_and_restarts_slot_ids() { + fn take_all_retains_capacity_and_restarts_slot_ids() { let mut arena = Arena::with_capacity(3); let first = arena.insert(1); let second = arena.insert(2); @@ -225,7 +224,7 @@ mod tests { let capacity = arena.slots.capacity(); arena.remove(second); - assert_eq!(arena.drain().collect::>(), vec![1, 3]); + assert_eq!(arena.take_all::>(), vec![1, 3]); assert_eq!(arena.len(), 0); assert_eq!(arena.slots.capacity(), capacity); assert_eq!( @@ -235,15 +234,13 @@ mod tests { } #[test] - fn take_all_releases_the_backing_allocation() { + fn into_values_skips_vacant_slots() { let mut arena = Arena::new(); arena.insert(1); let removed = arena.insert(2); arena.insert(3); arena.remove(removed); - let values = arena.take_all(); - assert_eq!(arena.slots.capacity(), 0); - assert_eq!(values.collect::>(), vec![1, 3]); + assert_eq!(arena.into_values().collect::>(), vec![1, 3]); } } diff --git a/asyncband/src/internal/countdown.rs b/asyncband/src/internal/countdown.rs index 92ec06b1..9c5b4aeb 100644 --- a/asyncband/src/internal/countdown.rs +++ b/asyncband/src/internal/countdown.rs @@ -57,7 +57,7 @@ impl CountdownState { pub fn wake_all(&self) { let wakers = { let mut waiters = self.waiters.lock(); - waiters.take_all() + waiters.take_all_and_release() }; wakers.for_each(Waker::wake); diff --git a/asyncband/src/internal/mod.rs b/asyncband/src/internal/mod.rs index da1423e7..702f7320 100644 --- a/asyncband/src/internal/mod.rs +++ b/asyncband/src/internal/mod.rs @@ -113,10 +113,16 @@ pub(crate) mod waitlist; #[cfg(any( feature = "barrier", feature = "broadcast", + feature = "completion", feature = "event", + feature = "latch", feature = "mutex", + feature = "once", + feature = "phaser", feature = "rwlock", feature = "semaphore", + feature = "waitgroup", + feature = "watch", ))] // Only the semaphore refills a batch and asks whether it will spill, so other feature subsets // leave that method unused. diff --git a/asyncband/src/internal/waker_batch.rs b/asyncband/src/internal/waker_batch.rs index 2a41105f..d77fb64b 100644 --- a/asyncband/src/internal/waker_batch.rs +++ b/asyncband/src/internal/waker_batch.rs @@ -24,6 +24,7 @@ use std::task::Waker; /// /// Only pushed entries are initialized. Consuming all inline entries makes their storage /// available for the semaphore's next notification batch. +/// Iterate by reference to avoid moving the inline storage into iterator adapters. pub struct WakerBatch { /// The initialized entries are exactly `start..end`. inline: [MaybeUninit; Self::STACK_SIZE], @@ -41,8 +42,8 @@ pub struct WakerBatch { impl WakerBatch { /// Wakers kept on the stack before the batch spills to the heap. /// - /// This is also the most wakers the semaphore collects per lock acquisition, so a drain that - /// wakes a typical waiter set never allocates; larger sets pay one allocation for the overflow. + /// This is also the most wakers the semaphore collects per lock acquisition. Larger batches + /// allocate overflow storage. pub const STACK_SIZE: usize = 32; pub const fn new() -> Self { @@ -62,6 +63,7 @@ impl WakerBatch { self.end == Self::STACK_SIZE || !self.spilled.is_empty() } + #[inline] pub fn push(&mut self, waker: Waker) { if self.end < Self::STACK_SIZE && self.spilled.is_empty() { self.inline[self.end].write(waker); @@ -72,9 +74,21 @@ impl WakerBatch { } } +impl FromIterator for WakerBatch { + #[inline] + fn from_iter>(iter: T) -> Self { + let mut batch = Self::new(); + for waker in iter { + batch.push(waker); + } + batch + } +} + impl Iterator for WakerBatch { type Item = Waker; + #[inline] fn next(&mut self) -> Option { if self.start < self.end { let index = self.start; @@ -91,6 +105,7 @@ impl Iterator for WakerBatch { } impl Drop for WakerBatch { + #[inline] fn drop(&mut self) { let initialized = ptr::slice_from_raw_parts_mut( self.inline[self.start..self.end] @@ -160,8 +175,7 @@ mod tests { fn yields_in_push_order_across_the_spill() { let log = log(); let count = STACK_SIZE + 8; - let mut batch = WakerBatch::new(); - push_wakers(&mut batch, &log, count); + let mut batch: WakerBatch = (0..count).map(|id| waker(&log, id)).collect(); assert!(batch.will_spill()); for waker in &mut batch { diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index 658dbd28..b992a97b 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -26,12 +26,13 @@ use std::task::Waker; use crate::internal::arena::Arena; use crate::internal::arena::SlotId; +use crate::internal::waker_batch::WakerBatch; /// An exclusive handle to one waker slot in a [`WakerSet`]. /// /// This token deliberately does not implement `Clone` or `Copy`. Its owner must not pass it back -/// to the set after the registration has been detached by [`WakerSet::drain`] or -/// [`WakerSet::take_all`]. +/// to the set after the registration has been detached by [`WakerSet::take_all`] or +/// [`WakerSet::take_all_and_release`]. #[derive(Debug)] pub struct WakerToken(SlotId); @@ -56,22 +57,25 @@ impl WakerSet { } } - /// Drains registered wakers while retaining slot capacity for reuse. + /// Collects registered wakers into an owned batch, retaining slot capacity for reuse. /// - /// The caller must invalidate outstanding tokens and collect the wakers under the set's lock, - /// then wake or drop them after releasing it. + /// Collection moves each waker under the set's lock. The caller must invalidate outstanding + /// tokens and wake or drop the batch after releasing the lock. #[inline] - pub fn drain(&mut self) -> impl Iterator + '_ { - self.wakers.drain() + pub fn take_all(&mut self) -> WakerBatch { + if self.wakers.is_empty() { + return WakerBatch::new(); + } + self.wakers.take_all() } - /// Takes all registered wakers together with the set's backing allocation. + /// Transfers registered wakers and the backing allocation when capacity is no longer needed. /// - /// The caller must invalidate every outstanding token and consume or drop the iterator after - /// releasing the lock that protects this set. + /// No wakers are moved individually under the lock. The caller must invalidate outstanding + /// tokens and consume or drop the iterator after unlocking, which also frees the allocation. #[inline] - pub fn take_all(&mut self) -> impl Iterator + 'static { - self.wakers.take_all() + pub fn take_all_and_release(&mut self) -> impl Iterator + 'static { + mem::replace(&mut self.wakers, Arena::new()).into_values() } /// Registers or updates a waker. diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 1be948a8..51c0c54a 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -131,6 +131,7 @@ use std::task::Poll; use std::task::Waker; use crate::internal::mutex::Mutex; +use crate::internal::waker_batch::WakerBatch; use crate::internal::wakerset::WakerSet; use crate::internal::wakerset::WakerToken; @@ -172,15 +173,13 @@ struct State { impl State { /// Advances a completed phase and returns its waiters for waking outside the lock. - fn advance_if_ready(&mut self) -> impl Iterator + 'static { - let wakers = if self.closed || self.unarrived != 0 { - None - } else { - self.phase = self.phase.wrapping_add(1); - self.unarrived = self.registered; - Some(self.waiters.take_all()) - }; - wakers.into_iter().flatten() + fn advance_if_ready(&mut self) -> WakerBatch { + if self.closed || self.unarrived != 0 { + return WakerBatch::new(); + } + self.phase = self.phase.wrapping_add(1); + self.unarrived = self.registered; + self.waiters.take_all() } fn completion(&self, observed: u64) -> Poll> { @@ -248,7 +247,7 @@ impl Phaser { return; } state.closed = true; - state.waiters.take_all() + state.waiters.take_all_and_release() }; wakers.for_each(Waker::wake); } @@ -391,9 +390,9 @@ impl Drop for PhaserParticipants { state.registered -= self.remaining; state.unarrived -= self.remaining; self.remaining = 0; - let wakers = state.advance_if_ready(); + let mut wakers = state.advance_if_ready(); drop(state); - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); } } @@ -441,9 +440,9 @@ impl PhaserParticipant { state.unarrived -= 1; } self.pending = Some(phase); - let wakers = state.advance_if_ready(); + let mut wakers = state.advance_if_ready(); drop(state); - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); Ok(phase) } @@ -489,9 +488,9 @@ impl PhaserParticipant { } else { Ok(state.phase) }; - let wakers = state.advance_if_ready(); + let mut wakers = state.advance_if_ready(); drop(state); - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); result } } diff --git a/asyncband/src/waitgroup/mod.rs b/asyncband/src/waitgroup/mod.rs index 79238f30..72042185 100644 --- a/asyncband/src/waitgroup/mod.rs +++ b/asyncband/src/waitgroup/mod.rs @@ -104,7 +104,7 @@ impl State { let wakers = { let mut waiters = self.waiters.lock(); - waiters.take_all() + waiters.take_all_and_release() }; wakers.for_each(Waker::wake); } diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index 989e0ea7..f9e5a338 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -153,7 +153,7 @@ impl Drop for Sender { if state.senders != 0 { return; } - state.waiters.take_all() + state.waiters.take_all_and_release() }; wakers.for_each(Waker::wake); } @@ -178,11 +178,11 @@ impl Sender { .expect("watch channel version counter overflowed"); let replaced = mem::replace(&mut state.value, value); state.version = version; - let wakers = state.waiters.take_all(); + let mut wakers = state.waiters.take_all(); drop(state); // Waker callbacks and the replaced value's destructor may reenter this channel. - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); drop(replaced); Ok(()) } @@ -203,9 +203,9 @@ impl Sender { .expect("watch channel version counter overflowed"); let replaced = mem::replace(&mut state.value, value); state.version = version; - let wakers = state.waiters.take_all(); + let mut wakers = state.waiters.take_all(); drop(state); - wakers.for_each(Waker::wake); + wakers.by_ref().for_each(Waker::wake); replaced } From 32ce1f03e725583c9f5f3d153bd81e2ff06a51f6 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 11:02:19 +0800 Subject: [PATCH 10/12] fixup Signed-off-by: tison --- asyncband/src/broadcast/mpmc/common.rs | 8 ++++---- asyncband/src/completion/mod.rs | 5 +++-- asyncband/src/internal/arena.rs | 12 +++++++++--- asyncband/src/internal/countdown.rs | 4 ++-- asyncband/src/internal/wakerset.rs | 8 ++++---- asyncband/src/phaser/mod.rs | 3 ++- asyncband/src/waitgroup/mod.rs | 3 ++- asyncband/src/watch/mod.rs | 2 +- 8 files changed, 27 insertions(+), 18 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index 64285fde..5b2cafc5 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -320,7 +320,7 @@ impl Backlog { /// subscription for the slowest cursor, so advancing the head costs the messages released /// instead of the receivers subscribed. /// - /// A receive releases exactly one message: its cursor is counted at the next version before + /// A `receive` releases exactly one message: its cursor is counted at the next version before /// it leaves `head`, so the zero-count prefix ends there. Only removing a lagging subscription /// can release more. /// @@ -434,12 +434,12 @@ impl Inner { pub fn disconnect(inner: &Mutex>) { let wakers = { let mut inner = inner.lock(); - inner.waiters.take_all_and_release() + mem::take(&mut inner.waiters).into_iter() }; wakers.for_each(Waker::wake); } -/// Releases a cancelled receive's waker registration, dropping the waker unlocked. +/// Releases a cancelled `receive`'s waker registration, dropping the waker unlocked. pub fn unregister( inner: &Mutex>, senders: &AtomicUsize, @@ -582,7 +582,7 @@ mod tests { assert_eq!(log.remove_receiver(b).len(), 3); log.assert_cursor_accounting(); - // A receive that catches up to the tail releases exactly one message and moves the cursor + // A `receive` that catches up to the tail releases exactly one message and moves the cursor // back to `at_tail`. assert!(log.publish(Arc::new(4)).is_none()); let (msg, reclaimed) = log.receive(a).unwrap(); diff --git a/asyncband/src/completion/mod.rs b/asyncband/src/completion/mod.rs index b91fe87e..9b8e0f62 100644 --- a/asyncband/src/completion/mod.rs +++ b/asyncband/src/completion/mod.rs @@ -50,6 +50,7 @@ use std::fmt; use std::future::Future; +use std::mem; use std::pin::Pin; use std::sync::Arc; use std::sync::OnceLock; @@ -127,7 +128,7 @@ impl Completer { }; let wakers = { let mut waiters = shared.waiters.lock(); - let wakers = waiters.take_all_and_release(); + let wakers = mem::take(&mut *waiters).into_iter(); // The single completer publishes only after every waiter token has been invalidated. assert!(shared.result.set(Some(value)).is_ok()); wakers @@ -147,7 +148,7 @@ impl Drop for Completer { }; let wakers = { let mut waiters = shared.waiters.lock(); - let wakers = waiters.take_all_and_release(); + let wakers = mem::take(&mut *waiters).into_iter(); assert!(shared.result.set(None).is_ok()); wakers }; diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index d4c41c6d..d0c06b28 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -60,6 +60,12 @@ pub struct Arena { len: usize, } +impl Default for Arena { + fn default() -> Self { + Self::new() + } +} + #[derive(Debug)] enum Slot { Occupied(T), @@ -184,7 +190,7 @@ impl Arena { /// Consumes the arena, yielding occupied values in slot order and releasing its allocation. #[inline] - pub fn into_values(self) -> impl Iterator { + pub fn into_iter(self) -> impl Iterator { self.slots.into_iter().filter_map(|slot| match slot { Slot::Occupied(value) => Some(value), Slot::Vacant { .. } => None, @@ -234,13 +240,13 @@ mod tests { } #[test] - fn into_values_skips_vacant_slots() { + fn into_iter_skips_vacant_slots() { let mut arena = Arena::new(); arena.insert(1); let removed = arena.insert(2); arena.insert(3); arena.remove(removed); - assert_eq!(arena.into_values().collect::>(), vec![1, 3]); + assert_eq!(arena.into_iter().collect::>(), vec![1, 3]); } } diff --git a/asyncband/src/internal/countdown.rs b/asyncband/src/internal/countdown.rs index 9c5b4aeb..fd4679c8 100644 --- a/asyncband/src/internal/countdown.rs +++ b/asyncband/src/internal/countdown.rs @@ -15,6 +15,7 @@ // specific language governing permissions and limitations // under the License. +use std::mem; use std::sync::atomic::AtomicU32; use std::sync::atomic::Ordering; use std::task::Context; @@ -57,9 +58,8 @@ impl CountdownState { pub fn wake_all(&self) { let wakers = { let mut waiters = self.waiters.lock(); - waiters.take_all_and_release() + mem::take(&mut *waiters).into_iter() }; - wakers.for_each(Waker::wake); } diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index b992a97b..0c6c28c3 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -32,12 +32,12 @@ use crate::internal::waker_batch::WakerBatch; /// /// This token deliberately does not implement `Clone` or `Copy`. Its owner must not pass it back /// to the set after the registration has been detached by [`WakerSet::take_all`] or -/// [`WakerSet::take_all_and_release`]. +/// [`WakerSet::into_iter`]. #[derive(Debug)] pub struct WakerToken(SlotId); /// Cancellable waker storage without an implicit lifecycle or generation. -#[derive(Debug)] +#[derive(Default, Debug)] pub struct WakerSet { wakers: Arena, } @@ -74,8 +74,8 @@ impl WakerSet { /// No wakers are moved individually under the lock. The caller must invalidate outstanding /// tokens and consume or drop the iterator after unlocking, which also frees the allocation. #[inline] - pub fn take_all_and_release(&mut self) -> impl Iterator + 'static { - mem::replace(&mut self.wakers, Arena::new()).into_values() + pub fn into_iter(self) -> impl Iterator + 'static { + self.wakers.into_iter() } /// Registers or updates a waker. diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 51c0c54a..6d71207f 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -124,6 +124,7 @@ use std::fmt; use std::future::Future; use std::iter::FusedIterator; +use std::mem; use std::pin::Pin; use std::sync::Arc; use std::task::Context; @@ -247,7 +248,7 @@ impl Phaser { return; } state.closed = true; - state.waiters.take_all_and_release() + mem::take(&mut state.waiters).into_iter() }; wakers.for_each(Waker::wake); } diff --git a/asyncband/src/waitgroup/mod.rs b/asyncband/src/waitgroup/mod.rs index 72042185..5af07294 100644 --- a/asyncband/src/waitgroup/mod.rs +++ b/asyncband/src/waitgroup/mod.rs @@ -59,6 +59,7 @@ use std::fmt; use std::future::Future; use std::future::IntoFuture; +use std::mem; use std::pin::Pin; use std::sync::Arc; use std::sync::atomic::AtomicUsize; @@ -104,7 +105,7 @@ impl State { let wakers = { let mut waiters = self.waiters.lock(); - waiters.take_all_and_release() + mem::take(&mut *waiters).into_iter() }; wakers.for_each(Waker::wake); } diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index f9e5a338..344fdd7b 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -153,7 +153,7 @@ impl Drop for Sender { if state.senders != 0 { return; } - state.waiters.take_all_and_release() + mem::take(&mut state.waiters).into_iter() }; wakers.for_each(Waker::wake); } From b4bf7aff753442b52a48cff058a610ccf7193116 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 11:09:38 +0800 Subject: [PATCH 11/12] refactor: clarify waker ownership and notification paths Signed-off-by: tison --- asyncband/src/barrier/mod.rs | 4 +- asyncband/src/broadcast/mpmc/bounded/mod.rs | 15 +++---- asyncband/src/broadcast/mpmc/common.rs | 20 +++------ asyncband/src/broadcast/mpmc/unbounded/mod.rs | 5 +-- asyncband/src/completion/mod.rs | 29 +++++------- asyncband/src/internal/arena.rs | 4 +- asyncband/src/internal/countdown.rs | 11 ++--- asyncband/src/internal/mod.rs | 3 -- asyncband/src/internal/semaphore.rs | 2 +- asyncband/src/internal/waker_batch.rs | 44 +++++++++---------- asyncband/src/internal/wakerset.rs | 40 +++++++---------- asyncband/src/phaser/mod.rs | 4 +- asyncband/src/waitgroup/mod.rs | 7 +-- asyncband/src/watch/mod.rs | 5 +-- 14 files changed, 80 insertions(+), 113 deletions(-) diff --git a/asyncband/src/barrier/mod.rs b/asyncband/src/barrier/mod.rs index f864beda..e52db7cd 100644 --- a/asyncband/src/barrier/mod.rs +++ b/asyncband/src/barrier/mod.rs @@ -184,7 +184,7 @@ impl Barrier { state.arrived += 1; // The final arrival completes this generation. Advance the generation while holding - // the state lock, then wake the drained followers after releasing it. + // the state lock, then wake the detached followers after releasing it. if state.arrived == self.n { state.arrived = 0; state.generation += 1; @@ -237,7 +237,7 @@ impl Future for BarrierWait<'_> { let mut state = barrier.state.lock(); if *generation < state.generation { - // Completion advances the generation and drains its old waiters under this same lock, + // Completion advances the generation and detaches its old waiters under this same lock, // so no registration represented by this token remains in the waker set. *token = None; return Poll::Ready(()); diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index e3590897..48ec26c6 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -403,24 +403,19 @@ impl BoundedSender { self.publish(msg, |msg| msg) } - /// The publishing step both send paths share. + /// Publishes a message for both send paths. /// - /// `into_msg` is called only once this decides the message will actually be retained, which is - /// what lets `try_send` defer its allocation past the capacity check while `try_publish` hands - /// over an `Arc` it allocated with the channel unlocked. - /// - /// Publishing and draining the wait set share one critical section, so a receiver can never - /// observe an empty buffer and park after this message became visible. + /// Calls `into_msg` only when retaining the payload, so `try_send` can defer its allocation + /// until capacity is available while `try_publish` passes through its existing `Arc`. fn publish

(&self, payload: P, into_msg: impl FnOnce(P) -> Arc) -> Result<(), P> { let mut discarded = None; let mut inner = self.shared.inner.lock(); if !inner.log.has_receivers() { - // Nothing can read this message. The payload leaves the critical section with us - // and is dropped below, so `T::drop` never runs under the lock. + // Drop the discarded payload after unlocking: its destructor may reenter the channel. inner.log.publish_discarded(); discarded = Some(payload); } else if inner.log.retained() == self.shared.capacity { - // Nothing was published, so there is no wait set to drain. + // Leave waiters registered because no message was published. return Err(payload); } else { inner.log.publish_retained(into_msg(payload)); diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index 5b2cafc5..1f4d1ca8 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -406,10 +406,8 @@ impl Backlog { /// Buffer, receiver cursors, and parked receivers, all under one lock. /// -/// The wait set lives beside the backlog so that publishing a message and draining the waiters -/// happen in one critical section. That is what makes the park path race-free: a receiver that -/// finds no message and then registers still holds this lock, so a concurrent send cannot slip -/// between the two steps and skip the wake-up. +/// Checking the backlog and registering a waker share this lock with publishing and detaching +/// registrations, so a send cannot slip between an empty check and registration. pub struct Inner { pub log: Backlog, pub waiters: WakerSet, @@ -432,14 +430,11 @@ impl Inner { /// /// Both families call this from the last sender's `Drop`. pub fn disconnect(inner: &Mutex>) { - let wakers = { - let mut inner = inner.lock(); - mem::take(&mut inner.waiters).into_iter() - }; - wakers.for_each(Waker::wake); + let wakers = mem::take(&mut inner.lock().waiters); + wakers.into_iter().for_each(Waker::wake); } -/// Releases a cancelled `receive`'s waker registration, dropping the waker unlocked. +/// Removes the waker registration for a cancelled `receive`, dropping the waker unlocked. pub fn unregister( inner: &Mutex>, senders: &AtomicUsize, @@ -480,9 +475,8 @@ pub fn try_receive( /// The one poll step behind `recv` on both channels. /// -/// Checking the backlog and registering a waker under the same lock prevents a publication from -/// landing between those steps. Publication and disconnection detach all registrations, so their -/// ready paths clear the token without unregistering it. +/// Publication and disconnection detach all registrations, so their ready paths clear the token +/// without unregistering it. pub fn poll_receive( inner: &Mutex>, senders: &AtomicUsize, diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 91866ac1..17f9d111 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -174,15 +174,12 @@ impl UnboundedSender { pub fn send(&self, msg: T) { let msg = Arc::new(msg); - // Publishing and draining the wait set share one critical section, so a receiver can never - // observe an empty buffer and park after this message became visible. let mut inner = self.shared.inner.lock(); let unretained = inner.log.publish(msg); let mut wakers = inner.waiters.take_all(); drop(inner); - // Notify all waiting receivers. An unsent message is dropped here too, once the lock is - // released. + // Wake callbacks and payload destruction may reenter the channel. wakers.by_ref().for_each(Waker::wake); drop(unretained); } diff --git a/asyncband/src/completion/mod.rs b/asyncband/src/completion/mod.rs index 9b8e0f62..fb5f0d79 100644 --- a/asyncband/src/completion/mod.rs +++ b/asyncband/src/completion/mod.rs @@ -126,17 +126,14 @@ impl Completer { let Some(shared) = self.shared.upgrade() else { return Err(value); }; - let wakers = { - let mut waiters = shared.waiters.lock(); - let wakers = mem::take(&mut *waiters).into_iter(); - // The single completer publishes only after every waiter token has been invalidated. - assert!(shared.result.set(Some(value)).is_ok()); - wakers - }; - // `complete` consumes the only completer. Disarm its destructor before invoking arbitrary - // wake callbacks; the completed state no longer needs abandonment handling. + let mut waiters = shared.waiters.lock(); + // Detach registrations before publishing completion to lock-free observers. + let wakers = mem::take(&mut *waiters); + assert!(shared.result.set(Some(value)).is_ok()); + drop(waiters); + // Disarm abandonment handling before invoking wake callbacks. self.shared = Weak::new(); - wakers.for_each(Waker::wake); + wakers.into_iter().for_each(Waker::wake); Ok(()) } } @@ -146,13 +143,11 @@ impl Drop for Completer { let Some(shared) = self.shared.upgrade() else { return; }; - let wakers = { - let mut waiters = shared.waiters.lock(); - let wakers = mem::take(&mut *waiters).into_iter(); - assert!(shared.result.set(None).is_ok()); - wakers - }; - wakers.for_each(Waker::wake); + let mut waiters = shared.waiters.lock(); + let wakers = mem::take(&mut *waiters); + assert!(shared.result.set(None).is_ok()); + drop(waiters); + wakers.into_iter().for_each(Waker::wake); } } diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index d0c06b28..0cd8647f 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -188,7 +188,9 @@ impl Arena { .collect() } - /// Consumes the arena, yielding occupied values in slot order and releasing its allocation. + /// Consumes the arena, yielding occupied values in slot order. + /// + /// The returned iterator owns the backing allocation and releases it when dropped. #[inline] pub fn into_iter(self) -> impl Iterator { self.slots.into_iter().filter_map(|slot| match slot { diff --git a/asyncband/src/internal/countdown.rs b/asyncband/src/internal/countdown.rs index fd4679c8..0a80f9c3 100644 --- a/asyncband/src/internal/countdown.rs +++ b/asyncband/src/internal/countdown.rs @@ -54,13 +54,10 @@ impl CountdownState { .map(|_| ()) } - /// Drains the waiter set under its lock, then wakes every waiter after releasing the lock. + /// Detaches the waiter set under its lock, then wakes every waiter after releasing the lock. pub fn wake_all(&self) { - let wakers = { - let mut waiters = self.waiters.lock(); - mem::take(&mut *waiters).into_iter() - }; - wakers.for_each(Waker::wake); + let wakers = mem::take(&mut *self.waiters.lock()); + wakers.into_iter().for_each(Waker::wake); } /// Polls for zero, registering the current waker if the countdown is still active. @@ -74,7 +71,7 @@ impl CountdownState { let mut waiters = self.waiters.lock(); if self.state() == 0 { - // A concurrent zero transition will drain after this lock is released. + // A concurrent zero transition detaches the registrations under this same lock. *token = None; return Poll::Ready(()); } diff --git a/asyncband/src/internal/mod.rs b/asyncband/src/internal/mod.rs index 702f7320..9e38c5af 100644 --- a/asyncband/src/internal/mod.rs +++ b/asyncband/src/internal/mod.rs @@ -124,9 +124,6 @@ pub(crate) mod waitlist; feature = "waitgroup", feature = "watch", ))] -// Only the semaphore refills a batch and asks whether it will spill, so other feature subsets -// leave that method unused. -#[allow(dead_code)] pub(crate) mod waker_batch; #[cfg(any( diff --git a/asyncband/src/internal/semaphore.rs b/asyncband/src/internal/semaphore.rs index 4ce4ad71..7d45152d 100644 --- a/asyncband/src/internal/semaphore.rs +++ b/asyncband/src/internal/semaphore.rs @@ -414,7 +414,7 @@ mod tests { #[test] fn release_distributes_permits_across_multiple_batches() { - const WAITER_COUNT: usize = WakerBatch::STACK_SIZE * 2 + 1; + const WAITER_COUNT: usize = WakerBatch::INLINE_CAPACITY * 2 + 1; let semaphore = Semaphore::new(0); let counters = (0..WAITER_COUNT) diff --git a/asyncband/src/internal/waker_batch.rs b/asyncband/src/internal/waker_batch.rs index d77fb64b..77526917 100644 --- a/asyncband/src/internal/waker_batch.rs +++ b/asyncband/src/internal/waker_batch.rs @@ -20,14 +20,14 @@ use std::mem::MaybeUninit; use std::ptr; use std::task::Waker; -/// An owning FIFO of wakers that stores the first [`Self::STACK_SIZE`] entries without allocating. +/// An owning FIFO of wakers that stores the first [`Self::INLINE_CAPACITY`] entries without +/// allocating. /// -/// Only pushed entries are initialized. Consuming all inline entries makes their storage -/// available for the semaphore's next notification batch. -/// Iterate by reference to avoid moving the inline storage into iterator adapters. +/// Only pushed entries are initialized. Iterate by reference to avoid moving the inline storage +/// into iterator adapters and to reuse the batch, as the semaphore does between notifications. pub struct WakerBatch { /// The initialized entries are exactly `start..end`. - inline: [MaybeUninit; Self::STACK_SIZE], + inline: [MaybeUninit; Self::INLINE_CAPACITY], /// The next inline entry to yield. start: usize, /// The next inline slot to push into. @@ -40,15 +40,14 @@ pub struct WakerBatch { } impl WakerBatch { - /// Wakers kept on the stack before the batch spills to the heap. + /// Number of wakers stored inline before using overflow storage. /// - /// This is also the most wakers the semaphore collects per lock acquisition. Larger batches - /// allocate overflow storage. - pub const STACK_SIZE: usize = 32; + /// The semaphore's permit-release loop uses this as its batch limit before unlocking. + pub const INLINE_CAPACITY: usize = 32; pub const fn new() -> Self { Self { - inline: [const { MaybeUninit::uninit() }; Self::STACK_SIZE], + inline: [const { MaybeUninit::uninit() }; Self::INLINE_CAPACITY], start: 0, end: 0, spilled: VecDeque::new(), @@ -59,17 +58,18 @@ impl WakerBatch { /// /// The semaphore stops filling a batch here so it can release its lock and wake what it has /// before collecting more. + #[inline] pub fn will_spill(&self) -> bool { - self.end == Self::STACK_SIZE || !self.spilled.is_empty() + self.end == Self::INLINE_CAPACITY || !self.spilled.is_empty() } #[inline] pub fn push(&mut self, waker: Waker) { - if self.end < Self::STACK_SIZE && self.spilled.is_empty() { + if self.will_spill() { + self.spilled.push_back(waker); + } else { self.inline[self.end].write(waker); self.end += 1; - } else { - self.spilled.push_back(waker); } } } @@ -127,7 +127,7 @@ mod tests { use super::WakerBatch; - const STACK_SIZE: usize = WakerBatch::STACK_SIZE; + const INLINE_CAPACITY: usize = WakerBatch::INLINE_CAPACITY; /// The ids of the wakers woken so far, in order. /// @@ -174,7 +174,7 @@ mod tests { #[test] fn yields_in_push_order_across_the_spill() { let log = log(); - let count = STACK_SIZE + 8; + let count = INLINE_CAPACITY + 8; let mut batch: WakerBatch = (0..count).map(|id| waker(&log, id)).collect(); assert!(batch.will_spill()); @@ -190,8 +190,8 @@ mod tests { #[test] fn drops_unconsumed_entries_exactly_once() { let log = log(); - let count = STACK_SIZE + 8; - for consumed in [0, 5, STACK_SIZE, STACK_SIZE + 3, count] { + let count = INLINE_CAPACITY + 8; + for consumed in [0, 5, INLINE_CAPACITY, INLINE_CAPACITY + 3, count] { let mut batch = WakerBatch::new(); push_wakers(&mut batch, &log, count); for _ in 0..consumed { @@ -209,9 +209,9 @@ mod tests { let log = log(); let mut batch = WakerBatch::new(); for _ in 0..3 { - push_wakers(&mut batch, &log, STACK_SIZE); + push_wakers(&mut batch, &log, INLINE_CAPACITY); assert!(batch.will_spill()); - assert_eq!(batch.by_ref().count(), STACK_SIZE); + assert_eq!(batch.by_ref().count(), INLINE_CAPACITY); assert!(!batch.will_spill()); assert_eq!(alive(&log), 0); } @@ -221,7 +221,7 @@ mod tests { fn keeps_push_order_while_spilled() { let log = log(); let mut batch = WakerBatch::new(); - push_wakers(&mut batch, &log, STACK_SIZE + 1); + push_wakers(&mut batch, &log, INLINE_CAPACITY + 1); // Free inline room; the spilled entry must still come out before anything pushed now. for _ in 0..4 { batch.next().unwrap().wake(); @@ -233,7 +233,7 @@ mod tests { waker.wake(); } - let mut expected = (0..STACK_SIZE + 1).collect::>(); + let mut expected = (0..INLINE_CAPACITY + 1).collect::>(); expected.push(999); assert_eq!(woken(&log), expected); assert_eq!(alive(&log), 0); diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index 0c6c28c3..1190fa87 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -15,11 +15,11 @@ // specific language governing permissions and limitations // under the License. -//! Cancellable storage for task wakers whose lifecycle is owned by the caller. +//! Cancellable task wakers protected by the owning primitive's state lock. //! -//! A `WakerSet` is protected by the state lock of its owning primitive. The owner must clear a -//! token instead of unregistering it after an operation that detached the set. This lets each -//! primitive use its existing terminal state or generation to recognize stale registrations. +//! The primitive uses its generation or terminal state to recognize detached registrations; +//! their tokens must be cleared rather than passed back to the set. Returned wakers must be woken +//! or dropped after releasing the lock. use std::mem; use std::task::Waker; @@ -30,9 +30,7 @@ use crate::internal::waker_batch::WakerBatch; /// An exclusive handle to one waker slot in a [`WakerSet`]. /// -/// This token deliberately does not implement `Clone` or `Copy`. Its owner must not pass it back -/// to the set after the registration has been detached by [`WakerSet::take_all`] or -/// [`WakerSet::into_iter`]. +/// Removing the registration, calling [`WakerSet::take_all`], or replacing the set invalidates it. #[derive(Debug)] pub struct WakerToken(SlotId); @@ -59,8 +57,8 @@ impl WakerSet { /// Collects registered wakers into an owned batch, retaining slot capacity for reuse. /// - /// Collection moves each waker under the set's lock. The caller must invalidate outstanding - /// tokens and wake or drop the batch after releasing the lock. + /// Moves each waker into the batch; up to [`WakerBatch::INLINE_CAPACITY`] fit without + /// allocating. #[inline] pub fn take_all(&mut self) -> WakerBatch { if self.wakers.is_empty() { @@ -69,27 +67,26 @@ impl WakerSet { self.wakers.take_all() } - /// Transfers registered wakers and the backing allocation when capacity is no longer needed. + /// Consumes the set, transferring its backing allocation to an iterator. /// - /// No wakers are moved individually under the lock. The caller must invalidate outstanding - /// tokens and consume or drop the iterator after unlocking, which also frees the allocation. + /// This avoids collecting a separate batch when capacity is no longer needed. The iterator + /// releases the allocation when dropped. #[inline] - pub fn into_iter(self) -> impl Iterator + 'static { + pub fn into_iter(self) -> impl Iterator { self.wakers.into_iter() } /// Registers or updates a waker. /// - /// If an existing waker is replaced, it is returned so the caller can drop it after releasing - /// the lock that protects this set. + /// Returns the previous waker only if it was replaced. #[inline] #[must_use = "drop the returned waker after releasing the waker set's state lock"] pub fn register(&mut self, token: &mut Option, waker: &Waker) -> Option { - if let Some(current) = token.as_ref().map(|token| { - self.wakers + if let Some(token) = token { + let current = self + .wakers .get_mut(token.0) - .expect("waker token must refer to an occupied slot") - }) { + .expect("waker token must refer to an occupied slot"); if current.will_wake(waker) { return None; } @@ -100,10 +97,7 @@ impl WakerSet { None } - /// Removes the waker identified by `token`. - /// - /// The owner must clear stale tokens without calling this method after detaching the set. The - /// returned waker must be dropped after releasing the lock that protects this set. + /// Removes and returns the waker identified by `token`, clearing the token. #[inline] #[must_use = "drop the returned waker after releasing the waker set's state lock"] pub fn unregister(&mut self, token: &mut Option) -> Option { diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 6d71207f..33964351 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -248,9 +248,9 @@ impl Phaser { return; } state.closed = true; - mem::take(&mut state.waiters).into_iter() + mem::take(&mut state.waiters) }; - wakers.for_each(Waker::wake); + wakers.into_iter().for_each(Waker::wake); } /// Returns an instantaneous count of registered participants, including those already arrived. diff --git a/asyncband/src/waitgroup/mod.rs b/asyncband/src/waitgroup/mod.rs index 5af07294..5be2927f 100644 --- a/asyncband/src/waitgroup/mod.rs +++ b/asyncband/src/waitgroup/mod.rs @@ -103,11 +103,8 @@ impl State { return; } - let wakers = { - let mut waiters = self.waiters.lock(); - mem::take(&mut *waiters).into_iter() - }; - wakers.for_each(Waker::wake); + let wakers = mem::take(&mut *self.waiters.lock()); + wakers.into_iter().for_each(Waker::wake); } fn poll_wait(&self, token: &mut Option, cx: &mut Context<'_>) -> Poll<()> { diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index 344fdd7b..001c7523 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -146,16 +146,15 @@ impl fmt::Debug for Sender { impl Drop for Sender { fn drop(&mut self) { - // Only the final sender detaches the parked receivers; their wake callbacks run unlocked. let wakers = { let mut state = self.shared.state.lock(); state.senders -= 1; if state.senders != 0 { return; } - mem::take(&mut state.waiters).into_iter() + mem::take(&mut state.waiters) }; - wakers.for_each(Waker::wake); + wakers.into_iter().for_each(Waker::wake); } } From 05632efca5f0259dce0ad528c1b522af8635128b Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 11:16:22 +0800 Subject: [PATCH 12/12] refactor: encapsulate waker set notification Replace WakerSet::into_iter with wake_all and restore WakerBatch::STACK_SIZE across code, tests, and documentation. Signed-off-by: tison --- asyncband/src/broadcast/mpmc/common.rs | 3 +-- asyncband/src/completion/mod.rs | 5 ++--- asyncband/src/internal/countdown.rs | 3 +-- asyncband/src/internal/mod.rs | 3 +-- asyncband/src/internal/semaphore.rs | 2 +- asyncband/src/internal/waker_batch.rs | 27 +++++++++++++------------- asyncband/src/internal/wakerset.rs | 12 +++++------- asyncband/src/phaser/mod.rs | 2 +- asyncband/src/waitgroup/mod.rs | 3 +-- asyncband/src/watch/mod.rs | 2 +- 10 files changed, 27 insertions(+), 35 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index 1f4d1ca8..73d96e0f 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -31,7 +31,6 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; -use std::task::Waker; use super::error::RecvError; use super::error::TryRecvError; @@ -431,7 +430,7 @@ impl Inner { /// Both families call this from the last sender's `Drop`. pub fn disconnect(inner: &Mutex>) { let wakers = mem::take(&mut inner.lock().waiters); - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } /// Removes the waker registration for a cancelled `receive`, dropping the waker unlocked. diff --git a/asyncband/src/completion/mod.rs b/asyncband/src/completion/mod.rs index fb5f0d79..95ce3c15 100644 --- a/asyncband/src/completion/mod.rs +++ b/asyncband/src/completion/mod.rs @@ -57,7 +57,6 @@ use std::sync::OnceLock; use std::sync::Weak; use std::task::Context; use std::task::Poll; -use std::task::Waker; use crate::internal::mutex::Mutex; use crate::internal::wakerset::WakerSet; @@ -133,7 +132,7 @@ impl Completer { drop(waiters); // Disarm abandonment handling before invoking wake callbacks. self.shared = Weak::new(); - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); Ok(()) } } @@ -147,7 +146,7 @@ impl Drop for Completer { let wakers = mem::take(&mut *waiters); assert!(shared.result.set(None).is_ok()); drop(waiters); - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } } diff --git a/asyncband/src/internal/countdown.rs b/asyncband/src/internal/countdown.rs index 0a80f9c3..eb37cb93 100644 --- a/asyncband/src/internal/countdown.rs +++ b/asyncband/src/internal/countdown.rs @@ -20,7 +20,6 @@ use std::sync::atomic::AtomicU32; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; -use std::task::Waker; use crate::internal::mutex::Mutex; use crate::internal::wakerset::WakerSet; @@ -57,7 +56,7 @@ impl CountdownState { /// Detaches the waiter set under its lock, then wakes every waiter after releasing the lock. pub fn wake_all(&self) { let wakers = mem::take(&mut *self.waiters.lock()); - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } /// Polls for zero, registering the current waker if the countdown is still active. diff --git a/asyncband/src/internal/mod.rs b/asyncband/src/internal/mod.rs index 9e38c5af..f9dceb04 100644 --- a/asyncband/src/internal/mod.rs +++ b/asyncband/src/internal/mod.rs @@ -136,7 +136,6 @@ pub(crate) mod waker_batch; feature = "waitgroup", feature = "watch", ))] -// Reusable waker sets and terminal primitives use different lifecycle policies, so some feature -// subsets leave one constructor or detach operation unused. +// Some feature subsets use only one constructor or one of the two notification paths. #[allow(dead_code)] pub(crate) mod wakerset; diff --git a/asyncband/src/internal/semaphore.rs b/asyncband/src/internal/semaphore.rs index 7d45152d..4ce4ad71 100644 --- a/asyncband/src/internal/semaphore.rs +++ b/asyncband/src/internal/semaphore.rs @@ -414,7 +414,7 @@ mod tests { #[test] fn release_distributes_permits_across_multiple_batches() { - const WAITER_COUNT: usize = WakerBatch::INLINE_CAPACITY * 2 + 1; + const WAITER_COUNT: usize = WakerBatch::STACK_SIZE * 2 + 1; let semaphore = Semaphore::new(0); let counters = (0..WAITER_COUNT) diff --git a/asyncband/src/internal/waker_batch.rs b/asyncband/src/internal/waker_batch.rs index 77526917..f1ec29e0 100644 --- a/asyncband/src/internal/waker_batch.rs +++ b/asyncband/src/internal/waker_batch.rs @@ -20,14 +20,13 @@ use std::mem::MaybeUninit; use std::ptr; use std::task::Waker; -/// An owning FIFO of wakers that stores the first [`Self::INLINE_CAPACITY`] entries without -/// allocating. +/// An owning FIFO of wakers that stores the first [`Self::STACK_SIZE`] entries without allocating. /// /// Only pushed entries are initialized. Iterate by reference to avoid moving the inline storage /// into iterator adapters and to reuse the batch, as the semaphore does between notifications. pub struct WakerBatch { /// The initialized entries are exactly `start..end`. - inline: [MaybeUninit; Self::INLINE_CAPACITY], + inline: [MaybeUninit; Self::STACK_SIZE], /// The next inline entry to yield. start: usize, /// The next inline slot to push into. @@ -43,11 +42,11 @@ impl WakerBatch { /// Number of wakers stored inline before using overflow storage. /// /// The semaphore's permit-release loop uses this as its batch limit before unlocking. - pub const INLINE_CAPACITY: usize = 32; + pub const STACK_SIZE: usize = 32; pub const fn new() -> Self { Self { - inline: [const { MaybeUninit::uninit() }; Self::INLINE_CAPACITY], + inline: [const { MaybeUninit::uninit() }; Self::STACK_SIZE], start: 0, end: 0, spilled: VecDeque::new(), @@ -60,7 +59,7 @@ impl WakerBatch { /// before collecting more. #[inline] pub fn will_spill(&self) -> bool { - self.end == Self::INLINE_CAPACITY || !self.spilled.is_empty() + self.end == Self::STACK_SIZE || !self.spilled.is_empty() } #[inline] @@ -127,7 +126,7 @@ mod tests { use super::WakerBatch; - const INLINE_CAPACITY: usize = WakerBatch::INLINE_CAPACITY; + const STACK_SIZE: usize = WakerBatch::STACK_SIZE; /// The ids of the wakers woken so far, in order. /// @@ -174,7 +173,7 @@ mod tests { #[test] fn yields_in_push_order_across_the_spill() { let log = log(); - let count = INLINE_CAPACITY + 8; + let count = STACK_SIZE + 8; let mut batch: WakerBatch = (0..count).map(|id| waker(&log, id)).collect(); assert!(batch.will_spill()); @@ -190,8 +189,8 @@ mod tests { #[test] fn drops_unconsumed_entries_exactly_once() { let log = log(); - let count = INLINE_CAPACITY + 8; - for consumed in [0, 5, INLINE_CAPACITY, INLINE_CAPACITY + 3, count] { + let count = STACK_SIZE + 8; + for consumed in [0, 5, STACK_SIZE, STACK_SIZE + 3, count] { let mut batch = WakerBatch::new(); push_wakers(&mut batch, &log, count); for _ in 0..consumed { @@ -209,9 +208,9 @@ mod tests { let log = log(); let mut batch = WakerBatch::new(); for _ in 0..3 { - push_wakers(&mut batch, &log, INLINE_CAPACITY); + push_wakers(&mut batch, &log, STACK_SIZE); assert!(batch.will_spill()); - assert_eq!(batch.by_ref().count(), INLINE_CAPACITY); + assert_eq!(batch.by_ref().count(), STACK_SIZE); assert!(!batch.will_spill()); assert_eq!(alive(&log), 0); } @@ -221,7 +220,7 @@ mod tests { fn keeps_push_order_while_spilled() { let log = log(); let mut batch = WakerBatch::new(); - push_wakers(&mut batch, &log, INLINE_CAPACITY + 1); + push_wakers(&mut batch, &log, STACK_SIZE + 1); // Free inline room; the spilled entry must still come out before anything pushed now. for _ in 0..4 { batch.next().unwrap().wake(); @@ -233,7 +232,7 @@ mod tests { waker.wake(); } - let mut expected = (0..INLINE_CAPACITY + 1).collect::>(); + let mut expected = (0..STACK_SIZE + 1).collect::>(); expected.push(999); assert_eq!(woken(&log), expected); assert_eq!(alive(&log), 0); diff --git a/asyncband/src/internal/wakerset.rs b/asyncband/src/internal/wakerset.rs index 1190fa87..75649cba 100644 --- a/asyncband/src/internal/wakerset.rs +++ b/asyncband/src/internal/wakerset.rs @@ -57,8 +57,7 @@ impl WakerSet { /// Collects registered wakers into an owned batch, retaining slot capacity for reuse. /// - /// Moves each waker into the batch; up to [`WakerBatch::INLINE_CAPACITY`] fit without - /// allocating. + /// Moves each waker into the batch; up to [`WakerBatch::STACK_SIZE`] fit without allocating. #[inline] pub fn take_all(&mut self) -> WakerBatch { if self.wakers.is_empty() { @@ -67,13 +66,12 @@ impl WakerSet { self.wakers.take_all() } - /// Consumes the set, transferring its backing allocation to an iterator. + /// Consumes the set, waking every registered waker and releasing its allocation. /// - /// This avoids collecting a separate batch when capacity is no longer needed. The iterator - /// releases the allocation when dropped. + /// Call after releasing the owning primitive's state lock. #[inline] - pub fn into_iter(self) -> impl Iterator { - self.wakers.into_iter() + pub fn wake_all(self) { + self.wakers.into_iter().for_each(Waker::wake); } /// Registers or updates a waker. diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 33964351..5d97a51f 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -250,7 +250,7 @@ impl Phaser { state.closed = true; mem::take(&mut state.waiters) }; - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } /// Returns an instantaneous count of registered participants, including those already arrived. diff --git a/asyncband/src/waitgroup/mod.rs b/asyncband/src/waitgroup/mod.rs index 5be2927f..4fb7c403 100644 --- a/asyncband/src/waitgroup/mod.rs +++ b/asyncband/src/waitgroup/mod.rs @@ -66,7 +66,6 @@ use std::sync::atomic::AtomicUsize; use std::sync::atomic::Ordering; use std::task::Context; use std::task::Poll; -use std::task::Waker; use crate::internal::mutex::Mutex; use crate::internal::wakerset::WakerSet; @@ -104,7 +103,7 @@ impl State { } let wakers = mem::take(&mut *self.waiters.lock()); - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } fn poll_wait(&self, token: &mut Option, cx: &mut Context<'_>) -> Poll<()> { diff --git a/asyncband/src/watch/mod.rs b/asyncband/src/watch/mod.rs index 001c7523..c2f20721 100644 --- a/asyncband/src/watch/mod.rs +++ b/asyncband/src/watch/mod.rs @@ -154,7 +154,7 @@ impl Drop for Sender { } mem::take(&mut state.waiters) }; - wakers.into_iter().for_each(Waker::wake); + wakers.wake_all(); } }