cancel.rs (3434B)
1 //! Composable cooperative cancellation without process-global signal handling. 2 3 use core::fmt; 4 5 /// Cloneable cooperative cancellation shared by supervisor-owned tasks. 6 /// 7 /// Child cancellation propagates to descendants but never back to its parent 8 /// or sideways to siblings. This type does not install or interpret operating 9 /// system signals. 10 #[derive(Clone, Default)] 11 pub struct CancellationToken { 12 inner: tokio_util::sync::CancellationToken, 13 } 14 15 impl CancellationToken { 16 #[must_use] 17 pub fn new() -> Self { 18 Self::default() 19 } 20 21 /// Creates a child that is cancelled with this token while retaining its own authority. 22 #[must_use] 23 pub fn child_token(&self) -> Self { 24 Self { 25 inner: self.inner.child_token(), 26 } 27 } 28 29 /// Requests cancellation. Repeated requests have no additional effect. 30 pub fn cancel(&self) { 31 self.inner.cancel(); 32 } 33 34 #[must_use] 35 pub fn is_cancelled(&self) -> bool { 36 self.inner.is_cancelled() 37 } 38 39 /// Completes after cancellation and remains immediately ready thereafter. 40 pub async fn cancelled(&self) { 41 self.inner.cancelled().await; 42 } 43 } 44 45 impl fmt::Debug for CancellationToken { 46 fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { 47 formatter 48 .debug_struct("CancellationToken") 49 .field("cancelled", &self.is_cancelled()) 50 .finish() 51 } 52 } 53 54 #[cfg(test)] 55 mod tests { 56 use super::*; 57 58 #[tokio::test] 59 async fn parent_child_propagation_is_directional() { 60 let parent = CancellationToken::new(); 61 let child = parent.child_token(); 62 let grandchild = child.child_token(); 63 let sibling = parent.child_token(); 64 65 child.cancel(); 66 child.cancelled().await; 67 grandchild.cancelled().await; 68 assert!(!parent.is_cancelled()); 69 assert!(!sibling.is_cancelled()); 70 71 parent.cancel(); 72 parent.cancelled().await; 73 sibling.cancelled().await; 74 assert!(parent.child_token().is_cancelled()); 75 } 76 77 #[tokio::test] 78 async fn cancellation_is_idempotent_and_observable_after_the_fact() { 79 let token = CancellationToken::new(); 80 token.cancel(); 81 token.cancel(); 82 assert!(token.is_cancelled()); 83 token.cancelled().await; 84 assert_eq!( 85 format!("{token:?}"), 86 "CancellationToken { cancelled: true }" 87 ); 88 } 89 90 #[tokio::test] 91 async fn dropping_a_waiter_does_not_consume_cancellation() { 92 let token = CancellationToken::new(); 93 let abandoned = token.cancelled(); 94 drop(abandoned); 95 96 token.cancel(); 97 token.cancelled().await; 98 assert!(token.is_cancelled()); 99 } 100 101 #[tokio::test] 102 async fn concurrent_waiter_registration_and_cancel_has_no_lost_wakeup() { 103 for round in 0..64 { 104 let token = CancellationToken::new(); 105 let waiter_token = token.clone(); 106 let waiter = tokio::spawn(async move { 107 if round % 2 == 0 { 108 tokio::task::yield_now().await; 109 } 110 waiter_token.cancelled().await; 111 waiter_token.is_cancelled() 112 }); 113 if round % 2 == 1 { 114 tokio::task::yield_now().await; 115 } 116 token.cancel(); 117 assert!(waiter.await.unwrap()); 118 } 119 } 120 }