Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
24 commits
Select commit Hold shift + click to select a range
ecb998c
refactor(internal): remove the dead Arena::values iterator
tisonkun Oct 3, 2026
67ed45e
refactor(internal): remove the unused Default impl for Arena
tisonkun Oct 3, 2026
04e063f
refactor(internal): gate Arena::len to test builds
tisonkun Oct 3, 2026
bf5d522
style: declare test_support::poll_once with plain pub
tisonkun Oct 3, 2026
a0b741c
docs: remove comments that restate what the code declares
tisonkun Oct 3, 2026
3d7258d
docs(mutex,rwlock): state the invariant behind the Arc extraction
tisonkun Oct 3, 2026
5ccd6fe
refactor(rwlock): drop unneeded pub(super) from write guard fields
tisonkun Oct 3, 2026
dc79f1e
docs(rwlock): explain the default reader cap in one accurate sentence
tisonkun Oct 3, 2026
65e0528
refactor(broadcast): collapse the sender disconnect match into a cond…
tisonkun Oct 3, 2026
ce3f6c5
docs(broadcast): name publish_retained as the backlog growth point
tisonkun Oct 3, 2026
a8d4933
refactor(spmc): remove the unreachable send-waiter handoff in Send::drop
tisonkun Oct 3, 2026
8d37e8e
docs(spmc): describe the single internal mutex accurately
tisonkun Oct 3, 2026
48229b9
docs(once-cell): align OnceCell docs with take() and task terminology
tisonkun Oct 3, 2026
f179369
refactor(once-cell): remove the vestigial permit rebind in set_value
tisonkun Oct 3, 2026
2b86eea
refactor(latch): implement wait_owned by awaiting wait()
tisonkun Oct 3, 2026
07de3a4
style(phaser): backtick .await in the must_use message
tisonkun Oct 3, 2026
75f10d1
refactor(pool): remove the dead default type parameter from UnreadyOb…
tisonkun Oct 3, 2026
5a5ad83
style(pool): narrow ObjectStatus mutation methods to pub(super)
tisonkun Oct 3, 2026
7afb9c7
docs(xtask): describe cargo x test as running workspace tests
tisonkun Oct 3, 2026
49ee5e1
chore(examples): remove the vestigial cargo-release metadata
tisonkun Oct 3, 2026
c2618fa
docs: link the changelogs absolutely from MIGRATE.md
tisonkun Oct 3, 2026
9c1967c
chore: drop redundant default-features from workspace-inherited deps
tisonkun Oct 3, 2026
826ba72
Revert "docs: link the changelogs absolutely from MIGRATE.md"
tisonkun Oct 3, 2026
0b3552f
docs: clarify channel operation terminology in comments
tisonkun Oct 3, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
4 changes: 1 addition & 3 deletions asyncband/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -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"] }
Expand Down
79 changes: 40 additions & 39 deletions asyncband/src/broadcast/mpmc/bounded/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
//!
Expand All @@ -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
//!
Expand Down Expand Up @@ -167,7 +167,6 @@ pub fn bounded<T: Clone>(capacity: usize) -> (BoundedSender<T>, BoundedReceiver<
}

struct Shared<T> {
/// Buffer, receiver cursors, and parked receivers, all under a single lock.
inner: Mutex<Inner<T>>,
/// Number of active senders.
senders: AtomicUsize,
Expand All @@ -181,17 +180,19 @@ struct Shared<T> {
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<T> Shared<T> {
/// 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));
Expand All @@ -200,7 +201,7 @@ impl<T> Shared<T> {

/// 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();
Expand Down Expand Up @@ -245,13 +246,10 @@ impl<T> fmt::Debug for BoundedSender<T> {

impl<T> Drop for BoundedSender<T> {
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);
}
}
}
Expand All @@ -267,13 +265,14 @@ impl<T> BoundedSender<T> {
///
/// 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
///
Expand All @@ -295,8 +294,8 @@ impl<T> BoundedSender<T> {

struct SendState<'a, T> {
sender: &'a BoundedSender<T>,
// 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<Arc<T>>,
Expand Down Expand Up @@ -386,18 +385,20 @@ impl<T> BoundedSender<T> {
/// ```
pub fn try_send(&self, value: T) -> Result<(), TrySendError<T>> {
// `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<T>) -> Result<(), Arc<T>> {
self.publish(msg, |msg| msg)
}
Expand Down Expand Up @@ -433,7 +434,7 @@ impl<T> BoundedSender<T> {
/// 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
///
Expand Down Expand Up @@ -547,9 +548,9 @@ impl<T: Clone> BoundedReceiver<T> {
///
/// # 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
///
Expand Down Expand Up @@ -594,8 +595,8 @@ impl<T: Clone> BoundedReceiver<T> {
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))
}
Expand Down Expand Up @@ -639,7 +640,7 @@ impl<T> BoundedReceiver<T> {
/// 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
///
Expand Down Expand Up @@ -668,7 +669,7 @@ struct Recv<'a, T> {

impl<T> 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;
}
Expand Down Expand Up @@ -701,7 +702,7 @@ impl<T: Clone> 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)))
}
Expand Down
2 changes: 1 addition & 1 deletion asyncband/src/broadcast/mpmc/bounded/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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();
Expand Down
Loading
Loading