lc_core/runnables/
cancellation.rs1use std::sync::atomic::{AtomicBool, Ordering};
26use std::sync::Arc;
27
28#[derive(Debug, Clone)]
33pub struct CancellationToken {
34 inner: Arc<AtomicBool>,
35}
36
37impl CancellationToken {
38 pub fn new() -> Self {
40 Self {
41 inner: Arc::new(AtomicBool::new(false)),
42 }
43 }
44
45 pub fn cancel(&self) {
49 self.inner.store(true, Ordering::SeqCst);
50 }
51
52 pub fn is_cancelled(&self) -> bool {
54 self.inner.load(Ordering::SeqCst)
55 }
56
57 pub async fn cancelled(&self) {
62 while !self.inner.load(Ordering::SeqCst) {
65 tokio::task::yield_now().await;
66 }
67 }
68}
69
70impl Default for CancellationToken {
71 fn default() -> Self {
72 Self::new()
73 }
74}
75
76#[cfg(test)]
77mod tests {
78 use super::*;
79
80 #[test]
81 fn test_new_token_is_not_cancelled() {
82 let token = CancellationToken::new();
83 assert!(!token.is_cancelled());
84 }
85
86 #[test]
87 fn test_cancel_sets_is_cancelled() {
88 let token = CancellationToken::new();
89 token.cancel();
90 assert!(token.is_cancelled());
91 }
92
93 #[test]
94 fn test_clone_shares_cancellation_state() {
95 let token = CancellationToken::new();
96 let clone = token.clone();
97
98 assert!(!token.is_cancelled());
99 assert!(!clone.is_cancelled());
100
101 clone.cancel();
102
103 assert!(token.is_cancelled());
104 assert!(clone.is_cancelled());
105 }
106
107 #[test]
108 fn test_multiple_clones_all_cancelled() {
109 let token = CancellationToken::new();
110 let c1 = token.clone();
111 let c2 = token.clone();
112 let c3 = token.clone();
113
114 token.cancel();
115
116 assert!(c1.is_cancelled());
117 assert!(c2.is_cancelled());
118 assert!(c3.is_cancelled());
119 }
120
121 #[tokio::test]
122 async fn test_cancelled_future_resolves() {
123 let token = CancellationToken::new();
124
125 let cloned = token.clone();
127 tokio::spawn(async move {
128 cloned.cancel();
129 });
130
131 token.cancelled().await;
133 assert!(token.is_cancelled());
134 }
135}