diff --git a/asyncband/Cargo.toml b/asyncband/Cargo.toml index a1ab5d9c..0e6f3e38 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/asyncband/src/broadcast/mpmc/bounded/mod.rs b/asyncband/src/broadcast/mpmc/bounded/mod.rs index f91faa24..0bfbc8b5 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 //! @@ -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, @@ -181,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)); @@ -200,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,13 +246,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 + // 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); } } } @@ -267,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 /// @@ -295,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>, @@ -386,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) } @@ -433,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 /// @@ -547,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 /// @@ -594,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)) } @@ -639,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 /// @@ -668,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; } @@ -701,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 5fdd7957..07c16df1 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 ab7f0fff..74a04a9a 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,12 +319,12 @@ 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`], 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 { @@ -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 636f8958..b5eb4afd 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 15366ef5..de2eab06 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. //! @@ -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, @@ -137,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); } } } @@ -272,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 /// @@ -393,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/internal/arena.rs b/asyncband/src/internal/arena.rs index 0cd8647f..68a6b4e1 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), @@ -127,6 +121,7 @@ impl Arena { } } + #[cfg(test)] pub fn len(&self) -> usize { self.len } @@ -135,13 +130,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 diff --git a/asyncband/src/latch/mod.rs b/asyncband/src/latch/mod.rs index 77893e22..feae9c16 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); - } -} diff --git a/asyncband/src/mpmc/bounded.rs b/asyncband/src/mpmc/bounded.rs index cd93c398..31dc1d9b 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 4cbe177d..200f013e 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 5ceda024..ebb03a0c 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 327a411c..211bab46 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 7fc71ac6..20bf3591 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 60557b77..a1e5abfc 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 c09ffee5..bd4f6110 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 2781ff6c..586ec9c1 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 69e545e7..9737e2f2 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 42802cfa..4955c27e 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 90617be4..6874b4c2 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/mutex/mod.rs b/asyncband/src/mutex/mod.rs index ed049641..ffed4559 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, } @@ -516,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 { @@ -561,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 { @@ -618,9 +619,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 +721,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 +773,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 +817,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 +837,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,11 +915,11 @@ 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); - // 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 { @@ -990,13 +982,13 @@ 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); 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/once/once_cell/mod.rs b/asyncband/src/once/once_cell/mod.rs index 39ee75f1..edab814f 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 @@ -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) } } diff --git a/asyncband/src/oneshot/mod.rs b/asyncband/src/oneshot/mod.rs index 310d4fc4..fd7339cf 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 521eaed2..4b397079 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 ad4287b6..af161d40 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/phaser/mod.rs b/asyncband/src/phaser/mod.rs index 70c524d3..b96cba5f 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, diff --git a/asyncband/src/pool/common.rs b/asyncband/src/pool/common.rs index 894b777b..49fbb1d8 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()); } } diff --git a/asyncband/src/pool/unbounded.rs b/asyncband/src/pool/unbounded.rs index 927530fd..ac0baffb 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, diff --git a/asyncband/src/rwlock/mod.rs b/asyncband/src/rwlock/mod.rs index 346e38b5..13ec0a2b 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, } @@ -119,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()) } diff --git a/asyncband/src/rwlock/owned_mapped_read_guard.rs b/asyncband/src/rwlock/owned_mapped_read_guard.rs index ae96acdc..e9a7365c 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>, } diff --git a/asyncband/src/rwlock/owned_read_guard.rs b/asyncband/src/rwlock/owned_read_guard.rs index 5c2d1be7..b9a2ae5d 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)) diff --git a/asyncband/src/rwlock/owned_write_guard.rs b/asyncband/src/rwlock/owned_write_guard.rs index 8d7f5584..99f1342b 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 896ed607..a5dd89f2 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> {} diff --git a/asyncband/src/spmc/bounded.rs b/asyncband/src/spmc/bounded.rs index f164013b..b4d7e529 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. /// @@ -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 18451db2..fe63ad2a 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 613985ea..dd79fbe9 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 82b77187..ad47bddf 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 { @@ -231,8 +232,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 +275,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); } } diff --git a/asyncband/src/spmc/unbounded.rs b/asyncband/src/spmc/unbounded.rs index 60e89473..b2b4851e 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) { @@ -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/asyncband/src/test_support.rs b/asyncband/src/test_support.rs index a450f64e..4bf84413 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())) } diff --git a/benchmarks/benches/ecosystem/mpsc/unbounded.rs b/benchmarks/benches/ecosystem/mpsc/unbounded.rs index b739f98a..60f36117 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 a9fcbca1..4d1e0a0e 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 1b8237eb..0be7228e 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 6d5cdb3f..dd03bdf5 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 440db89b..7728c17b 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 096b63f9..d7582f36 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/examples/Cargo.toml b/examples/Cargo.toml index 61f060e6..a76d3e35 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", diff --git a/tests-integration/tests/broadcast_test/bounded.rs b/tests-integration/tests/broadcast_test/bounded.rs index a278c837..5423783e 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 afd1d779..3324e4d9 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 5b8845c4..f81f8953 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 1ea57a78..42143233 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 f8b5d733..fa6e1120 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 69968b0f..e77d0d71 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 bd3eba01..9fdd122f 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); } diff --git a/xtask/Cargo.toml b/xtask/Cargo.toml index 5170bc1b..8ec8366d 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] diff --git a/xtask/src/main.rs b/xtask/src/main.rs index cacfe8cc..a047a4f4 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), }