Skip to main content

fusor_async/
cancellation.rs

1use 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/// Cancellation of a single operation. Connect transport cancellation here; stopping
16/// polling alone does not stop browser Fetch or undo a server operation.
17#[derive(Clone, Default)]
18pub struct RequestContext(Rc<Cancellation>);
19
20/// Owner-side cancellation capability, also usable by explicit write adapters.
21/// Cancellation stops local work; it cannot establish whether a server committed.
22/// Dropping this source cancels its context. Receiving contexts cannot cancel it.
23pub 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    /// The operation completed. Release the source without signalling cancellation.
39    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/// Unregisters its callback on drop. Keep it alive while the operation is pending.
50#[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}