saddle-runtime 0.3.30

Saddle managed asynchronous runtime and lifecycle
Documentation
//! Startup-reserved, one-shot process failure. No request permission is held.
use super::*;
use crate::request_task::reserved::Retained;
use saddle_observability::root_diagnostic::RecordedTerminalFailure;

struct StopState {
    primary: Mutex<Option<(RecordedTerminalFailure, Option<saddle_core::request_context::RequestContextSnapshot>)>>,
    wake: tokio::sync::Notify,
    task_storage_issued: std::sync::atomic::AtomicBool,
    health: Mutex<Option<crate::ApplicationHealth>>,
    _storage: saddle_admission::StoragePermit,
}

#[derive(Clone)]
#[doc(hidden)]
pub struct ProfuseGwFailureStop(Retained<StopState>);

impl ProfuseGwFailureStop {
    pub(super) fn task_storage(&self) -> Result<crate::request_task::reserved_set::ReservedProcessTaskStorage, AdmissionError> {
        // One formal listener startup may claim this source. Its children keep
        // their original permits; dropping a temporary wrapper must not let a
        // second listener acquire the same startup role while backing survives.
        if self.0.task_storage_issued.compare_exchange(false, true,
            std::sync::atomic::Ordering::AcqRel, std::sync::atomic::Ordering::Acquire).is_err() {
            return Err(AdmissionError::CapacityRejected);
        }
        match crate::request_task::reserved_set::ReservedProcessTaskStorage::from_process_parent(&self.0._storage) {
            Ok(storage) => Ok(storage),
            Err(error) => {
                self.0.task_storage_issued.store(false, std::sync::atomic::Ordering::Release);
                Err(error)
            }
        }
    }
    pub(super) fn prepare(process: &ProfuseGwRuntimeProcess) -> Result<Self, AdmissionError> {
        let storage = process.admission.try_process_storage(saddle_admission::StorageDemand::separate(&[
            (std::alloc::Layout::new::<StopState>(), 1),
            (std::alloc::Layout::new::<[std::sync::atomic::AtomicUsize; 2]>(), 1),
        ])?)?;
        Ok(Self(Retained::new(StopState {
            task_storage_issued: std::sync::atomic::AtomicBool::new(false),
            primary: Mutex::new(None), wake: tokio::sync::Notify::new(),
            health: Mutex::new(None), _storage: storage,
        })))
    }

    pub(super) fn bind_health(&self, health: crate::ApplicationHealth) {
        *self.0.health.lock().unwrap_or_else(|p| p.into_inner()) = Some(health.clone());
        if self.is_requested() { health.fail_closed(); }
    }

    pub fn is_requested(&self) -> bool {
        self.0.primary.lock().unwrap_or_else(|p| p.into_inner()).is_some()
    }

    pub fn report(
        &self,
        failure: saddle_admission::ProfuseGwTerminalAuditFailure,
        context: Option<saddle_core::request_context::RequestContextSnapshot>,
        source: Option<saddle_core::DiagnosticOccurrence>,
        output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
    ) {
        let recorded = RecordedTerminalFailure::audit(output, context.as_ref(), source, failure);
        let first = {
            let mut primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
            if primary.is_some() { false } else { *primary = Some((recorded, context)); true }
        };
        if first {
            if let Some(health) = self.0.health.lock().unwrap_or_else(|p| p.into_inner()).as_ref() {
                health.fail_closed();
            }
            self.0.wake.notify_one();
        }
    }

    pub(super) async fn wait(&self) -> saddle_core::Result<()> {
        loop {
            let notified = self.0.wake.notified();
            let primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
            if let Some((recorded, _)) = primary.as_ref() {
                return Err(recorded.project());
            }
            drop(primary);
            notified.await;
        }
    }

    pub fn record_cleanup(&self, snapshot: &saddle_admission::ProcessSnapshot,
        output: Option<&saddle_observability::EmergencyDiagnosticHandle>) {
        let primary = self.0.primary.lock().unwrap_or_else(|p| p.into_inner());
        if let Some((recorded, context)) = primary.as_ref() {
            let _ = saddle_observability::root_diagnostic::request_terminal_audit_cleanup(
                output, context.as_ref(), recorded.diagnostic().occurrence(), snapshot);
        }
    }
}

#[cfg(test)]
mod tests {
    use super::*;

    fn failure() -> saddle_admission::ProfuseGwTerminalAuditFailure {
        saddle_admission::ProfuseGwTerminalAuditFailure {
            error: AdmissionError::AccountClosed,
            audit: None,
            database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed,
        }
    }

    #[test]
    fn audit_stop_survives_missing_output_late_health_and_repeated_reports() {
        let mut process = ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process());
        drop(process.db_startup.take());
        let baseline = process.admission.resource_snapshot().framework_charged;
        let stop = ProfuseGwFailureStop::prepare(&process).unwrap();
        stop.report(failure(), None, None, None);
        let first = stop.0.primary.lock().unwrap().as_ref().unwrap().0.diagnostic().occurrence();
        let health = Application::new().health();
        stop.bind_health(health.clone());
        assert_eq!(health.snapshot().phase(), crate::ApplicationPhase::Draining);
        std::thread::scope(|scope| {
            for _ in 0..4 {
                let stop = stop.clone();
                scope.spawn(move || stop.report(failure(), None, None, None));
            }
        });
        assert!(first.matches_bounded_diagnostic(stop.0.primary.lock().unwrap().as_ref().unwrap().0.diagnostic()));
        let runtime = tokio::runtime::Builder::new_current_thread().enable_all().build().unwrap();
        runtime.block_on(async {
            // Every waiter sees the latched error, including after the first
            // notification was consumed. No second shutdown is required.
            for _ in 0..2 {
                let error = tokio::time::timeout(Duration::from_millis(100), stop.wait()).await.unwrap().unwrap_err();
                assert_eq!(error.code(), "runtime.request_audit_failed");
                assert!(error.source_unavailable());
                assert!(first.matches_diagnostic(error.diagnostic().unwrap()));
            }
        });
        drop(stop);
        assert_eq!(process.admission.resource_snapshot().framework_charged, baseline);
        process.finish().ok().expect("stop state released startup storage");
    }

    #[test]
    #[ignore = "isolated owned-runtime startup storage failure subprocesses"]
    fn terminal_storage_failure_preserves_real_original() {
        const CHILD: &str = "SADDLE_TERMINAL_STORAGE_CHILD";
        if let Ok(mode) = std::env::var(CHILD) {
            let mut process = ProfuseGwRuntimeProcess::new(crate::request_task::reserved::tests::process());
            drop(process.db_startup.take());
            // A real live process storage owner makes the next claim reject.
            // It remains live throughout finalization; do not fake its refund.
            let held = process.admission.try_process_storage(saddle_admission::StorageDemand::separate(&[]).unwrap()).unwrap();
            let factory = |_| async { panic!("storage rejection must precede Application construction"); #[allow(unreachable_code)] Ok(Application::new()) };
            let root = std::env::temp_dir().join(format!("saddle-terminal-storage-{}", std::process::id()));
            let terminal = if mode == "written" {
                std::fs::create_dir(&root).unwrap();
                let output = saddle_observability::EmergencyDiagnostics::start_checked(&saddle_observability::FileLoggingConfig::new(&root, saddle_observability::Rotation::Daily)).unwrap();
                run_profusegw_owned_application_with_diagnostics(process, output, factory).0
            } else { run_profusegw_owned_application(process, factory) };
            let failure = match terminal.finish() {
                Err(ProfuseGwCoordinatorFailure::Lifecycle(primary)) |
                Err(ProfuseGwCoordinatorFailure::LifecycleWithCleanup { primary, .. }) => primary,
                _ => panic!("real storage source must remain lifecycle primary"),
            };
            assert_eq!(failure.code(), "runtime.terminal_storage_unavailable");
            if mode == "written" {
                let receipt = failure.source_receipt::<saddle_observability::root_diagnostic::WrittenTerminalOriginal>().unwrap();
                assert!(receipt.occurrence().matches_diagnostic(failure.diagnostic().unwrap()));
                assert!(!failure.source_unavailable());
                let rows = std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap();
                assert!(rows.contains("CapacityRejected"));
                if let Some(oracle) = std::env::var_os("SADDLE_AQ_TERMINAL_ORACLE") {
                    let checked = std::process::Command::new("python3").arg("-B").arg(oracle)
                        .arg("--log").arg(root.join("saddle.emergency.log"))
                        .args(["--kind", "storage"]).status().unwrap();
                    assert!(checked.success(), "independent actual Admission original required");
                }
                assert!(!rows.contains("process terminal storage unavailable"), "safe summary is never recaptured as the original");
                std::fs::remove_dir_all(root).unwrap();
            } else {
                assert!(failure.source_unavailable());
                assert!(failure.source_receipt::<saddle_observability::root_diagnostic::WrittenTerminalOriginal>().is_none());
                assert_eq!(failure.unconfirmed_original().unwrap().to_string(), AdmissionError::CapacityRejected.to_string());
            }
            drop(held);
            println!("TERMINAL_STORAGE_ASSERTIONS_COMPLETE");
            return;
        }
        for mode in ["written", "unconfirmed"] {
            let mut child = std::process::Command::new(std::env::current_exe().unwrap())
                .args(["--exact", "profusegw::failure_stop::tests::terminal_storage_failure_preserves_real_original", "--ignored", "--nocapture"])
                .env(CHILD, mode).stdout(std::process::Stdio::piped()).stderr(std::process::Stdio::piped()).spawn().unwrap();
            let deadline = std::time::Instant::now() + Duration::from_secs(10);
            while child.try_wait().unwrap().is_none() {
                if std::time::Instant::now() >= deadline { child.kill().unwrap(); panic!("terminal storage child exceeded deadline"); }
                std::thread::sleep(Duration::from_millis(10));
            }
            let output = child.wait_with_output().unwrap();
            assert!(output.status.success(), "{}", String::from_utf8_lossy(&output.stderr));
            assert!(String::from_utf8_lossy(&output.stdout).contains("TERMINAL_STORAGE_ASSERTIONS_COMPLETE"));
            assert!(!String::from_utf8_lossy(&output.stderr).contains("panicked at"));
        }
    }
}