use std::sync::atomic::{AtomicBool, AtomicUsize, Ordering};
use std::sync::Arc;
#[derive(Debug)]
struct Inner {
cancelled: AtomicBool,
polls: AtomicUsize,
budget: usize,
}
#[derive(Clone, Debug, Default)]
pub struct CancelToken {
inner: Option<Arc<Inner>>,
}
impl CancelToken {
#[must_use]
pub fn new() -> Self {
Self::with_budget(usize::MAX)
}
#[must_use]
pub fn never() -> Self {
Self { inner: None }
}
#[must_use]
pub fn with_budget(budget: usize) -> Self {
Self {
inner: Some(Arc::new(Inner {
cancelled: AtomicBool::new(false),
polls: AtomicUsize::new(0),
budget,
})),
}
}
pub fn cancel(&self) {
if let Some(inner) = &self.inner {
inner.cancelled.store(true, Ordering::Release);
}
}
#[must_use]
pub fn is_cancelled(&self) -> bool {
match &self.inner {
None => false,
Some(inner) => {
let seen = inner.polls.fetch_add(1, Ordering::Relaxed);
inner.cancelled.load(Ordering::Acquire) || seen >= inner.budget
},
}
}
#[must_use]
pub fn peek_cancelled(&self) -> bool {
match &self.inner {
None => false,
Some(inner) => inner.cancelled.load(Ordering::Acquire),
}
}
#[must_use]
pub fn polls(&self) -> usize {
match &self.inner {
None => 0,
Some(inner) => inner.polls.load(Ordering::Relaxed),
}
}
#[must_use]
pub fn cancel_on_drop(&self) -> CancelOnDrop {
CancelOnDrop {
token: self.clone(),
armed: true,
}
}
}
#[derive(Debug)]
pub struct CancelOnDrop {
token: CancelToken,
armed: bool,
}
impl CancelOnDrop {
#[must_use]
pub fn token(&self) -> &CancelToken {
&self.token
}
pub fn disarm(&mut self) {
self.armed = false;
}
}
impl Drop for CancelOnDrop {
fn drop(&mut self) {
if self.armed {
self.token.cancel();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn never_token_is_free_and_never_cancels() {
let t = CancelToken::never();
t.cancel();
assert!(!t.is_cancelled(), "never() must not be cancellable");
assert!(!t.peek_cancelled());
assert_eq!(t.polls(), 0, "never() must hold no state");
}
#[test]
fn default_is_never() {
let t = CancelToken::default();
t.cancel();
assert!(
!t.is_cancelled(),
"Default must be the zero-cost never() token so adding a cancel field \
changes no existing caller's behaviour"
);
}
#[test]
fn cancel_is_observed_by_clones() {
let t = CancelToken::new();
let clone = t.clone();
assert!(!clone.peek_cancelled());
t.cancel();
assert!(
clone.peek_cancelled(),
"clones must share state or a spawn_blocking worker cannot see the guard fire"
);
}
#[test]
fn polls_counts_only_is_cancelled() {
let t = CancelToken::new();
assert_eq!(t.polls(), 0);
let _ = t.is_cancelled();
let _ = t.is_cancelled();
assert_eq!(t.polls(), 2);
let _ = t.peek_cancelled();
assert_eq!(t.polls(), 2, "peek must not consume budget");
}
#[test]
fn budget_trips_after_exactly_n_polls() {
let t = CancelToken::with_budget(3);
assert!(!t.is_cancelled(), "poll 0 within budget");
assert!(!t.is_cancelled(), "poll 1 within budget");
assert!(!t.is_cancelled(), "poll 2 within budget");
assert!(t.is_cancelled(), "poll 3 exhausts a budget of 3");
assert!(t.is_cancelled(), "and stays cancelled");
}
#[test]
fn zero_budget_trips_immediately() {
let t = CancelToken::with_budget(0);
assert!(t.is_cancelled(), "a budget of 0 permits no work at all");
}
#[test]
fn new_token_has_no_budget() {
let t = CancelToken::new();
for i in 0..10_000 {
assert!(
!t.is_cancelled(),
"poll {i} must not trip an unbudgeted token"
);
}
}
#[test]
fn guard_cancels_token_on_drop() {
let t = CancelToken::new();
{
let _guard = t.cancel_on_drop();
assert!(
!t.peek_cancelled(),
"the guard must not cancel while it is still alive"
);
}
assert!(
t.peek_cancelled(),
"dropping the guard is what turns an axum client-disconnect into a stop signal"
);
}
#[test]
fn guard_exposes_the_token_it_will_cancel() {
let t = CancelToken::new();
let guard = t.cancel_on_drop();
assert!(!guard.token().peek_cancelled());
drop(guard);
assert!(t.peek_cancelled());
}
#[test]
fn disarmed_guard_does_not_cancel_on_drop() {
let t = CancelToken::new();
{
let mut guard = t.cancel_on_drop();
guard.disarm();
}
assert!(
!t.peek_cancelled(),
"a disarmed guard must leave the token alive: work handed to a \
background task outlives the future that armed the guard"
);
}
#[test]
fn guard_is_armed_until_disarmed() {
let t = CancelToken::new();
let mut guard = t.cancel_on_drop();
drop(t.cancel_on_drop()); assert!(t.peek_cancelled());
guard.disarm();
}
}