saddle-runtime 0.3.24

Saddle managed asynchronous runtime and lifecycle
Documentation
//! External DB-shaped composition example. Physical return is SIMULATED here;
//! Database must actually finish its owning connection/phase before signing it.
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();
    // Failed preparation leaves that exact physical owner in its original slot.
    let (supervisor, observation) = driver.prepare_reserved(context)?;
    let stage = observation
        .start_database_stage(observer, output)
        .expect("one interval");
    // DB may use observation.diagnostic_view() to capture its actual Error,
    // then retain_database_failure(original) without leaking/clearing raw slots.
    // This future may borrow the original driver for step checkpoints.
    let outcome = supervisor.supervise(output, body).await;
    let supervision = outcome.completion.supervision();
    let operation = match &outcome.result {
        Ok(_) => OperationOutcome::Succeeded, // Ready T is never rewritten by cleanup.
        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); // End the original execution borrow BEFORE physical ownership move.
    let (physical, execution) = owner.take().unwrap().into_physical_finalization();
    // NO database I/O here. This is explicitly a simulated Database disposition.
    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;
    // Supervision stop remains independently available even with Ready(T).
    let _stop = supervision.stop;
    stage
        .finish(outcome.completion, facts)
        .ok()
        .expect("same invocation, not merely same request");
    // Original failures remain under the task ticket through its final recover.
    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");
}

// This compiles the full transaction borrowing topology, not a production DB
// replacement. D supplies its real kernel, connection, phase and Unknown facts.
#[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)?;
        // Actual DB must consume the checkpoint result, never run SQL after Stop.
        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,
    };
    // One supervisor covers ALL business waits; session borrows the immutable
    // invocation capability, not a mutable task context or nested supervisor.
    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
            // D's owning kernel executes rollback/commit here, retaining CommitUnknown.
        })
        .await;
    drop(driver);
    let (completion, execution) = owner.take().unwrap().into_physical_finalization();
    // SIMULATION ONLY. D performs real physical finalization before this receipt.
    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)
}