use std::{
future::{Future, poll_fn},
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;
#[path = "profusegw_serial_scope.rs"]
mod serial_scope;
#[cfg(test)]
#[path = "profusegw_ar_tests.rs"]
mod ar_tests;
use crate::alpha1_ingress::{
AbsoluteDeadlineError, AbsoluteDeadlineOwner, verify_absolute_deadline,
};
use crate::application::{ShutdownSignal, claim_owned_runtime};
pub use serial_scope::*;
#[doc(hidden)]
pub struct ProfuseGwManagedDeadline {
unix_ms: i64,
timer: Pin<Box<tokio::time::Sleep>>,
reserved_storage: Option<crate::request_task::reserved::ReservedRequestLifetime>,
}
#[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 {
result: Result<(), ProfuseGwCoordinatorFailure>,
}
#[doc(hidden)]
pub struct ProfuseGwManagedDispatch {
execution: ProfuseGwLightweightExecutionOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
db_execution: DbPhysicalExecutionHalf,
observation: Option<ProfuseGwRequestObservation>,
scope_stop: Option<ProfuseGwScopeStop>,
}
struct ProfuseGwRequestObservation {
observer: Observer,
context: CallContext,
event_context: EventContext,
diagnostic: Mutex<Option<saddle_core::DiagnosticOccurrence>>,
request_diagnostic: Mutex<Option<crate::diagnostics::RequestCaptured>>,
cleanup_failed: std::sync::atomic::AtomicBool,
}
impl ProfuseGwManagedDispatch {
pub(crate) fn retain_reserved_deadline(&mut self,storage:crate::request_task::reserved::ReservedRequestLifetime) {
self.deadline.reserved_storage=Some(storage);
}
}
impl ProfuseGwConcreteDbRequestLease {
pub fn reserved_diagnostic_view(&self, view: &crate::request_task::reserved::ReservedRequestView)
-> Result<crate::request_task::reserved::ReservedRequestView,crate::request_task::reserved::ReservedContextError> {
view.validate_execution(&self.execution).map_err(crate::request_task::reserved::ReservedContextError::Storage)?;
view.in_scope(&self.db_request,&self.db_execution)
}
}
impl ProfuseGwRequestObservation {
fn occurrence(&self) -> Option<saddle_core::DiagnosticOccurrence> {
*self.diagnostic.lock().unwrap_or_else(|e| e.into_inner())
}
}
fn record_request_stop(
observation: Option<&ProfuseGwRequestObservation>,
stop: ProfuseGwScopeStop,
physical_return: bool,
stage: saddle_core::DiagnosticStage,
) {
let occurrence = crate::diagnostics::request_stop_source(
stop == ProfuseGwScopeStop::TimedOut,
physical_return,
stage,
observation.map(|o| (&o.context, &o.event_context)),
);
if let Some(observation) = observation {
observation
.diagnostic
.lock()
.unwrap_or_else(|e| e.into_inner())
.get_or_insert(occurrence);
}
}
#[doc(hidden)]
pub struct ProfuseGwAdmissionEvent {
receipt: ProfuseGwAdmissionObservationReceipt,
}
#[doc(hidden)]
pub type ReservedDispatchOwner<X> = (
Option<ProfuseGwManagedDispatch>,
X,
Option<ProfuseGwAdmissionEvent>,
);
#[doc(hidden)]
pub struct ProfuseGwConcreteDbRequestLease {
execution: ProfuseGwLightweightExecutionOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
db_execution: DbPhysicalExecutionHalf,
observation: Option<ProfuseGwRequestObservation>,
scope_stop: Option<ProfuseGwScopeStop>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[doc(hidden)]
pub enum ProfuseGwScopeStop {
Cancelled,
TimedOut,
}
#[doc(hidden)]
pub struct ProfuseGwDatabaseFinalizationCompletion {
finalization: ProfuseGwLightweightDbFinalizationOwner,
deadline: ProfuseGwManagedDeadline,
db_request: DbPhysicalRequestHalf,
observation: Option<ProfuseGwRequestObservation>,
scope_stop: Option<ProfuseGwScopeStop>,
}
impl ProfuseGwDatabaseFinalizationCompletion {
#[doc(hidden)]
pub fn poll_physical_deadline(&mut self, context: &mut Context<'_>) -> Poll<()> {
if self.scope_stop.is_some()
|| tokio::time::Instant::now() >= self.deadline.timer.deadline()
|| self.deadline.timer.as_mut().poll(context).is_ready()
{
if self.scope_stop.is_none() {
self.scope_stop = Some(ProfuseGwScopeStop::TimedOut);
record_request_stop(
self.observation.as_ref(),
ProfuseGwScopeStop::TimedOut,
true,
saddle_core::DiagnosticStage::FinalizerResource,
);
}
Poll::Ready(())
} else {
Poll::Pending
}
}
}
#[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 ProfuseGwPostDatabaseManagedOwner<T> {
terminal: ProfuseGwPostDatabaseRequestTerminal,
value: T,
}
#[doc(hidden)]
pub struct ProfuseGwPostDatabaseRequestTerminal {
admission: saddle_admission::ProfuseGwLightweightDispositionReceipt,
deadline: ProfuseGwManagedDeadline,
observation: Option<ProfuseGwRequestObservation>,
scope_stop: Option<ProfuseGwScopeStop>,
next_scope: serial_scope::NextScope,
}
impl<T> ProfuseGwPostDatabaseManagedOwner<T> {
#[doc(hidden)]
pub fn into_response_parts(self) -> (T, ProfuseGwPostDatabaseRequestTerminal) {
(self.value, self.terminal)
}
}
#[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_unrooted(
self,
observer: &Observer,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
application: &saddle_core::ContextLabel,
lifecycle: saddle_core::RequestViewPhase,
construction: saddle_observability::AdmissionConstruction,
) -> saddle_observability::AdmissionCapacitySubmission {
self.record_capacity_context(observer, output,
saddle_observability::AdmissionEventContext::Unrooted { application, lifecycle },
construction)
}
pub(crate) fn record_capacity_context(
self,
observer: &Observer,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
context: saddle_observability::AdmissionEventContext<'_>,
construction: saddle_observability::AdmissionConstruction,
) -> saddle_observability::AdmissionCapacitySubmission {
let facts = saddle_observability::AdmissionCapacityFacts::from_receipt(&self.receipt);
observer.record_admission_capacity(output, context, &facts, construction)
}
#[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 {
record_post_release_capacity_snapshot(
&observation.observer,
&observation.context,
&observation.event_context,
receipt.post_release(),
receipt.elapsed_ms(),
);
observation.observer.record_resource_finalization(
&observation.context,
&observation.event_context,
disposition,
receipt.elapsed_ms(),
);
let axes = saddle_core::DiagnosticOutcomeAxes {
physical: match receipt.database() {
AdmissionDatabaseDisposition::NotUsed => {
saddle_core::PhysicalDispositionFact::NotUsed
}
AdmissionDatabaseDisposition::Returned => {
saddle_core::PhysicalDispositionFact::Returned
}
AdmissionDatabaseDisposition::Discarded => {
saddle_core::PhysicalDispositionFact::Discarded
}
},
cleanup: if observation
.cleanup_failed
.load(std::sync::atomic::Ordering::Relaxed)
{
saddle_core::CleanupOutcome::Failed
} else {
saddle_core::CleanupOutcome::Succeeded
},
..Default::default()
};
let captured = observation
.request_diagnostic
.lock()
.unwrap_or_else(|e| e.into_inner());
if let Some(captured) = captured.as_ref() {
captured.record(&axes);
} else {
crate::diagnostics::boundary(
observation.occurrence(),
&axes,
Some((&observation.context, &observation.event_context)),
);
}
}
}
fn record_post_release_capacity_snapshot(
observer: &Observer,
context: &CallContext,
event: &EventContext,
usage: saddle_admission::ProfuseGwCapacityUsageSnapshot,
elapsed_ms: u64,
) {
let capacity = usage.capacity();
let bottleneck = match capacity.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,
capacity.cpu_budget(),
usage.cpu_used(),
),
(
CapacityDimension::Memory,
capacity.memory_budget_bytes(),
usage.memory_used(),
),
(
CapacityDimension::Database,
capacity.database_budget(),
usage.database_used(),
),
(
CapacityDimension::ProfuseContract,
capacity.profusecontract_budget(),
usage.profusecontract_used(),
),
];
let limit = u64::try_from(capacity.active_limit()).unwrap_or(u64::MAX);
for (dimension, budget, used) in dimensions {
observer.record_capacity(
context,
event,
CapacityObservation::accepted_dimension(
dimension,
u64::try_from(budget).unwrap_or(u64::MAX),
limit,
u64::try_from(used).unwrap_or(u64::MAX),
bottleneck.clone(),
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(self.admission.request_timeout()) {
Ok((unix_ms, deadline)) => ProfuseGwManagedDeadline {
unix_ms,
timer: deadline.into_sleep(),
reserved_storage: None,
},
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,
scope_stop: 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 try_reserved_rejection<O:Send+'static,T:Send+'static,F>(
&self,application:saddle_core::request_context::ContextLabel,
output:Option<saddle_observability::EmergencyDiagnosticHandle>,body:std::alloc::Layout,
indirect:&[(std::alloc::Layout,usize)],owner:O,make:F,
)->Result<(crate::request_task::reserved::ReservedRequestRoot,crate::request_task::reserved::ReservedTaskFuture<O,T>,crate::request_task::reserved::ReservedTaskTicket<O>),
(O,F,crate::request_task::reserved::ReservedContextError)>
where F:for<'a> FnOnce(&'a mut O,crate::request_task::reserved::ReservedTaskContext)
->crate::request_task::reserved::ReservedBorrowedFuture<'a,T>+Send+'static {
let process=self.shared.process.lock().unwrap_or_else(|p|p.into_inner());
let Some(process)=process.as_ref() else {
return Err((owner,make,crate::request_task::reserved::ReservedContextError::Storage(AdmissionError::AccountClosed)));
};
crate::request_task::reserved::try_prepare_rejection_task(&process.admission.verified_profile(),application,output,body,indirect,owner,make)
}
#[doc(hidden)]
pub fn try_reserved_task_storage(&self)->Result<crate::request_task::reserved_set::ReservedProcessTaskStorage,AdmissionError> {
let process=self.shared.process.lock().unwrap_or_else(|p|p.into_inner());
let process=process.as_ref().ok_or(AdmissionError::AccountClosed)?;
crate::request_task::reserved_set::ReservedProcessTaskStorage::acquire(&process.admission)
}
#[doc(hidden)]
pub fn try_reserved_dispatch<X:Send+'static,T:Send+'static,F>(
&self,application:saddle_core::request_context::ContextLabel,
output:Option<saddle_observability::EmergencyDiagnosticHandle>,
body:std::alloc::Layout,indirect:&[(std::alloc::Layout,usize)],input:X,make:F,
)->ReservedDispatchOutcome<X,T,F>
where F:for<'a> FnOnce(&'a mut ReservedDispatchOwner<X>,crate::request_task::reserved::ReservedTaskContext)
->crate::request_task::reserved::ReservedBorrowedFuture<'a,T>+Send+'static {
use crate::request_task::reserved::ReservedTaskAllocation;
let process=self.shared.process.lock().unwrap_or_else(|p|p.into_inner());
let Some(process)=process.as_ref() else {
return ReservedDispatchOutcome::Rejected{input,make,reason:ReservedDispatchRejection::Admission(AdmissionError::AccountClosed)};
};
let (permit,receipt)=match process.admission.verified_profile().try_admit_observed() {
ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit,receipt)=>(permit,receipt),
ProfuseGwLightweightObservedAdmissionOutcome::CapacityRejected(receipt)=>return ReservedDispatchOutcome::Rejected{
input,make,reason:ReservedDispatchRejection::Capacity(ProfuseGwAdmissionEvent{receipt})},
ProfuseGwLightweightObservedAdmissionOutcome::Stop(error)=>return ReservedDispatchOutcome::Rejected{
input,make,reason:ReservedDispatchRejection::Admission(error)},
};
let admission=ProfuseGwAdmissionEvent{receipt};
let storage=match ReservedTaskAllocation::for_dispatch::<X,T,F>(&permit,body,indirect) {
Ok(storage)=>storage,
Err(error)=>{drop(permit);return ReservedDispatchOutcome::Rejected{input,make,reason:ReservedDispatchRejection::Admitted{admission,failure:ReservedDispatchConstructionFailure::Storage(error)}};}
};
let (unix_ms,deadline)=match default_profusegw_deadline(process.admission.request_timeout()) {
Ok(deadline)=>deadline,
Err(error)=>{drop((storage,permit));return ReservedDispatchOutcome::Rejected{input,make,reason:ReservedDispatchRejection::Admitted{admission,failure:ReservedDispatchConstructionFailure::Deadline(error)}};}
};
let Some((db_request,db_execution))=process.db_requests.issue_request() else {
drop((storage,permit));return ReservedDispatchOutcome::Rejected{input,make,reason:ReservedDispatchRejection::Admitted{admission,failure:ReservedDispatchConstructionFailure::RequestIdentity(AdmissionError::InvalidConfiguration)}};
};
let dispatch=ProfuseGwManagedDispatch{execution:permit.into_execution(),
deadline:ProfuseGwManagedDeadline{unix_ms,timer:deadline.into_sleep(),reserved_storage:None},db_request,db_execution,observation:None,scope_stop:None};
match storage.construct(dispatch,input,admission,application,output,make) {
Ok((root,future,ticket))=>ReservedDispatchOutcome::Ready{root,future,ticket},
Err((dispatch,input,admission,make,error))=>{
dispatch.cancel();
ReservedDispatchOutcome::Rejected{input,make,reason:ReservedDispatchRejection::Admitted{admission,failure:ReservedDispatchConstructionFailure::Context(error)}}
}
}
}
#[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),
}
}
}
#[doc(hidden)]
pub enum ReservedDispatchRejection {
Capacity(ProfuseGwAdmissionEvent),
Admission(AdmissionError),
Admitted { admission: ProfuseGwAdmissionEvent, failure: ReservedDispatchConstructionFailure },
}
#[doc(hidden)]
pub enum ReservedDispatchConstructionFailure {
Storage(AdmissionError),
RequestIdentity(AdmissionError),
Deadline(AbsoluteDeadlineError),
Context(crate::request_task::reserved::ReservedContextError),
}
#[doc(hidden)]
pub enum ReservedDispatchOutcome<X,T,F> {
Ready {
root:crate::request_task::reserved::ReservedRequestRoot,
future:crate::request_task::reserved::ReservedTaskFuture<ReservedDispatchOwner<X>,T>,
ticket:crate::request_task::reserved::ReservedTaskTicket<ReservedDispatchOwner<X>>,
},
Rejected {input:X,make:F,reason:ReservedDispatchRejection},
}
impl ProfuseGwTerminalProcess {
fn finalize(
process: Option<ProfuseGwRuntimeProcess>,
lifecycle: saddle_core::Result<()>,
) -> Self {
let finalization = match process {
Some(process) => process
.finish()
.map_err(ProfuseGwCoordinatorFailure::Finalization),
None => Err(ProfuseGwCoordinatorFailure::Lifecycle(
lifecycle_recovery_error(),
)),
};
if let Err(failure) = &finalization {
let code = match failure {
ProfuseGwCoordinatorFailure::Lifecycle(_) => "runtime.process_owner_missing",
ProfuseGwCoordinatorFailure::Finalization(AdmissionError::InvalidConfiguration) => {
"runtime.process_finalization_invalid_configuration"
}
ProfuseGwCoordinatorFailure::Finalization(
AdmissionError::ShutdownNotReconciled { .. },
) => "runtime.process_shutdown_not_reconciled",
ProfuseGwCoordinatorFailure::Finalization(AdmissionError::ProcessUnhealthy(_)) => {
"runtime.process_ledger_unhealthy"
}
ProfuseGwCoordinatorFailure::Finalization(
AdmissionError::LedgerOwnersOutstanding { .. },
) => "runtime.process_owners_outstanding",
_ => "runtime.process_resource_finalization_failed",
};
crate::diagnostics::cleanup(
crate::diagnostics::attach(
saddle_core::SaddleError::new(
saddle_core::ErrorKind::Infrastructure,
code,
"the process resource check did not complete cleanly",
),
saddle_core::DiagnosticStage::FinalizerResource,
code,
),
lifecycle.as_ref().err(),
saddle_core::DiagnosticStage::FinalizerResource,
);
}
Self::from_results(lifecycle, finalization)
}
fn from_results(
lifecycle: saddle_core::Result<()>,
finalization: Result<(), ProfuseGwCoordinatorFailure>,
) -> Self {
let result = match lifecycle {
Ok(()) => finalization,
Err(error) => Err(ProfuseGwCoordinatorFailure::Lifecycle(error)),
};
Self { result }
}
fn finish(self) -> Result<(), ProfuseGwCoordinatorFailure> {
self.result
}
}
impl ProfuseGwManagedDispatch {
#[doc(hidden)]
pub fn bind_ingress_deadline(
&mut self,
deadline_unix_ms: i64,
) -> Result<(), AbsoluteDeadlineError> {
let deadline = u64::try_from(deadline_unix_ms)
.map_err(|_| AbsoluteDeadlineError::OutOfRange)
.and_then(verify_absolute_deadline)?;
self.deadline.unix_ms = deadline_unix_ms;
self.deadline.timer.as_mut().reset(deadline.into_instant());
Ok(())
}
#[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.install_observation(observer, context, event_context);
self
}
#[doc(hidden)]
pub fn bind_reserved_observation(
mut self,
observer: Observer,
context: CallContext,
event_context: EventContext,
admission: ProfuseGwAdmissionEvent,
view: &crate::request_task::reserved::ReservedRequestView,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) -> (Self, saddle_observability::AdmissionCapacitySubmission) {
let submission = view.record_admission(admission, &observer, output,
saddle_observability::AdmissionConstruction::Ready);
self.install_observation(observer, context, event_context);
(self, submission)
}
fn install_observation(&mut self, observer: Observer, context: CallContext, event_context: EventContext) {
self.observation = Some(ProfuseGwRequestObservation {
observer,
context,
event_context,
diagnostic: Mutex::new(None),
request_diagnostic: Mutex::new(None),
cleanup_failed: std::sync::atomic::AtomicBool::new(false),
});
}
#[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_request_stop(
self.observation.as_ref(),
ProfuseGwScopeStop::Cancelled,
false,
saddle_core::DiagnosticStage::RequestHandler,
);
drop(self.deadline);
record_terminal(self.observation, self.execution.cancel_observed());
}
#[doc(hidden)]
pub fn timeout(self) {
record_request_stop(
self.observation.as_ref(),
ProfuseGwScopeStop::TimedOut,
false,
saddle_core::DiagnosticStage::RequestHandler,
);
drop(self.deadline);
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,
scope_stop: self.scope_stop,
}
}
}
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,
scope_stop,
} = 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,
scope_stop,
},
value,
});
}
};
let value = receipt.into_value();
drop(deadline);
record_terminal(observation, execution.cancel_observed());
Ok(value)
}
impl ProfuseGwConcreteDbRequestLease {
#[doc(hidden)]
pub async fn scope_checkpoint<C: Future<Output = ()> + ?Sized>(
&mut self,
cancel: Pin<&mut C>,
) -> Result<(), ProfuseGwScopeStop> {
scope_checkpoint(&mut self.deadline, &mut self.scope_stop, cancel).await
}
#[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,
scope_stop: self.scope_stop,
}
}
#[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,
scope_stop: self.scope_stop,
},
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,
}
async fn scope_checkpoint<C: Future<Output = ()> + ?Sized>(
deadline: &mut ProfuseGwManagedDeadline,
stop: &mut Option<ProfuseGwScopeStop>,
mut cancel: Pin<&mut C>,
) -> Result<(), ProfuseGwScopeStop> {
let mut yielded = false;
poll_fn(|cx| {
if let Some(reason) = *stop {
return Poll::Ready(Err(reason));
}
let reason = if cancel.as_mut().poll(cx).is_ready() {
Some(ProfuseGwScopeStop::Cancelled)
} else if tokio::time::Instant::now() >= deadline.timer.deadline()
|| deadline.timer.as_mut().poll(cx).is_ready()
{
Some(ProfuseGwScopeStop::TimedOut)
} else {
None
};
if let Some(reason) = reason {
*stop = Some(reason);
return Poll::Ready(Err(reason));
}
if !yielded {
yielded = true;
cx.waker().wake_by_ref();
return Poll::Pending;
}
Poll::Ready(Ok(()))
})
.await
}
#[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 Some(stop) = lease.scope_stop {
return Poll::Ready(match stop {
ProfuseGwScopeStop::Cancelled => OperationPoll::Cancelled,
ProfuseGwScopeStop::TimedOut => OperationPoll::TimedOut,
});
}
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() {
lease.scope_stop = Some(ProfuseGwScopeStop::Cancelled);
record_request_stop(
lease.observation.as_ref(),
ProfuseGwScopeStop::Cancelled,
false,
saddle_core::DiagnosticStage::RequestDb,
);
return Poll::Ready(OperationPoll::Cancelled);
}
if lease.deadline.timer.as_mut().poll(context).is_ready() {
lease.scope_stop = Some(ProfuseGwScopeStop::TimedOut);
record_request_stop(
lease.observation.as_ref(),
ProfuseGwScopeStop::TimedOut,
false,
saddle_core::DiagnosticStage::RequestDb,
);
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<ProfuseGwPostDatabaseManagedOwner<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 (continuation, disposition, value) = receipt.into_scope_continuation();
let admission = match disposition {
DbPhysicalDisposition::Returned => completion.finalization.connection_returned(),
DbPhysicalDisposition::Discarded => completion.finalization.connection_discarded(),
};
Ok(ProfuseGwPostDatabaseManagedOwner {
terminal: ProfuseGwPostDatabaseRequestTerminal {
admission,
deadline: completion.deadline,
observation: completion.observation,
scope_stop: completion.scope_stop,
next_scope: serial_scope::NextScope::Continuation(continuation),
},
value,
})
}
#[doc(hidden)]
pub fn finish_profusegw_after_database(terminal: ProfuseGwPostDatabaseRequestTerminal) {
let ProfuseGwPostDatabaseRequestTerminal {
admission,
deadline,
observation,
scope_stop: _,
next_scope: _,
} = terminal;
drop(deadline);
record_terminal(observation, admission.finish_observed());
}
#[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>>,
{
run_profusegw_owned_application_inner(process, factory, None).0
}
#[doc(hidden)]
pub fn run_profusegw_owned_application_with_diagnostics<Factory, FactoryFuture>(
process: ProfuseGwRuntimeProcess,
output: saddle_observability::EmergencyDiagnostics,
factory: Factory,
) -> (
ProfuseGwTerminalProcess,
crate::diagnostics::RuntimeDiagnosticExit,
)
where
Factory: FnOnce(ProfuseGwProcessLease) -> FactoryFuture,
FactoryFuture: Future<Output = saddle_core::Result<Application>>,
{
let (terminal, exit) = run_profusegw_owned_application_inner(process, factory, Some(output));
(
terminal,
exit.expect("the input diagnostic owner is always closed"),
)
}
fn run_profusegw_owned_application_inner<Factory, FactoryFuture>(
process: ProfuseGwRuntimeProcess,
factory: Factory,
output: Option<saddle_observability::EmergencyDiagnostics>,
) -> (
ProfuseGwTerminalProcess,
Option<crate::diagnostics::RuntimeDiagnosticExit>,
)
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::finalize(Some(process), Err(error)),
output.map(|o| crate::diagnostics::close_output(o, None)),
);
}
};
let shared = Arc::new(SharedRuntimeProcess {
process: Mutex::new(Some(process)),
});
let lease = ProfuseGwProcessLease {
shared: Arc::clone(&shared),
};
let outcome = crate::diagnostics::catching(
saddle_core::DiagnosticStage::StartupListener,
"runtime.application",
None,
|| {
runtime.block_on(async move {
let signal = ShutdownSignal::register()?;
let application = factory(lease).await?;
let finalizer = application.pending_driver_finalizer();
let shutdown_deadline = application.shutdown_deadline_handle();
let lifecycle_observer = application.lifecycle_observer_handle();
let result = application.run_until_shutdown(signal.wait()).await;
Ok::<_, saddle_core::SaddleError>((
finalizer,
shutdown_deadline,
lifecycle_observer,
result,
))
})
},
);
let mut output_deadline = None;
let lifecycle = match outcome {
Ok(Ok((finalizer, shutdown_deadline, lifecycle_observer, result))) => {
output_deadline = *shutdown_deadline
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner());
finalizer.finish(runtime, result, output_deadline, lifecycle_observer)
}
Ok(Err(error)) => {
drop(runtime);
Err(error)
}
Err(diagnostic) => {
drop(runtime);
Err(saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal,
"runtime.profusegw_lifecycle_panicked",
"the ProfuseGW Application lifecycle panicked before terminal recovery",
)
.with_diagnostic(diagnostic))
}
};
let process = shared
.process
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
.take();
let terminal = ProfuseGwTerminalProcess::finalize(process, lifecycle);
let exit = output.map(|o| crate::diagnostics::close_output(o, output_deadline));
(terminal, exit)
}
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(
timeout: Duration,
) -> 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(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 {
#[test]
fn reserved_scope_cancel_retains_original_owner_and_not_used() {
use super::*;
use crate::request_task::reserved::{ReservedTaskContext,ReservedBorrowedFuture,ReservedRequestFailure};
type Owner=ReservedDispatchOwner<Option<ProfuseGwSerialScope<std::future::Ready<()>>>>;
fn body(owner:&mut Owner,context:ReservedTaskContext)->impl Future<Output=Result<u32,ReservedRequestFailure>>+Send+'_ {
async move {
owner.1=Some(owner.0.take().unwrap().into_parameter_scope(std::future::ready(())));
let construction=owner.1.as_mut().unwrap().parameter_construction().ok().unwrap();
assert!(matches!(construction.supervise_reserved(&context,None,std::future::ready(41)).await,
Err(ProfuseGwReservedScopeFailure::Execution(ProfuseGwScopeFailure::Stopped(ProfuseGwScopeStop::Cancelled)))));
Ok(41)
}
}
fn factory<'a>(owner:&'a mut Owner,context:ReservedTaskContext)->ReservedBorrowedFuture<'a,u32>{Box::pin(body(owner,context))}
fn layout<I,R>(_:impl FnOnce(I)->R)->std::alloc::Layout {std::alloc::Layout::new::<R>()}
tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap().block_on(async {
let shared=Arc::new(SharedRuntimeProcess{process:Mutex::new(Some(ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process())))});
let lease=ProfuseGwProcessLease{shared:shared.clone()};
let startup=lease.take_database_startup_half().unwrap();
let demand=layout(|(owner,context):(&'static mut Owner,ReservedTaskContext)|body(owner,context));
let outcome=lease.try_reserved_dispatch(saddle_core::request_context::ContextLabel::checked("scope-test").unwrap(),None,demand,&[(std::alloc::Layout::new::<std::future::Ready<()>>(),1)],None,factory);
let ReservedDispatchOutcome::Ready{root,future,mut ticket,..}=outcome else {panic!("isolated scope fit")};
let mut tasks=tokio::task::JoinSet::new();let handle=tasks.spawn(future);ticket.bind(handle.id()).unwrap();
let joined=ticket.complete(tasks.join_next().await.unwrap()).ok().unwrap();
drop(root);
let mut recovered=joined.recover(Default::default()).ok().unwrap();
assert_eq!(recovered.result.take().unwrap().ok(),Some(41));
assert!(recovered.primary.is_some());
recovered.owner.1.take().unwrap().finish_unentered_response(()).ok().unwrap();
drop(recovered);
match lease.try_admit(){ProfuseGwCoordinatorAdmissionOutcome::Ready(d,_)=>d.cancel(),_=>panic!("scope storage and account returned")}
drop((startup,lease));shared.process.lock().unwrap().take().unwrap().finish().unwrap();
});
}
#[test]
fn reserved_dispatch_keeps_original_deadline_owner_until_matching_join() {
use super::*;
use crate::request_task::reserved::tests::process;
use saddle_core::request_context::ContextLabel;
tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap().block_on(async {
let shared=Arc::new(SharedRuntimeProcess{process:Mutex::new(Some(ProfuseGwRuntimeProcess::new(process())))});
let lease=ProfuseGwProcessLease{shared:shared.clone()};
let startup=lease.take_database_startup_half().unwrap();
let entered=Arc::new(std::sync::atomic::AtomicBool::new(false));
let seen=entered.clone();
let outcome=lease.try_reserved_dispatch(ContextLabel::checked("app").unwrap(),None,
std::alloc::Layout::new::<std::future::Pending<Result<u32,crate::request_task::reserved::ReservedRequestFailure>>>(),&[],(),
move |owner,_| {
assert!(owner.0.as_ref().unwrap().deadline_unix_ms()>0);
seen.store(true,std::sync::atomic::Ordering::SeqCst);
Box::pin(std::future::pending::<Result<u32,crate::request_task::reserved::ReservedRequestFailure>>())
});
let ReservedDispatchOutcome::Ready{root,future,mut ticket}=outcome else {panic!("isolated real dispatch fit")};
let mut tasks=tokio::task::JoinSet::new();
let abort=tasks.spawn(future);
ticket.bind(abort.id()).unwrap();
while !entered.load(std::sync::atomic::Ordering::SeqCst){tokio::task::yield_now().await;}
abort.abort();
let joined=ticket.complete(tasks.join_next().await.unwrap()).ok().unwrap();
drop(root);
let mut recovered=joined.recover(saddle_observability::root_diagnostic::RootOutcomeFacts::default()).ok().unwrap();
assert!(recovered.primary.is_some());
let dispatch=recovered.owner.0.take().unwrap();
assert!(matches!(lease.try_admit(),ProfuseGwCoordinatorAdmissionOutcome::CapacityRejected(_)));
finish_profusegw_without_database(dispatch,41).ok().unwrap();
match lease.try_admit(){ProfuseGwCoordinatorAdmissionOutcome::Ready(dispatch,_)=>dispatch.cancel(),_=>panic!("same profile recovers")}
drop((recovered,startup,lease));
shared.process.lock().unwrap().take().unwrap().finish().unwrap();
});
}
use std::sync::{
Arc,
atomic::{AtomicUsize, Ordering},
};
use super::default_profusegw_deadline;
#[test]
fn terminal_source_precedes_output_close() {
const CHILD: &str = "RUNTIME_TERMINAL_SOURCE_CHILD";
if let Some(path) = std::env::var_os(CHILD) {
let output = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(
std::path::PathBuf::from(&path),
saddle_observability::Rotation::Daily,
),
)
.unwrap();
assert!(crate::diagnostics::install_output(output.handle()).is_ok());
let primary = crate::diagnostics::attach(
saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal,
"runtime.test_primary",
"primary",
),
saddle_core::DiagnosticStage::ShutdownComponent,
"runtime.test_primary",
);
let primary_id = primary.diagnostic().unwrap().id();
crate::diagnostics::report(&primary);
let terminal = super::ProfuseGwTerminalProcess::finalize(None, Err(primary));
let exit = crate::diagnostics::close_output(
output,
Some(std::time::Instant::now() + std::time::Duration::from_secs(2)),
);
assert_eq!(
exit.shutdown,
saddle_observability::DiagnosticShutdown::Finished
);
assert_eq!(exit.snapshot.enqueued, exit.snapshot.written);
assert_eq!(exit.snapshot.dropped, 0);
assert!(matches!(terminal.finish(),
Err(super::ProfuseGwCoordinatorFailure::Lifecycle(e))
if e.code() == "runtime.test_primary"));
let log = std::fs::read_to_string(
std::path::PathBuf::from(path).join("saddle.emergency.log"),
)
.unwrap();
let rows: Vec<serde_json::Value> = log
.lines()
.map(|line| serde_json::from_str(line).unwrap())
.collect();
let cleanup = rows
.iter()
.find(|v| v["diagnostic"]["causes"][0]["code"] == "runtime.process_owner_missing")
.expect("resource source was written before close");
assert_eq!(cleanup["diagnostic"]["primary_diagnostic_id"], primary_id);
return;
}
let path =
std::env::temp_dir().join(format!("saddle-terminal-source-{}", std::process::id()));
std::fs::create_dir(&path).unwrap();
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"profusegw::tests::terminal_source_precedes_output_close",
])
.env(CHILD, &path)
.spawn()
.unwrap();
let started = std::time::Instant::now();
let status = loop {
if let Some(status) = child.try_wait().unwrap() {
break status;
}
if started.elapsed() > std::time::Duration::from_secs(10) {
child.kill().unwrap();
child.wait().unwrap();
panic!("terminal child exceeded test watchdog");
}
std::thread::sleep(std::time::Duration::from_millis(10));
};
assert!(status.success());
std::fs::remove_file(path.join("saddle.emergency.log")).unwrap();
std::fs::remove_dir(path).unwrap();
}
#[test]
fn terminal_resource_result_never_overwrites_primary() {
let primary = saddle_core::SaddleError::new(
saddle_core::ErrorKind::Internal,
"runtime.test_primary",
"primary",
);
let terminal = super::ProfuseGwTerminalProcess::from_results(
Err(primary),
Err(super::ProfuseGwCoordinatorFailure::Finalization(
saddle_admission::AdmissionError::AccountClosed,
)),
);
assert!(matches!(terminal.finish(),
Err(super::ProfuseGwCoordinatorFailure::Lifecycle(e))
if e.code() == "runtime.test_primary"));
let terminal = super::ProfuseGwTerminalProcess::from_results(
Ok(()),
Err(super::ProfuseGwCoordinatorFailure::Finalization(
saddle_admission::AdmissionError::AccountClosed,
)),
);
assert!(matches!(
terminal.finish(),
Err(super::ProfuseGwCoordinatorFailure::Finalization(
saddle_admission::AdmissionError::AccountClosed
))
));
}
fn checkpoint_deadline(at: tokio::time::Instant) -> super::ProfuseGwManagedDeadline {
super::ProfuseGwManagedDeadline {
unix_ms: 123,
timer: Box::pin(tokio::time::sleep_until(at)),
reserved_storage:None,
}
}
#[tokio::test]
async fn scope_checkpoint_expired_before_commit_is_sticky() {
let at = tokio::time::Instant::now();
let mut deadline = checkpoint_deadline(at);
let mut stop = None;
let mut cancel = Box::pin(std::future::pending());
assert_eq!(
super::scope_checkpoint(&mut deadline, &mut stop, cancel.as_mut()).await,
Err(super::ProfuseGwScopeStop::TimedOut)
);
assert_eq!(deadline.timer.deadline(), at);
let mut never_poll = Box::pin(std::future::poll_fn(|_| -> std::task::Poll<()> {
panic!("sticky stop repolled cancellation")
}));
assert_eq!(
super::scope_checkpoint(&mut deadline, &mut stop, never_poll.as_mut()).await,
Err(super::ProfuseGwScopeStop::TimedOut)
);
}
#[tokio::test]
async fn scope_checkpoint_drop_preserves_deadline_and_rechecks_cancel() {
use std::future::Future;
let at = tokio::time::Instant::now() + std::time::Duration::from_secs(60);
let mut deadline = checkpoint_deadline(at);
let mut stop = None;
let mut cancel = Box::pin(std::future::pending());
{
let checkpoint = super::scope_checkpoint(&mut deadline, &mut stop, cancel.as_mut());
let mut checkpoint = Box::pin(checkpoint);
let mut cx = std::task::Context::from_waker(std::task::Waker::noop());
assert!(checkpoint.as_mut().poll(&mut cx).is_pending());
}
assert_eq!(deadline.timer.deadline(), at);
assert_eq!(stop, None);
let mut cancelled = Box::pin(std::future::ready(()));
assert_eq!(
super::scope_checkpoint(&mut deadline, &mut stop, cancelled.as_mut()).await,
Err(super::ProfuseGwScopeStop::Cancelled)
);
assert_eq!(deadline.unix_ms, 123);
}
#[tokio::test]
async fn scope_checkpoint_ready_loop_yields_and_observes_expiry() {
let (sender, receiver) = tokio::sync::oneshot::channel();
tokio::spawn(async move {
let _ = sender.send(());
});
let at = tokio::time::Instant::now() + std::time::Duration::from_secs(60);
let mut deadline = checkpoint_deadline(at);
let mut stop = None;
let mut cancel = Box::pin(async move {
let _ = receiver.await;
});
let mut cancelled = false;
for _ in 0..10_000 {
if super::scope_checkpoint(&mut deadline, &mut stop, cancel.as_mut())
.await
.is_err()
{
cancelled = true;
break;
}
std::future::ready(()).await;
}
assert!(cancelled, "Ready loop starved the cancellation task");
assert_eq!(stop, Some(super::ProfuseGwScopeStop::Cancelled));
let at = tokio::time::Instant::now() + std::time::Duration::from_millis(2);
let mut deadline = checkpoint_deadline(at);
let mut stop = None;
let mut cancel = Box::pin(std::future::pending());
for _ in 0..1_000_000 {
if super::scope_checkpoint(&mut deadline, &mut stop, cancel.as_mut())
.await
.is_err()
{
break;
}
std::future::ready(()).await;
}
assert_eq!(stop, Some(super::ProfuseGwScopeStop::TimedOut));
assert_eq!(deadline.timer.deadline(), at);
}
#[test]
fn budget_default_deadline_is_runtime_owned_and_live() {
let (unix_ms, _owner) =
default_profusegw_deadline(std::time::Duration::from_millis(5000)).unwrap();
assert!(unix_ms > 0);
}
#[test]
fn configured_timeout_drives_existing_absolute_deadline() {
use std::time::{Duration, SystemTime, UNIX_EPOCH};
for timeout in [1500u64, 5000, 6000] {
let before = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis();
let (unix_ms, _owner) =
default_profusegw_deadline(Duration::from_millis(timeout)).unwrap();
let after = SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis();
assert!(
(before + timeout as u128..=after + timeout as u128).contains(&(unix_ms as u128))
);
}
assert!(default_profusegw_deadline(Duration::ZERO).is_err());
assert!(default_profusegw_deadline(Duration::MAX).is_err());
assert!(default_profusegw_deadline(Duration::from_millis(u64::MAX)).is_err());
}
#[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);
}
}