Skip to main content

lora_executor/
cancel.rs

1//! Cooperative query cancellation.
2//!
3//! Queries already carry an optional deadline that the executor checks at
4//! operator boundaries and inside hot loops. Cancellation rides on that
5//! mechanism instead of threading a second token through every operator:
6//! a cancellable query is given a *unique* deadline instant, and cancelling
7//! it records that instant in a small process-wide set. The deadline check
8//! consults the set, so a cancelled query stops at its next check point on
9//! whichever thread is running it (including parallel workers), and then
10//! unwinds through the normal timeout path: the error propagates, WAL
11//! transactions abort, and locks are released.
12//!
13//! The set is only consulted while at least one cancellation is pending,
14//! so queries that never use cancellation pay one relaxed atomic load.
15
16use std::collections::HashSet;
17use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
18use std::sync::{Mutex, OnceLock};
19use std::time::Duration;
20
21use web_time::Instant;
22
23/// How far out an "untimed" cancellable deadline sits. Far enough that
24/// it never fires as a timeout, near enough that `Instant` arithmetic
25/// cannot overflow on any platform.
26const NO_TIMEOUT: Duration = Duration::from_secs(60 * 60 * 24 * 365);
27
28static PENDING: AtomicUsize = AtomicUsize::new(0);
29static SEQ: AtomicU64 = AtomicU64::new(0);
30
31fn cancelled() -> &'static Mutex<HashSet<Instant>> {
32    static SET: OnceLock<Mutex<HashSet<Instant>>> = OnceLock::new();
33    SET.get_or_init(|| Mutex::new(HashSet::new()))
34}
35
36/// A deadline that can also be cancelled before it expires.
37///
38/// Pass [`Self::deadline`] wherever a query deadline is accepted. Call
39/// [`Self::cancel`] from any thread to stop the query at its next check
40/// point; it then fails exactly as if its deadline had passed. Dropping
41/// the handle forgets the cancellation.
42#[derive(Debug)]
43pub struct CancellableDeadline {
44    deadline: Instant,
45    registered: bool,
46}
47
48impl CancellableDeadline {
49    /// A deadline `timeout` from now (or effectively never when `None`)
50    /// that can be cancelled early.
51    pub fn new(timeout: Option<Duration>) -> Self {
52        let base = Instant::now()
53            .checked_add(timeout.unwrap_or(NO_TIMEOUT))
54            .unwrap_or_else(Instant::now);
55        // Offset by a per-handle count of nanoseconds so two queries that
56        // start in the same instant never share a deadline value.
57        let nudge = Duration::from_nanos(SEQ.fetch_add(1, Ordering::Relaxed) % 1_000_000);
58        Self {
59            deadline: base.checked_add(nudge).unwrap_or(base),
60            registered: false,
61        }
62    }
63
64    pub fn deadline(&self) -> Instant {
65        self.deadline
66    }
67
68    /// Request cancellation. Idempotent.
69    pub fn cancel(&mut self) {
70        if self.registered {
71            return;
72        }
73        let mut set = cancelled().lock().unwrap_or_else(|p| p.into_inner());
74        if set.insert(self.deadline) {
75            PENDING.fetch_add(1, Ordering::Release);
76        }
77        self.registered = true;
78    }
79
80    pub fn is_cancelled(&self) -> bool {
81        self.registered
82    }
83}
84
85impl Drop for CancellableDeadline {
86    fn drop(&mut self) {
87        if !self.registered {
88            return;
89        }
90        let mut set = cancelled().lock().unwrap_or_else(|p| p.into_inner());
91        if set.remove(&self.deadline) {
92            PENDING.fetch_sub(1, Ordering::Release);
93        }
94    }
95}
96
97/// Whether the query owning `deadline` has been cancelled.
98#[inline]
99pub(crate) fn is_cancelled(deadline: Instant) -> bool {
100    if PENDING.load(Ordering::Acquire) == 0 {
101        return false;
102    }
103    cancelled()
104        .lock()
105        .unwrap_or_else(|p| p.into_inner())
106        .contains(&deadline)
107}
108
109thread_local! {
110    static ACTIVE: std::cell::Cell<Option<Instant>> = const { std::cell::Cell::new(None) };
111}
112
113/// Makes a query's deadline visible to everything that runs on this
114/// thread while the query executes, including pull pipelines that
115/// operators build lazily mid-query (`OPTIONAL MATCH`, `CALL {}`) and
116/// buffered sub-executors. Executors enter one at each public entry
117/// point; nesting keeps the outer deadline when the inner has none, and
118/// dropping the scope restores whatever was active before.
119pub(crate) struct DeadlineScope {
120    prev: Option<Instant>,
121}
122
123impl DeadlineScope {
124    pub(crate) fn enter(deadline: Option<Instant>) -> Self {
125        let prev = ACTIVE.with(|a| a.get());
126        ACTIVE.with(|a| a.set(deadline.or(prev)));
127        Self { prev }
128    }
129}
130
131impl Drop for DeadlineScope {
132    fn drop(&mut self) {
133        let prev = self.prev;
134        ACTIVE.with(|a| a.set(prev));
135    }
136}
137
138/// The deadline of the query currently executing on this thread, if any.
139pub(crate) fn active_deadline() -> Option<Instant> {
140    ACTIVE.with(|a| a.get())
141}
142
143/// True once `deadline` has passed or its query was cancelled.
144#[inline]
145pub fn deadline_reached(deadline: Instant) -> bool {
146    Instant::now() >= deadline || is_cancelled(deadline)
147}
148
149#[cfg(test)]
150mod tests {
151    use super::*;
152
153    #[test]
154    fn cancel_is_visible_and_forgotten_on_drop() {
155        let mut a = CancellableDeadline::new(None);
156        let b = CancellableDeadline::new(None);
157        assert_ne!(a.deadline(), b.deadline());
158        assert!(!deadline_reached(a.deadline()));
159        a.cancel();
160        assert!(deadline_reached(a.deadline()));
161        assert!(!deadline_reached(b.deadline()));
162        let d = a.deadline();
163        drop(a);
164        assert!(!is_cancelled(d));
165    }
166
167    #[test]
168    fn timeout_still_fires() {
169        let a = CancellableDeadline::new(Some(Duration::ZERO));
170        std::thread::sleep(Duration::from_millis(1));
171        assert!(deadline_reached(a.deadline()));
172    }
173}