saddle-runtime 0.3.29

Saddle managed asynchronous runtime and lifecycle
Documentation
//! Application-owned fixed post-driver finalizer slot.
//!
//! This is framework assembly surface only. It is public because Rust has no
//! cross-crate friend visibility and is not re-exported by the Saddle facade.

use std::{
    sync::{Arc, Mutex, MutexGuard, mpsc},
    time::Instant,
};

use saddle_admission::VerifiedPostDriverInstallBinding;
use saddle_core::{ErrorKind, Result, SaddleError};

use crate::compiled_route::OfficialCompiledDriverFinalizer;

enum InstallState {
    Unarmed,
    MustSubmit,
    Consumed,
}

struct State {
    install: InstallState,
    submitted: Option<OfficialCompiledDriverFinalizer>,
    consumed: bool,
}

/// Cloneable access to one Application-owned fixed slot. Clones do not own the
/// finalizer; they can only arm the contract or linearly submit its opaque
/// value during component shutdown.
#[derive(Clone)]
pub struct PendingDriverFinalizerSlot {
    state: Arc<Mutex<State>>,
    #[cfg(test)]
    worker_result: Arc<std::sync::atomic::AtomicU8>,
}

impl PendingDriverFinalizerSlot {
    pub(crate) fn new() -> Self {
        Self {
            state: Arc::new(Mutex::new(State {
                install: InstallState::Unarmed,
                submitted: None,
                consumed: false,
            })),
            #[cfg(test)]
            worker_result: Arc::new(std::sync::atomic::AtomicU8::new(0)),
        }
    }

    #[cfg(test)]
    pub(crate) fn is_unarmed_for_test(&self) -> bool {
        let state = lock(&self.state);
        matches!(state.install, InstallState::Unarmed)
            && state.submitted.is_none()
            && !state.consumed
    }

    /// Performs the final synchronous, infallible authority commit. All
    /// fallible/async work completed while the slot was still unarmed.
    #[doc(hidden)]
    pub fn commit_verified_install(&self, binding: VerifiedPostDriverInstallBinding) {
        let mut state = lock(&self.state);
        if !matches!(state.install, InstallState::Unarmed)
            || state.submitted.is_some()
            || state.consumed
        {
            crate::diagnostics::finalizer_contract_abort("runtime.finalizer_duplicate_install");
        }
        let _binding = binding;
        state.install = InstallState::MustSubmit;
    }

    pub(crate) fn reserved_submit_handle(&self) -> MustSubmitDriverFinalizer {
        MustSubmitDriverFinalizer { slot: self.clone() }
    }

    /// Submits the sole opaque finalizer. Duplicate, unarmed or post-consume
    /// submission is a finite fail-closed contract violation.
    pub fn submit(&self, finalizer: OfficialCompiledDriverFinalizer) {
        let mut state = lock(&self.state);
        if !matches!(state.install, InstallState::MustSubmit)
            || state.submitted.is_some()
            || state.consumed
        {
            crate::diagnostics::finalizer_contract_abort("runtime.finalizer_invalid_submission");
        }
        state.submitted = Some(finalizer);
    }

    pub(crate) fn finish(
        &self,
        runtime: tokio::runtime::Runtime,
        application_result: Result<()>,
        shutdown_deadline: Option<Instant>,
        lifecycle_observer: Option<(saddle_observability::Observer, String)>,
    ) -> Result<()> {
        let pending = {
            let mut state = lock(&self.state);
            if state.consumed {
                crate::diagnostics::finalizer_contract_abort("runtime.finalizer_already_consumed");
            }
            state.consumed = true;
            let must_submit = matches!(state.install, InstallState::MustSubmit);
            state.install = InstallState::Consumed;
            if must_submit {
                Some(state.submitted.take().unwrap_or_else(|| {
                    crate::diagnostics::finalizer_contract_abort(
                        "runtime.finalizer_missing_submission",
                    )
                }))
            } else {
                None
            }
        };

        let finalizer_result = match pending {
            Some(pending) => {
                let Some(deadline) = shutdown_deadline else {
                    let failed = crate::diagnostics::post_driver_source(
                        finalization_failure("runtime.finalizer_missing_shutdown_deadline"),
                        application_result.as_ref().err(),
                        lifecycle_observer.as_ref().map(|(_, application)| application.as_str()),
                    );
                    drop(runtime);
                    return crate::application::combine_lifecycle_results(application_result, Err(failed));
                };
                let (sender, receiver) = mpsc::sync_channel(1);
                let primary_occurrence = application_result.as_ref().err()
                    .and_then(|error| error.diagnostic().map(|diagnostic| diagnostic.occurrence()));
                let worker_application = lifecycle_observer.as_ref()
                    .map(|(_, application)| application.clone());
                #[cfg(test)]
                let worker_result = Arc::clone(&self.worker_result);
                std::thread::spawn(move || {
                    let result = crate::diagnostics::catching(
                        saddle_core::DiagnosticStage::FinalizerResource,
                        "runtime.post_driver",
                        None,
                        || {
                            pending
                                .bind_runtime(runtime)
                                .finish()
                                .map_err(|_| finalization_failure("runtime.finalizer_bind_failed"))
                                .and_then(|report| {
                                    #[cfg(test)]
                                    if std::env::var_os("RUNTIME_DIAGNOSTIC_FINALIZER_FAULT")
                                        .is_some()
                                    {
                                        // Actual post-driver worker, after physical
                                        // driver/ledger finalization, before result delivery.
                                        panic!("DIAGNOSTIC_PRIVATE_SENTINEL");
                                    }
                                    if report.ledger.healthy
                                        && !report.watermark.breached
                                        && !report.task_failed
                                    {
                                        Ok(())
                                    } else {
                                        Err(finalization_failure(if !report.ledger.healthy {
                                            "runtime.finalizer_ledger_unhealthy"
                                        } else if report.watermark.breached {
                                            "runtime.finalizer_watermark_breached"
                                        } else {
                                            "runtime.finalizer_task_failed"
                                        }))
                                    }
                                })
                        },
                    )
                    .unwrap_or_else(|d| Err(finalization_error().with_diagnostic(d)))
                    .map_err(|error| {
                        let raw = SaddleError::new(error.kind(), error.code(), error.message());
                        crate::diagnostics::post_driver_source_with_occurrence(
                            error, primary_occurrence, worker_application.as_deref(), &raw,
                        )
                    });
                    #[cfg(test)]
                    worker_result.store(
                        if result.is_ok() { 1 } else { 2 },
                        std::sync::atomic::Ordering::SeqCst,
                    );
                    let _ = sender.send(result);
                });
                let wait_started = Instant::now();
                match receiver.recv_timeout(deadline.saturating_duration_since(wait_started)) {
                    Ok(result) => result.map_err(|error| (error, None, true)),
                    Err(error) => {
                        let failed = receive_failure(&error);
                        if matches!(error, mpsc::RecvTimeoutError::Timeout) {
                            if let Some((observer, application)) = lifecycle_observer.as_ref() {
                                observer.record_lifecycle_timeout(
                                application.as_str(),
                                saddle_observability::LifecycleTimeoutStage::PostDriverFinalization,
                                u64::try_from(wait_started.elapsed().as_millis())
                                    .unwrap_or(u64::MAX),
                            );
                            }
                        }
                        // The existing worker still owns Runtime and the
                        // prepaid ledger until it actually finishes. A
                        // timeout is an explicit unreconciled failure, not a
                        // physical-destruction receipt. Keep the original
                        // application error through the final `and` below.
                        Err((failed, Some(error), false))
                    }
                }
            }
            None => {
                drop(runtime);
                Ok(())
            }
        };
        let finalizer_result = finalizer_result.map_err(|(error, raw, worker_recorded)| {
            if worker_recorded { return error; }
            let primary = application_result.as_ref().err();
            let application = lifecycle_observer.as_ref().map(|(_, application)| application.as_str());
            match raw {
                Some(raw) => crate::diagnostics::post_driver_source_with_raw(error, primary, application, &raw),
                None => crate::diagnostics::post_driver_source(error, primary, application),
            }
        });
        crate::application::combine_lifecycle_results(application_result, finalizer_result)
    }
}

fn finalization_failure(code: &'static str) -> SaddleError {
    crate::diagnostics::attach(
        finalization_error(),
        saddle_core::DiagnosticStage::FinalizerResource,
        code,
    )
}

fn receive_failure(error: &mpsc::RecvTimeoutError) -> SaddleError {
    finalization_failure(match error {
        mpsc::RecvTimeoutError::Timeout => "runtime.finalizer_timeout",
        mpsc::RecvTimeoutError::Disconnected => "runtime.finalizer_disconnected",
    })
}

/// Component-facing submission half. It conveys no Admission or identity
/// facts and can only submit into the already committed slot.
#[doc(hidden)]
pub struct MustSubmitDriverFinalizer {
    slot: PendingDriverFinalizerSlot,
}

impl MustSubmitDriverFinalizer {
    #[doc(hidden)]
    pub fn submit(self, finalizer: OfficialCompiledDriverFinalizer) {
        self.slot.submit(finalizer);
    }
}

fn lock(state: &Mutex<State>) -> MutexGuard<'_, State> {
    state.lock().unwrap_or_else(|_| {
        crate::diagnostics::finalizer_contract_abort("runtime.finalizer_slot_poisoned")
    })
}

fn finalization_error() -> SaddleError {
    SaddleError::new(
        ErrorKind::Infrastructure,
        "runtime.driver_finalization_failed",
        "the managed Runtime driver did not finalize cleanly",
    )
}

#[cfg(test)]
pub(crate) fn diagnostic_test_finish(
    runtime: tokio::runtime::Runtime,
    pending: OfficialCompiledDriverFinalizer,
    primary: SaddleError,
) -> SaddleError {
    diagnostic_test_finish_with_witness(runtime, pending, primary).0
}

#[cfg(test)]
pub(crate) fn diagnostic_test_finish_with_witness(
    runtime: tokio::runtime::Runtime,
    pending: OfficialCompiledDriverFinalizer,
    primary: SaddleError,
) -> (SaddleError, Arc<std::sync::atomic::AtomicU8>) {
    diagnostic_test_finish_with_deadline(runtime, pending, primary, std::time::Duration::from_secs(2))
}

#[cfg(test)]
pub(crate) fn diagnostic_test_finish_with_deadline(
    runtime: tokio::runtime::Runtime,
    pending: OfficialCompiledDriverFinalizer,
    primary: SaddleError,
    deadline: std::time::Duration,
) -> (SaddleError, Arc<std::sync::atomic::AtomicU8>) {
    let slot = PendingDriverFinalizerSlot::new();
    let worker_result = Arc::clone(&slot.worker_result);
    lock(&slot.state).install = InstallState::MustSubmit;
    slot.submit(pending);
    let result = slot.finish(
        runtime,
        Err(primary),
        Some(Instant::now() + deadline),
        None,
    );
    (result.expect_err("primary error must remain the returned error"), worker_result)
}

#[cfg(test)]
mod diagnostic_tests {
    use super::*;
    #[test]
    fn disconnected_is_not_timeout() {
        for (error, expected) in [
            (
                mpsc::RecvTimeoutError::Disconnected,
                "runtime.finalizer_disconnected",
            ),
            (mpsc::RecvTimeoutError::Timeout, "runtime.finalizer_timeout"),
        ] {
            let error = receive_failure(&error);
            let json = serde_json::to_value(error.diagnostic().unwrap()).unwrap();
            assert_eq!(json["causes"][0]["code"], expected);
        }
    }
}