use crate::{DiagnosticSubmission,EmergencyDiagnosticHandle,Observer};
use saddle_core::request_context::ContextIdentity;
use serde::Serialize;
impl Observer {
pub fn record_standard_rpc_shutdown(&self, output: Option<&EmergencyDiagnosticHandle>,
joined: bool, rpc: &saddle_admission::RpcStageSnapshot,
) -> DiagnosticSubmission {
#[derive(Serialize)]
struct Record {
schema: &'static str, event: &'static str, source: &'static str,
process_id: u32, process_generation: u128, driver_joined: bool,
pool_physical_owners: usize, pool_storage_bytes: usize,
request_physical_references: usize, request_backing_bytes: usize,
}
let Some(output) = output else { return DiagnosticSubmission::OutputUnavailable; };
output.submit_fixed_record(&Record {
schema: "standard-rpc-shutdown/1", event: "framework.standard_rpc_shutdown",
source: "post-driver-join-and-pool-reclaim-original-ledger",
process_id: std::process::id(), process_generation: crate::diagnostic::finalization_process_instance(),
driver_joined: joined, pool_physical_owners: rpc.standard_pool_owners,
pool_storage_bytes: rpc.standard_pool_storage_bytes,
request_physical_references: rpc.physical_references, request_backing_bytes: rpc.backing_bytes,
})
}
}
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 {
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})
}
}
#[doc(hidden)]
pub fn ingress_fifo_process_generation() -> u128 {
crate::diagnostic::finalization_process_instance()
}
impl Observer {
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 {
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,
})
}
}