Skip to main content

oxilite_core/
job.rs

1//! Step machines: the sans-IO execution model.
2//!
3//! A [`Job`] is driven by repeatedly calling [`Job::step`] with the response to the previous
4//! request. Drivers are trivial loops, see [`run_sync`].
5//!
6// @lat: [[architecture#Sans-IO core]]
7
8use crate::error::Result;
9use crate::sql::{Request, Response};
10
11/// What a job wants next.
12#[derive(Debug)]
13pub enum Step<T> {
14    /// Run this request and resume with its response.
15    Execute(Request),
16    /// Finished.
17    Done(T),
18}
19
20/// A resumable operation.
21pub trait Job {
22    type Output;
23
24    /// Advances the job. The first call gets `None`; later calls get the response to the
25    /// previously returned request.
26    fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>>;
27}
28
29/// A backend that runs requests synchronously.
30pub trait SyncBackend {
31    fn execute(&self, request: &Request) -> Result<Response>;
32    fn capabilities(&self) -> &crate::sql::Capabilities;
33
34    /// Starts an interactive transaction (savepoint). Only for backends with
35    /// `interactive_transactions`.
36    fn begin(&self) -> Result<()> {
37        Err(crate::Error::unsupported("interactive transactions"))
38    }
39    fn commit(&self) -> Result<()> {
40        Err(crate::Error::unsupported("interactive transactions"))
41    }
42    fn rollback(&self) -> Result<()> {
43        Err(crate::Error::unsupported("interactive transactions"))
44    }
45}
46
47impl<B: SyncBackend + ?Sized> SyncBackend for &B {
48    fn execute(&self, request: &Request) -> Result<Response> {
49        (**self).execute(request)
50    }
51    fn capabilities(&self) -> &crate::sql::Capabilities {
52        (**self).capabilities()
53    }
54    fn begin(&self) -> Result<()> {
55        (**self).begin()
56    }
57    fn commit(&self) -> Result<()> {
58        (**self).commit()
59    }
60    fn rollback(&self) -> Result<()> {
61        (**self).rollback()
62    }
63}
64
65/// A backend that runs requests asynchronously (D1, network engines).
66///
67/// Futures are not required to be `Send`: Workers are single-threaded.
68#[allow(async_fn_in_trait)]
69pub trait AsyncBackend {
70    async fn execute(&self, request: &Request) -> Result<Response>;
71    fn capabilities(&self) -> &crate::sql::Capabilities;
72}
73
74/// Runs a job to completion on a sync backend.
75pub fn run_sync<J: Job>(backend: &impl SyncBackend, mut job: J) -> Result<J::Output> {
76    let mut response = None;
77    loop {
78        match job.step(response.take())? {
79            Step::Execute(request) => response = Some(backend.execute(&request)?),
80            Step::Done(out) => return Ok(out),
81        }
82    }
83}
84
85/// Runs a job to completion on an async backend.
86pub async fn run_async<J: Job>(backend: &impl AsyncBackend, mut job: J) -> Result<J::Output> {
87    let mut response = None;
88    loop {
89        match job.step(response.take())? {
90            Step::Execute(request) => response = Some(backend.execute(&request).await?),
91            Step::Done(out) => return Ok(out),
92        }
93    }
94}
95
96/// A job made of one request and a decoder.
97pub struct OneShot<T> {
98    request: Option<Request>,
99    decode: Option<Box<dyn FnOnce(Response) -> Result<T>>>,
100}
101
102impl<T> OneShot<T> {
103    pub fn new(request: Request, decode: impl FnOnce(Response) -> Result<T> + 'static) -> Self {
104        Self {
105            request: Some(request),
106            decode: Some(Box::new(decode)),
107        }
108    }
109}
110
111impl<T> Job for OneShot<T> {
112    type Output = T;
113
114    fn step(&mut self, response: Option<Response>) -> Result<Step<T>> {
115        match response {
116            None => Ok(Step::Execute(
117                self.request.take().expect("OneShot started twice"),
118            )),
119            Some(r) => Ok(Step::Done((self
120                .decode
121                .take()
122                .expect("OneShot resumed twice"))(
123                r
124            )?)),
125        }
126    }
127}
128
129/// A job running a sequence of requests (e.g. chunked bulk loads), returning total changes
130/// of the statements flagged in `count`.
131pub struct Sequence {
132    requests: std::vec::IntoIter<Request>,
133    changes: u64,
134}
135
136impl Sequence {
137    pub fn new(requests: Vec<Request>) -> Self {
138        Self {
139            requests: requests.into_iter(),
140            changes: 0,
141        }
142    }
143}
144
145impl Job for Sequence {
146    type Output = u64;
147
148    fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
149        if let Some(r) = response {
150            self.changes += r.iter().map(|rs| rs.changes).sum::<u64>();
151        }
152        Ok(match self.requests.next() {
153            Some(r) => Step::Execute(r),
154            None => Step::Done(self.changes),
155        })
156    }
157}
158
159/// Runs `first`, then the job `next` builds from its output (e.g. resolve versions, then
160/// query), as one job.
161/// Builds the second job of a [`Then`] from the first one's output.
162type Continuation<A, B> = Box<dyn FnOnce(<A as Job>::Output) -> Result<B>>;
163
164pub struct Then<A: Job, B> {
165    first: A,
166    next: Option<Continuation<A, B>>,
167    second: Option<B>,
168}
169
170impl<A: Job, B: Job> Then<A, B> {
171    pub fn new(first: A, next: impl FnOnce(A::Output) -> Result<B> + 'static) -> Self {
172        Self {
173            first,
174            next: Some(Box::new(next)),
175            second: None,
176        }
177    }
178}
179
180impl<A: Job, B: Job> Job for Then<A, B> {
181    type Output = B::Output;
182    fn step(&mut self, response: Option<Response>) -> Result<Step<B::Output>> {
183        if let Some(b) = self.second.as_mut() {
184            return b.step(response);
185        }
186        match self.first.step(response)? {
187            Step::Execute(r) => Ok(Step::Execute(r)),
188            Step::Done(out) => {
189                let make = self.next.take().expect("Then continued twice");
190                self.second.insert(make(out)?).step(None)
191            }
192        }
193    }
194}