From 64b2cafd919d1a87775b5a942ada7574e2195945 Mon Sep 17 00:00:00 2001 From: tison Date: Sun, 4 Oct 2026 00:40:11 +0800 Subject: [PATCH] feat(cancellation): explore independent cancellation signals --- CHANGELOG.md | 1 + README.md | 1 + asyncband/Cargo.toml | 1 + asyncband/src/cancellation/mod.rs | 138 +++++++++++++++ asyncband/src/lib.rs | 3 + examples/Cargo.toml | 5 + examples/src/cancel_lookups.rs | 65 +++++++ tests-integration/Cargo.toml | 1 + tests-integration/tests/cancellation_test.rs | 173 +++++++++++++++++++ 9 files changed, 388 insertions(+) create mode 100644 asyncband/src/cancellation/mod.rs create mode 100644 examples/src/cancel_lookups.rs create mode 100644 tests-integration/tests/cancellation_test.rs diff --git a/CHANGELOG.md b/CHANGELOG.md index a59b9983..057a333b 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -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 diff --git a/README.md b/README.md index 17ee4426..b143ece2 100644 --- a/README.md +++ b/README.md @@ -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. | diff --git a/asyncband/Cargo.toml b/asyncband/Cargo.toml index 0e6f3e38..9dc6d347 100644 --- a/asyncband/Cargo.toml +++ b/asyncband/Cargo.toml @@ -47,6 +47,7 @@ default = [] barrier = [] blocking = [] broadcast = [] +cancellation = ["latch"] completion = [] condvar = ["mutex"] event = [] diff --git a/asyncband/src/cancellation/mod.rs b/asyncband/src/cancellation/mod.rs new file mode 100644 index 00000000..20c9f3a2 --- /dev/null +++ b/asyncband/src/cancellation/mod.rs @@ -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, +} + +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, +} + +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 + 'static { + self.signal.clone().wait_owned() + } +} diff --git a/asyncband/src/lib.rs b/asyncband/src/lib.rs index f2003d85..7055fe14 100644 --- a/asyncband/src/lib.rs +++ b/asyncband/src/lib.rs @@ -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. | @@ -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")] diff --git a/examples/Cargo.toml b/examples/Cargo.toml index a76d3e35..b3e7fd9c 100644 --- a/examples/Cargo.toml +++ b/examples/Cargo.toml @@ -24,6 +24,7 @@ rust-version.workspace = true [dependencies] asyncband = { workspace = true, features = [ + "cancellation", "completion", "event", "lazy-cell", @@ -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" diff --git a/examples/src/cancel_lookups.rs b/examples/src/cancel_lookups.rs new file mode 100644 index 00000000..07e7dd78 --- /dev/null +++ b/examples/src/cancel_lookups.rs @@ -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"); +} diff --git a/tests-integration/Cargo.toml b/tests-integration/Cargo.toml index e93720d7..e2167246 100644 --- a/tests-integration/Cargo.toml +++ b/tests-integration/Cargo.toml @@ -30,6 +30,7 @@ asyncband = { workspace = true, features = [ "barrier", "blocking", "broadcast", + "cancellation", "completion", "condvar", "event", diff --git a/tests-integration/tests/cancellation_test.rs b/tests-integration/tests/cancellation_test.rs new file mode 100644 index 00000000..90b1ae62 --- /dev/null +++ b/tests-integration/tests/cancellation_test.rs @@ -0,0 +1,173 @@ +// 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. + +use std::pin::pin; +use std::sync::Arc; +use std::sync::Barrier; + +use asyncband::cancellation::CancellationSource; +use tests_integration::WakeCounter; +use tests_integration::poll_once; +use tests_integration::poll_with; + +#[test] +fn request_reaches_current_and_late_waits() { + let source = CancellationSource::new(); + let token = source.token(); + let observer = token.clone(); + let (first_waker, first_wakes) = WakeCounter::new(); + let (second_waker, second_wakes) = WakeCounter::new(); + let mut first = pin!(token.cancelled()); + let mut second = pin!(observer.cancelled()); + assert!(!token.is_cancelled()); + assert!(poll_with(first.as_mut(), &first_waker).is_pending()); + assert!(poll_with(second.as_mut(), &second_waker).is_pending()); + + source.cancel(); + source.cancel(); + let late = source.token(); + drop(source); + + assert_eq!(first_wakes.count(), 1); + assert_eq!(second_wakes.count(), 1); + assert!(poll_once(first.as_mut()).is_ready()); + assert!(poll_once(second.as_mut()).is_ready()); + assert!(poll_once(pin!(late.cancelled())).is_ready()); + assert!(poll_once(pin!(token.cancelled())).is_ready()); + assert!(late.is_cancelled()); +} + +#[test] +fn dropping_source_neither_cancels_nor_wakes() { + let source = CancellationSource::new(); + let token = source.token(); + let (waker, wakes) = WakeCounter::new(); + let baseline = Arc::strong_count(&wakes); + let mut wait = Box::pin(token.cancelled_owned()); + assert!(poll_with(wait.as_mut(), &waker).is_pending()); + + drop(source); + + assert!(!token.is_cancelled()); + assert_eq!(wakes.count(), 0); + assert!(poll_with(wait.as_mut(), &waker).is_pending()); + assert!(poll_once(pin!(token.clone().cancelled())).is_pending()); + assert_eq!(Arc::strong_count(&wakes), baseline + 1); + drop(wait); + assert_eq!(Arc::strong_count(&wakes), baseline); +} + +#[test] +fn source_can_issue_tokens_after_all_observers_are_dropped() { + let source = CancellationSource::new(); + drop(source.token()); + assert!(!source.token().is_cancelled()); + + source.cancel(); + + let token = source.token(); + assert!(token.is_cancelled()); + assert!(poll_once(pin!(token.cancelled())).is_ready()); +} + +#[test] +fn dropping_wait_only_removes_its_registration() { + let source = CancellationSource::new(); + let token = source.token(); + let (retired_waker, retired_wakes) = WakeCounter::new(); + let (remaining_waker, remaining_wakes) = WakeCounter::new(); + let (retry_waker, retry_wakes) = WakeCounter::new(); + let baseline = Arc::strong_count(&retired_wakes); + let mut retired = Box::pin(token.cancelled()); + let mut remaining = pin!(token.cancelled()); + assert!(poll_with(retired.as_mut(), &retired_waker).is_pending()); + assert!(poll_with(remaining.as_mut(), &remaining_waker).is_pending()); + + drop(retired); + assert_eq!(Arc::strong_count(&retired_wakes), baseline); + assert!(!token.is_cancelled()); + let mut retry = pin!(token.cancelled()); + assert!(poll_with(retry.as_mut(), &retry_waker).is_pending()); + source.cancel(); + + assert_eq!(retired_wakes.count(), 0); + assert_eq!(remaining_wakes.count(), 1); + assert_eq!(retry_wakes.count(), 1); + assert!(poll_once(remaining.as_mut()).is_ready()); + assert!(poll_once(retry.as_mut()).is_ready()); +} + +#[test] +fn owned_wait_outlives_its_token() { + let source = CancellationSource::new(); + let token = source.token(); + let mut wait = pin!(token.cancelled_owned()); + drop(token); + assert!(poll_once(wait.as_mut()).is_pending()); + + source.cancel(); + + assert!(poll_once(wait.as_mut()).is_ready()); +} + +#[test] +fn registration_racing_with_concurrent_requests_does_not_lose_wake() { + for _ in 0..64 { + let source = CancellationSource::new(); + let token = source.token(); + let start = Barrier::new(3); + let (waker, wakes) = WakeCounter::new(); + let mut wait = pin!(token.cancelled()); + + let first_poll = std::thread::scope(|scope| { + for _ in 0..2 { + scope.spawn(|| { + start.wait(); + source.cancel(); + }); + } + start.wait(); + poll_with(wait.as_mut(), &waker) + }); + + if first_poll.is_pending() { + // A subsequent ready poll alone would not detect a missed executor notification. + assert!(wakes.count() > 0); + assert!(poll_once(wait.as_mut()).is_ready()); + } + assert!(token.is_cancelled()); + } +} + +#[tokio::test] +async fn cancellation_does_not_finish_a_task() { + let source = CancellationSource::new(); + let wait = source.token().cancelled_owned(); + let (observed_tx, observed_rx) = tokio::sync::oneshot::channel(); + let (cleanup_tx, cleanup_rx) = tokio::sync::oneshot::channel(); + let worker = tokio::spawn(async move { + wait.await; + observed_tx.send(()).unwrap(); + cleanup_rx.await.unwrap(); + }); + + source.cancel(); + observed_rx.await.unwrap(); + assert!(!worker.is_finished()); + cleanup_tx.send(()).unwrap(); + worker.await.unwrap(); +}