saddle-observability 0.3.33

Saddle structured logging and trace correlation
Documentation
//! Closed adapter for process terminal audit and startup terminal storage.
//! A receipt can only be issued by the original capture in this module.
use super::{original, OriginalCaptureState, UnrootedCaptureFacts};
use crate::{DiagnosticSubmission, EmergencyDiagnosticHandle};
use saddle_core::{BoundedDiagnostic, BoundedDiagnosticCause, CaptureSite,
    DiagnosticCategory, DiagnosticCode, DiagnosticOccurrence, DiagnosticStage,
    request_context::RequestContextSnapshot};

/// Confirmation of this terminal original, never of its disposition or cleanup.
/// Private fields and no public constructor prevent promotion of output facts.
#[derive(Clone, Copy)]
pub struct WrittenTerminalOriginal {
    occurrence: DiagnosticOccurrence,
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::sync::atomic::{AtomicUsize, Ordering};
    static RUN: AtomicUsize = AtomicUsize::new(0);
    static TEST_OUTPUT: std::sync::Mutex<()> = std::sync::Mutex::new(());
    fn close(mut output: crate::EmergencyDiagnostics) {
        let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2);
        while output.shutdown() == crate::DiagnosticShutdown::Pending {
            assert!(std::time::Instant::now() < deadline);
            std::thread::yield_now();
        }
    }
    fn output() -> (std::path::PathBuf, crate::EmergencyDiagnostics) {
        let root = std::env::temp_dir().join(format!("saddle-terminal-{}-{}", std::process::id(), RUN.fetch_add(1, Ordering::Relaxed)));
        std::fs::create_dir(&root).unwrap();
        let output = crate::EmergencyDiagnostics::start_checked(&crate::FileLoggingConfig::new(&root, crate::Rotation::Daily)).unwrap();
        (root, output)
    }
    fn failure() -> saddle_admission::ProfuseGwTerminalAuditFailure {
        saddle_admission::ProfuseGwTerminalAuditFailure {
            error: saddle_admission::AdmissionError::RequestEscaped(saddle_admission::AuditReport {
                escape_allocations: 1, escape_bytes: 40, ..Default::default()
            }), audit: None, database: saddle_admission::ProfuseGwDatabaseDisposition::NotUsed,
        }
    }
    #[test]
    fn original_precedes_repeated_same_identity_projection() {
        let _exclusive = TEST_OUTPUT.lock().unwrap();
        let (root, output) = output();
        let recorded = RecordedTerminalFailure::audit(Some(&output.handle()), None, None, failure());
        let log = root.join("saddle.emergency.log");
        // Actual read before either projection, not an exit-time inference.
        let before = std::fs::read_to_string(&log).unwrap();
        assert!(before.contains("RequestEscaped"));
        assert!(before.contains("escape_bytes: 40"));
        assert!(recorded.original_confirmed());
        assert_eq!(recorded.disposition_submission(), Some(DiagnosticSubmission::Written));
        for _ in 0..2 {
            let safe = recorded.project();
            let receipt = safe.source_receipt::<WrittenTerminalOriginal>().unwrap();
            assert!(receipt.occurrence().matches_diagnostic(safe.diagnostic().unwrap()));
            assert!(receipt.occurrence().matches_bounded_diagnostic(recorded.diagnostic()));
            assert!(!safe.source_unavailable());
            assert!(safe.unconfirmed_original().is_none());
        }
        assert_eq!(before, std::fs::read_to_string(&log).unwrap(), "projection never rewrites source");
        close(output); std::fs::remove_dir_all(root).unwrap();
    }
    #[test]
    fn missing_original_and_failed_disposition_are_independent() {
        if std::env::var_os("SADDLE_TERMINAL_OUTPUT_FAILURE_CHILD").is_none() {
            let mut child = std::process::Command::new("bash");
            child.args(["-c", "trap '' XFSZ; exec \"$@\"", "terminal-output-failure"])
                .arg(std::env::current_exe().unwrap())
                .args(["--exact", "root_diagnostic::terminal::tests::missing_original_and_failed_disposition_are_independent", "--nocapture"])
                .env("SADDLE_TERMINAL_OUTPUT_FAILURE_CHILD", "1")
                .stdout(std::process::Stdio::piped()).stderr(std::process::Stdio::piped());
            let mut child = child.spawn().unwrap();
            let deadline = std::time::Instant::now() + std::time::Duration::from_secs(10);
            while child.try_wait().unwrap().is_none() {
                if std::time::Instant::now() >= deadline { child.kill().unwrap(); panic!("terminal output child exceeded deadline"); }
                std::thread::sleep(std::time::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_OUTPUT_NEGATIVE_COMPLETE"));
            assert!(String::from_utf8_lossy(&output.stderr).contains("FileTooLarge"));
            return;
        }
        let _exclusive = TEST_OUTPUT.lock().unwrap();
        let absent = RecordedTerminalFailure::audit(None, None, None, failure());
        let safe = absent.project();
        assert!(safe.source_unavailable());
        assert!(safe.source_receipt::<WrittenTerminalOriginal>().is_none());
        assert!(safe.unconfirmed_original().unwrap().to_string().contains("escape_bytes: 40"));
        let (root, mut output) = output();
        let handle = output.handle();
        let mut recorded = RecordedTerminalFailure::audit_original(Some(&handle), None, None, failure());
        assert!(recorded.original_confirmed());
        // Restrict only this child after its real original was read back.
        assert!(std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap().contains("RequestEscaped"));
        assert!(std::process::Command::new("prlimit").args(["--pid", &std::process::id().to_string(), "--fsize=0:0"]).status().unwrap().success());
        recorded.record_disposition(Some(&handle), None, None);
        assert_ne!(recorded.disposition_submission(), Some(DiagnosticSubmission::Written));
        let safe = recorded.project();
        assert!(safe.source_unavailable());
        assert!(safe.source_receipt::<WrittenTerminalOriginal>().is_some(), "only original is confirmed");
        assert!(safe.unconfirmed_original().is_none());
        let failed = RecordedTerminalFailure::audit(Some(&handle), None, None, failure());
        assert!(!failed.original_confirmed());
        assert!(failed.project().source_receipt::<WrittenTerminalOriginal>().is_none());
        close(output); std::fs::remove_dir_all(root).unwrap();
        println!("TERMINAL_OUTPUT_NEGATIVE_COMPLETE");
    }
    #[test]
    fn storage_preserves_actual_admission_source_or_unconfirmed_original() {
        let _exclusive = TEST_OUTPUT.lock().unwrap();
        let error = saddle_admission::AdmissionError::CapacityRejected;
        let absent = RecordedTerminalFailure::storage(None, error);
        let safe = absent.project();
        assert_eq!(safe.code(), "runtime.terminal_storage_unavailable");
        assert!(safe.source_unavailable());
        assert_eq!(safe.unconfirmed_original().unwrap().to_string(), error.to_string());
        let (root, output) = output();
        let recorded = RecordedTerminalFailure::storage(Some(&output.handle()), error);
        let before = std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap();
        assert!(before.contains("CapacityRejected"));
        let safe = recorded.project();
        assert!(!safe.source_unavailable());
        assert!(safe.source_receipt::<WrittenTerminalOriginal>().unwrap().occurrence().matches_diagnostic(safe.diagnostic().unwrap()));
        assert_eq!(before, std::fs::read_to_string(root.join("saddle.emergency.log")).unwrap());
        close(output); std::fs::remove_dir_all(root).unwrap();
    }
}
impl WrittenTerminalOriginal {
    pub fn occurrence(&self) -> DiagnosticOccurrence { self.occurrence }
}

enum Original {
    Audit(saddle_admission::ProfuseGwTerminalAuditFailure),
    Storage(saddle_admission::AdmissionError),
}
enum Capture {
    Written(WrittenTerminalOriginal),
    Unconfirmed,
}

/// Bounded retained source, not request permission. Stored in process reserve.
pub struct RecordedTerminalFailure {
    diagnostic: BoundedDiagnostic,
    original: Original,
    capture: Capture,
    disposition: Option<DiagnosticSubmission>,
    facts: Option<UnrootedCaptureFacts>,
}

impl RecordedTerminalFailure {
    pub fn audit(
        output: Option<&EmergencyDiagnosticHandle>,
        context: Option<&RequestContextSnapshot>,
        source: Option<DiagnosticOccurrence>,
        failure: saddle_admission::ProfuseGwTerminalAuditFailure,
    ) -> Self {
        let mut result = Self::audit_original(output, context, source, failure);
        result.record_disposition(output, context, source);
        result
    }

    fn audit_original(
        output: Option<&EmergencyDiagnosticHandle>,
        context: Option<&RequestContextSnapshot>,
        source: Option<DiagnosticOccurrence>,
        failure: saddle_admission::ProfuseGwTerminalAuditFailure,
    ) -> Self {
        let mut diagnostic = BoundedDiagnostic::capture(DiagnosticCategory::UnexpectedError,
            CaptureSite::FirstObserved, BoundedDiagnosticCause::new(
                DiagnosticStage::FinalizerResource, DiagnosticCode::new("runtime.request_audit_failed").unwrap()));
        if let Some(source) = source { diagnostic = diagnostic.during_cleanup_occurrence(source); }
        let facts = original::request_terminal_audit_error(output, context, &diagnostic, source, &failure);
        let capture = if facts.original_capture() == OriginalCaptureState::CompleteWritten
            && facts.occurrence().matches_bounded_diagnostic(&diagnostic) {
            Capture::Written(WrittenTerminalOriginal { occurrence: facts.occurrence() })
        } else { Capture::Unconfirmed };
        Self { diagnostic, original: Original::Audit(failure), capture, disposition: None, facts: Some(facts) }
    }

    fn record_disposition(&mut self, output: Option<&EmergencyDiagnosticHandle>,
        context: Option<&RequestContextSnapshot>, source: Option<DiagnosticOccurrence>) {
        let Original::Audit(failure) = &self.original else { unreachable!("audit disposition only"); };
        self.disposition = Some(original::request_terminal_audit_disposition(
            output, context, source, self.facts.as_ref().expect("audit capture facts"), failure.audit));
    }

    pub fn storage(output: Option<&EmergencyDiagnosticHandle>, error: saddle_admission::AdmissionError) -> Self {
        let diagnostic = BoundedDiagnostic::capture(DiagnosticCategory::UnexpectedError,
            CaptureSite::FirstObserved, BoundedDiagnosticCause::new(
                DiagnosticStage::StartupListener, DiagnosticCode::new("runtime.terminal_storage_unavailable").unwrap()));
        let projected = saddle_core::Diagnostic::project_simple_bounded(&diagnostic).unwrap();
        let facts = original::process_component_start_error(output,
            saddle_core::ContextFact::NotEstablished, &projected, &error);
        let capture = if facts.original_capture() == OriginalCaptureState::CompleteWritten
            && facts.occurrence().matches_bounded_diagnostic(&diagnostic) {
            Capture::Written(WrittenTerminalOriginal { occurrence: facts.occurrence() })
        } else { Capture::Unconfirmed };
        Self { diagnostic, original: Original::Storage(error), capture, disposition: None, facts: None }
    }

    pub fn diagnostic(&self) -> &BoundedDiagnostic { &self.diagnostic }
    pub fn original_confirmed(&self) -> bool { matches!(self.capture, Capture::Written(_)) }
    pub fn disposition_submission(&self) -> Option<DiagnosticSubmission> { self.disposition }

    /// Repeated projection neither writes the source nor creates an occurrence.
    pub fn project(&self) -> saddle_core::SaddleError {
        match &self.capture {
            Capture::Written(receipt) if receipt.occurrence.matches_bounded_diagnostic(&self.diagnostic) =>
                self.project_written(receipt),
            _ => self.project_unconfirmed(),
        }
    }

    // The source inventory recognizes only these private, sealed projection
    // leaves and pins their signing/encoding dependencies. No caller diagnostic
    // or public output facts can reach this positive pairing.
    fn project_written(&self, receipt: &WrittenTerminalOriginal) -> saddle_core::SaddleError {
        let diagnostic = saddle_core::Diagnostic::project_simple_bounded(&self.diagnostic).unwrap();
        let (code, message) = self.classification();
        let error = saddle_core::SaddleError::new(saddle_core::ErrorKind::Infrastructure, code, message)
            .with_diagnostic(diagnostic).with_source_receipt(*receipt);
        if self.disposition.is_some_and(|submission| submission != DiagnosticSubmission::Written) {
            error.with_unconfirmed_source()
        } else { error }
    }

    fn project_unconfirmed(&self) -> saddle_core::SaddleError {
        let diagnostic = saddle_core::Diagnostic::project_simple_bounded(&self.diagnostic).unwrap();
        let (code, message) = self.classification();
        saddle_core::SaddleError::new(saddle_core::ErrorKind::Infrastructure, code, message)
            .with_diagnostic(diagnostic).with_unconfirmed_original(self.copy_original())
    }

    fn classification(&self) -> (&'static str, &'static str) {
        match self.original {
            Original::Audit(_) => ("runtime.request_audit_failed", "request resource audit failed; process stopped"),
            Original::Storage(_) => ("runtime.terminal_storage_unavailable", "process terminal storage unavailable"),
        }
    }

    fn copy_original(&self) -> Box<dyn std::error::Error + Send + Sync> {
        match &self.original {
            Original::Audit(failure) => Box::new(saddle_admission::ProfuseGwTerminalAuditFailure {
                error: failure.error, audit: failure.audit, database: failure.database,
            }),
            Original::Storage(error) => Box::new(*error),
        }
    }
}