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
1 change: 1 addition & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ All notable changes to this project will be documented in this file.

### New features

* Add opt-in cooperative cancellation with a `CancellationSource` that explicitly requests cancellation and cloneable, read-only `CancellationToken` observers; borrowed and owned waits observe a persistent signal without tracking task completion, and dropping the source does not request cancellation.
* Add bounded MPMC `reserve` and `try_reserve` methods returning a borrowed `Permit`, so callers can wait for capacity before constructing a value; sends and reservations receive capacity in wait-queue order, and unused permits release it.

### Notable changes
Expand Down
1 change: 1 addition & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,7 @@ Runnable examples live in the [`examples`](examples) workspace crate. They demon
| | [`Latch`](https://docs.rs/asyncband/*/asyncband/latch/struct.Latch.html) | `latch` | Wait until a fixed one-way countdown reaches zero. |
| | [`Phaser`](https://docs.rs/asyncband/*/asyncband/phaser/struct.Phaser.html) | `phaser` | Coordinate repeated phases with a dynamic participant set. |
| | [`WaitGroup`](https://docs.rs/asyncband/*/asyncband/waitgroup/struct.WaitGroup.html) | `waitgroup` | Dynamically register participants and wait until all have completed. |
| | [`cancellation`](https://docs.rs/asyncband/*/asyncband/cancellation/) | `cancellation` | Request cooperative cancellation without tracking task completion. |
| | [`Shutdown`](https://docs.rs/asyncband/*/asyncband/shutdown/struct.Shutdown.html) | `shutdown` | Request shutdown and wait until all completion guards are dropped. |
| Work coalescing | [`Once`](https://docs.rs/asyncband/*/asyncband/once/struct.Once.html) | `once` | Complete one asynchronous initialization; cancelled or panicked attempts may be retried. |
| | [`OnceCell`](https://docs.rs/asyncband/*/asyncband/once/struct.OnceCell.html) | `once-cell` | Store one value from an access-time initializer; failed, cancelled, or panicked attempts may be retried. |
Expand Down
1 change: 1 addition & 0 deletions asyncband/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -47,6 +47,7 @@ default = []
barrier = []
blocking = []
broadcast = []
cancellation = ["latch"]
completion = []
condvar = ["mutex"]
event = []
Expand Down
138 changes: 138 additions & 0 deletions asyncband/src/cancellation/mod.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,138 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

//! Request cooperative cancellation without tracking task completion.
//!
//! A [`CancellationSource`] controls one signal. Its cloneable [`CancellationToken`] observers
//! can query or wait for that signal, but cannot request cancellation themselves. The request is
//! sticky: every current and future wait completes once cancellation has been requested.
//!
//! Cancellation is advisory. Observing it does not mean that a task has stopped or finished its
//! cleanup. Callers choose whether and when to stop work, and retain responsibility for joining
//! spawned tasks. Dropping a wait removes that wait's registration; it does not cancel other work.
//!
//! # Source lifetime
//!
//! Dropping the source does not request cancellation or wake waiters. If it was dropped without
//! calling [`cancel`](CancellationSource::cancel), tokens remain uncancelled and their waits stay
//! pending indefinitely. Use the wait alongside work that may finish normally, or retain the
//! source and explicitly cancel it when a stop request is required. A request already issued
//! remains observable after the source is dropped.
//!
//! # Examples
//!
//! ```
//! use asyncband::cancellation::CancellationSource;
//!
//! # #[tokio::main]
//! # async fn main() {
//! let source = CancellationSource::new();
//! let token = source.token();
//! let observer = token.clone();
//!
//! source.cancel();
//! tokio::join!(token.cancelled(), observer.cancelled());
//! assert!(token.is_cancelled());
//! # }
//! ```

use std::future::Future;
use std::sync::Arc;

use crate::latch::Latch;

/// The authority to request cancellation for a set of [`CancellationToken`] observers.
///
/// This type is not cloneable. Share a reference, or explicitly put the source in an [`Arc`],
/// when multiple controllers need cancellation authority. Give workers tokens instead.
///
/// Dropping the source does not request cancellation. See the [module documentation](self) for
/// the behavior of tokens that outlive it.
#[derive(Debug)]
pub struct CancellationSource {
signal: Arc<Latch>,
}

impl Default for CancellationSource {
fn default() -> Self {
Self::new()
}
}

impl CancellationSource {
/// Creates a source with no cancellation request.
pub fn new() -> Self {
Self {
signal: Arc::new(Latch::new(1)),
}
}

/// Returns a read-only observer of this source's signal.
///
/// A token created after cancellation observes the same request immediately.
pub fn token(&self) -> CancellationToken {
CancellationToken {
signal: self.signal.clone(),
}
}

/// Requests cancellation and wakes registered waits.
///
/// The request is persistent and this method is idempotent, including concurrent calls.
/// It does not wait for tasks to observe the request, finish work, or complete cleanup.
pub fn cancel(&self) {
self.signal.count_down();
}
}

/// A cloneable, read-only observer of a cooperative cancellation request.
///
/// Tokens do not keep any task running and do not delay its completion. Every wait registers
/// independently, including multiple waits using the same token. Cloning or dropping a token
/// never changes the cancellation state.
#[derive(Debug, Clone)]
pub struct CancellationToken {
signal: Arc<Latch>,
}

impl CancellationToken {
/// Returns whether cancellation has been requested.
///
/// Once true, this remains true. A false result is only a snapshot. Dropping an uncancelled
/// source leaves the result false; source destruction is not reported as cancellation.
pub fn is_cancelled(&self) -> bool {
self.signal.try_wait().is_ok()
}

/// Waits until cancellation is requested without consuming the signal.
///
/// This method is cancel safe: dropping a pending wait unregisters only that wait. Later
/// waits still observe any request, including one made before they are first polled.
///
/// If the source is dropped before cancellation, this wait remains pending indefinitely.
pub async fn cancelled(&self) {
self.signal.wait().await;
}

/// Returns a cancellation wait that can outlive this token and move into a spawned task.
///
/// The future owns an observer of the signal. Like [`cancelled`](Self::cancelled), it is
/// cancel safe and remains pending if the source is dropped without requesting cancellation.
pub fn cancelled_owned(&self) -> impl Future<Output = ()> + 'static {
self.signal.clone().wait_owned()
}
}
3 changes: 3 additions & 0 deletions asyncband/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -66,6 +66,7 @@
//! | | [`Latch`](latch::Latch) | `latch` | Wait until a fixed one-way countdown reaches zero. |
//! | | [`Phaser`](phaser::Phaser) | `phaser` | Coordinate repeated phases with a dynamic participant set. |
//! | | [`WaitGroup`](waitgroup::WaitGroup) | `waitgroup` | Dynamically register participants and wait until all have completed. |
//! | | [`cancellation`] | `cancellation` | Request cooperative cancellation without tracking task completion. |
//! | | [`Shutdown`](shutdown::Shutdown) | `shutdown` | Request shutdown and wait until all completion guards are dropped. |
//! | Work coalescing | [`Once`](once::Once) | `once` | Complete one asynchronous initialization; cancelled or panicked attempts may be retried. |
//! | | [`OnceCell`](once::OnceCell) | `once-cell` | Store one value from an access-time initializer; failed, cancelled, or panicked attempts may be retried. |
Expand Down Expand Up @@ -131,6 +132,8 @@ pub mod barrier;
pub mod blocking;
#[cfg(feature = "broadcast")]
pub mod broadcast;
#[cfg(feature = "cancellation")]
pub mod cancellation;
#[cfg(feature = "completion")]
pub mod completion;
#[cfg(feature = "condvar")]
Expand Down
5 changes: 5 additions & 0 deletions examples/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -24,6 +24,7 @@ rust-version.workspace = true

[dependencies]
asyncband = { workspace = true, features = [
"cancellation",
"completion",
"event",
"lazy-cell",
Expand All @@ -43,6 +44,10 @@ tokio = { workspace = true, features = [
[lints]
workspace = true

[[example]]
name = "cancel_lookups"
path = "src/cancel_lookups.rs"

[[example]]
name = "coalesced_worker"
path = "src/coalesced_worker.rs"
Expand Down
65 changes: 65 additions & 0 deletions examples/src/cancel_lookups.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,65 @@
// Licensed to the Apache Software Foundation (ASF) under one
// or more contributor license agreements. See the NOTICE file
// distributed with this work for additional information
// regarding copyright ownership. The ASF licenses this file
// to you under the Apache License, Version 2.0 (the
// "License"); you may not use this file except in compliance
// with the License. You may obtain a copy of the License at
//
// http://www.apache.org/licenses/LICENSE-2.0
//
// Unless required by applicable law or agreed to in writing,
// software distributed under the License is distributed on an
// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
// KIND, either express or implied. See the License for the
// specific language governing permissions and limitations
// under the License.

//! Cancel redundant replica lookups after the first answer arrives.
//!
//! Cancellation belongs to this request, not the service's lifetime. Workers receive only an
//! observer; the caller owns cancellation authority and task handles. The signal does not wait
//! for workers to exit, so the caller joins them separately. A lookup that finishes before
//! observing cancellation may still return an answer.

use std::time::Duration;

use asyncband::cancellation::CancellationSource;
use asyncband::cancellation::CancellationToken;
use tokio::task::JoinSet;

async fn lookup(
replica: &'static str,
latency: Duration,
token: CancellationToken,
) -> Option<&'static str> {
// The timer stands in for caller-provided I/O; cancellation itself has no runtime dependency.
tokio::select! {
_ = token.cancelled() => None,
_ = tokio::time::sleep(latency) => Some(replica),
}
}

#[tokio::main(flavor = "current_thread")]
async fn main() {
let source = CancellationSource::new();
let mut lookups = JoinSet::new();
for (replica, millis) in [("nearby", 5), ("regional", 50), ("remote", 100)] {
lookups.spawn(lookup(
replica,
Duration::from_millis(millis),
source.token(),
));
}

let winner = lookups.join_next().await.unwrap().unwrap().unwrap();
source.cancel();

let mut cancelled = 0;
while let Some(result) = lookups.join_next().await {
if result.unwrap().is_none() {
cancelled += 1;
}
}
println!("Answer from {winner}; {cancelled} other lookups observed cancellation");
}
1 change: 1 addition & 0 deletions tests-integration/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,7 @@ asyncband = { workspace = true, features = [
"barrier",
"blocking",
"broadcast",
"cancellation",
"completion",
"condvar",
"event",
Expand Down
Loading
Loading