use std::{
future::{Future, poll_fn},
panic::{AssertUnwindSafe, catch_unwind},
pin::Pin,
sync::{Arc, Mutex},
task::{Context, Poll},
time::{Duration, SystemTime, UNIX_EPOCH},
};
use saddle_admission::{
AdmissionError, DeploymentResourceBudget, ProfuseGwLightweightAdmissionOutcome,
ProfuseGwLightweightDbFinalizationOwner, ProfuseGwLightweightExecutionOwner,
ProfuseGwLightweightProcessOwner, ProfuseGwLightweightStartupFailure,
prepare_profusegw_lightweight_profile,
};
use saddle_core::{
DbPhysicalDisposition, DbPhysicalDispositionIssuer, DbPhysicalDispositionOwner,
DbPhysicalExecutionHalf, DbPhysicalRequestHalf, DbPhysicalRequestIssuer, DbPhysicalStartupHalf,
pair_db_physical_disposition, seal_db_request_not_used,
};
use crate::Application;
use crate::alpha1_ingress::{
AbsoluteDeadlineError, AbsoluteDeadlineOwner, verify_absolute_deadline,
};
use crate::application::{ShutdownSignal, claim_owned_runtime};
const PROFUSEGW_DEFAULT_REQUEST_TIMEOUT: Duration = Duration::from_millis(5_000);
#[doc(hidden)]
pub struct ProfuseGwManagedDeadline {
unix_ms: i64,
timer: Pin<Box<tokio::time::Sleep>>,
}
#[doc(hidden)]
pub struct ProfuseGwRuntimeProcess {
admission: ProfuseGwLightweightProcessOwner,
db_requests: DbPhysicalRequestIssuer,
db_startup: Option<DbPhysicalStartupHalf>,
}
struct SharedRuntimeProcess {
process: Mutex<Option<ProfuseGwRuntimeProcess>>,
}
#[doc(hidden)]
pub struct ProfuseGwProcessLease {
shared: Arc<SharedRuntimeProcess>,
}
#[doc(hidden)]
pub struct ProfuseGwTerminalProcess {
process: Option<ProfuseGwRuntimeProcess>,
lifecycle: saddle_core::Result<()>,
}
#[doc(hidden)]
pub struct ProfuseGwManagedDispatch {
execution: ProfuseGwLightweightExecutionOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
db_execution: DbPhysicalExecutionHalf,
}
#[doc(hidden)]
pub struct ProfuseGwConcreteDbRequestLease {
execution: ProfuseGwLightweightExecutionOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
db_execution: DbPhysicalExecutionHalf,
}
#[doc(hidden)]
pub struct ProfuseGwDatabaseFinalizationCompletion {
finalization: ProfuseGwLightweightDbFinalizationOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
}
impl ProfuseGwDatabaseFinalizationCompletion {
#[doc(hidden)]
pub fn poll_physical_deadline(&mut self, context: &mut Context<'_>) -> Poll<()> {
self.deadline.timer.as_mut().poll(context)
}
}
#[doc(hidden)]
pub struct ProfuseGwDatabasePhysicalFinalizationHandoff {
completion: ProfuseGwDatabaseFinalizationCompletion,
execution: DbPhysicalExecutionHalf,
}
#[doc(hidden)]
pub enum ProfuseGwDatabaseOperationOutcome<T> {
Ready(ProfuseGwConcreteDbRequestLease, T),
TimedOut(ProfuseGwConcreteDbRequestLease),
Cancelled(ProfuseGwConcreteDbRequestLease),
}
#[doc(hidden)]
pub struct ProfuseGwDatabaseDispositionFailure<T> {
completion: ProfuseGwDatabaseFinalizationCompletion,
physical: DbPhysicalDispositionOwner<T>,
}
#[doc(hidden)]
pub struct ProfuseGwUnusedDatabaseFailure<T> {
dispatch: ProfuseGwManagedDispatch,
value: T,
}
#[doc(hidden)]
pub enum ProfuseGwCoordinatorAdmissionOutcome {
Ready(ProfuseGwManagedDispatch),
CapacityRejected,
Stop(AdmissionError),
DeadlineUnavailable(AbsoluteDeadlineError),
}
#[doc(hidden)]
pub enum ProfuseGwCoordinatorFailure {
Startup(ProfuseGwLightweightStartupFailure),
Lifecycle(saddle_core::SaddleError),
Finalization(AdmissionError),
}
impl ProfuseGwRuntimeProcess {
fn new(admission: ProfuseGwLightweightProcessOwner) -> Self {
let (db_startup, db_requests) =
DbPhysicalDispositionIssuer::issue().into_startup_and_request_issuer();
Self {
admission,
db_requests,
db_startup: Some(db_startup),
}
}
#[doc(hidden)]
pub fn try_admit(&self) -> ProfuseGwCoordinatorAdmissionOutcome {
let deadline = match default_profusegw_deadline() {
Ok((unix_ms, deadline)) => ProfuseGwManagedDeadline {
unix_ms,
timer: deadline.into_sleep(),
},
Err(error) => {
return ProfuseGwCoordinatorAdmissionOutcome::DeadlineUnavailable(error);
}
};
match self.admission.verified_profile().try_admit() {
ProfuseGwLightweightAdmissionOutcome::Ready(permit) => {
let Some((db_request, db_execution)) = self.db_requests.issue_request() else {
return ProfuseGwCoordinatorAdmissionOutcome::Stop(
AdmissionError::InvalidConfiguration,
);
};
ProfuseGwCoordinatorAdmissionOutcome::Ready(ProfuseGwManagedDispatch {
execution: permit.into_execution(),
deadline,
db_request,
db_execution,
})
}
ProfuseGwLightweightAdmissionOutcome::CapacityRejected => {
ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected
}
ProfuseGwLightweightAdmissionOutcome::Stop(error) => {
ProfuseGwCoordinatorAdmissionOutcome::Stop(error)
}
}
}
fn finish(self) -> Result<(), AdmissionError> {
let startup_consumed = self.db_startup.is_none();
let result = self.admission.finish();
if !startup_consumed {
return Err(AdmissionError::InvalidConfiguration);
}
result
}
}
impl ProfuseGwProcessLease {
#[doc(hidden)]
pub fn take_database_startup_half(&self) -> Option<DbPhysicalStartupHalf> {
self.shared
.process
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.as_mut()?
.db_startup
.take()
}
#[doc(hidden)]
pub fn try_admit(&self) -> ProfuseGwCoordinatorAdmissionOutcome {
let process = self
.shared
.process
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
match process.as_ref() {
Some(process) => process.try_admit(),
None => ProfuseGwCoordinatorAdmissionOutcome::Stop(AdmissionError::AccountClosed),
}
}
}
impl ProfuseGwTerminalProcess {
fn finish(self) -> Result<(), ProfuseGwCoordinatorFailure> {
let finalization = match self.process {
Some(process) => process
.finish()
.map_err(ProfuseGwCoordinatorFailure::Finalization),
None => Err(ProfuseGwCoordinatorFailure::Lifecycle(
lifecycle_recovery_error(),
)),
};
match self.lifecycle {
Ok(()) => finalization,
Err(error) => {
let _finalization = finalization;
Err(ProfuseGwCoordinatorFailure::Lifecycle(error))
}
}
}
}
impl ProfuseGwManagedDispatch {
#[doc(hidden)]
pub fn deadline_unix_ms(&self) -> i64 {
self.deadline.unix_ms
}
#[doc(hidden)]
pub async fn deadline_elapsed(&mut self) {
self.deadline.timer.as_mut().await;
}
#[doc(hidden)]
pub fn cancel(self) {
self.execution.cancel();
}
#[doc(hidden)]
pub fn timeout(self) {
self.execution.timeout();
}
#[doc(hidden)]
pub fn into_database_request(self) -> ProfuseGwConcreteDbRequestLease {
ProfuseGwConcreteDbRequestLease {
execution: self.execution,
deadline: self.deadline,
db_request: self.db_request,
db_execution: self.db_execution,
}
}
}
impl<T> ProfuseGwUnusedDatabaseFailure<T> {
#[doc(hidden)]
pub fn into_inputs(self) -> (ProfuseGwManagedDispatch, T) {
(self.dispatch, self.value)
}
}
#[doc(hidden)]
pub fn finish_profusegw_without_database<T>(
dispatch: ProfuseGwManagedDispatch,
value: T,
) -> Result<T, ProfuseGwUnusedDatabaseFailure<T>> {
let ProfuseGwManagedDispatch {
execution,
deadline,
db_request,
db_execution,
} = dispatch;
let receipt = match seal_db_request_not_used(db_request, db_execution, value) {
Ok(receipt) => receipt,
Err((db_request, db_execution, value)) => {
return Err(ProfuseGwUnusedDatabaseFailure {
dispatch: ProfuseGwManagedDispatch {
execution,
deadline,
db_request,
db_execution,
},
value,
});
}
};
let value = receipt.into_value();
execution.cancel();
drop(deadline);
Ok(value)
}
impl ProfuseGwConcreteDbRequestLease {
#[doc(hidden)]
pub fn restore_dispatch(self) -> ProfuseGwManagedDispatch {
ProfuseGwManagedDispatch {
execution: self.execution,
deadline: self.deadline,
db_request: self.db_request,
db_execution: self.db_execution,
}
}
#[doc(hidden)]
pub fn into_physical_finalization(self) -> ProfuseGwDatabasePhysicalFinalizationHandoff {
ProfuseGwDatabasePhysicalFinalizationHandoff {
completion: ProfuseGwDatabaseFinalizationCompletion {
finalization: self.execution.begin_database_finalization(),
deadline: self.deadline,
db_request: self.db_request,
},
execution: self.db_execution,
}
}
}
impl ProfuseGwDatabasePhysicalFinalizationHandoff {
#[doc(hidden)]
pub fn into_database_execution(
self,
) -> (
ProfuseGwDatabaseFinalizationCompletion,
DbPhysicalExecutionHalf,
) {
(self.completion, self.execution)
}
}
impl<T> ProfuseGwDatabaseDispositionFailure<T> {
#[doc(hidden)]
pub fn into_inputs(
self,
) -> (
ProfuseGwDatabaseFinalizationCompletion,
DbPhysicalDispositionOwner<T>,
) {
(self.completion, self.physical)
}
}
enum OperationPoll<T> {
Ready(T),
TimedOut,
Cancelled,
}
#[doc(hidden)]
pub async fn poll_profusegw_database_operation<T, O, C>(
mut lease: ProfuseGwConcreteDbRequestLease,
operation: O,
cancel: C,
) -> ProfuseGwDatabaseOperationOutcome<T>
where
O: Future<Output = T>,
C: Future<Output = ()>,
{
tokio::pin!(operation);
tokio::pin!(cancel);
let outcome = poll_fn(|context| {
if let Poll::Ready(value) = lease
.execution
.poll_database_query(operation.as_mut(), context)
{
return Poll::Ready(OperationPoll::Ready(value));
}
if cancel.as_mut().poll(context).is_ready() {
return Poll::Ready(OperationPoll::Cancelled);
}
if lease.deadline.timer.as_mut().poll(context).is_ready() {
return Poll::Ready(OperationPoll::TimedOut);
}
Poll::Pending
})
.await;
match outcome {
OperationPoll::Ready(value) => ProfuseGwDatabaseOperationOutcome::Ready(lease, value),
OperationPoll::TimedOut => ProfuseGwDatabaseOperationOutcome::TimedOut(lease),
OperationPoll::Cancelled => ProfuseGwDatabaseOperationOutcome::Cancelled(lease),
}
}
#[doc(hidden)]
pub fn finish_profusegw_database_disposition<T>(
completion: ProfuseGwDatabaseFinalizationCompletion,
physical: DbPhysicalDispositionOwner<T>,
) -> Result<T, ProfuseGwDatabaseDispositionFailure<T>> {
let receipt = match pair_db_physical_disposition(physical, completion.db_request) {
Ok(receipt) => receipt,
Err((physical, db_request)) => {
return Err(ProfuseGwDatabaseDispositionFailure {
completion: ProfuseGwDatabaseFinalizationCompletion {
db_request,
..completion
},
physical,
});
}
};
let (disposition, value) = receipt.into_outcome();
let admission = match disposition {
DbPhysicalDisposition::Returned => completion.finalization.connection_returned(),
DbPhysicalDisposition::Discarded => completion.finalization.connection_discarded(),
};
admission.finish();
drop(completion.deadline);
Ok(value)
}
#[doc(hidden)]
pub async fn coordinate_profusegw_app_run<Process, ProcessFuture>(
budget: DeploymentResourceBudget,
process: Process,
) -> Result<(), ProfuseGwCoordinatorFailure>
where
Process: FnOnce(ProfuseGwRuntimeProcess) -> ProcessFuture,
ProcessFuture: Future<Output = ProfuseGwTerminalProcess>,
{
let admission = prepare_profusegw_lightweight_profile(budget)
.map_err(ProfuseGwCoordinatorFailure::Startup)?;
process(ProfuseGwRuntimeProcess::new(admission))
.await
.finish()
}
#[doc(hidden)]
pub fn run_profusegw_owned_application<Factory, FactoryFuture>(
process: ProfuseGwRuntimeProcess,
factory: Factory,
) -> ProfuseGwTerminalProcess
where
Factory: FnOnce(ProfuseGwProcessLease) -> FactoryFuture,
FactoryFuture: Future<Output = saddle_core::Result<Application>>,
{
let runtime = match claim_owned_runtime() {
Ok(runtime) => runtime,
Err(error) => {
return ProfuseGwTerminalProcess {
process: Some(process),
lifecycle: Err(error),
};
}
};
let shared = Arc::new(SharedRuntimeProcess {
process: Mutex::new(Some(process)),
});
let lease = ProfuseGwProcessLease {
shared: Arc::clone(&shared),
};
let outcome = catch_unwind(AssertUnwindSafe(|| {
runtime.block_on(async move {
let signal = ShutdownSignal::register()?;
let application = factory(lease).await?;
let finalizer = application.pending_driver_finalizer();
let result = application.run_until_shutdown(signal.wait()).await;
Ok::<_, saddle_core::SaddleError>((finalizer, result))
})
}));
let lifecycle = match outcome {
Ok(Ok((finalizer, result))) => finalizer.finish(runtime, result),
Ok(Err(error)) => {
drop(runtime);
Err(error)
}
Err(_) => {
drop(runtime);
Err(saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal,
"runtime.profusegw_lifecycle_panicked",
"the ProfuseGW Application lifecycle panicked before terminal recovery",
))
}
};
let process = shared
.process
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.take();
ProfuseGwTerminalProcess { process, lifecycle }
}
fn lifecycle_recovery_error() -> saddle_core::SaddleError {
saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal,
"runtime.profusegw_owner_recovery_failed",
"the ProfuseGW process authority was unavailable at lifecycle terminal",
)
}
fn default_profusegw_deadline() -> Result<(i64, AbsoluteDeadlineOwner), AbsoluteDeadlineError> {
let now = SystemTime::now()
.duration_since(UNIX_EPOCH)
.map_err(|_| AbsoluteDeadlineError::ClockBeforeUnixEpoch)?;
let now_ms = u64::try_from(now.as_millis()).map_err(|_| AbsoluteDeadlineError::OutOfRange)?;
let timeout_ms = u64::try_from(PROFUSEGW_DEFAULT_REQUEST_TIMEOUT.as_millis())
.map_err(|_| AbsoluteDeadlineError::OutOfRange)?;
let deadline_unix_ms = now_ms
.checked_add(timeout_ms)
.ok_or(AbsoluteDeadlineError::OutOfRange)?;
let deadline = verify_absolute_deadline(deadline_unix_ms)?;
let deadline_unix_ms =
i64::try_from(deadline_unix_ms).map_err(|_| AbsoluteDeadlineError::OutOfRange)?;
Ok((deadline_unix_ms, deadline))
}
#[cfg(test)]
mod tests {
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use super::default_profusegw_deadline;
#[test]
fn budget_default_deadline_is_runtime_owned_and_live() {
let (unix_ms, _owner) = default_profusegw_deadline().unwrap();
assert!(unix_ms > 0);
}
#[tokio::test]
async fn process_adapter_shape_is_once_only() {
async fn invoke_once<P, F, Fut>(owner: P, process: F) -> P
where
F: FnOnce(P) -> Fut,
Fut: Future<Output = P>,
{
process(owner).await
}
let calls = Arc::new(AtomicUsize::new(0));
let observed = Arc::clone(&calls);
let owner = Box::new(41_u64);
let address = (&*owner) as *const u64;
let returned = invoke_once(owner, move |owner| async move {
observed.fetch_add(1, Ordering::SeqCst);
owner
})
.await;
assert_eq!(calls.load(Ordering::SeqCst), 1);
assert_eq!((&*returned) as *const u64, address);
}
}