use saddle_core::{DbPhysicalProcessCapability, OperationOutcome, PhysicalDispositionFact};
use saddle_observability::root_diagnostic::RootOutcomeFacts;
use saddle_observability::{EmergencyDiagnosticHandle, Observer};
use saddle_runtime::profusegw::*;
use saddle_runtime::request_task::reserved::*;
use std::future::Future;
#[allow(dead_code)]
async fn database_scope<C: Future<Output = ()>, F: Future>(
owner: &mut Option<ProfuseGwSerialScope<C>>,
context: &ReservedTaskContext,
observer: &Observer,
output: Option<&EmergencyDiagnosticHandle>,
db_physical: &DbPhysicalProcessCapability,
body: F,
) -> Result<
ProfuseGwSuspendedScope<C, Result<F::Output, ProfuseGwReservedScopeFailure>>,
ProfuseGwReservedScopeFailure,
> {
let driver = owner
.as_mut()
.expect("original externally retained owner")
.driver();
let (supervisor, observation) = driver.prepare_reserved(context)?;
let stage = observation
.start_database_stage(observer, output)
.expect("one interval");
let outcome = supervisor.supervise(output, body).await;
let supervision = outcome.completion.supervision();
let operation = match &outcome.result {
Ok(_) => OperationOutcome::Succeeded, Err(ProfuseGwReservedScopeFailure::Preparation(_)) => OperationOutcome::Rejected,
Err(ProfuseGwReservedScopeFailure::Execution(ProfuseGwScopeFailure::AlreadySupervised)) => {
OperationOutcome::Rejected
}
Err(ProfuseGwReservedScopeFailure::Execution(ProfuseGwScopeFailure::Panicked)) => {
OperationOutcome::Panicked
}
Err(ProfuseGwReservedScopeFailure::Execution(ProfuseGwScopeFailure::Stopped(
ProfuseGwScopeStop::TimedOut,
))) => OperationOutcome::TimedOut,
Err(_) => OperationOutcome::Cancelled,
};
drop(driver); let (physical, execution) = owner.take().unwrap().into_physical_finalization();
let receipt = db_physical
.connection_discarded(execution, outcome.result)
.ok()
.expect("same process");
let suspended = physical.complete(receipt).ok().expect("same original pair");
let mut facts = RootOutcomeFacts::default();
facts.axes.operation = operation;
facts.axes.physical = PhysicalDispositionFact::Discarded;
let _stop = supervision.stop;
stage
.finish(outcome.completion, facts)
.ok()
.expect("same invocation, not merely same request");
Ok(suspended)
}
fn main() {
for (name, layout) in supervised_scope_layouts() {
println!("{name}: size={} align={}", layout.size(), layout.align());
}
println!(
"supervisor={} outcome_u32={}",
std::mem::size_of::<
ProfuseGwReservedScopeSupervisor<'static, 'static, std::future::Pending<()>>,
>(),
std::mem::size_of::<ProfuseGwReservedScopeOutcome<u32>>()
);
println!("RG_SCOPE_CONSUMER COMPILED physical=SIMULATED formal_DB=NOT_RUN");
}
#[allow(dead_code)]
enum DemoError {
OriginalRecorded,
RetentionUnavailable(ReservedRequestFailure),
Preparation(ReservedContextError),
}
struct DemoSession<'a, 'scope, C> {
driver: &'a ProfuseGwScopeDriver<'scope, C>,
observation: &'a ReservedScopeObservation,
output: Option<&'a EmergencyDiagnosticHandle>,
}
impl<C: Future<Output = ()>> DemoSession<'_, '_, C> {
async fn step(
&mut self,
name: saddle_core::request_context::RegisteredContextOperation,
fail: bool,
) -> Result<u32, DemoError> {
let operation = self
.observation
.operation(name)
.map_err(DemoError::Preparation)?;
if self.driver.checkpoint().await.is_err() {
return Err(DemoError::OriginalRecorded);
}
if fail {
let error = std::io::Error::new(
std::io::ErrorKind::PermissionDenied,
"original driver error",
);
let diagnostic = saddle_core::Diagnostic::capture(
saddle_core::DiagnosticCategory::ExpectedRejection,
saddle_core::CaptureSite::FirstObserved,
saddle_core::DiagnosticCause::new(
saddle_core::DiagnosticStage::RequestDb,
saddle_core::DiagnosticCode::new("db.operation").unwrap(),
),
);
operation
.source_existing_error(
&error,
diagnostic,
saddle_core::DiagnosticCode::new("db.operation").unwrap(),
self.output,
Default::default(),
)
.map_err(DemoError::RetentionUnavailable)?;
return Err(DemoError::OriginalRecorded);
}
Ok(41)
}
}
#[allow(dead_code)]
async fn transaction_shape<C: Future<Output = ()>>(
owner: &mut Option<ProfuseGwSerialScope<C>>,
context: &ReservedTaskContext,
observer: &Observer,
output: Option<&EmergencyDiagnosticHandle>,
physical: &DbPhysicalProcessCapability,
fail_write: bool,
) -> Result<
ProfuseGwSuspendedScope<C, Result<Result<u32, DemoError>, ProfuseGwReservedScopeFailure>>,
ProfuseGwReservedScopeFailure,
> {
use saddle_core::request_context::RegisteredContextOperation as Op;
let driver = owner.as_mut().unwrap().driver();
let (begin, begin_obs) = driver.prepare_reserved(context)?;
let begin_stage = begin_obs.start_database_stage(observer, output).unwrap();
let begun = begin
.supervise(output, async { driver.checkpoint().await })
.await;
let (body, body_obs) = driver.prepare_reserved(context)?;
let body_stage = body_obs.start_database_stage(observer, output).unwrap();
let mut session = DemoSession {
driver: &driver,
observation: &body_obs,
output,
};
let result = body
.supervise(output, async {
let read = session
.step(Op::checked("db.registered.read").unwrap(), false)
.await?;
let chosen = if read == 41 {
"db.registered.write"
} else {
"db.registered.other"
};
tokio::task::yield_now().await;
session.step(Op::checked(chosen).unwrap(), fail_write).await
})
.await;
drop(session);
let rollback = !matches!(&result.result, Ok(Ok(_)));
let (end, end_obs) = driver.prepare_reserved(context)?;
let end_stage = end_obs.start_database_stage(observer, output).unwrap();
let ended = end
.supervise(output, async {
let _operation = end_obs.operation(
Op::checked(if rollback { "db.rollback" } else { "db.commit" }).unwrap(),
);
driver.checkpoint().await
})
.await;
drop(driver);
let (completion, execution) = owner.take().unwrap().into_physical_finalization();
let receipt = physical
.connection_discarded(execution, result.result)
.ok()
.unwrap();
let suspended = completion.complete(receipt).ok().unwrap();
let mut facts = RootOutcomeFacts::default();
facts.axes.operation = if rollback {
OperationOutcome::Rejected
} else {
OperationOutcome::Succeeded
};
facts.axes.physical = PhysicalDispositionFact::Discarded;
begin_stage.finish(begun.completion, facts).ok().unwrap();
body_stage.finish(result.completion, facts).ok().unwrap();
end_stage.finish(ended.completion, facts).ok().unwrap();
Ok(suspended)
}