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, ProfuseGwAdmissionObservationReceipt,
ProfuseGwCapacityBottleneck, ProfuseGwCapacityDecision, ProfuseGwCapacityRejectReason,
ProfuseGwCreditChange, ProfuseGwDatabaseDisposition as AdmissionDatabaseDisposition,
ProfuseGwLightweightDbFinalizationOwner, ProfuseGwLightweightExecutionOwner,
ProfuseGwLightweightObservedAdmissionOutcome, ProfuseGwLightweightProcessOwner,
ProfuseGwLightweightStartupFailure, ProfuseGwTerminalObservationReceipt,
prepare_profusegw_lightweight_profile,
};
use saddle_core::{
DbPhysicalDisposition, DbPhysicalDispositionIssuer, DbPhysicalDispositionOwner,
DbPhysicalExecutionHalf, DbPhysicalRequestHalf, DbPhysicalRequestIssuer, DbPhysicalStartupHalf,
pair_db_physical_disposition, seal_db_request_not_used,
};
use saddle_observability::{
BottleneckIdentity, CallContext, CapacityDimension, CapacityObservation, DatabaseDisposition,
EventContext, Observer, RejectReason,
};
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,
observation: Option<ProfuseGwRequestObservation>,
}
struct ProfuseGwRequestObservation {
observer: Observer,
context: CallContext,
event_context: EventContext,
}
#[doc(hidden)]
pub struct ProfuseGwAdmissionEvent {
receipt: ProfuseGwAdmissionObservationReceipt,
}
#[doc(hidden)]
pub struct ProfuseGwConcreteDbRequestLease {
execution: ProfuseGwLightweightExecutionOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
db_execution: DbPhysicalExecutionHalf,
observation: Option<ProfuseGwRequestObservation>,
}
#[doc(hidden)]
pub struct ProfuseGwDatabaseFinalizationCompletion {
finalization: ProfuseGwLightweightDbFinalizationOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
observation: Option<ProfuseGwRequestObservation>,
}
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, ProfuseGwAdmissionEvent),
CapacityRejected(ProfuseGwAdmissionEvent),
Stop(AdmissionError),
DeadlineUnavailable(AbsoluteDeadlineError),
}
#[doc(hidden)]
pub enum ProfuseGwCoordinatorFailure {
Startup(ProfuseGwLightweightStartupFailure),
Lifecycle(saddle_core::SaddleError),
Finalization(AdmissionError),
}
impl ProfuseGwAdmissionEvent {
#[doc(hidden)]
pub fn record(self, observer: &Observer, context: &CallContext, event: &EventContext) {
let receipt = self.receipt;
record_capacity_snapshot(
observer,
context,
event,
receipt.snapshot(),
receipt.used(),
receipt.decision(),
receipt.reason(),
receipt.elapsed_ms(),
);
}
}
#[doc(hidden)]
pub fn record_profusegw_startup_capacity(
observer: &Observer,
context: &CallContext,
event: &EventContext,
snapshot: saddle_admission::ProfuseGwCapacitySnapshot,
) {
record_capacity_snapshot(
observer,
context,
event,
snapshot,
0,
ProfuseGwCapacityDecision::Accepted,
None,
0,
);
}
fn record_capacity_snapshot(
observer: &Observer,
context: &CallContext,
event: &EventContext,
snapshot: saddle_admission::ProfuseGwCapacitySnapshot,
used: usize,
decision: ProfuseGwCapacityDecision,
reason: Option<ProfuseGwCapacityRejectReason>,
elapsed_ms: u64,
) {
let bottleneck = match snapshot.bottleneck() {
ProfuseGwCapacityBottleneck::Cpu => "cpu",
ProfuseGwCapacityBottleneck::Memory => "memory",
ProfuseGwCapacityBottleneck::Database => "database",
ProfuseGwCapacityBottleneck::ProfuseContract => "profuse_contract",
};
let Ok(bottleneck) = BottleneckIdentity::new(bottleneck) else {
return;
};
let dimensions = [
(CapacityDimension::Cpu, snapshot.cpu_budget()),
(CapacityDimension::Memory, snapshot.memory_budget_bytes()),
(CapacityDimension::Database, snapshot.database_budget()),
(
CapacityDimension::ProfuseContract,
snapshot.profusecontract_budget(),
),
];
let limit = u64::try_from(snapshot.active_limit()).unwrap_or(u64::MAX);
let used = u64::try_from(used).unwrap_or(u64::MAX);
for (dimension, budget) in dimensions {
let budget = u64::try_from(budget).unwrap_or(u64::MAX);
let observation = if decision == ProfuseGwCapacityDecision::CapacityRejected {
let Ok(reason) = RejectReason::new(match reason {
Some(ProfuseGwCapacityRejectReason::AtLimit) | None => "at_limit",
}) else {
return;
};
CapacityObservation::rejected_dimension(
dimension,
budget,
limit,
used,
bottleneck.clone(),
reason,
elapsed_ms,
)
} else {
CapacityObservation::accepted_dimension(
dimension,
budget,
limit,
used,
bottleneck.clone(),
elapsed_ms,
)
};
observer.record_capacity(context, event, observation);
}
}
fn record_terminal(
observation: Option<ProfuseGwRequestObservation>,
receipt: ProfuseGwTerminalObservationReceipt,
) {
if receipt.credit() != ProfuseGwCreditChange::Released {
return;
}
let disposition = match receipt.database() {
AdmissionDatabaseDisposition::NotUsed => DatabaseDisposition::NotUsed,
AdmissionDatabaseDisposition::Returned => DatabaseDisposition::Returned,
AdmissionDatabaseDisposition::Discarded => DatabaseDisposition::Discarded,
};
if let Some(observation) = observation {
observation.observer.record_resource_finalization(
&observation.context,
&observation.event_context,
disposition,
receipt.elapsed_ms(),
);
}
}
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_observed() {
ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, receipt) => {
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,
observation: None,
},
ProfuseGwAdmissionEvent { receipt },
)
}
ProfuseGwLightweightObservedAdmissionOutcome::CapacityRejected(receipt) => {
ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(ProfuseGwAdmissionEvent {
receipt,
})
}
ProfuseGwLightweightObservedAdmissionOutcome::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 capacity_snapshot(&self) -> Option<saddle_admission::ProfuseGwCapacitySnapshot> {
self.shared
.process
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.as_ref()
.map(|process| process.admission.capacity_snapshot())
}
#[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 bind_observation(
mut self,
observer: Observer,
context: CallContext,
event_context: EventContext,
admission: ProfuseGwAdmissionEvent,
) -> Self {
admission.record(&observer, &context, &event_context);
self.observation = Some(ProfuseGwRequestObservation {
observer,
context,
event_context,
});
self
}
#[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) {
record_terminal(self.observation, self.execution.cancel_observed());
}
#[doc(hidden)]
pub fn timeout(self) {
record_terminal(self.observation, self.execution.timeout_observed());
}
#[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,
observation: self.observation,
}
}
}
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,
observation,
} = 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,
observation,
},
value,
});
}
};
let value = receipt.into_value();
record_terminal(observation, execution.cancel_observed());
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,
observation: self.observation,
}
}
#[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,
observation: self.observation,
},
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(),
};
record_terminal(completion.observation, admission.finish_observed());
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);
}
}