1//! Work a judgement hands to the backend, and resumes from.
2//!
3//! A judgement runs on the executor and may not call Postgres, yet a
4//! judgement over several rounds (a label-tree search) needs the cache
5//! between them: each round's keys are looked up before it is sent.
6//! [`on_backend`] queues a closure and returns a future of its result.
7//! The scan runs the queue on the backend between waits, in its settling
8//! context, and the judgement resumes on the next poll.
9//!
10//! The queue is the backend's, not one scan's: whichever jev scan is
11//! waiting runs it, as it polls every scan's tasks. A closure whose
12//! future was dropped (a cancel, a rescan, a LIMIT) is skipped unrun.
13
14use std::cell::RefCell;
15use std::future::Future;
16use std::rc::{Rc, Weak};
17use std::task::{Poll, Waker};
18
19type Slot<T> = RefCell<(Option<T>, Option<Waker>)>;
20
21thread_local! {
22    static QUEUE: RefCell<Vec<Box<dyn FnOnce()>>> = const { RefCell::new(Vec::new()) };
23    /// The waiting scan's wait, woken when work is queued.
24    static WAITING: RefCell<Option<Waker>> = const { RefCell::new(None) };
25}
26
27/// Runs `work` on the backend, where Postgres may be called, and resolves
28/// to its result. It may raise an ERROR, which aborts the statement as
29/// any other on the backend does.
30pub fn on_backend<T: 'static>(work: impl FnOnce() -> T + 'static) -> impl Future<Output = T> {
31    let slot: Rc<Slot<T>> = Rc::new(RefCell::new((None, None)));
32    let weak: Weak<Slot<T>> = Rc::downgrade(&slot);
33    QUEUE.with(|q| {
34        q.borrow_mut().push(Box::new(move || {
35            // Dropped unanswered: nothing is waiting for the result.
36            let Some(slot) = weak.upgrade() else { return };
37            let out = work();
38            let waker = {
39                let mut s = slot.borrow_mut();
40                s.0 = Some(out);
41                s.1.take()
42            };
43            if let Some(w) = waker {
44                w.wake();
45            }
46        }))
47    });
48    WAITING.with(|w| {
49        if let Some(w) = w.borrow_mut().take() {
50            w.wake();
51        }
52    });
53    std::future::poll_fn(move |cx| {
54        let mut s = slot.borrow_mut();
55        match s.0.take() {
56            Some(out) => Poll::Ready(out),
57            None => {
58                s.1 = Some(cx.waker().clone());
59                Poll::Pending
60            }
61        }
62    })
63}
64
65/// Whether work is queued; the scan's wait also registers `waker` to be
66/// woken when some is.
67pub(super) fn pending(waker: &Waker) -> bool {
68    WAITING.with(|w| *w.borrow_mut() = Some(waker.clone()));
69    QUEUE.with(|q| !q.borrow().is_empty())
70}
71
72/// Runs what is queued, including what the work itself queues.
73pub(super) fn run() {
74    loop {
75        let batch = QUEUE.with(|q| std::mem::take(&mut *q.borrow_mut()));
76        if batch.is_empty() {
77            return;
78        }
79        for work in batch {
80            work();
81        }
82    }
83}