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    /// Iterations of expression-level loops since the last clock read.
112    static EVAL_TICK: std::cell::Cell<u32> = const { std::cell::Cell::new(0) };
113    /// Set once an expression-level check saw the active deadline pass
114    /// (or its query cancelled). Sticky for the rest of the query on this
115    /// thread so every enclosing loop unwinds at its next iteration.
116    static EVAL_TRIPPED: std::cell::Cell<bool> = const { std::cell::Cell::new(false) };
117}
118
119/// Expression-level loop iterations between two clock reads. One
120/// iteration is at least a row clone plus a storage lookup, so 256 of them
121/// cost tens of microseconds: the check stays far below 1% of the work
122/// while a deadline is still noticed within a millisecond or so.
123const EVAL_CHECK_STRIDE: u32 = 256;
124
125/// Makes a query's deadline visible to everything that runs on this
126/// thread while the query executes, including pull pipelines that
127/// operators build lazily mid-query (`OPTIONAL MATCH`, `CALL {}`) and
128/// buffered sub-executors. Executors enter one at each public entry
129/// point; nesting keeps the outer deadline when the inner has none, and
130/// dropping the scope restores whatever was active before.
131pub(crate) struct DeadlineScope {
132    prev: Option<Instant>,
133}
134
135impl DeadlineScope {
136    pub(crate) fn enter(deadline: Option<Instant>) -> Self {
137        let prev = ACTIVE.with(|a| a.get());
138        if prev.is_none() {
139            // A new query on this thread: forget a trip left by the last.
140            EVAL_TRIPPED.with(|t| t.set(false));
141        }
142        ACTIVE.with(|a| a.set(deadline.or(prev)));
143        Self { prev }
144    }
145}
146
147impl Drop for DeadlineScope {
148    fn drop(&mut self) {
149        let prev = self.prev;
150        ACTIVE.with(|a| a.set(prev));
151        if prev.is_none() {
152            EVAL_TRIPPED.with(|t| t.set(false));
153        }
154    }
155}
156
157/// Cheap per-iteration check for loops that run inside one expression
158/// evaluation (pattern and list comprehensions, `reduce`, list
159/// quantifiers, pattern expansion). Those loops sit below every operator
160/// boundary the pipeline checks, so without this a single `WHERE` could
161/// run arbitrarily far past the query's deadline. Reads the clock every
162/// [`EVAL_CHECK_STRIDE`] calls; once it fires it stays fired for the rest
163/// of the query, and `eval_expr_result` then reports the timeout.
164#[inline]
165pub(crate) fn eval_deadline_hit() -> bool {
166    if EVAL_TRIPPED.with(|t| t.get()) {
167        return true;
168    }
169    let tick = EVAL_TICK.with(|t| {
170        let n = t.get().wrapping_add(1);
171        t.set(n);
172        n
173    });
174    if !tick.is_multiple_of(EVAL_CHECK_STRIDE) {
175        return false;
176    }
177    match active_deadline() {
178        Some(deadline) if deadline_reached(deadline) => {
179            EVAL_TRIPPED.with(|t| t.set(true));
180            true
181        }
182        _ => false,
183    }
184}
185
186/// Whether an expression on this thread stopped early because the active
187/// query's deadline passed: its value is incomplete and must not be used.
188#[inline]
189pub(crate) fn eval_tripped() -> bool {
190    EVAL_TRIPPED.with(|t| t.get())
191}
192
193/// The deadline of the query currently executing on this thread, if any.
194pub(crate) fn active_deadline() -> Option<Instant> {
195    ACTIVE.with(|a| a.get())
196}
197
198/// True once `deadline` has passed or its query was cancelled.
199#[inline]
200pub fn deadline_reached(deadline: Instant) -> bool {
201    Instant::now() >= deadline || is_cancelled(deadline)
202}
203
204#[cfg(test)]
205mod tests {
206    use super::*;
207
208    #[test]
209    fn cancel_is_visible_and_forgotten_on_drop() {
210        let mut a = CancellableDeadline::new(None);
211        let b = CancellableDeadline::new(None);
212        assert_ne!(a.deadline(), b.deadline());
213        assert!(!deadline_reached(a.deadline()));
214        a.cancel();
215        assert!(deadline_reached(a.deadline()));
216        assert!(!deadline_reached(b.deadline()));
217        let d = a.deadline();
218        drop(a);
219        assert!(!is_cancelled(d));
220    }
221
222    #[test]
223    fn timeout_still_fires() {
224        let a = CancellableDeadline::new(Some(Duration::ZERO));
225        std::thread::sleep(Duration::from_millis(1));
226        assert!(deadline_reached(a.deadline()));
227    }
228}