use crate::error::Result;
use crate::sql::{Request, Response};
#[derive(Debug)]
pub enum Step<T> {
Execute(Request),
Done(T),
}
pub trait Job {
type Output;
fn step(&mut self, response: Option<Response>) -> Result<Step<Self::Output>>;
}
pub trait SyncBackend {
fn execute(&self, request: &Request) -> Result<Response>;
fn capabilities(&self) -> &crate::sql::Capabilities;
fn begin(&self) -> Result<()> {
Err(crate::Error::unsupported("interactive transactions"))
}
fn commit(&self) -> Result<()> {
Err(crate::Error::unsupported("interactive transactions"))
}
fn rollback(&self) -> Result<()> {
Err(crate::Error::unsupported("interactive transactions"))
}
}
impl<B: SyncBackend + ?Sized> SyncBackend for &B {
fn execute(&self, request: &Request) -> Result<Response> {
(**self).execute(request)
}
fn capabilities(&self) -> &crate::sql::Capabilities {
(**self).capabilities()
}
fn begin(&self) -> Result<()> {
(**self).begin()
}
fn commit(&self) -> Result<()> {
(**self).commit()
}
fn rollback(&self) -> Result<()> {
(**self).rollback()
}
}
#[allow(async_fn_in_trait)]
pub trait AsyncBackend {
async fn execute(&self, request: &Request) -> Result<Response>;
fn capabilities(&self) -> &crate::sql::Capabilities;
}
pub fn run_sync<J: Job>(backend: &impl SyncBackend, mut job: J) -> Result<J::Output> {
let mut response = None;
loop {
match job.step(response.take())? {
Step::Execute(request) => response = Some(backend.execute(&request)?),
Step::Done(out) => return Ok(out),
}
}
}
pub async fn run_async<J: Job>(backend: &impl AsyncBackend, mut job: J) -> Result<J::Output> {
let mut response = None;
loop {
match job.step(response.take())? {
Step::Execute(request) => response = Some(backend.execute(&request).await?),
Step::Done(out) => return Ok(out),
}
}
}
pub struct OneShot<T> {
request: Option<Request>,
decode: Option<Box<dyn FnOnce(Response) -> Result<T>>>,
}
impl<T> OneShot<T> {
pub fn new(request: Request, decode: impl FnOnce(Response) -> Result<T> + 'static) -> Self {
Self {
request: Some(request),
decode: Some(Box::new(decode)),
}
}
}
impl<T> Job for OneShot<T> {
type Output = T;
fn step(&mut self, response: Option<Response>) -> Result<Step<T>> {
match response {
None => Ok(Step::Execute(
self.request.take().expect("OneShot started twice"),
)),
Some(r) => Ok(Step::Done((self
.decode
.take()
.expect("OneShot resumed twice"))(
r
)?)),
}
}
}
pub struct Sequence {
requests: std::vec::IntoIter<Request>,
changes: u64,
}
impl Sequence {
pub fn new(requests: Vec<Request>) -> Self {
Self {
requests: requests.into_iter(),
changes: 0,
}
}
}
impl Job for Sequence {
type Output = u64;
fn step(&mut self, response: Option<Response>) -> Result<Step<u64>> {
if let Some(r) = response {
self.changes += r.iter().map(|rs| rs.changes).sum::<u64>();
}
Ok(match self.requests.next() {
Some(r) => Step::Execute(r),
None => Step::Done(self.changes),
})
}
}
type Continuation<A, B> = Box<dyn FnOnce(<A as Job>::Output) -> Result<B>>;
pub struct Then<A: Job, B> {
first: A,
next: Option<Continuation<A, B>>,
second: Option<B>,
}
impl<A: Job, B: Job> Then<A, B> {
pub fn new(first: A, next: impl FnOnce(A::Output) -> Result<B> + 'static) -> Self {
Self {
first,
next: Some(Box::new(next)),
second: None,
}
}
}
impl<A: Job, B: Job> Job for Then<A, B> {
type Output = B::Output;
fn step(&mut self, response: Option<Response>) -> Result<Step<B::Output>> {
if let Some(b) = self.second.as_mut() {
return b.step(response);
}
match self.first.step(response)? {
Step::Execute(r) => Ok(Step::Execute(r)),
Step::Done(out) => {
let make = self.next.take().expect("Then continued twice");
self.second.insert(make(out)?).step(None)
}
}
}
}