From ecb998c9744abdb937c8d54bf66b25d82de9f811 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:15:31 +0800 Subject: [PATCH 01/24] refactor(internal): remove the dead Arena::values iterator No caller exists in any feature combination: the method was added for a broadcast consumer that has since been removed, and the module-level #[allow(dead_code)] on the arena module hid that. Iteration over occupied values is covered by take_all and into_iter. --- asyncband/src/internal/arena.rs | 7 ------- 1 file changed, 7 deletions(-) diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index 0cd8647..8c76b93 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -135,13 +135,6 @@ impl Arena { self.len == 0 } - pub fn values(&self) -> impl Iterator { - self.slots.iter().filter_map(|slot| match slot { - Slot::Occupied(value) => Some(value), - Slot::Vacant { .. } => None, - }) - } - /// Removes the value stored at `id`. /// /// # Panics From 67ed45e250da88d55b491a22f43f30fd0e746aa5 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:16:39 +0800 Subject: [PATCH 02/24] refactor(internal): remove the unused Default impl for Arena Nothing constructs an Arena through Default: every mem::take site in the crate takes a WakerSet, which has its own Default impl built on WakerSet::new. WaitList already sets the precedent of a const fn new without a Default impl, so new_without_default stays quiet. --- asyncband/src/internal/arena.rs | 6 ------ 1 file changed, 6 deletions(-) diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index 8c76b93..6975af5 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -60,12 +60,6 @@ pub struct Arena { len: usize, } -impl Default for Arena { - fn default() -> Self { - Self::new() - } -} - #[derive(Debug)] enum Slot { Occupied(T), From 04e063f484507a504d532d3d589c0043d2c60208 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:17:24 +0800 Subject: [PATCH 03/24] refactor(internal): gate Arena::len to test builds Its only callers are WaitList::occupied_len, which is itself #[cfg(test)], and the arena unit tests. Production code observes occupancy through is_empty, which stays available in all builds. --- asyncband/src/internal/arena.rs | 1 + 1 file changed, 1 insertion(+) diff --git a/asyncband/src/internal/arena.rs b/asyncband/src/internal/arena.rs index 6975af5..68a6b4e 100644 --- a/asyncband/src/internal/arena.rs +++ b/asyncband/src/internal/arena.rs @@ -121,6 +121,7 @@ impl Arena { } } + #[cfg(test)] pub fn len(&self) -> usize { self.len } From bf5d522963877bbe95dd7404fff193b251cf20a7 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:17:47 +0800 Subject: [PATCH 04/24] style: declare test_support::poll_once with plain pub The module itself is a private mod declaration, so crate visibility is already restricted at the module boundary. Repository style keeps restricted visibility on the boundary and plain pub on the module's items, the same pattern the internal module follows. --- asyncband/src/test_support.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/asyncband/src/test_support.rs b/asyncband/src/test_support.rs index a450f64..4bf8441 100644 --- a/asyncband/src/test_support.rs +++ b/asyncband/src/test_support.rs @@ -21,6 +21,6 @@ use std::task::Context; use std::task::Poll; use std::task::Waker; -pub(crate) fn poll_once(future: Pin<&mut F>) -> Poll { +pub fn poll_once(future: Pin<&mut F>) -> Poll { future.poll(&mut Context::from_waker(Waker::noop())) } From a0b741c6f15b241cd5c54ada14498b5bd879c09a Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:20:27 +0800 Subject: [PATCH 05/24] docs: remove comments that restate what the code declares These comments narrate the very line they precede instead of explaining a rationale: field docs that paraphrase the field's type ("Container storing the protected data" on UnsafeCell, "Non-null pointer to the mapped data" on NonNull), the "ownership certificate" voice-over on Arc lock fields, "Release the lock by calling release on the semaphore" above s.release(1), and the "Use DerefMut ..." meta commentary. The broadcast Shared::inner field doc additionally duplicated the Inner type's own documentation. Comments carrying real invariants (variance, non-zero max_readers) are kept. --- asyncband/src/broadcast/mpmc/bounded/mod.rs | 1 - asyncband/src/broadcast/mpmc/unbounded/mod.rs | 1 - asyncband/src/mutex/mod.rs | 13 ------------- asyncband/src/rwlock/mod.rs | 2 -- asyncband/src/rwlock/owned_mapped_read_guard.rs | 4 ---- 5 files changed, 21 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index f91faa2..8fcb79a 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -167,7 +167,6 @@ pub fn bounded(capacity: usize) -> (BoundedSender, BoundedReceiver< } struct Shared { - /// Buffer, receiver cursors, and parked receivers, all under a single lock. inner: Mutex>, /// Number of active senders. senders: AtomicUsize, diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 15366ef..f5e8e99 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -103,7 +103,6 @@ pub fn unbounded() -> (UnboundedSender, UnboundedReceiver) { } struct Shared { - /// Buffer, receiver cursors, and parked receivers, all under a single lock. inner: Mutex>, /// Number of active senders. senders: AtomicUsize, diff --git a/asyncband/src/mutex/mod.rs b/asyncband/src/mutex/mod.rs index ed04964..9b61e1e 100644 --- a/asyncband/src/mutex/mod.rs +++ b/asyncband/src/mutex/mod.rs @@ -78,9 +78,7 @@ use crate::internal::semaphore; /// /// See the [module level documentation](self) for more. pub struct Mutex { - /// Semaphore used to control access to protected data, ensuring mutual exclusion s: semaphore::Semaphore, - /// Container storing the protected data, allowing interior mutability c: UnsafeCell, } @@ -618,9 +616,7 @@ impl OwnedMutexGuard { /// ``` #[must_use = "dropping the guard releases the mutex immediately"] pub struct MappedMutexGuard<'a, T: ?Sized> { - /// Non-null pointer to the mapped data d: NonNull, - /// Reference to the original mutex's semaphore, used for releasing the lock s: &'a semaphore::Semaphore, // Mutable access requires invariance over T. variance: PhantomData<&'a mut T>, @@ -722,7 +718,6 @@ impl<'a, T: ?Sized> MappedMutexGuard<'a, T> { F: FnOnce(&mut T) -> &mut U, U: ?Sized, { - // Use DerefMut to safely get mutable reference, avoiding explicit unsafe block let d = NonNull::from(f(&mut *orig)); let orig = ManuallyDrop::new(orig); MappedMutexGuard { @@ -775,7 +770,6 @@ impl<'a, T: ?Sized> MappedMutexGuard<'a, T> { F: FnOnce(&mut T) -> Option<&mut U>, U: ?Sized, { - // Use DerefMut to safely get mutable reference, avoiding explicit unsafe block match f(&mut *orig) { Some(d) => { let d = NonNull::from(d); @@ -820,11 +814,7 @@ impl<'a, T: ?Sized> MappedMutexGuard<'a, T> { /// ``` #[must_use = "dropping the guard releases the mutex immediately"] pub struct OwnedMappedMutexGuard { - // This Arc acts as an ownership certificate, ensuring the Mutex remains valid - // and the lock is not released lock: Arc>, - // This NonNull pointer precisely points to the subfield U, telling us which - // memory location we can operate on, with compile-time guarantee of non-null d: NonNull, // Mutable access requires invariance over U. variance: PhantomData<*mut U>, @@ -844,7 +834,6 @@ unsafe impl Sync for OwnedMapp impl Drop for OwnedMappedMutexGuard { fn drop(&mut self) { - // Release the lock by calling release on the semaphore self.lock.s.release(1); } } @@ -923,7 +912,6 @@ impl OwnedMappedMutexGuard { F: FnOnce(&mut U) -> &mut V, V: ?Sized, { - // Use DerefMut to maintain consistency with other map implementations let d = NonNull::from(f(&mut *orig)); let orig = ManuallyDrop::new(orig); @@ -990,7 +978,6 @@ impl OwnedMappedMutexGuard { F: FnOnce(&mut U) -> Option<&mut V>, V: ?Sized, { - // Use DerefMut to maintain consistency with other filter_map implementations match f(&mut *orig) { Some(d) => { let d = NonNull::from(d); diff --git a/asyncband/src/rwlock/mod.rs b/asyncband/src/rwlock/mod.rs index 346e38b..2c2ee23 100644 --- a/asyncband/src/rwlock/mod.rs +++ b/asyncband/src/rwlock/mod.rs @@ -76,9 +76,7 @@ pub struct RwLock { /// /// This is ensured to be non-zero. max_readers: usize, - /// Semaphore to coordinate read and write access to T s: Semaphore, - /// The inner data. c: UnsafeCell, } diff --git a/asyncband/src/rwlock/owned_mapped_read_guard.rs b/asyncband/src/rwlock/owned_mapped_read_guard.rs index ae96acd..e9a7365 100644 --- a/asyncband/src/rwlock/owned_mapped_read_guard.rs +++ b/asyncband/src/rwlock/owned_mapped_read_guard.rs @@ -62,11 +62,7 @@ use crate::rwlock::RwLock; /// ``` #[must_use = "dropping the guard releases its read access immediately"] pub struct OwnedMappedRwLockReadGuard { - // This Arc acts as an ownership certificate, ensuring the RwLock remains valid - // and the lock is not released lock: Arc>, - // This NonNull pointer precisely points to the subfield U, telling us which - // memory location we can operate on d: NonNull, variance: PhantomData U>, } From 3d7258d179fae5c59fe623e2e9d946fae25adbf0 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:24:03 +0800 Subject: [PATCH 06/24] docs(mutex,rwlock): state the invariant behind the Arc extraction "We safely extract the Arc from the ManuallyDrop guard" just asserts safety in the imperative mood without giving a reason, and OwnedMutexGuard::map had no comment at all while its siblings did. Every extraction site now states the actual invariant, matching the wording the rwlock write guards already used: the source guard is ManuallyDrop and will not be dropped, so moving the Arc out transfers lock ownership instead of releasing it. --- asyncband/src/mutex/mod.rs | 11 ++++++++--- asyncband/src/rwlock/owned_read_guard.rs | 6 ++++-- 2 files changed, 12 insertions(+), 5 deletions(-) diff --git a/asyncband/src/mutex/mod.rs b/asyncband/src/mutex/mod.rs index 9b61e1e..ffed455 100644 --- a/asyncband/src/mutex/mod.rs +++ b/asyncband/src/mutex/mod.rs @@ -514,6 +514,8 @@ impl OwnedMutexGuard { let guard = ManuallyDrop::new(orig); + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&guard.lock) }; OwnedMappedMutexGuard { @@ -559,7 +561,8 @@ impl OwnedMutexGuard { let d = NonNull::from(d); let guard = ManuallyDrop::new(orig); - // SAFETY: We safely extract the Arc from the ManuallyDrop guard + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&guard.lock) }; Ok(OwnedMappedMutexGuard { @@ -915,7 +918,8 @@ impl OwnedMappedMutexGuard { let d = NonNull::from(f(&mut *orig)); let orig = ManuallyDrop::new(orig); - // SAFETY: We safely extract the Arc from the ManuallyDrop guard + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&orig.lock) }; OwnedMappedMutexGuard { @@ -983,7 +987,8 @@ impl OwnedMappedMutexGuard { let d = NonNull::from(d); let orig = ManuallyDrop::new(orig); - // SAFETY: We safely extract the Arc from the ManuallyDrop guard + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&orig.lock) }; Ok(OwnedMappedMutexGuard { diff --git a/asyncband/src/rwlock/owned_read_guard.rs b/asyncband/src/rwlock/owned_read_guard.rs index 5c2d1be..b9a2ae5 100644 --- a/asyncband/src/rwlock/owned_read_guard.rs +++ b/asyncband/src/rwlock/owned_read_guard.rs @@ -162,7 +162,8 @@ impl OwnedRwLockReadGuard { let d = std::ptr::NonNull::from(f(unsafe { &*orig.lock.c.get() })); let orig = std::mem::ManuallyDrop::new(orig); - // Safely extract the Arc from the guard + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&orig.lock) }; OwnedMappedRwLockReadGuard::new(d, lock) @@ -216,7 +217,8 @@ impl OwnedRwLockReadGuard { let d = std::ptr::NonNull::from(d); let orig = std::mem::ManuallyDrop::new(orig); - // Safely extract the Arc from the guard + // SAFETY: The guard is wrapped in `ManuallyDrop` and will not be dropped, + // so the `Arc` can be moved out to transfer lock ownership to the new guard. let lock = unsafe { std::ptr::read(&orig.lock) }; Ok(OwnedMappedRwLockReadGuard::new(d, lock)) From 5ccd6fef4e3bbaba1050ba487b61c839c9fc9dc7 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:24:23 +0800 Subject: [PATCH 07/24] refactor(rwlock): drop unneeded pub(super) from write guard fields Every access to RwLockWriteGuard and OwnedRwLockWriteGuard fields happens inside the file that defines the guard, so the pub(super) widens visibility without a consumer. The read guard fields keep theirs: downgrade constructs read guards from the write guard modules. --- asyncband/src/rwlock/owned_write_guard.rs | 4 ++-- asyncband/src/rwlock/write_guard.rs | 4 ++-- 2 files changed, 4 insertions(+), 4 deletions(-) diff --git a/asyncband/src/rwlock/owned_write_guard.rs b/asyncband/src/rwlock/owned_write_guard.rs index 8d7f558..99f1342 100644 --- a/asyncband/src/rwlock/owned_write_guard.rs +++ b/asyncband/src/rwlock/owned_write_guard.rs @@ -85,8 +85,8 @@ impl RwLock { /// alive without borrowing it and releases the lock when dropped. #[must_use = "dropping the guard releases its write access immediately"] pub struct OwnedRwLockWriteGuard { - pub(super) permits_acquired: usize, - pub(super) lock: Arc>, + permits_acquired: usize, + lock: Arc>, } unsafe impl Send for OwnedRwLockWriteGuard {} diff --git a/asyncband/src/rwlock/write_guard.rs b/asyncband/src/rwlock/write_guard.rs index 896ed60..a5dd89f 100644 --- a/asyncband/src/rwlock/write_guard.rs +++ b/asyncband/src/rwlock/write_guard.rs @@ -77,8 +77,8 @@ impl RwLock { /// [`RwLock::write`] and [`RwLock::try_write`] create this guard. Dropping it releases the lock. #[must_use = "dropping the guard releases its write access immediately"] pub struct RwLockWriteGuard<'a, T: ?Sized> { - pub(super) permits_acquired: usize, - pub(super) lock: &'a RwLock, + permits_acquired: usize, + lock: &'a RwLock, } unsafe impl Send for RwLockWriteGuard<'_, T> {} From dc79f1eaaa3ea5f6d3b6b96d22c3dfd34c4a494b Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:24:50 +0800 Subject: [PATCH 08/24] docs(rwlock): explain the default reader cap in one accurate sentence "large enough while not touch the edge" was ungrammatical and named no constraint: the cap exists to keep semaphore permit arithmetic away from usize::MAX while remaining unreachable in practice. --- asyncband/src/rwlock/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/asyncband/src/rwlock/mod.rs b/asyncband/src/rwlock/mod.rs index 2c2ee23..13ec0a2 100644 --- a/asyncband/src/rwlock/mod.rs +++ b/asyncband/src/rwlock/mod.rs @@ -117,7 +117,7 @@ impl RwLock { /// let rwlock = RwLock::new(5); /// ``` pub const fn new(t: T) -> RwLock { - // large enough while not touch the edge + // Effectively unlimited, while keeping permit arithmetic far from usize::MAX. RwLock::with_max_readers(t, NonZeroUsize::new(usize::MAX >> 1).unwrap()) } From 65e0528397a100846a42946ff5edfe4d14cedd69 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:25:34 +0800 Subject: [PATCH 09/24] refactor(broadcast): collapse the sender disconnect match into a conditional The Drop impls matched on the fetch_sub result with a no-op fallback arm whose only content was a comment saying there is nothing to do. An equality check expresses the last-sender test directly, the same shape the mpsc sender and waitgroup disconnect paths already use. The bounded sender keeps its comment explaining why only receivers need waking. --- asyncband/src/broadcast/mpmc/bounded/mod.rs | 11 ++++------- asyncband/src/broadcast/mpmc/unbounded/mod.rs | 7 ++----- 2 files changed, 6 insertions(+), 12 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index 8fcb79a..d81a014 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -244,13 +244,10 @@ impl fmt::Debug for BoundedSender { impl Drop for BoundedSender { fn drop(&mut self) { - match self.shared.senders.fetch_sub(1, Ordering::AcqRel) { - // Only parked receivers need waking. A parked producer borrows a live sender for the - // duration of its `send`, so the last sender cannot be dropping while one exists. - 1 => common::disconnect(&self.shared.inner), - _ => { - // there are still other senders left, do nothing - } + // Only parked receivers need waking. A parked producer borrows a live sender for the + // duration of its `send`, so the last sender cannot be dropping while one exists. + if self.shared.senders.fetch_sub(1, Ordering::AcqRel) == 1 { + common::disconnect(&self.shared.inner); } } } diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index f5e8e99..0d95f7c 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -136,11 +136,8 @@ impl fmt::Debug for UnboundedSender { impl Drop for UnboundedSender { fn drop(&mut self) { - match self.shared.senders.fetch_sub(1, Ordering::AcqRel) { - 1 => common::disconnect(&self.shared.inner), - _ => { - // there are still other senders left, do nothing - } + if self.shared.senders.fetch_sub(1, Ordering::AcqRel) == 1 { + common::disconnect(&self.shared.inner); } } } From ce3f6c530323ccc31e39dd1efc512b0b05d2e916 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:25:55 +0800 Subject: [PATCH 10/24] docs(broadcast): name publish_retained as the backlog growth point The reclaim_vacated comment claimed the buffer grows only in Backlog::publish, but the growth is the push_back in publish_retained, which the bounded channel calls directly, bypassing publish entirely. The accounting argument only holds if the reference names the actual growth point. --- asyncband/src/broadcast/mpmc/common.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index ab7f0ff..08bbcd8 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -323,8 +323,8 @@ impl Backlog { /// it leaves `head`, so the zero-count prefix ends there. Only removing a lagging subscription /// can release more. /// - /// `buffer` shrinks here and grows only in [`Backlog::publish`], so this is the one place - /// `retained()` can fall. A bounded channel therefore accounts for released capacity at + /// `buffer` shrinks here and grows only in [`Backlog::publish_retained`], so this is the one + /// place `retained()` can fall. A bounded channel therefore accounts for released capacity at /// exactly the two call sites that reach this: [`Backlog::receive`] and /// [`Backlog::remove_receiver`]. fn reclaim_vacated(&mut self) -> Reclaimed { From a8d49335aa3ea413e8c070f7193e223919d42bb6 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:27:05 +0800 Subject: [PATCH 11/24] refactor(spmc): remove the unreachable send-waiter handoff in Send::drop The queue is single-producer: send borrows the sender exclusively, so at most one Send future exists and send_waiters can hold only its own node. Once Drop removes that node, notify_one always finds an empty list and the "another blocked sender" scenario the field comment describes cannot occur. The mpmc queue keeps its handoff because multiple producers make it live; the receiver-side handoff here stays for the same reason. The removed branch could only fire for a mem::forget-leaked send future, where waking the stale waker is a pointless poll. --- asyncband/src/spmc/queue.rs | 16 ++-------------- 1 file changed, 2 insertions(+), 14 deletions(-) diff --git a/asyncband/src/spmc/queue.rs b/asyncband/src/spmc/queue.rs index 82b7718..36e9b38 100644 --- a/asyncband/src/spmc/queue.rs +++ b/asyncband/src/spmc/queue.rs @@ -231,8 +231,6 @@ impl Shared { struct Send<'a, T> { shared: &'a Shared, waiter: Option, - // `Drop` passes an unconsumed notification on before this value is destroyed, because its - // destructor may depend on another blocked sender making progress. value: Option, } @@ -276,23 +274,13 @@ impl Drop for Send<'_, T> { let Some(id) = self.waiter.take() else { return; }; - let (retired, waker) = { + let retired = { let mut state = self.shared.state.lock(); if state.receivers == 0 { return; } - let retired = state.send_waiters.remove_waiter(id); - // Hand an unconsumed notification to the next sender while the slot is still free. - let waker = if matches!(retired, Waiter::Notified) && state.has_capacity() { - state.send_waiters.notify_one() - } else { - None - }; - (retired, waker) + state.send_waiters.remove_waiter(id) }; - if let Some(waker) = waker { - waker.wake(); - } drop(retired); } } From 8d37e8e93ace59d9d185d75412555f9767542fd3 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:27:42 +0800 Subject: [PATCH 12/24] docs(spmc): describe the single internal mutex accurately Both constructors claimed operations acquire "internal mutexes", but the queue state lives behind one Mutex, and the unbounded paragraph even contradicted itself two sentences later. The mpmc and mpsc modules already say "an internal mutex". --- asyncband/src/spmc/bounded.rs | 2 +- asyncband/src/spmc/unbounded.rs | 2 +- 2 files changed, 2 insertions(+), 2 deletions(-) diff --git a/asyncband/src/spmc/bounded.rs b/asyncband/src/spmc/bounded.rs index f164013..4ff7cda 100644 --- a/asyncband/src/spmc/bounded.rs +++ b/asyncband/src/spmc/bounded.rs @@ -29,7 +29,7 @@ use super::queue::Shared; /// The queue stores at most `capacity` values. Sending waits for a receiver to free capacity when /// the queue is full. /// -/// Operations briefly acquire internal mutexes. No lock is held across an await point, while +/// Operations briefly acquire an internal mutex. No lock is held across an await point, while /// waking tasks, or while dropping messages. The `try_*` methods do not wait for capacity or /// messages, but may wait to acquire a mutex. /// diff --git a/asyncband/src/spmc/unbounded.rs b/asyncband/src/spmc/unbounded.rs index 60e8947..5e97498 100644 --- a/asyncband/src/spmc/unbounded.rs +++ b/asyncband/src/spmc/unbounded.rs @@ -28,7 +28,7 @@ use super::queue::Shared; /// /// Sends are synchronous and values may be buffered until available memory is exhausted. /// -/// Operations briefly acquire internal mutexes. No lock is held across an await point, while +/// Operations briefly acquire an internal mutex. No lock is held across an await point, while /// waking tasks, or while dropping messages. Sending and trying to receive may wait to acquire /// a mutex, but never wait for capacity or new messages. pub fn unbounded() -> (UnboundedSender, UnboundedReceiver) { From 48229b90e6a01261e947d206c89146dae36eebb5 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:28:15 +0800 Subject: [PATCH 13/24] docs(once-cell): align OnceCell docs with take() and task terminology The struct summary promised std-style "initialized at most once" semantics, but take() deliberately moves the cell back to an uninitialized state so it can be initialized again; the module-level docs already describe the actual contract. The set() docs also said "thread" where every sibling method says "task". --- asyncband/src/once/once_cell/mod.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/asyncband/src/once/once_cell/mod.rs b/asyncband/src/once/once_cell/mod.rs index 39ee75f..f084b6f 100644 --- a/asyncband/src/once/once_cell/mod.rs +++ b/asyncband/src/once/once_cell/mod.rs @@ -31,7 +31,7 @@ use crate::internal::value_cell::ValueCell; use crate::semaphore::Semaphore; use crate::semaphore::SemaphorePermit; -/// A thread-safe cell whose value is asynchronously initialized at most once. +/// A thread-safe cell that stores one value from an initializer supplied at access time. /// /// Callers provide an initializer when accessing an empty cell. An initializer that returns an /// error, panics, or is cancelled leaves the cell empty so a later caller can retry. Use @@ -297,7 +297,7 @@ impl OnceCell { /// Initializes the contents of the cell to `value`. /// - /// May wait if another thread is currently attempting to initialize the cell. The cell is + /// May wait if another task is currently attempting to initialize the cell. The cell is /// guaranteed to contain a value when `set` returns, though not necessarily the one provided. /// /// Returns `Ok(())` if the cell was uninitialized and `Err(value)` if the cell was already From f1793695afbe232c39822af87304f087b678419c Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:28:38 +0800 Subject: [PATCH 14/24] refactor(once-cell): remove the vestigial permit rebind in set_value let _permit = permit extended nothing: a by-value parameter already lives until the end of the function. It is scaffolding from the pre-ValueCell implementation, when the body held the permit across several statements; naming the parameter _permit expresses the held-but-unused intent directly and the SAFETY comment already documents the permit invariant. --- asyncband/src/once/once_cell/mod.rs | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/asyncband/src/once/once_cell/mod.rs b/asyncband/src/once/once_cell/mod.rs index f084b6f..edab814 100644 --- a/asyncband/src/once/once_cell/mod.rs +++ b/asyncband/src/once/once_cell/mod.rs @@ -380,8 +380,7 @@ impl OnceCell { self.value.take() } - fn set_value(&self, value: T, permit: SemaphorePermit<'_>) -> &T { - let _permit = permit; + fn set_value(&self, value: T, _permit: SemaphorePermit<'_>) -> &T { // SAFETY: Holding the only semaphore permit serializes initialization. unsafe { self.value.set(value) } } From 2b86eead3c9374136bfd87a07008fdc26bd0d067 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:29:15 +0800 Subject: [PATCH 15/24] refactor(latch): implement wait_owned by awaiting wait() OwnedLatchWait duplicated LatchWait field-for-field and poll-for-poll, differing only in holding Arc instead of &Latch, and the intern_poll helper existed only to be shared by the two. Both event primitives already implement the identical method as self.wait().await: the Arc lives inside the async block while wait() borrows from it, and cancel-safety is unchanged because the same LatchWait future drives the wait. With a single caller left, intern_poll is inlined back into LatchWait::poll. --- asyncband/src/latch/mod.rs | 35 ++--------------------------------- 1 file changed, 2 insertions(+), 33 deletions(-) diff --git a/asyncband/src/latch/mod.rs b/asyncband/src/latch/mod.rs index 77893e2..feae9c1 100644 --- a/asyncband/src/latch/mod.rs +++ b/asyncband/src/latch/mod.rs @@ -202,17 +202,7 @@ impl Latch { /// # } /// ``` pub async fn wait_owned(self: Arc) { - let fut = OwnedLatchWait { - token: None, - latch: self, - }; - fut.await - } -} - -impl Latch { - fn intern_poll(&self, token: &mut Option, cx: &mut Context<'_>) -> Poll<()> { - self.state.poll_wait(token, cx) + self.wait().await } } @@ -227,7 +217,7 @@ impl Future for LatchWait<'_> { fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { let Self { token, latch } = self.get_mut(); - latch.intern_poll(token, cx) + latch.state.poll_wait(token, cx) } } @@ -236,24 +226,3 @@ impl Drop for LatchWait<'_> { self.latch.state.unregister(&mut self.token); } } - -#[must_use = "futures do nothing unless you `.await` or poll them"] -struct OwnedLatchWait { - token: Option, - latch: Arc, -} - -impl Future for OwnedLatchWait { - type Output = (); - - fn poll(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { - let Self { token, latch } = self.get_mut(); - latch.intern_poll(token, cx) - } -} - -impl Drop for OwnedLatchWait { - fn drop(&mut self) { - self.latch.state.unregister(&mut self.token); - } -} From 07de3a4acd2d645cb858ff6957393bdf25036083 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:29:33 +0800 Subject: [PATCH 16/24] style(phaser): backtick .await in the must_use message Every other future type in the crate uses the message with `.await` in backticks; this one diverged by the missing markup alone. --- asyncband/src/phaser/mod.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/asyncband/src/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 70c524d..b96cba5 100644 --- a/asyncband/src/phaser/mod.rs +++ b/asyncband/src/phaser/mod.rs @@ -503,7 +503,7 @@ impl Drop for PhaserParticipant { } } -#[must_use = "futures do nothing unless you .await or poll them"] +#[must_use = "futures do nothing unless you `.await` or poll them"] struct PhaserWait<'a> { phaser: &'a Phaser, observed: u64, From 75f10d1669320f241f17d89bfae5d6fdd9cb1897 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:30:03 +0800 Subject: [PATCH 17/24] refactor(pool): remove the dead default type parameter from UnreadyObject The struct is private and constructed exactly once, inside Pool::get where M is inferred from the pool; neither impl block relies on the default. It was copied from the public Object> signature, where the default backs the Pool convenience API, and the bounded pool's UnreadyObject never had one. --- asyncband/src/pool/unbounded.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/asyncband/src/pool/unbounded.rs b/asyncband/src/pool/unbounded.rs index 927530f..ac0baff 100644 --- a/asyncband/src/pool/unbounded.rs +++ b/asyncband/src/pool/unbounded.rs @@ -571,7 +571,7 @@ impl> Object { /// If the check fails, `detach()` should be called to permanently remove the object /// from the pool. If dropped without calling either method (due to being cancelled), /// the behavior depends on the pool's [`RecycleCancelledStrategy`] configuration. -struct UnreadyObject = NeverManageObject> { +struct UnreadyObject> { state: Option>, pool: Weak>, recycle_cancelled_strategy: RecycleCancelledStrategy, From 5a5ad83fc0154f92aa4c9b348d373280d684ff7d Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:30:24 +0800 Subject: [PATCH 18/24] style(pool): narrow ObjectStatus mutation methods to pub(super) Both call sites live in the pool subtree (bounded.rs, unbounded.rs), so crate-wide visibility overshoots the consumers. pub would leak the methods into the public API because ObjectStatus is re-exported, and pub(super) is the idiom sibling modules already use for exactly this case, e.g. watch's SendError::new. --- asyncband/src/pool/common.rs | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/asyncband/src/pool/common.rs b/asyncband/src/pool/common.rs index 894b777..49fbb1d 100644 --- a/asyncband/src/pool/common.rs +++ b/asyncband/src/pool/common.rs @@ -57,11 +57,11 @@ impl ObjectStatus { self.recycle_count } - pub(crate) fn mark_recycled(&mut self) { + pub(super) fn mark_recycled(&mut self) { self.recycle_count += 1; } - pub(crate) fn mark_returned(&mut self) { + pub(super) fn mark_returned(&mut self) { self.last_returned = Some(Instant::now()); } } From 7afb9c78cd0cc8a48a93defd95198f24e22c0edc Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:30:40 +0800 Subject: [PATCH 19/24] docs(xtask): describe cargo x test as running workspace tests The subcommand invokes cargo test --workspace, which runs unit tests, the tests-integration suites, and doc tests; "Run unit tests" under-reports the workflow that AGENTS.md designates as the source of truth. The sibling entries already say "workspace" explicitly. --- xtask/src/main.rs | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/xtask/src/main.rs b/xtask/src/main.rs index cacfe8c..a047a4f 100644 --- a/xtask/src/main.rs +++ b/xtask/src/main.rs @@ -63,7 +63,7 @@ enum SubCommand { Miri(CommandMiri), #[clap(about = "Verify API compatibility for a planned release.")] Semver(CommandSemver), - #[clap(about = "Run unit tests.")] + #[clap(about = "Run workspace tests.")] Test(CommandTest), } From 49ee5e12f25f9829901a9fea5194a079f1f34a31 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:31:07 +0800 Subject: [PATCH 20/24] chore(examples): remove the vestigial cargo-release metadata The table is read only by cargo-release, which this repository does not use: the release skill publishes with cargo publish through ATR and Trusted Publishing, and nothing in the repo consumes the key. The package is already unpublished via publish = false, which is what the other private workspace members rely on. --- examples/Cargo.toml | 3 --- 1 file changed, 3 deletions(-) diff --git a/examples/Cargo.toml b/examples/Cargo.toml index 61f060e..a76d3e3 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -22,9 +22,6 @@ publish = false edition.workspace = true rust-version.workspace = true -[package.metadata.release] -release = false - [dependencies] asyncband = { workspace = true, features = [ "completion", From c2618fae268f6172bb3cd3c59b9a893523282a4d Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:31:55 +0800 Subject: [PATCH 21/24] docs: link the changelogs absolutely from MIGRATE.md The asyncband/ copy of MIGRATE.md ships inside the published crate and is also what the GitHub directory view renders, but its relative CHANGELOG.md and CHANGELOG-OLD.md links resolve against asyncband/, where no changelogs exist, and 404 there as well as in the packaged crate. Absolute links work in every context; the two copies stay byte-identical. --- MIGRATE.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/MIGRATE.md b/MIGRATE.md index 7604dce..e91636d 100644 --- a/MIGRATE.md +++ b/MIGRATE.md @@ -23,7 +23,7 @@ Asyncband continues the codebase formerly published as [`mea`](https://crates.io ## Recommended migration path -First, upgrade the existing dependency to `mea` 0.6.7 and resolve any changes required by earlier MEA releases. The [historical changelog](CHANGELOG-OLD.md) documents those releases. +First, upgrade the existing dependency to `mea` 0.6.7 and resolve any changes required by earlier MEA releases. The [historical changelog](https://github.com/apache/asyncband/blob/main/CHANGELOG-OLD.md) documents those releases. Next, switch from `mea` 0.6.7 to `asyncband` 0.6.7 without changing the dependency's feature configuration, and replace Rust paths from `mea::` to `asyncband::`: @@ -37,6 +37,6 @@ asyncband = "0.6.7" Asyncband 0.6.7 is the compatibility point for the rename, so the dependency name and Rust paths are the only changes expected in this step. No compatibility package or re-export keeps the `mea` crate name available; downstream crates must update those names directly. -Once the project builds with Asyncband 0.6.7, follow the [Asyncband changelog](CHANGELOG.md) when upgrading to later releases. +Once the project builds with Asyncband 0.6.7, follow the [Asyncband changelog](https://github.com/apache/asyncband/blob/main/CHANGELOG.md) when upgrading to later releases. For the background to the rename, see the [Asyncband proposal discussion](https://lists.apache.org/thread/f31qd3jm3odomjwy3lqkk21coyqsr9xs). From 9c1967ce9b1055f1590b17962708f98fc1c49885 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 14:33:00 +0800 Subject: [PATCH 22/24] chore: drop redundant default-features from workspace-inherited deps The workspace declarations for hashbrown and ureq already set default-features = false, so repeating it in the member manifests is a no-op; cargo tree confirms the resolved feature sets are unchanged. --- asyncband/Cargo.toml | 4 +--- xtask/Cargo.toml | 5 +---- 2 files changed, 2 insertions(+), 7 deletions(-) diff --git a/asyncband/Cargo.toml b/asyncband/Cargo.toml index a1ab5d9..0e6f3e3 100644 --- a/asyncband/Cargo.toml +++ b/asyncband/Cargo.toml @@ -70,9 +70,7 @@ waitgroup = [] watch = [] [dependencies] -hashbrown = { workspace = true, default-features = false, features = [ - "inline-more", -], optional = true } +hashbrown = { workspace = true, features = ["inline-more"], optional = true } [dev-dependencies] tokio = { workspace = true, features = ["full"] } diff --git a/xtask/Cargo.toml b/xtask/Cargo.toml index 5170bc1..8ec8366 100644 --- a/xtask/Cargo.toml +++ b/xtask/Cargo.toml @@ -31,10 +31,7 @@ cargo_metadata = { workspace = true } clap = { workspace = true, features = ["derive"] } semver = { workspace = true } serde = { workspace = true, features = ["derive"] } -ureq = { workspace = true, default-features = false, features = [ - "json", - "rustls", -] } +ureq = { workspace = true, features = ["json", "rustls"] } which = { workspace = true } [lints] From 826ba729b75155484a3c789a68691aa1899b6743 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 15:30:27 +0800 Subject: [PATCH 23/24] Revert "docs: link the changelogs absolutely from MIGRATE.md" This reverts commit c2618fae268f6172bb3cd3c59b9a893523282a4d. --- MIGRATE.md | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/MIGRATE.md b/MIGRATE.md index e91636d..7604dce 100644 --- a/MIGRATE.md +++ b/MIGRATE.md @@ -23,7 +23,7 @@ Asyncband continues the codebase formerly published as [`mea`](https://crates.io ## Recommended migration path -First, upgrade the existing dependency to `mea` 0.6.7 and resolve any changes required by earlier MEA releases. The [historical changelog](https://github.com/apache/asyncband/blob/main/CHANGELOG-OLD.md) documents those releases. +First, upgrade the existing dependency to `mea` 0.6.7 and resolve any changes required by earlier MEA releases. The [historical changelog](CHANGELOG-OLD.md) documents those releases. Next, switch from `mea` 0.6.7 to `asyncband` 0.6.7 without changing the dependency's feature configuration, and replace Rust paths from `mea::` to `asyncband::`: @@ -37,6 +37,6 @@ asyncband = "0.6.7" Asyncband 0.6.7 is the compatibility point for the rename, so the dependency name and Rust paths are the only changes expected in this step. No compatibility package or re-export keeps the `mea` crate name available; downstream crates must update those names directly. -Once the project builds with Asyncband 0.6.7, follow the [Asyncband changelog](https://github.com/apache/asyncband/blob/main/CHANGELOG.md) when upgrading to later releases. +Once the project builds with Asyncband 0.6.7, follow the [Asyncband changelog](CHANGELOG.md) when upgrading to later releases. For the background to the rename, see the [Asyncband proposal discussion](https://lists.apache.org/thread/f31qd3jm3odomjwy3lqkk21coyqsr9xs). From 0b3552ff65851c689348b3321adfb71897a80455 Mon Sep 17 00:00:00 2001 From: tison Date: Sat, 3 Oct 2026 15:46:02 +0800 Subject: [PATCH 24/24] docs: clarify channel operation terminology in comments --- asyncband/src/broadcast/mpmc/bounded/mod.rs | 69 ++++++++++--------- asyncband/src/broadcast/mpmc/bounded/tests.rs | 2 +- asyncband/src/broadcast/mpmc/common.rs | 49 ++++++------- asyncband/src/broadcast/mpmc/mod.rs | 11 +-- asyncband/src/broadcast/mpmc/unbounded/mod.rs | 10 +-- asyncband/src/mpmc/bounded.rs | 16 ++--- asyncband/src/mpmc/error.rs | 4 +- asyncband/src/mpmc/mod.rs | 2 +- asyncband/src/mpmc/queue.rs | 12 ++-- asyncband/src/mpmc/unbounded.rs | 4 +- asyncband/src/mpsc/bounded/mod.rs | 4 +- asyncband/src/mpsc/bounded/receiver.rs | 8 +-- asyncband/src/mpsc/bounded/sender.rs | 13 ++-- asyncband/src/mpsc/error.rs | 17 ++--- asyncband/src/mpsc/unbounded/mod.rs | 4 +- asyncband/src/mpsc/unbounded/receiver.rs | 8 +-- asyncband/src/oneshot/mod.rs | 2 +- asyncband/src/oneshot/receiver.rs | 24 +++---- asyncband/src/oneshot/sender.rs | 6 +- asyncband/src/spmc/bounded.rs | 13 ++-- asyncband/src/spmc/error.rs | 4 +- asyncband/src/spmc/mod.rs | 6 +- asyncband/src/spmc/queue.rs | 15 ++-- asyncband/src/spmc/unbounded.rs | 4 +- .../benches/ecosystem/mpsc/unbounded.rs | 4 +- benchmarks/benches/ecosystem/spmc/mod.rs | 2 +- benchmarks/benches/ecosystem/watch/paths.rs | 3 +- .../primitives/broadcast/mpmc/bounded.rs | 8 +-- .../primitives/broadcast/mpmc/unbounded.rs | 6 +- benchmarks/src/channels/mpmc.rs | 4 +- .../tests/broadcast_test/bounded.rs | 22 +++--- .../tests/broadcast_test/unbounded.rs | 10 +-- .../tests/mpmc_test/notification.rs | 2 +- .../tests/mpsc_test/backpressure.rs | 2 +- .../tests/mpsc_test/callbacks.rs | 2 +- .../tests/mpsc_test/concurrency.rs | 2 +- tests-integration/tests/oneshot_test/main.rs | 2 +- 37 files changed, 194 insertions(+), 182 deletions(-) diff --git a/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index d81a014..0bfbc8b 100644 --- a/asyncband/src/broadcast/mpmc/bounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/bounded/mod.rs @@ -18,7 +18,7 @@ //! A multi-producer multi-consumer broadcast channel with a bounded buffer. //! //! This channel supports multiple senders and multiple receivers. Each message sent by any sender -//! is received by all active receivers. Nothing is ever displaced to make room, so a receive never +//! is received by all active receivers. Nothing is ever displaced to make room, so receiving never //! reports lag; instead the channel retains at most `capacity` messages and makes producers wait. //! //! # Capacity @@ -32,11 +32,11 @@ //! subscription exerts backpressure" means, and it is the trade a lossless bounded broadcast //! makes. Drop a receiver that will not drain, and its backlog is released immediately. //! -//! If no receivers are active the channel retains nothing, so a send never waits. +//! If no receivers are active the channel retains nothing, so sending never waits. //! -//! A successful receive releases its subscription's claim before returning the value; processing -//! that value afterward does not hold capacity. The capacity limit excludes pending sends and -//! values already handed to application code. +//! Receiving a value releases the subscription's claim before returning that value; processing +//! it afterward does not hold capacity. The capacity limit excludes values held by pending `send` +//! futures and values already handed to application code. //! //! # Receivers //! @@ -49,7 +49,7 @@ //! Waiting producers are woken as capacity frees, but capacity is not reserved for them: a //! producer calling [`BoundedSender::try_send`] can take a slot that a woken producer was about to //! use, and that producer then waits again. Publication itself is one indivisible step, so -//! cancelling a send can never leave a gap in the committed order. +//! cancelling a `send` future can never leave a gap in the committed order. //! //! # Examples //! @@ -180,17 +180,19 @@ struct Shared { tx_permits: Semaphore, /// Number of producers in the waiting path of [`BoundedSender::send`]. /// - /// This upper bound lets receives skip the semaphore lock when no producer is waiting. + /// This upper bound lets receivers skip the semaphore lock when no producer is waiting. blocked_senders: AtomicUsize, } impl Shared { - /// Notifies blocked producers after a receive or subscription removal frees capacity. + /// Notifies waiting producers after receiving a message or removing a subscription frees + /// capacity. /// /// Call with the channel unlocked, before payload cloning or destruction can panic. fn release_reclaimed(&self, freed: usize) { - // 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. + // Cap permits at the number of blocked producers. Surplus permits make later `send` futures + // retry a full channel instead of waiting, 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)); @@ -199,7 +201,7 @@ impl Shared { /// Wakes every blocked producer when the last subscription leaves. /// - /// Sends now discard payloads without consuming capacity, so every producer can proceed. + /// Sending now discards payloads without consuming capacity, so every producer can proceed. fn release_all(&self) { if self.waiting_senders() > 0 { self.tx_permits.notify_all(); @@ -245,7 +247,7 @@ impl fmt::Debug for BoundedSender { impl Drop for BoundedSender { fn drop(&mut self) { // Only parked receivers need waking. A parked producer borrows a live sender for the - // duration of its `send`, so the last sender cannot be dropping while one exists. + // lifetime of its `send` future, so the last sender cannot be dropping while one exists. if self.shared.senders.fetch_sub(1, Ordering::AcqRel) == 1 { common::disconnect(&self.shared.inner); } @@ -263,13 +265,14 @@ impl BoundedSender { /// /// This method is cancel safe in the sense that matters for a lossless log: the value is /// either published to every active receiver or not published at all. Publication happens in - /// one indivisible step, so a cancelled send cannot leave a reserved but unfilled position in - /// the committed order. A send cancelled before it published drops the value with the future. + /// one indivisible step, so cancelling a `send` future cannot leave a reserved but unfilled + /// position in the committed order. Dropping the future before publication drops the unsent + /// value. /// /// # Panics /// - /// Panics if the internal message version counter overflows. After `u64::MAX` successful sends - /// on one channel instance, the next send panics. + /// Panics if the internal message version counter overflows. After `u64::MAX` successful + /// publications on one channel instance, the next attempt to publish a message panics. /// /// # Examples /// @@ -291,8 +294,8 @@ impl BoundedSender { struct SendState<'a, T> { sender: &'a BoundedSender, - // Declared before `value` so a cancelled send hands its registration back to the next - // waiting producer before running the payload's destructor. + // Declared before `value` so cancellation can wake the next waiting producer before + // the payload's destructor runs. acquire: Acquire<'a>, // Boxed once, out of the critical section, and reused by every retry. value: Option>, @@ -382,18 +385,20 @@ impl BoundedSender { /// ``` pub fn try_send(&self, value: T) -> Result<(), TrySendError> { // `Arc::new` runs inside the critical section, but only after the capacity check, so a - // rejected send never allocates. Unlike `T::clone` and `T::drop` it cannot run user code + // rejected attempt never allocates. Unlike `T::clone` and `T::drop` it cannot run user code // that reenters this channel, so it is safe to hold the lock across it. Hoisting it out // measured no faster even with eight producers contending — the allocator's thread-local - // cache already makes it cheap — and it measured slower wherever sends block, because a - // rejected send would then allocate and free once before `send` boxes the value for real. + // cache already makes it cheap — and it measured slower when producers wait for capacity, + // because a rejected attempt would then allocate and free once before the `send` future + // boxes the value for real. self.publish(value, Arc::new).map_err(TrySendError::Full) } /// Publishes a message that is already boxed, handing it back if the channel is still full. /// - /// This is the retry step of a waiting `send`, which boxes once with the channel unlocked and - /// then reuses that `Arc` for every attempt rather than reallocating per retry. + /// This is the retry step of a pending `send` future, which boxes once with the channel + /// unlocked and then reuses that `Arc` for every attempt rather than reallocating per + /// retry. fn try_publish(&self, msg: Arc) -> Result<(), Arc> { self.publish(msg, |msg| msg) } @@ -429,7 +434,7 @@ impl BoundedSender { /// against its [`capacity`](BoundedSender::capacity). /// /// The returned value is an instantaneous snapshot. It is suitable for diagnostics and soft - /// flow-control decisions, but concurrent sends and receives may change it immediately. + /// flow-control decisions, but other tasks may change it immediately by sending or receiving. /// /// # Examples /// @@ -543,9 +548,9 @@ impl BoundedReceiver { /// /// # Cancel safety /// - /// This method is cancel safe. If `recv` is used as the event in a `select` statement and some - /// other branch completes first, it is guaranteed that no messages were received on this - /// channel. + /// Dropping a pending `recv` future leaves this receiver's cursor unchanged. A subsequent + /// `recv` future can still return the same next value, so these futures may safely be raced + /// with other futures in a selection construct. /// /// # Examples /// @@ -590,8 +595,8 @@ impl BoundedReceiver { common::try_receive(&self.shared.inner, &self.shared.senders, self.key)?; // Release before taking the payload: `take_msg` runs `T::clone` and `T::drop`, and if - // either panics the slots this receive already freed would otherwise never be handed to a - // parked producer, stalling it permanently. + // either panics the reclaimed slots would otherwise never be handed to a waiting producer, + // stalling it permanently. self.shared.release_reclaimed(reclaimed.len()); Ok(common::take_msg(msg, reclaimed)) } @@ -635,7 +640,7 @@ impl BoundedReceiver { /// slowest active receiver. /// /// The returned value is an instantaneous snapshot. It is suitable for detecting that this - /// receiver is falling behind, but concurrent sends may change it immediately. + /// receiver is falling behind, but other tasks may publish more messages immediately. /// /// # Examples /// @@ -664,7 +669,7 @@ struct Recv<'a, T> { impl Drop for Recv<'_, T> { fn drop(&mut self) { - // Ready paths clear the token, so only a cancelled pending receive takes this lock. + // Ready paths clear the token, so only dropping a pending `Recv` future takes this lock. if self.token.is_none() { return; } @@ -697,7 +702,7 @@ impl Future for Recv<'_, T> { }; // Release before taking the payload, for the same reason as `try_recv`: a panicking - // `T::clone` must not strand producers on slots this receive already freed. + // `T::clone` must not strand producers waiting for the reclaimed slots. receiver.shared.release_reclaimed(reclaimed.len()); Poll::Ready(Ok(common::take_msg(msg, reclaimed))) } diff --git a/asyncband/src/broadcast/mpmc/bounded/tests.rs b/asyncband/src/broadcast/mpmc/bounded/tests.rs index 5fdd795..07c16df 100644 --- a/asyncband/src/broadcast/mpmc/bounded/tests.rs +++ b/asyncband/src/broadcast/mpmc/bounded/tests.rs @@ -63,7 +63,7 @@ fn buffer_is_preallocated_and_never_shrinks() { fn a_large_reclaim_leaves_no_permit_slack() { // Dropping a lagging subscription frees the whole backlog in one step, far more slots than the // single parked producer can use. Permits beyond that producer would sit in the semaphore, and - // the next send to block would burn each one on a publish attempt that cannot succeed. + // the next producer to wait would burn each one on a publication attempt that cannot succeed. let capacity = 64; let (tx, mut fast) = bounded(capacity); let lagging = tx.subscribe(); diff --git a/asyncband/src/broadcast/mpmc/common.rs b/asyncband/src/broadcast/mpmc/common.rs index 08bbcd8..74a04a9 100644 --- a/asyncband/src/broadcast/mpmc/common.rs +++ b/asyncband/src/broadcast/mpmc/common.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -//! Storage, cursors, and the receive step shared by the bounded and unbounded MPMC broadcast +//! Storage, cursors, and receiving logic shared by the bounded and unbounded MPMC broadcast //! channels. //! //! Both channels retain the same committed backlog and reclaim it the same way; they differ only @@ -43,7 +43,7 @@ use crate::internal::wakerset::WakerToken; /// Retained capacity below which an elastic backlog is never shrunk back. pub const MIN_RETAINED_CAPACITY: usize = 64; -/// A received message together with the retained prefix that the receive released. +/// A received message together with any buffer prefix reclaimed while advancing the cursor. /// /// The two travel together because the caller has to act on both with the channel unlocked, and a /// bounded channel has to hand the released capacity back before it touches the payload. @@ -52,7 +52,7 @@ pub type Received = (Arc, Reclaimed); /// Messages removed from the shared buffer and waiting to be dropped after it is unlocked. /// /// Keeping the first message out of the `Vec` avoids a heap allocation on the common path where -/// one receive reclaims exactly one message. +/// advancing a receiver's cursor reclaims exactly one message. pub struct Reclaimed { first: Option>, rest: Vec>, @@ -114,7 +114,7 @@ struct Slot { pub struct Backlog { /// Messages whose versions are in the range `[head, tail)`, each with its cursor count. /// - /// Each message is held behind an `Arc` so the receive path can move the payload out of the + /// Each message is held behind an `Arc` so a receiver can move the payload out of the /// critical section. Cloning the `Arc` under the lock keeps `T::clone` — and, for reclaimed /// messages, `T::drop` — outside it, which matters because both are arbitrary user code that /// may call back into this channel. @@ -300,7 +300,7 @@ impl Backlog { } else { Reclaimed::empty() }; - // A reclaim triggered by this receive always begins with this receiver's own message: the + // Reclaiming after advancing the cursor always begins with this receiver's own message: the // reclaim path runs only for a cursor leaving `head`, so the first slot drained is `msg`. // `take_msg` relies on this to recognize that it owns the payload. debug_assert!( @@ -319,9 +319,9 @@ 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 - /// it leaves `head`, so the zero-count prefix ends there. Only removing a lagging subscription - /// can release more. + /// A call to `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. /// /// `buffer` shrinks here and grows only in [`Backlog::publish_retained`], so this is the one /// place `retained()` can fall. A bounded channel therefore accounts for released capacity at @@ -333,9 +333,9 @@ impl Backlog { // Move reclaimed messages out so their Drop impls run after the channel is unlocked. Keep // the first one separate so the usual one-message reclaim does not allocate, and skip // building a `Drain` that would yield nothing: even an empty one costs a few nanoseconds - // on every receive. A bulk reclaim counts the zero-count prefix up front so it moves out - // in one drain instead of growing a vector geometrically, which measured 15% slower for a - // 32-message backlog. + // whenever a message is received. A bulk reclaim counts the zero-count prefix up front so + // it moves out in one drain instead of growing a vector geometrically, which + // measured 15% slower for a 32-message backlog. let first = self.buffer.pop_front().map(|slot| slot.msg); let extra = self .buffer @@ -406,7 +406,7 @@ impl Backlog { /// Buffer, receiver cursors, and parked receivers, all under one lock. /// /// 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. +/// registrations, so a message cannot be published between an empty check and registration. pub struct Inner { pub log: Backlog, pub waiters: WakerSet, @@ -433,7 +433,7 @@ pub fn disconnect(inner: &Mutex>) { wakers.wake_all(); } -/// Removes the waker registration for a cancelled `receive`, dropping the waker unlocked. +/// Removes a cancelled `recv` future's waker registration, dropping the waker after unlocking. pub fn unregister( inner: &Mutex>, senders: &AtomicUsize, @@ -452,7 +452,7 @@ pub fn unregister( drop(waker); } -/// Receives without waiting, yielding the message and the prefix the receive released. +/// Receives without waiting, returning the message and any buffer prefix reclaimed along with it. /// /// The caller owns what happens next: a bounded channel hands the released count back to blocked /// producers before it touches the payload. @@ -472,7 +472,7 @@ pub fn try_receive( } } -/// The one poll step behind `recv` on both channels. +/// Polls the shared receiving logic for both channels' `recv` futures. /// /// Publication and disconnection detach all registrations, so their ready paths clear the token /// without unregistering it. @@ -505,13 +505,14 @@ pub fn poll_receive( /// Drops the reclaimed backlog, then yields the received message, both with the channel unlocked. /// -/// A non-empty backlog means this receive drained `msg` from the buffer, so once the backlog is -/// dropped this receive holds the only reference and the payload can be moved out instead of -/// cloned. A channel with a single receiver therefore never clones a payload. +/// A non-empty backlog means receiving this message removed `msg` from the buffer. Once the +/// reclaimed backlog is dropped, the caller may hold the only remaining reference and can then +/// move the payload out instead of cloning it. A channel with a single receiver never clones a +/// payload. /// /// Ownership is decided from that bookkeeping rather than by probing the reference count. An -/// [`Arc::try_unwrap`] on every receive would fail under fan-out, and its failed compare-exchange -/// writes to a cache line that every receiver draining the message shares. +/// [`Arc::try_unwrap`] call for every received message would fail under fan-out, and its failed +/// compare-exchange writes to a cache line that every receiver draining the message shares. /// /// This runs `T::clone` and `T::drop`, either of which may panic, so a bounded channel must /// already have released the reclaimed capacity before calling it. @@ -534,8 +535,8 @@ mod tests { use super::Backlog; - /// Drives every cursor-count mutation point — subscribe, publish, receive, and drop — and - /// checks the accounting after each step. + /// Drives every cursor-count mutation point — subscribing, publishing, receiving, and dropping + /// — and checks the accounting after each step. #[test] fn cursor_accounting_holds_across_subscribe_receive_and_drop() { let mut log = Backlog::::elastic(); @@ -575,8 +576,8 @@ 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 - // back to `at_tail`. + // A call to `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(); assert_eq!(*msg, 4); diff --git a/asyncband/src/broadcast/mpmc/mod.rs b/asyncband/src/broadcast/mpmc/mod.rs index 636f895..b5eb4af 100644 --- a/asyncband/src/broadcast/mpmc/mod.rs +++ b/asyncband/src/broadcast/mpmc/mod.rs @@ -18,17 +18,18 @@ //! Multi-producer, multi-consumer broadcast channels. //! //! Both channels are lossless: every value a channel accepts stays readable by every subscription -//! that was active when it was accepted, so a receive never reports lag. They differ in what a +//! that was active when it was accepted, so receiving never reports lag. They differ in what a //! producer does when the slowest subscription stops reclaiming. [`bounded`] retains at most the //! capacity it was built with and makes producers wait for that subscription. [`unbounded`] never //! waits to send and lets the retained backlog grow instead. //! //! # Delivery and processing //! -//! A receive advances its subscription before returning the value. The channel tracks unread -//! messages, not application work: retaining a received value or processing it asynchronously -//! does not hold backlog capacity. There is no acknowledgement or processing-completion barrier. -//! If cloning a received value panics, that subscription has still advanced past the value. +//! Receiving a value advances the subscription before returning that value. The channel tracks +//! unread messages, not application work: retaining a received value or processing it +//! asynchronously does not hold backlog capacity. There is no acknowledgement or +//! processing-completion barrier. If cloning a received value panics, that subscription has still +//! advanced past the value. //! //! Sending with no subscriptions discards the value and succeeds. A later subscription starts //! with future publications; it does not replay discarded or previously retained values. diff --git a/asyncband/src/broadcast/mpmc/unbounded/mod.rs b/asyncband/src/broadcast/mpmc/unbounded/mod.rs index 0d95f7c..de2eab0 100644 --- a/asyncband/src/broadcast/mpmc/unbounded/mod.rs +++ b/asyncband/src/broadcast/mpmc/unbounded/mod.rs @@ -17,7 +17,7 @@ //! An unbounded fan-out channel with multiple senders and receivers. //! -//! A send publishes one value to every receiver that exists at that moment. Receivers advance +//! Sending publishes one value to every receiver that exists at that moment. Receivers advance //! independently, and a receiver created later starts with the next value rather than replaying //! earlier values. //! @@ -268,9 +268,9 @@ impl UnboundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` leaves this receiver's cursor unchanged. Its next call can still - /// return the same next value, so `recv` can be raced with other futures in a selection - /// construct. + /// Dropping a pending `recv` future leaves this receiver's cursor unchanged. A subsequent + /// `recv` future can still return the same next value, so these futures may safely be raced + /// with other futures in a selection construct. /// /// # Examples /// @@ -389,7 +389,7 @@ struct Recv<'a, T> { impl Drop for Recv<'_, T> { fn drop(&mut self) { - // Ready paths clear the token, so only a cancelled pending receive takes this lock. + // Ready paths clear the token, so only dropping a pending `Recv` future takes this lock. if self.token.is_none() { return; } diff --git a/asyncband/src/mpmc/bounded.rs b/asyncband/src/mpmc/bounded.rs index cd93c39..31dc1d9 100644 --- a/asyncband/src/mpmc/bounded.rs +++ b/asyncband/src/mpmc/bounded.rs @@ -28,8 +28,8 @@ use super::queue::Shared; /// Creates a bounded multi-producer, multi-consumer queue. /// /// Queued values, held permits, and capacity granted to waiting senders occupy at most `capacity` -/// slots. Pending sends and reservations receive capacity in wait-queue order. Sending waits for -/// a receiver to free capacity when none is available. +/// slots. Pending `send` and `reserve` operations receive capacity in wait-queue order. Sending +/// waits for a receiver to free capacity when none is available. /// /// The `try_*` methods do not wait for capacity or messages, but may briefly block on an internal /// mutex. @@ -84,9 +84,9 @@ impl BoundedSender { /// /// # Cancel safety /// - /// Dropping a pending `send` releases its waiting resources before dropping `value`, without - /// sending it or retaining capacity. Use [`reserve`](Self::reserve) to wait for capacity before - /// constructing a value. + /// Dropping a pending `send` future releases its waiting resources before dropping `value`, + /// without sending it or retaining capacity. Use [`reserve`](Self::reserve) to wait for + /// capacity before constructing a value. pub async fn send(&self, value: T) -> Result<(), SendError> { self.shared.send(value).await } @@ -100,7 +100,7 @@ impl BoundedSender { /// /// # Cancel safety /// - /// Dropping a pending reservation releases its place in the wait queue. If it was already + /// Dropping a pending `reserve` future releases its place in the wait queue. If it was already /// granted capacity, that capacity passes to the next waiter or becomes available again. /// /// # Examples @@ -176,8 +176,8 @@ impl BoundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not consume a value or prevent other receivers from receiving - /// it. + /// Dropping a pending `recv` future does not consume a value or prevent other receivers from + /// receiving it. pub async fn recv(&self) -> Result { self.shared.recv().await } diff --git a/asyncband/src/mpmc/error.rs b/asyncband/src/mpmc/error.rs index 4cbe177..200f013 100644 --- a/asyncband/src/mpmc/error.rs +++ b/asyncband/src/mpmc/error.rs @@ -103,7 +103,7 @@ impl fmt::Debug for TrySendError { impl std::error::Error for TrySendError {} -/// Error returned by a receive operation. +/// Error returned when receiving a value fails. #[derive(Debug, Clone, PartialEq, Eq)] pub enum RecvError { /// All senders have been dropped, and no buffered values remain. @@ -118,7 +118,7 @@ impl fmt::Display for RecvError { impl std::error::Error for RecvError {} -/// Error returned by a non-blocking receive operation. +/// Error returned by a non-blocking attempt to receive a value. #[derive(Debug, Clone, PartialEq, Eq)] pub enum TryRecvError { /// No value is currently available, but at least one sender remains. diff --git a/asyncband/src/mpmc/mod.rs b/asyncband/src/mpmc/mod.rs index 5ceda02..ebb03a0 100644 --- a/asyncband/src/mpmc/mod.rs +++ b/asyncband/src/mpmc/mod.rs @@ -20,7 +20,7 @@ //! Receivers compete for values: each value accepted by a sender is delivered to exactly one //! receiver while a receiver remains. Clone a receiver to distribute work across multiple //! asynchronous tasks. Dropping the final receiver releases any buffered values and makes later -//! sends return their value in an error. +//! attempts to send return their value in an error. mod bounded; mod error; diff --git a/asyncband/src/mpmc/queue.rs b/asyncband/src/mpmc/queue.rs index 327a411..211bab4 100644 --- a/asyncband/src/mpmc/queue.rs +++ b/asyncband/src/mpmc/queue.rs @@ -81,13 +81,13 @@ impl State { None } - /// Queues a value and selects the receiver to wake. + /// Queues a value and selects a waiting receiver's waker. fn push(&mut self, value: T) -> Option { self.values.push_back(value); self.recv_waiters.notify_one() } - /// Takes the next value and grants its capacity to the oldest waiting sender. + /// Takes the next value and grants its capacity to the first waiting operation. fn pop(&mut self) -> Result<(T, Option), TryRecvError> { if let Some(value) = self.values.pop_front() { Ok((value, self.release())) @@ -105,7 +105,7 @@ enum RecvWaiter { Notified, } -/// A grant transfers capacity to a detached sender waiter until it claims or cancels it. +/// A grant assigns capacity to a detached waiter node until its future claims or releases it. enum SendWaiter { Waiting(Waker), Granted, @@ -126,9 +126,9 @@ impl WaitList { Some(waker) } - /// Queues a blocked operation or refreshes the waker of a queued one. + /// Registers a pending operation's waker or refreshes an existing registration. /// - /// A notified receive that still found no value queues again at the back. + /// If the future finds no value after notification, its waiter rejoins the queue. #[must_use = "drop the replaced waker after releasing the queue lock"] fn register(&mut self, id: &mut Option, current: &Waker) -> Option { if let Some(queued) = *id { @@ -149,7 +149,7 @@ impl WaitList { } impl WaitList { - /// Queues a blocked sender or refreshes its waker without losing its place. + /// Registers an operation waiting for capacity or refreshes its waker without losing its place. #[must_use = "drop the replaced waker after releasing the queue lock"] fn register_waiter(&mut self, id: &mut Option, current: &Waker) -> Option { if let Some(queued) = *id { diff --git a/asyncband/src/mpmc/unbounded.rs b/asyncband/src/mpmc/unbounded.rs index 7fc71ac..20bf359 100644 --- a/asyncband/src/mpmc/unbounded.rs +++ b/asyncband/src/mpmc/unbounded.rs @@ -118,8 +118,8 @@ impl UnboundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not consume a value or prevent other receivers from receiving - /// it. + /// Dropping a pending `recv` future does not consume a value or prevent other receivers from + /// receiving it. pub async fn recv(&self) -> Result { self.shared.recv().await } diff --git a/asyncband/src/mpsc/bounded/mod.rs b/asyncband/src/mpsc/bounded/mod.rs index 60557b7..a1e5abf 100644 --- a/asyncband/src/mpsc/bounded/mod.rs +++ b/asyncband/src/mpsc/bounded/mod.rs @@ -37,8 +37,8 @@ pub use self::sender::Permit; /// Creates a bounded mpsc channel with room for `buffer` queued messages. /// /// [`BoundedSender::send`] waits for capacity when the buffer is full. Receiving a message releases -/// one slot for a waiting sender. Capacity is granted in the order that pending sends and -/// reservations enter the wait queue; new senders cannot take an already granted slot. +/// one slot for a waiting sender. Capacity is granted in the order that pending `send` and +/// `reserve` operations enter the wait queue; new senders cannot take an already granted slot. /// /// Message storage is preallocated for `buffer` values. Queued messages and outstanding /// reservations together occupy at most `buffer` capacity units. diff --git a/asyncband/src/mpsc/bounded/receiver.rs b/asyncband/src/mpsc/bounded/receiver.rs index c09ffee..bd4f611 100644 --- a/asyncband/src/mpsc/bounded/receiver.rs +++ b/asyncband/src/mpsc/bounded/receiver.rs @@ -32,7 +32,7 @@ use crate::mpsc::TryRecvError; /// The receiving endpoint of a bounded mpsc channel. /// /// Instances are created by the [`bounded`](crate::mpsc::bounded) function. Dropping the receiver -/// discards queued values and disconnects pending sends and reservations. +/// discards queued values and causes pending `send` and `reserve` operations to fail. pub struct BoundedReceiver { shared: Arc>>, } @@ -108,9 +108,9 @@ impl BoundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not remove a message from the channel. A later `recv` call - /// can still observe the next queued value, so `recv` may safely be raced with other futures - /// in a selection construct. + /// Dropping a pending `recv` future does not remove a message from the channel. A subsequent + /// `recv` future can still observe the next queued value, so these futures may safely be raced + /// with other futures in a selection construct. /// /// # Examples /// diff --git a/asyncband/src/mpsc/bounded/sender.rs b/asyncband/src/mpsc/bounded/sender.rs index 2781ff6..586ec9c 100644 --- a/asyncband/src/mpsc/bounded/sender.rs +++ b/asyncband/src/mpsc/bounded/sender.rs @@ -80,10 +80,11 @@ impl BoundedSender { /// /// # Cancel safety /// - /// Dropping a pending `send` loses its place waiting for capacity and drops `value`; a call - /// that has returned `Pending` has not sent the message. Use [`try_send`](Self::try_send) when - /// the caller must retain ownership if capacity is unavailable, or [`reserve`](Self::reserve) - /// to wait for capacity before constructing the message. + /// Dropping a pending `send` future loses its place waiting for capacity and drops `value`; + /// a future that has returned `Pending` has not sent the message. Use + /// [`try_send`](Self::try_send) when the caller must retain ownership if capacity is + /// unavailable, or [`reserve`](Self::reserve) to wait for capacity before constructing the + /// message. pub async fn send(&self, value: T) -> Result<(), SendError> { let value = match self.try_send(value) { Ok(()) => return Ok(()), @@ -108,8 +109,8 @@ impl BoundedSender { /// /// # Cancel safety /// - /// Dropping a pending reservation loses its place in the wait queue. If capacity has already - /// been granted, it is released to the next waiter or made available to a new sender. + /// Dropping a pending `reserve` future loses its place in the wait queue. If capacity has + /// already been granted, it is released to the next waiter or made available to a new sender. /// /// # Examples /// diff --git a/asyncband/src/mpsc/error.rs b/asyncband/src/mpsc/error.rs index 69e545e..9737e2f 100644 --- a/asyncband/src/mpsc/error.rs +++ b/asyncband/src/mpsc/error.rs @@ -18,14 +18,14 @@ use std::any::type_name; use std::fmt; -/// A send or capacity reservation failed because the receiver has been dropped. +/// An attempt to send a message or reserve capacity failed because the receiver has been dropped. /// /// Returned by [`UnboundedSender::send`], [`BoundedSender::send`], [`Permit::send`], and /// [`reserve`]. /// -/// A failed send retains the unsent message. A failed reservation carries `()` because no -/// message has been provided yet. Access the value with [`as_inner`](Self::as_inner) or -/// [`into_inner`](Self::into_inner). +/// If sending fails, the error retains the unsent message. If reserving capacity fails, the error +/// carries `()` because no message has been provided yet. Access the value with +/// [`as_inner`](Self::as_inner) or [`into_inner`](Self::into_inner). /// /// [`UnboundedSender::send`]: crate::mpsc::UnboundedSender::send /// [`BoundedSender::send`]: crate::mpsc::BoundedSender::send @@ -68,8 +68,9 @@ impl std::error::Error for SendError {} /// An attempt to send or reserve capacity failed. /// /// Returned by [`try_send`](crate::mpsc::BoundedSender::try_send) and -/// [`try_reserve`](crate::mpsc::BoundedSender::try_reserve). A failed send retains the unsent -/// message; a failed reservation carries `()` because no message has been provided yet. +/// [`try_reserve`](crate::mpsc::BoundedSender::try_reserve). If sending fails, the error retains +/// the unsent message; if reserving capacity fails, it carries `()` because no message has been +/// provided yet. #[derive(Clone, PartialEq, Eq)] pub enum TrySendError { /// No capacity is available for sending or reserving a message. @@ -115,7 +116,7 @@ impl fmt::Debug for TrySendError { impl std::error::Error for TrySendError {} -/// A receive operation cannot produce another value. +/// An attempt to receive a value failed because the channel is disconnected. #[derive(Debug, Clone, PartialEq, Eq)] pub enum RecvError { /// All senders have been dropped, and no buffered messages remain. @@ -130,7 +131,7 @@ impl fmt::Display for RecvError { impl std::error::Error for RecvError {} -/// A non-blocking receive did not produce a value. +/// A non-blocking attempt to receive a value failed. #[derive(Debug, Clone, PartialEq, Eq)] pub enum TryRecvError { /// No message is currently available, but at least one sender remains. diff --git a/asyncband/src/mpsc/unbounded/mod.rs b/asyncband/src/mpsc/unbounded/mod.rs index 42802cf..4955c27 100644 --- a/asyncband/src/mpsc/unbounded/mod.rs +++ b/asyncband/src/mpsc/unbounded/mod.rs @@ -31,7 +31,7 @@ mod sender; pub use self::receiver::UnboundedReceiver; pub use self::sender::UnboundedSender; -/// Creates an unbounded mpsc channel whose send operation never waits for capacity. +/// Creates an unbounded mpsc channel. Sending never waits for capacity. /// /// Pending messages can grow with producer demand and are limited only by successful memory /// allocation. Use a [`bounded`](crate::mpsc::bounded) channel or external admission control when @@ -59,7 +59,7 @@ pub fn unbounded() -> (UnboundedSender, UnboundedReceiver) { } // Queue contents, endpoint liveness, and wake registration share one lock. Only the receiver -// accesses its current batch; refilling that batch preserves the order of concurrent sends. +// accesses its current batch; refilling that batch preserves the order of concurrent `send` calls. struct State { buffer: Buffer, senders: usize, diff --git a/asyncband/src/mpsc/unbounded/receiver.rs b/asyncband/src/mpsc/unbounded/receiver.rs index 90617be..6874b4c 100644 --- a/asyncband/src/mpsc/unbounded/receiver.rs +++ b/asyncband/src/mpsc/unbounded/receiver.rs @@ -34,7 +34,7 @@ use crate::mpsc::TryRecvError; /// The receiving endpoint of an unbounded mpsc channel. /// /// Instances are created by the [`unbounded`](crate::mpsc::unbounded) function. Dropping the -/// receiver discards queued values and makes subsequent sends fail. +/// receiver discards queued values and makes subsequent attempts to send fail. pub struct UnboundedReceiver { shared: Arc>>, // Only accessed through `get_mut`; the mutex preserves Sync for Send-only payloads. @@ -120,9 +120,9 @@ impl UnboundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not remove a message from the channel. A later receive - /// operation can still observe the next queued value, so `recv` may safely be raced with other - /// futures in a selection construct. + /// Dropping a pending `recv` future does not remove a message from the channel. A subsequent + /// `recv` future can still observe the next queued value, so these futures may safely be raced + /// with other futures in a selection construct. /// /// # Examples /// diff --git a/asyncband/src/oneshot/mod.rs b/asyncband/src/oneshot/mod.rs index 310d4fc..fd7339c 100644 --- a/asyncband/src/oneshot/mod.rs +++ b/asyncband/src/oneshot/mod.rs @@ -84,7 +84,7 @@ //! Dropping either the receiver or its future disconnects the channel and discards any unread //! message. //! -//! To keep a pending receive alive when another branch wins, call [`Receiver::into_future`] +//! To keep a pending operation alive when another branch wins, call [`Receiver::into_future`] //! before selecting and borrow the resulting future as `&mut Recv`. mod receiver; diff --git a/asyncband/src/oneshot/receiver.rs b/asyncband/src/oneshot/receiver.rs index 521eaed..4b39707 100644 --- a/asyncband/src/oneshot/receiver.rs +++ b/asyncband/src/oneshot/receiver.rs @@ -47,7 +47,7 @@ use super::drop_message_and_deallocate_channel; /// Awaiting converts this receiver into a [`Recv`] future. Dropping either the receiver or its /// future disconnects the channel and discards any unread message. /// -/// To keep a pending receive alive when another branch wins, call [`Receiver::into_future`] +/// To keep a pending operation alive when another branch wins, call [`Receiver::into_future`] /// before selecting and borrow the resulting future as `&mut Recv`. The receiver itself does not /// implement [`Future`]. pub struct Receiver { @@ -62,8 +62,8 @@ impl fmt::Debug for Receiver { unsafe impl Send for Receiver {} -// Receiver must not be `Sync`: receive operations taking `&self` assume that no other receive -// operation runs concurrently. +// Receiver must not be `Sync`: `try_recv` takes `&self` but assumes exclusive access to the +// message. impl Unpin for Receiver {} @@ -86,7 +86,7 @@ impl Receiver { /// This occurs when the associated [`Sender`] is dropped without sending a message, or after /// the message is received. /// - /// If `true` is returned, all future receive operations are guaranteed to return an error. + /// If `true` is returned, any subsequent attempt to receive the message will return an error. pub fn is_disconnected(&self) -> bool { // SAFETY: The existence of `self` guarantees that the receiver is still alive. If the // sender was dropped, it observed the live receiver and left allocation cleanup to it, so @@ -102,8 +102,8 @@ impl Receiver { /// Returns true if there is a message in the channel, ready to be received. /// - /// If `true` is returned, the next call to receive the message is guaranteed to return - /// the message immediately. + /// If `true` is returned, the next attempt to receive the message is guaranteed to complete + /// immediately. pub fn has_message(&self) -> bool { // SAFETY: The existence of `self` guarantees that the receiver is still alive. If the // sender was dropped, it observed the live receiver and left allocation cleanup to it, so @@ -112,7 +112,7 @@ impl Receiver { // ORDERING: This method only observes the atomic state. MESSAGE is terminal for the sender, // and receiver operations cannot run concurrently, so atomic coherence preserves this - // observation for the next receive. Accessing the message synchronizes separately. + // observation until the message is received. Accessing the message synchronizes separately. matches!(channel.state.load(Ordering::Relaxed), MESSAGE) } @@ -123,9 +123,9 @@ impl Receiver { /// * `Err(TryRecvError::Disconnected)` if the [`Sender`] was dropped before sending anything or /// if the message has already been extracted by a previous `try_recv` call. /// - /// If a message is returned, the channel is disconnected and any subsequent receive operation - /// using this receiver will return an error: [`TryRecvError::Disconnected`] for `try_recv`, - /// or [`RecvError::Disconnected`] for [`recv`](Receiver::into_future). + /// If a message is returned, the channel is disconnected and any subsequent attempt to receive + /// a message will return an error: [`TryRecvError::Disconnected`] for `try_recv`, or + /// [`RecvError::Disconnected`] when awaiting the receiver. pub fn try_recv(&self) -> Result { // SAFETY: The channel will not be freed while this method is still running. let channel = unsafe { self.channel_ptr.as_ref() }; @@ -203,8 +203,8 @@ impl Drop for Receiver { /// Created by [`Receiver::into_future`], this future owns the receiving endpoint. Dropping it /// disconnects the channel and discards any unread message. /// -/// Select on `&mut Recv` to keep a pending receive alive when another branch wins. A completed -/// `Recv` must not be polled again. +/// Select on `&mut Recv` to keep the pending future alive when another branch wins. A completed +/// `Recv` future must not be polled again. pub struct Recv { channel_ptr: NonNull>, } diff --git a/asyncband/src/oneshot/sender.rs b/asyncband/src/oneshot/sender.rs index ad4287b..af161d4 100644 --- a/asyncband/src/oneshot/sender.rs +++ b/asyncband/src/oneshot/sender.rs @@ -95,8 +95,8 @@ impl Sender { if receiver_owns_allocation { waker.wake(); } else { - // The send remains successful because this sender owned the waker before the - // receiver cancelled. + // Sending still succeeds because this sender owned the waker before the + // receiving future was cancelled. // // SAFETY: Receiver cancellation transferred message and allocation cleanup to // this sender. The original pointer provenance may therefore be reclaimed as a @@ -137,7 +137,7 @@ impl Sender { let channel = unsafe { self.channel_ptr.as_ref() }; // ORDERING: Relaxed is sufficient for the method's contract: if this returns true, a - // future call to send is guaranteed to return an error. + // future call to `send` is guaranteed to return an error. // // Once true has been observed, it will remain true. However, if false is observed, the // receiver might just have been dropped without this thread observing it yet. diff --git a/asyncband/src/spmc/bounded.rs b/asyncband/src/spmc/bounded.rs index 4ff7cda..b4d7e52 100644 --- a/asyncband/src/spmc/bounded.rs +++ b/asyncband/src/spmc/bounded.rs @@ -75,10 +75,11 @@ impl BoundedSender { /// /// # Cancel safety /// - /// Dropping a pending `send` removes it from the wait queue and drops `value`; a call that has - /// returned `Pending` has not sent the value. Cancelling releases the exclusive sender borrow - /// and leaves available capacity usable by the next send. Use [`try_send`](Self::try_send) when - /// the caller must retain ownership if capacity is unavailable. + /// Dropping a pending `send` future removes its waiter and drops `value`. A future that has + /// returned `Pending` has not sent the value. Cancellation releases the exclusive sender borrow + /// and allows the next operation to use any available capacity. Use + /// [`try_send`](Self::try_send) when the caller must retain ownership if capacity is + /// unavailable. pub async fn send(&mut self, value: T) -> Result<(), SendError> { self.shared.send(value).await } @@ -129,8 +130,8 @@ impl BoundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not consume a value. Any selected value notification is - /// passed to another waiting receiver, so cancellation does not prevent it from receiving. + /// Dropping a pending `recv` future does not consume a value. Any unconsumed notification is + /// passed to another waiting task, so cancellation does not prevent it from receiving. pub async fn recv(&self) -> Result { self.shared.recv().await } diff --git a/asyncband/src/spmc/error.rs b/asyncband/src/spmc/error.rs index 18451db..fe63ad2 100644 --- a/asyncband/src/spmc/error.rs +++ b/asyncband/src/spmc/error.rs @@ -102,7 +102,7 @@ impl fmt::Debug for TrySendError { impl std::error::Error for TrySendError {} -/// Error returned by a receive operation. +/// Error returned when receiving a value fails. #[derive(Debug, Clone, PartialEq, Eq)] pub enum RecvError { /// The sender has been dropped, and no buffered values remain. @@ -117,7 +117,7 @@ impl fmt::Display for RecvError { impl std::error::Error for RecvError {} -/// Error returned by a non-blocking receive operation. +/// Error returned by a non-blocking attempt to receive a value. #[derive(Debug, Clone, PartialEq, Eq)] pub enum TryRecvError { /// No value is currently available, but the sender remains. diff --git a/asyncband/src/spmc/mod.rs b/asyncband/src/spmc/mod.rs index 613985e..dd79fbe 100644 --- a/asyncband/src/spmc/mod.rs +++ b/asyncband/src/spmc/mod.rs @@ -23,8 +23,8 @@ //! receiver while receivers remain. Values leave the queue in FIFO order, but consumer completion //! order and an equal distribution of work are not guaranteed. //! -//! The sender cannot be cloned, and every send operation requires `&mut self`, including for the -//! lifetime of a bounded send future. It can move between tasks, but shared references cannot send. +//! The sender cannot be cloned, and sending requires `&mut self`, including for the lifetime of a +//! bounded `send` future. The sender can move between tasks, but shared references cannot send. //! Dropping the sender lets receivers drain buffered messages before observing disconnection. //! Dropping the last receiver releases buffered messages and makes sending return the unsent value. //! @@ -75,7 +75,7 @@ //! } //! ``` //! -//! A bounded send future retains the exclusive borrow until completion or cancellation: +//! A bounded `send` future retains the exclusive borrow until completion or cancellation: //! //! ```compile_fail,E0499 //! let (mut sender, _receiver) = asyncband::spmc::bounded(1); diff --git a/asyncband/src/spmc/queue.rs b/asyncband/src/spmc/queue.rs index 36e9b38..ad47bdd 100644 --- a/asyncband/src/spmc/queue.rs +++ b/asyncband/src/spmc/queue.rs @@ -52,13 +52,13 @@ impl State { .is_none_or(|capacity| self.values.len() < capacity) } - /// Queues a value and selects the receiver to wake. + /// Queues a value and selects a waiting receiver's waker. fn push(&mut self, value: T) -> Option { self.values.push_back(value); self.recv_waiters.notify_one() } - /// Takes the next value and selects the sender to wake. + /// Takes the next value and selects a waiting sender's waker. fn pop(&mut self) -> Result<(T, Option), TryRecvError> { if let Some(value) = self.values.pop_front() { // Unbounded queues never block senders, so their sender queue is always empty. @@ -71,10 +71,11 @@ impl State { } } -/// A pending receive or bounded send. +/// Notification state for an operation waiting to send or receive a value. /// -/// Notification makes a waiter runnable; it does not reserve a value or slot. The detached node -/// remains owned by its future until it retries or is dropped. +/// Notification wakes the waiting task; it does not reserve a value or queue slot. +/// The future retains its waiter ID so it can reclaim the detached node when polled again or +/// dropped. Disconnection clears the waiter storage instead. enum Waiter { Waiting(Waker), Notified, @@ -95,9 +96,9 @@ impl WaitList { self.remove_unlinked_waiter(id) } - /// Queues a blocked operation or refreshes the waker of a queued one. + /// Registers a pending operation's waker or refreshes an existing registration. /// - /// A notified operation that still found no value or slot queues again at the back. + /// If the future finds no value or slot after notification, its waiter rejoins the queue. #[must_use = "drop the replaced waker after releasing the queue lock"] fn register(&mut self, id: &mut Option, current: &Waker) -> Option { if let Some(queued) = *id { diff --git a/asyncband/src/spmc/unbounded.rs b/asyncband/src/spmc/unbounded.rs index 5e97498..b2b4851 100644 --- a/asyncband/src/spmc/unbounded.rs +++ b/asyncband/src/spmc/unbounded.rs @@ -111,8 +111,8 @@ impl UnboundedReceiver { /// /// # Cancel safety /// - /// Dropping a pending `recv` does not consume a value. Any selected value notification is - /// passed to another waiting receiver, so cancellation does not prevent it from receiving. + /// Dropping a pending `recv` future does not consume a value. Any unconsumed notification is + /// passed to another waiting task, so cancellation does not prevent it from receiving. pub async fn recv(&self) -> Result { self.shared.recv().await } diff --git a/benchmarks/benches/ecosystem/mpsc/unbounded.rs b/benchmarks/benches/ecosystem/mpsc/unbounded.rs index b739f98..60f3611 100644 --- a/benchmarks/benches/ecosystem/mpsc/unbounded.rs +++ b/benchmarks/benches/ecosystem/mpsc/unbounded.rs @@ -126,7 +126,7 @@ fn repeated_bursts, T, F: Fn() -> T>( } }; // Measure recurring bursts after the initial allocation, including a deliberately retained - // backlog where requested. Do not require an extra empty receive to trigger reclamation. + // backlog where requested. Do not require an extra `try_recv` call to trigger reclamation. run(); bencher.counter(ItemsCount::new(messages)).bench_local(run); } @@ -228,7 +228,7 @@ fn scheduled_bursts_inline>( for _ in 0..BATCH_MESSAGES / burst_messages { let first = { // Exclude Tokio's cooperative-budget Pending from the initial empty probe. - // Retain the same receive future so its channel registration drives the + // Retain the same `recv_async` future so its channel registration drives the // wake. let mut receive = pin!(tokio::task::unconstrained(C::recv_async(&mut receiver))); diff --git a/benchmarks/benches/ecosystem/spmc/mod.rs b/benchmarks/benches/ecosystem/spmc/mod.rs index a9fcbca..4d1e0a0 100644 --- a/benchmarks/benches/ecosystem/spmc/mod.rs +++ b/benchmarks/benches/ecosystem/spmc/mod.rs @@ -15,7 +15,7 @@ // specific language governing permissions and limitations // under the License. -//! On a current-thread executor, unbounded sends can finish before consumers run. +//! On a current-thread executor, sending to an unbounded queue can finish before consumers run. mod bounded; mod unbounded; diff --git a/benchmarks/benches/ecosystem/watch/paths.rs b/benchmarks/benches/ecosystem/watch/paths.rs index 1b8237e..0be7228 100644 --- a/benchmarks/benches/ecosystem/watch/paths.rs +++ b/benchmarks/benches/ecosystem/watch/paths.rs @@ -15,7 +15,8 @@ // specific language governing permissions and limitations // under the License. -// Both adapters return an owned usize. Tokio's recv combines changed and borrow_and_update. +// Both adapters return an owned `usize`. The Tokio adapter's `recv` combines `changed` and +// `borrow_and_update`. use std::pin::pin; diff --git a/benchmarks/benches/primitives/broadcast/mpmc/bounded.rs b/benchmarks/benches/primitives/broadcast/mpmc/bounded.rs index 6d5cdb3..dd03bdf 100644 --- a/benchmarks/benches/primitives/broadcast/mpmc/bounded.rs +++ b/benchmarks/benches/primitives/broadcast/mpmc/bounded.rs @@ -64,9 +64,9 @@ fn try_send_and_drain_fanout(bencher: Bencher, receiver_count: usize) { fn full_channel_sender_handoff_cycle(bencher: Bencher, sender_count: usize) { let mut context = bench_context(); - // Register every producer on a full channel, complete one send after reclaiming a slot, - // cancel the remaining sends, and restore the original full backlog. Registration and - // cancellation are timed as part of this cycle. + // Register every producer on a full channel, complete one `send` future after reclaiming a + // slot, cancel the remaining futures, and restore the original full backlog. Registration + // and cancellation are timed as part of this cycle. bencher .with_inputs(|| { let (tx, rx) = mpmc::bounded(1); @@ -81,7 +81,7 @@ fn full_channel_sender_handoff_cycle(bencher: Bencher, sender_count: usize) { poll_pending(send.as_mut(), &mut context); } - // Poll queued sends until one republishes into the freed slot. + // Poll pending `send` futures until one publishes into the freed slot. black_box(rx.try_recv().unwrap()); for send in &mut sends { if send.as_mut().poll(&mut context).is_ready() { diff --git a/benchmarks/benches/primitives/broadcast/mpmc/unbounded.rs b/benchmarks/benches/primitives/broadcast/mpmc/unbounded.rs index 440db89..7728c17 100644 --- a/benchmarks/benches/primitives/broadcast/mpmc/unbounded.rs +++ b/benchmarks/benches/primitives/broadcast/mpmc/unbounded.rs @@ -56,7 +56,7 @@ const RECLAIM_FANOUTS: &[Fanout] = &[ Fanout { peak: 256, live: 1 }, ]; -// With the payload shared, each receive clones it and the second one reclaims the slot. +// Two receivers read the shared payload; the second reclaims the slot. #[divan::bench] fn send_and_try_recv_shared(bencher: Bencher) { let (sender, mut first) = mpmc::unbounded(); @@ -68,8 +68,8 @@ fn send_and_try_recv_shared(bencher: Bencher) { }); } -// The `usize` benchmarks above hide what a receive costs for a payload that owns memory: a clone -// there is an allocation, not a register move. +// The `usize` benchmarks above hide the cost of receiving a payload that owns memory: cloning it +// requires an allocation rather than a register move. fn payload() -> String { "x".repeat(64) } diff --git a/benchmarks/src/channels/mpmc.rs b/benchmarks/src/channels/mpmc.rs index 096b63f..d7582f3 100644 --- a/benchmarks/src/channels/mpmc.rs +++ b/benchmarks/src/channels/mpmc.rs @@ -73,8 +73,8 @@ pub const TOPOLOGIES: &[Topology] = &[ }, ]; -// The caller only coordinates the batch. All measured sends and receives run in spawned tasks, -// including on the current-thread runtime; no data is received by Runtime::block_on itself. +// The caller only coordinates the batch. Spawned tasks send and receive all measured messages, +// including on the current-thread runtime; `Runtime::block_on` itself receives no data. pub struct TaskBatch { start: Arc, workers: JoinSet<(usize, usize)>, diff --git a/tests-integration/tests/broadcast_test/bounded.rs b/tests-integration/tests/broadcast_test/bounded.rs index a278c83..5423783 100644 --- a/tests-integration/tests/broadcast_test/bounded.rs +++ b/tests-integration/tests/broadcast_test/bounded.rs @@ -52,7 +52,7 @@ impl Drop for Reentrant { } } -/// A payload that panics while a shared receive clones it. +/// A payload that panics when a receiver clones it. #[derive(Debug)] struct PanicOnClone { value: u64, @@ -342,8 +342,8 @@ fn dropping_the_last_receiver_wakes_every_blocked_sender() { tx.try_send(0).unwrap(); tx.try_send(1).unwrap(); - // More blocked producers than the drop will reclaim slots. Once no receiver remains every - // send succeeds unconditionally, so waking only `reclaimed` of them would strand the rest. + // More blocked producers than the drop will reclaim slots. Once no receiver remains, sending + // succeeds unconditionally, so waking only `reclaimed` of them would strand the rest. let trackers = (0..BLOCKED) .map(|_| Arc::new(WakeCounter::default())) .collect::>(); @@ -415,8 +415,8 @@ fn cancelled_send_publishes_nothing() { assert!(poll_once(send.as_mut()).is_pending()); drop(send); - // The cancelled value never entered the committed order, so the next receive sees only what - // was already published, and the one after it is a fresh send. + // The cancelled value never entered the committed order. The receiver sees the previously + // published value followed by a newly sent value. assert_eq!(rx.try_recv(), Ok(0)); tx.try_send(2).unwrap(); assert_eq!(rx.try_recv(), Ok(2)); @@ -571,8 +571,8 @@ fn large_reclaim_notifies_every_blocked_sender() { for tracker in trackers { assert_eq!(tracker.count(), 1); } - // Every send fits without another receive. All must have been notified, not merely made - // ready for a poll that an executor would otherwise have no reason to perform. + // Every pending value fits without receiving another message. All waiting tasks must have + // been notified, not merely made ready for a poll the executor has no reason to perform. for send in &mut sends { assert!(poll_once(send.as_mut()).is_ready()); } @@ -601,14 +601,14 @@ fn bounded_panicking_clone_leaves_the_channel_consistent() { }) .unwrap(); - // Two receivers share the payload, so this receive has to clone it. + // Two receivers share the payload, so `try_recv` has to clone it. let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { rx1.try_recv().map(|msg| msg.value) })); assert!(result.is_err()); - // The failed receive still consumed the message for `rx1`, and left the channel usable for - // both receivers. + // The panicking `try_recv` call still consumed the message for `rx1`, and left the channel + // usable for both receivers. assert_eq!(rx1.try_recv().unwrap().value, 2); assert_eq!(rx2.try_recv().unwrap().value, 1); assert_eq!(rx2.try_recv().unwrap().value, 2); @@ -659,7 +659,7 @@ fn bounded_message_destructors_run_outside_the_channel_lock() { .unwrap(); } - // Reclaim through a receive, and then through receiver drops. + // Exercise reclamation while receiving messages and dropping receivers. assert_eq!(rx1.try_recv().unwrap().value, 0); drop(rx2); assert_eq!(rx1.try_recv().unwrap().value, 1); diff --git a/tests-integration/tests/broadcast_test/unbounded.rs b/tests-integration/tests/broadcast_test/unbounded.rs index afd1d77..3324e4d 100644 --- a/tests-integration/tests/broadcast_test/unbounded.rs +++ b/tests-integration/tests/broadcast_test/unbounded.rs @@ -52,7 +52,7 @@ impl Drop for Reentrant { } } -/// A payload that panics while a shared receive clones it. +/// A payload that panics when a receiver clones it. #[derive(Debug)] struct PanicOnClone { value: u64, @@ -263,14 +263,14 @@ fn panicking_clone_leaves_the_channel_consistent() { panic: false, }); - // Two receivers share the payload, so this receive has to clone it. + // Two receivers share the payload, so `try_recv` has to clone it. let result = std::panic::catch_unwind(std::panic::AssertUnwindSafe(|| { rx1.try_recv().map(|msg| msg.value) })); assert!(result.is_err()); - // The failed receive still consumed the message for `rx1`, and left the channel usable for - // both receivers. + // The panicking `try_recv` call still consumed the message for `rx1`, and left the channel + // usable for both receivers. assert_eq!(rx1.try_recv().unwrap().value, 2); assert_eq!(rx2.try_recv().unwrap().value, 1); assert_eq!(rx2.try_recv().unwrap().value, 2); @@ -291,7 +291,7 @@ fn message_destructors_run_outside_the_channel_lock() { }); } - // Reclaim through a receive, and then through a receiver drop. + // Exercise reclamation while receiving messages and dropping receivers. assert_eq!(rx1.try_recv().unwrap().value, 0); drop(rx2); assert_eq!(rx1.try_recv().unwrap().value, 1); diff --git a/tests-integration/tests/mpmc_test/notification.rs b/tests-integration/tests/mpmc_test/notification.rs index 5b8845c..f81f895 100644 --- a/tests-integration/tests/mpmc_test/notification.rs +++ b/tests-integration/tests/mpmc_test/notification.rs @@ -71,7 +71,7 @@ fn receiver_notifications>( assert_eq!(receiver.try_recv(), Ok(1)); drop(first); assert_eq!(second_wakes.count(), 0); - // Cancellation must leave the next send able to notify this waiter. + // Sending another value must still notify this waiter after cancellation. send(&sender, 2); 2 } diff --git a/tests-integration/tests/mpsc_test/backpressure.rs b/tests-integration/tests/mpsc_test/backpressure.rs index 1ea57a7..4214323 100644 --- a/tests-integration/tests/mpsc_test/backpressure.rs +++ b/tests-integration/tests/mpsc_test/backpressure.rs @@ -66,7 +66,7 @@ fn cancelling_a_sender_preserves_capacity_and_notifies_the_next_waiter() { assert_eq!(second_wakes.count(), 0); } drop(first); - // Even a granted send must not publish its value until it is polled to completion. + // Even after capacity is granted, the `send` future must be polled to publish its value. if cancel_after_notification { assert_eq!(rx.try_recv(), Err(mpsc::TryRecvError::Empty)); } diff --git a/tests-integration/tests/mpsc_test/callbacks.rs b/tests-integration/tests/mpsc_test/callbacks.rs index f8b5d73..fa6e112 100644 --- a/tests-integration/tests/mpsc_test/callbacks.rs +++ b/tests-integration/tests/mpsc_test/callbacks.rs @@ -211,7 +211,7 @@ fn unbounded_disconnect_drops_partial_and_queued_batches_outside_lock() { let drops = drops.clone(); Arc::new(move || { // This exercises both a live receiver and disconnection. The marker has no - // callback, so destroying an unsuccessful send cannot recursively send again. + // callback, so dropping a rejected value cannot recursively send again. let _ = tx.send(Value(None)); drops.fetch_add(1, Ordering::Relaxed); }) diff --git a/tests-integration/tests/mpsc_test/concurrency.rs b/tests-integration/tests/mpsc_test/concurrency.rs index 69968b0..e77d0d7 100644 --- a/tests-integration/tests/mpsc_test/concurrency.rs +++ b/tests-integration/tests/mpsc_test/concurrency.rs @@ -182,7 +182,7 @@ fn bounded_try_recv_does_not_report_empty_after_completed_sends() { let mut received = 0; while received < PRODUCERS * MESSAGES_PER_PRODUCER { - // Once more sends have completed than messages received, Empty cannot be correct. + // Once more messages have been sent than received, `Empty` cannot be correct. let has_completed_send = completed.load(Ordering::Acquire) > received; match rx.try_recv() { Ok((producer, sequence)) => { diff --git a/tests-integration/tests/oneshot_test/main.rs b/tests-integration/tests/oneshot_test/main.rs index bd3eba0..9fdd122 100644 --- a/tests-integration/tests/oneshot_test/main.rs +++ b/tests-integration/tests/oneshot_test/main.rs @@ -295,7 +295,7 @@ fn poll_then_drop_receiver_during_send() { let sender_thread = spawn_named("sender", move || sender.send(message)); drop(receiver); - // Whether send or receiver drop wins, exactly one side owns and drops the message. + // Whether sending or dropping the receiver wins, exactly one side owns and drops the message. drop(sender_thread.join().unwrap()); assert_eq!(message_drop_count.load(Ordering::Relaxed), 1); }