saddle-observability 0.3.36

Saddle structured logging and trace correlation
Documentation
//! Read-only transitions at the formal ingress queue operations.
use crate::{DiagnosticSubmission,EmergencyDiagnosticHandle,Observer};
use saddle_core::request_context::ContextIdentity;
use serde::Serialize;

impl Observer {
    pub fn record_ingress_fifo(&self,output:Option<&EmergencyDiagnosticHandle>,request_id:Option<ContextIdentity>,
        queue_sequence:u64,transition:&'static str,reason:&'static str)->DiagnosticSubmission {
        #[derive(Serialize)]
        struct Record {
            schema:&'static str,event:&'static str,process_id:u32,process_generation:u128,
            request_id:Option<ContextIdentity>,queue_sequence:u64,transition:&'static str,reason:&'static str,
        }
        let Some(output)=output else{return DiagnosticSubmission::OutputUnavailable;};
        output.submit_fixed_record(&Record {schema:"ingress-fifo/1",event:"framework.ingress_fifo",
            process_id:std::process::id(),process_generation:crate::diagnostic::finalization_process_instance(),
            request_id,queue_sequence,transition,reason})
    }
}

impl Observer {
    /// Attempt identity only. Resource amounts remain on the original same-lock
    /// capacity receipt; input_peak is the already recognized input demand.
    pub fn record_ingress_fifo_attempt(&self, output: Option<&EmergencyDiagnosticHandle>,
        request_id: Option<ContextIdentity>, queue_sequence: u64, decision_attempt: u64,
        input_peak: usize, outcome: &'static str) -> DiagnosticSubmission {
        #[derive(Serialize)]
        struct Record {
            schema: &'static str, event: &'static str, process_id: u32,
            process_generation: u128, request_id: Option<ContextIdentity>,
            queue_sequence: u64, decision_attempt: u64, input_peak: usize, outcome: &'static str,
        }
        let Some(output) = output else { return DiagnosticSubmission::OutputUnavailable; };
        output.submit_fixed_record(&Record {schema: "ingress-fifo-attempt/1",
            event: "framework.ingress_fifo_attempt", process_id: std::process::id(),
            process_generation: crate::diagnostic::finalization_process_instance(),
            request_id, queue_sequence, decision_attempt, input_peak, outcome})
    }
}

/// The capacity receipt and FIFO marker use this same diagnostic process epoch.
#[doc(hidden)]
pub fn ingress_fifo_process_generation() -> u128 {
    crate::diagnostic::finalization_process_instance()
}

impl Observer {
    /// A post-join sample, not an admission decision or an atomic combined ledger receipt.
    pub fn record_resource_quiescent(&self, output: Option<&EmergencyDiagnosticHandle>,
        context: Option<saddle_core::request_context::RequestContextSnapshot>, completion_sequence: u64,
        resource: &saddle_admission::ProcessSnapshot, rpc: &saddle_admission::RpcStageSnapshot,
    ) -> DiagnosticSubmission {
        #[derive(Serialize)]
        struct Record {
            schema: &'static str, event: &'static str, source: &'static str,
            process_id: u32, process_generation: u128,
            context: Option<saddle_core::request_context::RequestContextSnapshot>, completion_sequence: u64,
            active_accounts: usize, charged: usize, task_charged: usize, framework_charged: usize,
            registered_waiters: usize, healthy: bool, rpc_capacity: usize,
            rpc_operations: usize, rpc_physical_references: usize,
        }
        let Some(output) = output else { return DiagnosticSubmission::OutputUnavailable; };
        output.submit_fixed_record(&Record {schema: "resource-quiescent/1", event: "framework.resource_quiescent",
            source: "postjoin-empty-collections-original-ledger-samples", process_id: std::process::id(),
            process_generation: crate::diagnostic::finalization_process_instance(), context, completion_sequence,
            active_accounts: resource.active_accounts, charged: resource.charged,
            task_charged: resource.task_charged, framework_charged: resource.framework_charged,
            registered_waiters: resource.registered_waiters, healthy: resource.healthy,
            rpc_capacity: rpc.capacity, rpc_operations: rpc.operations, rpc_physical_references: rpc.physical_references})
    }
}

impl EmergencyDiagnosticHandle {
    /// Scalar recording-point evidence after the caller's original queue mutation.
    /// This is not an exact physical removal timestamp or a new deadline owner.
    pub fn record_ingress_fifo_clock(&self, request_id: Option<ContextIdentity>,
        queue_sequence: u64, transition: &'static str, reason: &'static str,
        effective_deadline_unix_ms: i64) -> DiagnosticSubmission {
        #[derive(Serialize)]
        struct Record {
            schema: &'static str, event: &'static str, process_id: u32,
            process_generation: u128, request_id: Option<ContextIdentity>,
            queue_sequence: u64, transition: &'static str, reason: &'static str,
            effective_deadline_unix_ms: i64, recorded_unix_ns: Option<u128>,
        }
        let recorded_unix_ns = std::time::SystemTime::now()
            .duration_since(std::time::UNIX_EPOCH).ok().map(|elapsed| elapsed.as_nanos());
        self.submit_fixed_record(&Record {
            schema: "ingress-fifo-clock/1", event: "framework.ingress_fifo_clock",
            process_id: std::process::id(),
            process_generation: crate::diagnostic::finalization_process_instance(),
            request_id, queue_sequence, transition, reason,
            effective_deadline_unix_ms, recorded_unix_ns,
        })
    }
}