fusor_async/
cancellation.rs1use std::{
2 cell::{Cell, RefCell},
3 collections::BTreeMap,
4 rc::{Rc, Weak},
5};
6
7type Callback = Box<dyn FnOnce()>;
8#[derive(Default)]
9struct Cancellation {
10 cancelled: Cell<bool>,
11 next: Cell<u64>,
12 callbacks: RefCell<BTreeMap<u64, Callback>>,
13}
14
15#[derive(Clone, Default)]
18pub struct RequestContext(Rc<Cancellation>);
19
20pub struct CancellationSource(Option<RequestContext>);
24impl Default for CancellationSource {
25 fn default() -> Self {
26 Self(Some(RequestContext::default()))
27 }
28}
29impl CancellationSource {
30 pub fn context(&self) -> RequestContext {
31 self.0.as_ref().expect("live cancellation source").clone()
32 }
33 pub fn cancel(&self) {
34 if let Some(context) = &self.0 {
35 context.cancel();
36 }
37 }
38 pub fn complete(mut self) {
40 self.0.take();
41 }
42}
43impl Drop for CancellationSource {
44 fn drop(&mut self) {
45 self.cancel();
46 }
47}
48
49#[must_use = "retain the cancellation registration while the operation is pending"]
51pub struct CancelRegistration {
52 context: Weak<Cancellation>,
53 id: u64,
54}
55impl Drop for CancelRegistration {
56 fn drop(&mut self) {
57 if let Some(context) = self.context.upgrade() {
58 let callback = context.callbacks.borrow_mut().remove(&self.id);
59 drop(callback);
60 }
61 }
62}
63impl RequestContext {
64 pub fn is_cancelled(&self) -> bool {
65 self.0.cancelled.get()
66 }
67 pub fn on_cancel(&self, callback: impl FnOnce() + 'static) -> CancelRegistration {
68 let id = self
69 .0
70 .next
71 .get()
72 .checked_add(1)
73 .expect("cancellation registration overflow");
74 self.0.next.set(id);
75 if self.is_cancelled() {
76 callback();
77 } else {
78 self.0.callbacks.borrow_mut().insert(id, Box::new(callback));
79 }
80 CancelRegistration {
81 context: Rc::downgrade(&self.0),
82 id,
83 }
84 }
85 pub(crate) fn cancel(&self) {
86 if self.0.cancelled.replace(true) {
87 return;
88 }
89 let callbacks = self.0.callbacks.take();
90 for callback in callbacks.into_values() {
91 callback();
92 }
93 }
94}