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}