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}