Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
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: 4 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -31,6 +31,10 @@ All notable changes to this project will be documented in this file.

* Executor waker operations, including cloning, waking, and dropping, are expected not to panic; recovery from panicking waker callbacks is no longer supported.

### Improvements

* Reduce SPMC waiting and cancellation overhead. Bounded SPMC queues now preallocate storage for the requested capacity when created instead of growing it during sending.

## v0.7.3 (2026-09-29)

### New features
Expand Down
36 changes: 16 additions & 20 deletions asyncband/src/spmc/bounded.rs
Original file line number Diff line number Diff line change
Expand Up @@ -22,28 +22,29 @@ use super::RecvError;
use super::SendError;
use super::TryRecvError;
use super::TrySendError;
use super::queue::Producer;
use super::queue::Shared;
use super::queue::channel;

/// Creates a bounded single-producer, multi-consumer queue.
///
/// The queue stores at most `capacity` values. Sending waits for a receiver to free capacity when
/// the queue is full.
/// The queue stores at most `capacity` values in preallocated storage. Sending waits for a
/// receiver to free capacity when the queue is full.
///
/// 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.
///
/// # Panics
///
/// Panics if `capacity` is zero.
/// Panics if `capacity` is zero or the preallocated message storage exceeds the allocation size
/// limit.
#[track_caller]
pub fn bounded<T>(capacity: usize) -> (BoundedSender<T>, BoundedReceiver<T>) {
assert!(capacity > 0, "spmc bounded queue requires capacity > 0");
let shared = Arc::new(Shared::bounded(capacity));
let (producer, shared) = channel(capacity);
(
BoundedSender {
shared: shared.clone(),
},
BoundedSender { producer, capacity },
BoundedReceiver { shared },
)
}
Expand All @@ -53,7 +54,9 @@ pub fn bounded<T>(capacity: usize) -> (BoundedSender<T>, BoundedReceiver<T>) {
/// Instances are created by [`bounded`] and cannot be cloned. Sending requires exclusive access to
/// this endpoint.
pub struct BoundedSender<T> {
shared: Arc<Shared<T>>,
producer: Producer<T>,
// Only the producer can increase the queue length, so consumers need no capacity counter.
capacity: usize,
}

impl<T> fmt::Debug for BoundedSender<T> {
Expand All @@ -62,38 +65,31 @@ impl<T> fmt::Debug for BoundedSender<T> {
}
}

impl<T> Drop for BoundedSender<T> {
fn drop(&mut self) {
self.shared.drop_sender();
}
}

impl<T> BoundedSender<T> {
/// Sends a value, waiting until capacity is available if the queue is full.
///
/// If all receivers have been dropped, the value is returned in [`SendError`].
///
/// # Cancel safety
///
/// 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
/// Dropping a pending `send` future drops `value` without enqueueing it and releases the
/// exclusive sender borrow. Any available capacity remains usable by the next operation. 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<T>> {
self.shared.send(value).await
self.producer.send(value, self.capacity).await
}

/// Attempts to send a value without waiting for capacity.
///
/// Returns [`TrySendError::Full`] when the queue has reached its exact capacity and
/// [`TrySendError::Disconnected`] when all receivers have been dropped.
pub fn try_send(&mut self, value: T) -> Result<(), TrySendError<T>> {
self.shared.try_send(value)
self.producer.try_send(value, self.capacity)
}
}

/// Receives values from the associated [`BoundedSender`] handles.
/// Receives values from the associated [`BoundedSender`].
///
/// Cloned receivers compete for values, and every accepted value is returned by exactly one
/// receiver while a receiver remains. Dropping the final receiver releases buffered values.
Expand Down
Loading
Loading