use std::{
future::{Future, poll_fn},
pin::pin,
sync::atomic::Ordering,
task::Poll,
};
use saddle_core::{DbPhysicalDisposition, DbPhysicalExecutionHalf};
use saddle_runtime::profusegw::{
ProfuseGwDatabaseFinalizationCompletion, finish_profusegw_database_disposition,
};
use sqlx::{MySql, pool::PoolConnection};
use crate::startup_pool::StartupManagedDatabaseProcessCapability;
#[doc(hidden)]
pub struct DatabasePhysicalTerminal<T> {
disposition: DbPhysicalDisposition,
value: T,
}
impl<T> DatabasePhysicalTerminal<T> {
pub fn into_outcome(self) -> (DbPhysicalDisposition, T) {
(self.disposition, self.value)
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(crate) enum RequestedPhysicalDisposition {
Return,
Discard,
}
impl StartupManagedDatabaseProcessCapability {
pub(crate) fn ensure_live(&self) -> bool {
self.process_live.load(Ordering::Acquire)
}
pub(crate) async fn finish_physical<T>(
&self,
mut connection: Option<PoolConnection<MySql>>,
mut completion: ProfuseGwDatabaseFinalizationCompletion,
execution: DbPhysicalExecutionHalf,
requested: RequestedPhysicalDisposition,
value: T,
) -> DatabasePhysicalTerminal<T> {
let disposition = match (requested, connection.as_mut()) {
(RequestedPhysicalDisposition::Return, Some(connection_ref)) => {
let returned = {
let return_to_pool = connection_ref.return_to_pool();
let mut return_to_pool = pin!(return_to_pool);
poll_fn(|context| {
if return_to_pool.as_mut().poll(context).is_ready() {
return Poll::Ready(true);
}
if completion.poll_physical_deadline(context).is_ready() {
return Poll::Ready(false);
}
Poll::Pending
})
.await
};
if returned {
drop(connection.take());
DbPhysicalDisposition::Returned
} else {
drop(
connection
.take()
.expect("physical connection exists")
.detach(),
);
DbPhysicalDisposition::Discarded
}
}
(_, Some(_)) => {
drop(
connection
.take()
.expect("physical connection exists")
.detach(),
);
DbPhysicalDisposition::Discarded
}
_ => DbPhysicalDisposition::Discarded,
};
let physical = match disposition {
DbPhysicalDisposition::Returned => self.physical.connection_returned(execution, value),
DbPhysicalDisposition::Discarded => {
self.physical.connection_discarded(execution, value)
}
};
let physical = match physical {
Ok(physical) => physical,
Err(_) => std::process::abort(),
};
let value = match finish_profusegw_database_disposition(completion, physical) {
Ok(value) => value,
Err(_) => std::process::abort(),
};
DatabasePhysicalTerminal { disposition, value }
}
}