use crate::FileLoggingConfig;
use saddle_core::{CallContext, Diagnostic};
use serde::Serialize;
use std::{
fmt,
fs::{self, OpenOptions},
io::{self, Write},
path::PathBuf,
sync::{
Arc,
atomic::{AtomicBool, AtomicU64, Ordering},
mpsc::{self, SyncSender, TrySendError},
},
thread::{self, JoinHandle},
time::Duration,
};
pub const EMERGENCY_FILE_NAME: &str = "saddle.emergency.log";
const QUEUE_CAPACITY: usize = 64;
const RECORD_BYTES: usize = 32 * 1024;
static ACTIVE: AtomicBool = AtomicBool::new(false);
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum DiagnosticSubmission {
OutputUnavailable,
Enqueued,
Full,
Closed,
EncodingFailed,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum DiagnosticShutdown {
Pending,
Finished,
WorkerPanicked,
}
#[derive(Clone, Copy, Debug, Default, Eq, PartialEq)]
pub struct DiagnosticOutputSnapshot {
pub initialized: bool,
pub first_failure: Option<DiagnosticIoFailure>,
pub enqueued: u64,
pub written: u64,
pub dropped: u64,
pub output_failed: u64,
pub stderr_failed: u64,
pub stderr_suppressed: u64,
pub truncated: u64,
pub closed: bool,
}
#[derive(Default)]
struct Status {
initialized: AtomicBool,
failure: AtomicU64,
enqueued: AtomicU64,
written: AtomicU64,
dropped: AtomicU64,
output_failed: AtomicU64,
stderr_failed: AtomicU64,
stderr_suppressed: AtomicU64,
truncated: AtomicU64,
closed: AtomicBool,
submitting: AtomicU64,
}
impl Status {
fn fail(&self, stage: u64, error: &io::Error) {
let kind = match error.kind() {
io::ErrorKind::PermissionDenied => 1,
io::ErrorKind::NotFound => 2,
io::ErrorKind::NotADirectory => 3,
io::ErrorKind::InvalidInput => 4,
io::ErrorKind::StorageFull => 5,
io::ErrorKind::BrokenPipe => 6,
io::ErrorKind::WouldBlock => 7,
_ => 8,
};
let os = error.raw_os_error();
let packed = (stage << 56)
| (kind << 48)
| (u64::from(os.is_some()) << 40)
| u64::from(os.unwrap_or(0) as u32);
let _ = self
.failure
.compare_exchange(0, packed, Ordering::AcqRel, Ordering::Acquire);
}
fn snapshot(&self) -> DiagnosticOutputSnapshot {
let failure = self.failure.load(Ordering::Acquire);
DiagnosticOutputSnapshot {
initialized: self.initialized.load(Ordering::Acquire),
first_failure: (failure != 0).then(|| DiagnosticIoFailure {
stage: if failure >> 56 == 1 {
"initialize"
} else {
"write"
},
kind: [
"none",
"permission_denied",
"not_found",
"not_a_directory",
"invalid_input",
"storage_full",
"broken_pipe",
"would_block",
"other",
][((failure >> 48) & 255) as usize],
os_code: (failure & (1 << 40) != 0).then_some(failure as u32 as i32),
}),
enqueued: self.enqueued.load(Ordering::Acquire),
written: self.written.load(Ordering::Acquire),
dropped: self.dropped.load(Ordering::Acquire),
output_failed: self.output_failed.load(Ordering::Acquire),
stderr_failed: self.stderr_failed.load(Ordering::Acquire),
stderr_suppressed: self.stderr_suppressed.load(Ordering::Acquire),
truncated: self.truncated.load(Ordering::Acquire),
closed: self.closed.load(Ordering::Acquire),
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub struct DiagnosticIoFailure {
pub stage: &'static str,
pub kind: &'static str,
pub os_code: Option<i32>,
}
pub struct EmergencyDiagnostics {
handle: EmergencyDiagnosticHandle,
worker: Option<JoinHandle<()>>,
target: PathBuf,
}
#[derive(Clone)]
pub struct EmergencyDiagnosticHandle {
sender: SyncSender<DiagnosticPacket>,
status: Arc<Status>,
}
#[allow(clippy::large_enum_variant)] enum DiagnosticPacket {
Startup(Vec<u8>, Option<saddle_core::DeferredDiagnosticStack>),
Admitted(FixedDiagnosticBytes),
Original(FixedDiagnosticBytes),
}
impl DiagnosticPacket {
fn bytes(&self) -> &[u8] {
match self {
Self::Startup(v, _) => v,
Self::Admitted(v) | Self::Original(v) => &v.bytes[..v.len],
}
}
fn resolve_on_worker(&mut self, status: &Status) {
let Self::Startup(bytes, stack) = self else {
return;
};
let Some(stack) = stack.take() else {
return;
};
let resolved = stack.resolve_on_output_worker();
let Ok(mut record) = serde_json::from_slice::<serde_json::Value>(bytes) else {
status.dropped.fetch_add(1, Ordering::Relaxed);
return;
};
record["diagnostic"]["stack_status"] = resolved.stack_status.into();
record["diagnostic"]["stack"] = resolved.stack.into();
record["diagnostic"]["stack_truncated"] = resolved.stack_truncated.into();
if resolved.stack_truncated {
status.truncated.fetch_add(1, Ordering::Relaxed);
}
let mut encoded = LimitedBytes(Vec::with_capacity(RECORD_BYTES));
if serde_json::to_writer(&mut encoded, &record).is_err() {
encoded.0.clear();
record["diagnostic"]["stack"] = "".into();
record["diagnostic"]["stack_status"] = "omitted_output_limit".into();
record["diagnostic"]["stack_truncated"] = true.into();
record["encoding_truncated"] = true.into();
status.truncated.fetch_add(1, Ordering::Relaxed);
if serde_json::to_writer(&mut encoded, &record).is_err() {
return;
}
}
encoded.0.push(b'\n');
*bytes = encoded.0;
}
}
pub(crate) struct FixedDiagnosticBytes {
pub(crate) bytes: [u8; 8192],
pub(crate) len: usize,
}
pub(crate) fn root_frame_layouts() -> (std::alloc::Layout, std::alloc::Layout, usize) {
(
std::alloc::Layout::new::<FixedDiagnosticBytes>(),
std::alloc::Layout::new::<DiagnosticPacket>(),
QUEUE_CAPACITY,
)
}
impl Write for FixedDiagnosticBytes {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if bytes.len() > self.bytes.len().saturating_sub(self.len + 1) {
return Err(io::ErrorKind::WriteZero.into());
}
self.bytes[self.len..self.len + bytes.len()].copy_from_slice(bytes);
self.len += bytes.len();
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[derive(Debug)]
pub enum EmergencyInitError {
InvalidDirectory,
AlreadyActive,
Spawn(io::ErrorKind),
}
impl fmt::Display for EmergencyInitError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "emergency diagnostic initialization failed: {self:?}")
}
}
impl std::error::Error for EmergencyInitError {}
struct LimitedBytes(Vec<u8>);
impl Write for LimitedBytes {
fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
if self.0.len().saturating_add(bytes.len()) > RECORD_BYTES - 1 {
return Err(io::Error::new(
io::ErrorKind::InvalidData,
"diagnostic encoding bound",
));
}
self.0.extend_from_slice(bytes);
Ok(bytes.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[derive(Serialize)]
struct Envelope<'a> {
event: &'static str,
timestamp_unix_ms: u128,
diagnostic: &'a Diagnostic,
trace_id: Option<&'a str>,
rpc_id: Option<String>,
span_id: Option<String>,
operation: Option<String>,
request: Option<&'a str>,
route: Option<&'a str>,
}
impl<'a> Envelope<'a> {
fn new(
diagnostic: &'a Diagnostic,
context: Option<(&'a CallContext, &'a crate::EventContext)>,
) -> Self {
Self {
event: "framework.diagnostic",
timestamp_unix_ms: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis(),
diagnostic,
trace_id: context.map(|(call, _)| call.trace_correlation_id().as_str()),
rpc_id: context
.and_then(|(call, _)| call.rpc_correlation_id().map(|rpc| rpc.as_str().to_owned())),
span_id: context.map(|(call, _)| call.span_id().to_string()),
operation: context.map(|(call, _)| call.operation().to_string()),
request: context.map(|(_, event)| event.diagnostic_request()),
route: context.map(|(_, event)| event.diagnostic_route()),
}
}
}
impl EmergencyDiagnosticHandle {
pub fn submit_bounded(
&self,
diagnostic: Option<&saddle_core::BoundedDiagnostic>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &crate::EventContext)>,
) -> DiagnosticSubmission {
self.submit_bounded_parts(
diagnostic,
axes,
context.map(|(c, e)| (c, e.diagnostic_request(), e.diagnostic_route())),
None,
)
}
pub(crate) fn submit_bounded_parts(
&self,
diagnostic: Option<&saddle_core::BoundedDiagnostic>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &str, &str)>,
stage: Option<(&'static str, u64)>,
) -> DiagnosticSubmission {
self.encode_bounded(
if stage.is_some() { None } else { diagnostic },
diagnostic.map(|d| d.occurrence()),
axes,
context,
stage,
stage.is_some(),
)
}
pub fn submit_boundary(
&self,
occurrence: Option<saddle_core::DiagnosticOccurrence>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &crate::EventContext)>,
) -> DiagnosticSubmission {
self.encode_bounded(
None,
occurrence,
axes,
context.map(|(c, e)| (c, e.diagnostic_request(), e.diagnostic_route())),
None,
true,
)
}
pub(crate) fn submit_boundary_parts(
&self,
occurrence: Option<saddle_core::DiagnosticOccurrence>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &str, &str)>,
stage: Option<(&'static str, u64)>,
) -> DiagnosticSubmission {
self.encode_bounded(None, occurrence, axes, context, stage, true)
}
#[allow(clippy::too_many_arguments)]
fn encode_bounded(
&self,
diagnostic: Option<&saddle_core::BoundedDiagnostic>,
occurrence: Option<saddle_core::DiagnosticOccurrence>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &str, &str)>,
stage: Option<(&'static str, u64)>,
boundary: bool,
) -> DiagnosticSubmission {
self.status.submitting.fetch_add(1, Ordering::AcqRel);
struct Exit<'a>(&'a AtomicU64);
impl Drop for Exit<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::Release);
}
}
let _exit = Exit(&self.status.submitting);
if context.is_some_and(|(c, request, route)| {
c.operation().as_str().len() > 256 || request.len() > 128 || route.len() > 192
}) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::EncodingFailed;
}
if self.status.closed.load(Ordering::Acquire) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::Closed;
}
#[derive(Serialize)]
struct Record<'a> {
event: &'static str,
stage: Option<&'static str>,
elapsed_ms: Option<u64>,
timestamp_unix_ms: u128,
diagnostic: Option<&'a saddle_core::BoundedDiagnostic>,
diagnostic_reference: Option<saddle_core::DiagnosticOccurrence>,
axes: &'a saddle_core::DiagnosticOutcomeAxes,
trace_id: Option<&'a str>,
rpc_id: Option<&'a str>,
span_id: Option<&'a str>,
operation: Option<&'a str>,
request: Option<&'a str>,
route: Option<&'a str>,
}
let mut span = [b'0'; 16];
if let Some((call, _, _)) = context {
let value = call.span_id().as_u64();
for (i, byte) in span.iter_mut().enumerate() {
*byte = b"0123456789abcdef"[((value >> ((15 - i) * 4)) & 15) as usize];
}
}
let record = Record {
event: if stage.is_some() {
"framework.stage.finished"
} else if boundary {
"framework.boundary.outcome"
} else {
"framework.diagnostic"
},
stage: stage.map(|s| s.0),
elapsed_ms: stage.map(|s| s.1),
timestamp_unix_ms: std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis(),
diagnostic,
diagnostic_reference: occurrence,
axes,
trace_id: context.map(|(c, _, _)| c.trace_correlation_id().as_str()),
rpc_id: context.and_then(|(c, _, _)| c.rpc_correlation_id().map(|r| r.as_str())),
span_id: context.and_then(|_| std::str::from_utf8(&span).ok()),
operation: context.map(|(c, _, _)| c.operation().as_str()),
request: context.map(|(_, request, _)| request),
route: context.map(|(_, _, route)| route),
};
self.encode_and_queue(&record)
}
pub(crate) fn submit_fixed_record<T: Serialize>(&self, record: &T) -> DiagnosticSubmission {
self.status.submitting.fetch_add(1, Ordering::AcqRel);
let result = if self.status.closed.load(Ordering::Acquire) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Closed
} else {
self.encode_and_queue(record)
};
self.status.submitting.fetch_sub(1, Ordering::Release);
result
}
fn encode_and_queue<T: Serialize>(&self, record: &T) -> DiagnosticSubmission {
self.encode_and_queue_fallback(record, None, false)
}
pub(crate) fn submit_original_record<T: Serialize>(&self, record: &T) -> DiagnosticSubmission {
self.status.submitting.fetch_add(1, Ordering::AcqRel);
let result = if self.status.closed.load(Ordering::Acquire) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Closed
} else {
self.encode_and_queue_fallback(record, None, true)
};
self.status.submitting.fetch_sub(1, Ordering::Release);
result
}
pub(crate) fn submit_fixed_with_fallback<T: Serialize>(
&self,
record: &T,
fallback: &T,
) -> DiagnosticSubmission {
self.status.submitting.fetch_add(1, Ordering::AcqRel);
let result = if self.status.closed.load(Ordering::Acquire) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Closed
} else {
self.encode_and_queue_fallback(record, Some(fallback), false)
};
self.status.submitting.fetch_sub(1, Ordering::Release);
result
}
fn encode_and_queue_fallback<T: Serialize>(
&self,
record: &T,
fallback: Option<&T>,
original: bool,
) -> DiagnosticSubmission {
let mut bytes = FixedDiagnosticBytes {
bytes: [0; 8192],
len: 0,
};
if serde_json::to_writer(&mut bytes, record).is_err() {
bytes.len = 0;
let recovered = fallback
.is_some_and(|fallback| serde_json::to_writer(&mut bytes, fallback).is_ok());
if !recovered {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::EncodingFailed;
}
self.status.truncated.fetch_add(1, Ordering::Relaxed);
}
bytes.bytes[bytes.len] = b'\n';
bytes.len += 1;
let packet = if original {
DiagnosticPacket::Original(bytes)
} else {
DiagnosticPacket::Admitted(bytes)
};
match self.sender.try_send(packet) {
Ok(()) => {
self.status.enqueued.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Enqueued
}
Err(TrySendError::Full(_)) => {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Full
}
Err(TrySendError::Disconnected(_)) => {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Closed
}
}
}
pub fn snapshot(&self) -> DiagnosticOutputSnapshot {
self.status.snapshot()
}
pub fn submit(&self, diagnostic: &Diagnostic) -> DiagnosticSubmission {
self.submit_context(diagnostic, None)
}
pub fn submit_context(
&self,
diagnostic: &Diagnostic,
context: Option<(&CallContext, &crate::EventContext)>,
) -> DiagnosticSubmission {
self.submit_existing_record(
&Envelope::new(diagnostic, context),
diagnostic.deferred_stack(),
)
}
pub(crate) fn submit_existing_record<T: Serialize>(
&self,
record: &T,
stack: Option<saddle_core::DeferredDiagnosticStack>,
) -> DiagnosticSubmission {
let active = self.status.submitting.fetch_add(1, Ordering::AcqRel);
struct Exit<'a>(&'a AtomicU64);
impl Drop for Exit<'_> {
fn drop(&mut self) {
self.0.fetch_sub(1, Ordering::AcqRel);
}
}
let _exit = Exit(&self.status.submitting);
if active >= QUEUE_CAPACITY as u64 {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::Full;
}
if self.status.closed.load(Ordering::Acquire) {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::Closed;
}
let mut bytes = LimitedBytes(Vec::with_capacity(RECORD_BYTES));
if serde_json::to_writer(&mut bytes, record).is_err() {
bytes.0.clear();
let mut reduced = match serde_json::to_value(record) {
Ok(value) => value,
Err(_) => {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::EncodingFailed;
}
};
reduced["diagnostic"]["stack"] = serde_json::Value::String(String::new());
reduced["diagnostic"]["stack_truncated"] = true.into();
reduced["encoding_truncated"] = true.into();
self.status.truncated.fetch_add(1, Ordering::Relaxed);
if serde_json::to_writer(&mut bytes, &reduced).is_err() {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
return DiagnosticSubmission::EncodingFailed;
}
}
bytes.0.push(b'\n');
match self
.sender
.try_send(DiagnosticPacket::Startup(bytes.0, stack))
{
Ok(()) => {
self.status.enqueued.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Enqueued
}
Err(TrySendError::Full(_)) => {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Full
}
Err(TrySendError::Disconnected(_)) => {
self.status.dropped.fetch_add(1, Ordering::Relaxed);
DiagnosticSubmission::Closed
}
}
}
}
impl crate::Observer {
pub fn record_diagnostic(
&self,
diagnostic: &Diagnostic,
emergency: &EmergencyDiagnosticHandle,
context: Option<(&CallContext, &crate::EventContext)>,
) -> DiagnosticSubmission {
let result = emergency.submit_context(diagnostic, context);
self.mirror_existing_record(&Envelope::new(diagnostic, context), diagnostic.category());
result
}
pub(crate) fn mirror_existing_record<T: Serialize>(
&self,
record: &T,
category: saddle_core::DiagnosticCategory,
) {
if let Ok(serde_json::Value::Object(fields)) = serde_json::to_value(record) {
let level = if matches!(category, saddle_core::DiagnosticCategory::ExpectedRejection) {
crate::EventLevel::Info
} else {
crate::EventLevel::Error
};
let mut record = crate::logger::LogRecord::new(level, "framework.diagnostic");
record.data = fields;
record.data.remove("event");
record.data.remove("timestamp_unix_ms");
self.emit(record);
}
}
}
impl EmergencyDiagnostics {
pub fn start(config: &FileLoggingConfig) -> Result<Self, EmergencyInitError> {
let directory = config
.directory()
.to_str()
.ok_or(EmergencyInitError::InvalidDirectory)?;
if directory.is_empty()
|| directory.len() > 4096
|| directory.chars().any(char::is_control)
|| directory.contains("://")
|| directory.contains(['@', '?', '#'])
{
return Err(EmergencyInitError::InvalidDirectory);
}
if ACTIVE
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_err()
{
return Err(EmergencyInitError::AlreadyActive);
}
let target = config.directory().join(EMERGENCY_FILE_NAME);
let path = target.clone();
let (sender, receiver) = mpsc::sync_channel::<DiagnosticPacket>(QUEUE_CAPACITY);
let status = Arc::new(Status::default());
let worker_status = status.clone();
let worker = thread::Builder::new()
.name("saddle-emergency-writer".into())
.spawn(move || {
struct ActiveGuard;
impl Drop for ActiveGuard {
fn drop(&mut self) {
ACTIVE.store(false, Ordering::Release);
}
}
let _active = ActiveGuard;
let file = (|| {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
use std::os::unix::fs::OpenOptionsExt;
let file = OpenOptions::new()
.create(true)
.append(true)
.mode(0o600)
.custom_flags(
rustix::fs::OFlags::NOFOLLOW.bits() as i32
| rustix::fs::OFlags::NONBLOCK.bits() as i32,
)
.open(&path)?;
if !file.metadata()?.is_file() {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"non-regular emergency file",
));
}
rustix::fs::flock(&file, rustix::fs::FlockOperation::NonBlockingLockExclusive)
.map_err(io::Error::from)?;
Ok(file)
})();
if let Err(error) = &file {
worker_status.fail(1, error);
}
worker_status.initialized.store(true, Ordering::Release);
let mut output = file;
let mut stderr_attempted = false;
loop {
match receiver.recv_timeout(Duration::from_millis(10)) {
Ok(mut bytes) => {
bytes.resolve_on_worker(&worker_status);
deliver_with_disclosure(
&mut output,
&mut io::stderr(),
&path,
bytes.bytes(),
&worker_status,
&mut stderr_attempted,
!matches!(bytes, DiagnosticPacket::Original(_)),
);
}
Err(mpsc::RecvTimeoutError::Timeout) => {
if worker_status.closed.load(Ordering::Acquire)
&& worker_status.submitting.load(Ordering::Acquire) == 0
{
break;
}
}
Err(mpsc::RecvTimeoutError::Disconnected) => break,
}
}
})
.map_err(|error| {
ACTIVE.store(false, Ordering::Release);
EmergencyInitError::Spawn(error.kind())
})?;
Ok(Self {
handle: EmergencyDiagnosticHandle { sender, status },
worker: Some(worker),
target,
})
}
pub fn handle(&self) -> EmergencyDiagnosticHandle {
self.handle.clone()
}
pub fn target(&self) -> &std::path::Path {
&self.target
}
pub fn snapshot(&self) -> DiagnosticOutputSnapshot {
self.handle.snapshot()
}
pub fn shutdown(&mut self) -> DiagnosticShutdown {
self.handle.status.closed.store(true, Ordering::Release);
if self
.worker
.as_ref()
.is_some_and(|worker| !worker.is_finished())
{
return DiagnosticShutdown::Pending;
}
match self.worker.take().map(JoinHandle::join) {
Some(Err(_)) => DiagnosticShutdown::WorkerPanicked,
_ => DiagnosticShutdown::Finished,
}
}
}
impl Drop for EmergencyDiagnostics {
fn drop(&mut self) {
self.handle.status.closed.store(true, Ordering::Release);
}
}
#[cfg(test)]
fn deliver<W: Write, E: Write>(
output: &mut io::Result<W>,
stderr: &mut E,
target: &std::path::Path,
bytes: &[u8],
status: &Status,
stderr_attempted: &mut bool,
) {
deliver_with_disclosure(
output,
stderr,
target,
bytes,
status,
stderr_attempted,
true,
);
}
#[allow(clippy::too_many_arguments)] fn deliver_with_disclosure<W: Write, E: Write>(
output: &mut io::Result<W>,
stderr: &mut E,
target: &std::path::Path,
bytes: &[u8],
status: &Status,
stderr_attempted: &mut bool,
allow_stderr_content: bool,
) {
let result = match output {
Ok(writer) => writer.write_all(bytes),
Err(error) => Err(error
.raw_os_error()
.map(io::Error::from_raw_os_error)
.unwrap_or_else(|| io::Error::from(error.kind()))),
};
if let Err(error) = result {
status.fail(2, &error);
status.output_failed.fetch_add(1, Ordering::Relaxed);
if !*stderr_attempted {
*stderr_attempted = true;
let failure = serde_json::json!({"event":"framework.diagnostic.output_failed", "target":target.to_string_lossy(), "io_kind":format!("{:?}", error.kind()), "os_code":error.raw_os_error()});
if serde_json::to_writer(&mut *stderr, &failure)
.and_then(|_| stderr.write_all(b"\n").map_err(serde_json::Error::io))
.is_err()
|| (allow_stderr_content && stderr.write_all(bytes).is_err())
{
status.stderr_failed.fetch_add(1, Ordering::Relaxed);
}
} else {
status.stderr_suppressed.fetch_add(1, Ordering::Relaxed);
}
} else {
status.written.fetch_add(1, Ordering::Release);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn ae_paused_queue_partial_full_disconnect_and_admitted_encoding() {
use crate::{
AdmissionCapacityFacts, AdmissionConstruction, AdmissionEventContext,
CapacityDimension, Observer, ObserverConfig,
};
use saddle_admission::*;
use saddle_core::*;
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let pending = freeze_deployment_resource_budget(
2, 128, 5000, 1, 3, 1_000_000, 1_000_000, 1_000_000, 1_000_000,
)
.unwrap()
.reserve_known_database_memory(0, 0, 0)
.ok()
.unwrap();
let (application, listener) = BootstrapRendezvousIssuer::issue()
.freeze_application(GeneratedApplicationFreezeSource::new(
"app",
b"descriptor",
&["route"],
))
.unwrap();
let listener = listener
.freeze_listener(ListenerStartupFreezeSource::new(
"app",
"127.0.0.1:8000".parse().unwrap(),
"127.0.0.1:9000".parse().unwrap(),
Duration::from_millis(5000),
))
.ok()
.unwrap();
let (whole, receipt) = pair_bootstrap_rendezvous(application, listener)
.ok()
.unwrap();
let budget = bind_deployment_resource_budget_bootstrap(pending, whole, receipt)
.ok()
.unwrap();
let process = prepare_profusegw_lightweight_profile(budget).ok().unwrap();
let (permit, original) = match process.verified_profile().try_admit_observed() {
ProfuseGwLightweightObservedAdmissionOutcome::Ready(permit, receipt) => {
(permit, receipt)
}
_ => panic!("fixture admission failed"),
};
let facts = AdmissionCapacityFacts::from_receipt(&original);
let observer = Observer::with_writer(ObserverConfig::default(), io::sink()).unwrap();
let app = ContextLabel::checked("app").unwrap();
let context = AdmissionEventContext::Unrooted {
application: &app,
lifecycle: RequestViewPhase::Reading,
};
let (sender, receiver) = mpsc::sync_channel(QUEUE_CAPACITY);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
for _ in 0..63 {
assert_eq!(
handle.submit_fixed_record(&0),
DiagnosticSubmission::Enqueued
);
}
let ledger = ProcessLedger::new(ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
let result = observer.record_admission_capacity(
Some(&handle),
context,
&facts,
AdmissionConstruction::Cancelled,
);
assert_eq!(result.cpu, DiagnosticSubmission::Enqueued);
assert_eq!(result.memory, DiagnosticSubmission::Full);
assert_eq!(result.database, DiagnosticSubmission::Full);
assert_eq!(result.profuse_contract, DiagnosticSubmission::Full);
})
.unwrap();
let mut future = std::pin::pin!(future);
assert!(
matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Ready(Ok(_))
),
"actual unrooted encode must not escape the request allocation envelope"
);
for _ in 0..63 {
receiver.try_recv().unwrap();
}
let row: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(row["capacity_dimension"], "cpu");
assert_eq!(row["decision"], "accepted");
assert_eq!(row["construction"]["status"], "cancelled");
assert_eq!(handle.snapshot().enqueued, 64);
assert_eq!(handle.snapshot().dropped, 3);
assert_eq!(
handle.submit_fixed_record(&"x".repeat(8192)),
DiagnosticSubmission::EncodingFailed
);
assert!(receiver.try_recv().is_err());
drop(receiver);
let closed = observer.record_admission_capacity(
Some(&handle),
context,
&facts,
AdmissionConstruction::Ready,
);
assert_eq!(
[
closed.cpu,
closed.memory,
closed.database,
closed.profuse_contract
],
[DiagnosticSubmission::Closed; 4]
);
for dimension in [
CapacityDimension::Cpu,
CapacityDimension::Memory,
CapacityDimension::Database,
CapacityDimension::ProfuseContract,
] {
assert_eq!(
observer
.metrics_snapshot()
.capacity_accepted(Some(dimension)),
2
);
}
assert_eq!(original.decision(), ProfuseGwCapacityDecision::Accepted);
permit.into_execution().cancel();
process.finish().unwrap();
}
#[test]
fn root_original_partial_cycle_full_and_disclosure_are_explicit() {
use crate::root_diagnostic::OriginalCaptureState;
use crate::{RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent};
use saddle_core::*;
#[derive(Debug)]
struct BadFormat;
impl std::fmt::Display for BadFormat {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("exact-before-format-error")?;
Err(std::fmt::Error)
}
}
impl std::error::Error for BadFormat {}
#[derive(Debug)]
struct Cycle;
impl std::fmt::Display for Cycle {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.write_str("cycle-original")
}
}
impl std::error::Error for Cycle {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
Some(self)
}
}
let root = RequestRootPublisher::create(
ContextLabel::checked("test").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Handler));
let (sender, receiver) = mpsc::sync_channel(64);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let source = |error: &(dyn std::error::Error + 'static), output| {
RootDiagnosticScope::new(&view, output).source_error(
error,
DiagnosticStage::RequestHandler,
RootRequestEvent::Handler,
RootOutcomeFacts::default(),
)
};
let missing = source(&BadFormat, None);
assert_eq!(
missing.original_capture(),
OriginalCaptureState::OutputUnavailable
);
let partial = source(&BadFormat, Some(&handle));
assert_eq!(
partial.original_capture(),
OriginalCaptureState::FormattingFailed
);
assert_eq!(partial.submission(), DiagnosticSubmission::EncodingFailed);
let mut rows: Vec<serde_json::Value> = receiver
.try_iter()
.map(|p| serde_json::from_slice(p.bytes()).unwrap())
.collect();
assert_eq!(rows.last().unwrap()["payload"], "exact-before-format-error");
assert_eq!(rows.last().unwrap()["state"], "formatting_failed");
assert!(!rows.iter().any(|r| r["state"] == "exposed_chain_complete"));
let cycle = source(&Cycle, Some(&handle));
assert_eq!(cycle.original_capture(), OriginalCaptureState::CauseCycle);
rows = receiver
.try_iter()
.map(|p| serde_json::from_slice(p.bytes()).unwrap())
.collect();
assert_eq!(rows.last().unwrap()["state"], "cause_cycle");
for _ in 0..64 {
assert_eq!(
handle.submit_original_record(&0),
DiagnosticSubmission::Enqueued
);
}
let full = source(&Cycle, Some(&handle));
assert_eq!(
full.original_capture(),
OriginalCaptureState::OutputRejected
);
assert_eq!(full.submission(), DiagnosticSubmission::Full);
let _ = receiver.try_iter().count();
handle.status.closed.store(true, Ordering::Release);
assert_eq!(
source(&Cycle, Some(&handle)).submission(),
DiagnosticSubmission::Closed
);
let mut file: io::Result<Vec<u8>> = Err(io::ErrorKind::PermissionDenied.into());
let mut stderr = Vec::new();
deliver_with_disclosure(
&mut file,
&mut stderr,
std::path::Path::new("saddle.emergency.log"),
b"CT_ARTIFICIAL_SECRET_NOT_A_CREDENTIAL",
&Status::default(),
&mut false,
false,
);
let stderr = String::from_utf8(stderr).unwrap();
assert!(stderr.contains("output_failed"));
assert!(!stderr.contains("CT_ARTIFICIAL_SECRET"));
}
#[test]
fn root_original_admitted_and_existing_panic_keep_identity() {
use crate::root_diagnostic::OriginalCaptureState;
use crate::{RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent};
use saddle_core::*;
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let root = RequestRootPublisher::create(
ContextLabel::checked("test").unwrap(),
ContextFact::Unavailable,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Handler));
let (sender, receiver) = mpsc::sync_channel(64);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let diagnostic = Diagnostic::capture(
DiagnosticCategory::Panic,
CaptureSite::Origin,
DiagnosticCause::new(
DiagnosticStage::RequestHandler,
DiagnosticCode::new("panic.caught").unwrap(),
),
);
let id = diagnostic.id();
let error = io::Error::from(io::ErrorKind::TimedOut);
let ledger = saddle_admission::ProcessLedger::new(saddle_admission::ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
let scope = RootDiagnosticScope::new(&view, Some(&handle));
let failure = scope.source_error(
&error,
DiagnosticStage::RequestDb,
RootRequestEvent::Database,
RootOutcomeFacts::default(),
);
assert_eq!(
failure.original_capture(),
OriginalCaptureState::CompleteEnqueued
);
let panic = scope.source_existing_description(
&"raw panic: CT_ARTIFICIAL_SECRET_NOT_A_CREDENTIAL",
diagnostic,
DiagnosticCode::new("task.failed").unwrap(),
RootRequestEvent::Supervision,
RootOutcomeFacts::default(),
);
assert_eq!(
panic.original_capture(),
OriginalCaptureState::CompleteEnqueued
);
panic
.finish(&view, None, RootOutcomeFacts::default())
.ok()
.unwrap()
})
.unwrap();
let mut future = std::pin::pin!(future);
assert!(
matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Ready(Ok(_))
),
"source formatter/queue must not escape admission"
);
let rows: Vec<serde_json::Value> = receiver
.try_iter()
.map(|p| serde_json::from_slice(p.bytes()).unwrap())
.collect();
let panic: Vec<_> = rows
.iter()
.filter(|r| r["occurrence"]["diagnostic_id"] == id)
.collect();
assert_eq!(panic.len(), 4);
let header: serde_json::Value =
serde_json::from_str(panic[0]["payload"].as_str().unwrap()).unwrap();
assert_eq!(header["facts"]["diagnostic_id"], id);
assert_eq!(header["facts"]["category"], "panic");
assert_eq!(
header["source_interface"],
"unavailable_description_debug_only"
);
assert_eq!(
panic[1]["payload"],
"raw panic: CT_ARTIFICIAL_SECRET_NOT_A_CREDENTIAL"
);
assert_eq!(panic[2]["channel"], "debug");
assert_eq!(panic[3]["state"], "description_complete_source_unavailable");
}
use saddle_core::{
CaptureSite, DiagnosticCategory, DiagnosticCause, DiagnosticCode, DiagnosticStage,
};
#[test]
fn root_contract_oversize_detail_and_basic_encoding_failure_are_explicit() {
use crate::{RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent};
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(4);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let text = "m".repeat(256);
let protocol = "\\\"".repeat(128);
let call = CallContext::new(
ApplicationId::new(text.clone()),
ModuleId::new(text.clone()),
ServiceId::new(text.clone()),
OperationId::new(text.clone()),
TraceId::from_u128(1),
SpanId::from_u64(1),
)
.with_trace_correlation_id(TraceCorrelationId::new(protocol.clone()).unwrap())
.with_rpc_correlation_id(RpcCorrelationId::new(&protocol));
let mut root = RequestRootPublisher::create(
ContextLabel::checked(&text).unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
assert!(
root.publish(
RequestIdentityGroup::from_validated(
&call,
&protocol,
&text,
1,
ContextFact::Present(ContextLabel::checked(&text).unwrap())
)
.unwrap()
)
.is_ok()
);
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Database));
let metadata = "T".repeat(192);
let cause = BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("driver.type").unwrap(),
)
.with_object(InlineDiagnosticText::metadata(&metadata))
.with_types(
Some(InlineDiagnosticText::metadata(&metadata)),
Some(InlineDiagnosticText::metadata(&metadata)),
None,
);
let detail = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::Origin,
cause,
)
.wrap(cause)
.wrap(cause)
.wrap(cause);
let id = detail.id();
let receipt = RootDiagnosticScope::new(&view, Some(&handle)).source(
detail,
DiagnosticCode::new("driver.type").unwrap(),
RootRequestEvent::Database,
RootOutcomeFacts::default(),
);
assert_eq!(receipt.submission(), DiagnosticSubmission::Enqueued);
let row: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(row["detail_status"], "omitted_encoding_capacity");
assert_eq!(row["occurrence"]["diagnostic_id"], id);
assert!(row["diagnostic"].is_null());
assert_eq!(handle.snapshot().truncated, 1);
let (receipt, result) = receipt
.boundary(
&view,
Some(&handle),
RootRequestEvent::Database,
RootOutcomeFacts::default(),
)
.ok()
.unwrap();
assert_eq!(result, DiagnosticSubmission::EncodingFailed);
assert!(receiver.try_recv().is_err());
assert_eq!(receipt.submission(), DiagnosticSubmission::Enqueued);
assert_eq!(handle.snapshot().dropped, 1);
}
#[test]
fn root_contract_stage_terminal_is_single_and_admitted() {
use crate::{
Observer, ObserverConfig, RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent,
Stage, StageOutcome,
};
use saddle_core::*;
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let observer = Observer::with_writer(ObserverConfig::default(), std::io::sink()).unwrap();
let root = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Handler));
let foreign = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let foreign_view = foreign
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Handler));
let ledger = saddle_admission::ProcessLedger::new(saddle_admission::ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
let facts = |operation| RootOutcomeFacts {
axes: DiagnosticOutcomeAxes {
operation,
..Default::default()
},
..Default::default()
};
let stage = RootDiagnosticScope::new(&view, None)
.start_stage(&observer, RootRequestEvent::Response);
let stage = stage
.finish_nonfailure(facts(OperationOutcome::Failed))
.err()
.unwrap();
assert_eq!(
stage
.finish_nonfailure(facts(OperationOutcome::Succeeded))
.ok()
.unwrap(),
DiagnosticSubmission::OutputUnavailable
);
let source = RootDiagnosticScope::new(&view, None).source(
BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::Origin,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestHandler,
DiagnosticCode::new("handler.failed").unwrap(),
),
),
DiagnosticCode::new("handler.failed").unwrap(),
RootRequestEvent::Handler,
facts(OperationOutcome::Failed),
);
let foreign_stage = RootDiagnosticScope::new(&foreign_view, None)
.start_stage(&observer, RootRequestEvent::Response);
let (foreign_stage, source) = foreign_stage
.finish_failure(source, facts(OperationOutcome::Failed))
.err()
.unwrap();
assert_eq!(
foreign_stage
.finish_nonfailure(facts(OperationOutcome::Rejected))
.ok()
.unwrap(),
DiagnosticSubmission::OutputUnavailable
);
let stage = RootDiagnosticScope::new(&view, None)
.start_stage(&observer, RootRequestEvent::Response);
let (source, status) = stage
.finish_failure(source, facts(OperationOutcome::Failed))
.ok()
.unwrap();
assert_eq!(status, DiagnosticSubmission::OutputUnavailable);
source
.finish(&view, None, facts(OperationOutcome::Failed))
.ok()
.unwrap()
})
.unwrap();
let mut future = std::pin::pin!(future);
assert!(matches!(
future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop())),
Poll::Ready(Ok(_))
));
let metrics = observer.metrics_snapshot();
assert_eq!(metrics.requests(StageOutcome::Success), 1);
assert_eq!(metrics.requests(StageOutcome::Rejected), 1);
assert_eq!(metrics.requests(StageOutcome::Failure), 1);
assert_eq!(metrics.requests(StageOutcome::Cancelled), 0);
assert_eq!(
metrics.stage_latency(Stage::Response).iter().sum::<u64>(),
3
);
}
#[test]
fn root_contract_existing_source_cleanup_and_admitted_submission() {
use crate::{RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent};
use saddle_core::*;
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let (sender, receiver) = mpsc::sync_channel(4);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let root = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::Unavailable,
)
.unwrap();
let view = root.reference().view(
RequestLocalFacts::new(RequestViewPhase::Handler).with_task(ContextFact::Present(17)),
);
let diagnostic = Diagnostic::capture(
DiagnosticCategory::Panic,
CaptureSite::Origin,
DiagnosticCause::new(
DiagnosticStage::RequestHandler,
DiagnosticCode::new("panic.caught").unwrap(),
),
);
let id = diagnostic.id();
let first = RootDiagnosticScope::new(&view, Some(&handle)).source_existing(
diagnostic,
DiagnosticCode::new("task.failed").unwrap(),
RootRequestEvent::Supervision,
RootOutcomeFacts::default(),
);
let occurrence = first.occurrence();
let ledger = saddle_admission::ProcessLedger::new(saddle_admission::ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
let cleanup = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::Origin,
BoundedDiagnosticCause::new(
DiagnosticStage::FinalizerResource,
DiagnosticCode::new("cleanup.failed").unwrap(),
),
)
.during_cleanup_occurrence(occurrence);
let second = RootDiagnosticScope::new(&view, Some(&handle)).source(
cleanup,
DiagnosticCode::new("cleanup.failed").unwrap(),
RootRequestEvent::Finalization,
RootOutcomeFacts::default(),
);
assert_eq!(second.submission(), DiagnosticSubmission::Enqueued);
first
.finish(&view, None, RootOutcomeFacts::default())
.ok()
.unwrap()
})
.unwrap();
let mut future = std::pin::pin!(future);
let result = future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop()));
assert!(
matches!(result, Poll::Ready(Ok(_))),
"bounded receipt path must not allocate outside admitted contract: {result:?}"
);
let primary: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let cleanup: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(primary["diagnostic"]["diagnostic_id"], id);
assert_eq!(primary["category"], "panic");
assert_eq!(
primary["diagnostic"]["stack_status"],
"unavailable_deferred"
);
assert_eq!(cleanup["occurrence"]["primary_diagnostic_id"], id);
assert_ne!(cleanup["occurrence"]["diagnostic_id"], id);
assert_eq!(cleanup["context"], primary["context"]);
}
#[test]
fn root_contract_real_file_readback_preserves_safe_facts_without_secret() {
use crate::{RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent};
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(1);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let root = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let cause = BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.type").unwrap(),
)
.with_object(InlineDiagnosticText::metadata(
"mysql://SECRET_SENTINEL@host",
))
.with_sqlstate("22003")
.with_database_code(1264)
.with_column(Some(3), Some(4))
.with_types(
Some(InlineDiagnosticText::metadata("i32")),
Some(InlineDiagnosticText::metadata("BIGINT")),
None,
);
let failure = RootDiagnosticScope::new(&view, Some(&handle)).source(
BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::Origin,
cause,
),
DiagnosticCode::new("db.type").unwrap(),
RootRequestEvent::Database,
RootOutcomeFacts::default(),
);
let public = failure
.finish(&view, None, RootOutcomeFacts::default())
.ok()
.unwrap();
drop((view, root));
let path = std::env::temp_dir().join(format!(
"saddle-c-readback-{}-{}.jsonl",
std::process::id(),
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap()
.as_nanos()
));
let file = std::fs::OpenOptions::new()
.create_new(true)
.write(true)
.open(&path)
.unwrap();
let packet = receiver.try_recv().unwrap();
let mut output = Ok(file);
deliver(
&mut output,
&mut io::sink(),
&path,
packet.bytes(),
&handle.status,
&mut false,
);
output.as_mut().unwrap().flush().unwrap();
let raw = std::fs::read(&path).unwrap();
std::fs::remove_file(&path).unwrap();
assert!(!String::from_utf8_lossy(&raw).contains("SECRET_SENTINEL"));
let row: serde_json::Value = serde_json::from_slice(&raw).unwrap();
assert_eq!(row["diagnostic"]["causes"][0]["object"]["redacted"], true);
assert_eq!(row["diagnostic"]["causes"][0]["sqlstate"]["value"], "22003");
assert_eq!(row["diagnostic"]["causes"][0]["column_index"], 3);
assert_eq!(handle.snapshot().written, 1);
assert_eq!(public.source_submission(), DiagnosticSubmission::Enqueued);
assert_eq!(
public.terminal_submission(),
DiagnosticSubmission::OutputUnavailable
);
}
#[test]
fn root_contract_source_boundary_supervision_and_queue_degradation() {
use crate::{
RootDiagnosticScope, RootOutcomeFacts, RootRequestEvent, RootSupervisionReturn,
};
use saddle_core::{
BoundedDiagnostic, BoundedDiagnosticCause, ContextFact, ContextLabel,
RequestLocalFacts, RequestRootPublisher, RequestViewPhase,
};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let root = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap();
let view = root
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let make = |output| {
RootDiagnosticScope::new(&view, output).source(
BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::Origin,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDecode,
DiagnosticCode::new("read.failed").unwrap(),
),
),
DiagnosticCode::new("decode.failed").unwrap(),
RootRequestEvent::Ingress,
RootOutcomeFacts::default(),
)
};
let missing = make(None);
assert_eq!(
missing.submission(),
DiagnosticSubmission::OutputUnavailable
);
let receipt = make(Some(&handle));
assert_eq!(receipt.submission(), DiagnosticSubmission::Enqueued);
let receipt = RootSupervisionReturn::<()>::failed(receipt)
.consume()
.err()
.unwrap();
let foreign = RequestRootPublisher::create(
ContextLabel::checked("app").unwrap(),
ContextFact::NotEstablished,
)
.unwrap()
.reference()
.view(RequestLocalFacts::new(RequestViewPhase::Reading));
let receipt = receipt
.boundary(
&foreign,
Some(&handle),
RootRequestEvent::Supervision,
RootOutcomeFacts::default(),
)
.err()
.unwrap();
let public = receipt
.finish(&view, Some(&handle), RootOutcomeFacts::default())
.ok()
.unwrap();
assert_eq!(make(Some(&handle)).submission(), DiagnosticSubmission::Full);
handle.status.closed.store(true, Ordering::Release);
assert_eq!(
make(Some(&handle)).submission(),
DiagnosticSubmission::Closed
);
drop((missing, foreign, view, root));
let source: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let boundary: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let public = serde_json::to_value(public).unwrap();
assert_eq!(source["occurrence"], boundary["occurrence"]);
assert_eq!(source["occurrence"], public["occurrence"]);
assert_eq!(source["context"], boundary["source_context"]);
assert_eq!(source["context"]["trace_id"]["state"], "not_established");
assert!(source["diagnostic"].is_object());
assert!(boundary["diagnostic"].is_null());
assert!(public.get("context").is_none());
assert_eq!(handle.snapshot().written, 0);
assert_eq!(handle.snapshot().enqueued, 2);
assert_eq!(handle.snapshot().dropped, 2);
}
#[test]
fn outbound_child_queue_source_boundary_full_closed() {
use crate::{EventContext, RequestDiagnosticScope, RequestIdentity, RouteIdentity};
use saddle_core::{
BoundedDiagnosticCause, CallContext, DiagnosticOutcomeAxes, RpcCorrelationId, SpanId,
TraceId,
};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let call = |rpc: &str, span| {
CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"operation".into(),
TraceId::from_u128(1),
SpanId::from_u64(span),
)
.with_rpc_correlation_id(RpcCorrelationId::new(rpc))
};
let parent_call = call("0", 1);
let child_call = call("0.1", 2);
let event = EventContext::new(
RequestIdentity::new("request-1").unwrap(),
RouteIdentity::new("route").unwrap(),
1,
)
.unwrap();
let parent = RequestDiagnosticScope::established(&handle, &parent_call, &event);
let child = parent.derive_outbound_child(&child_call, &event).unwrap();
let fail = || {
child.fail(
7,
DiagnosticCategory::UnexpectedError,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDecode,
DiagnosticCode::new("outbound.failed").unwrap(),
),
)
};
let source = fail();
assert_eq!(source.submission(), DiagnosticSubmission::Enqueued);
let (retained, delivery) =
source.finish_boundary_retained(Some(&handle), &DiagnosticOutcomeAxes::default());
assert_eq!(*retained.error(), 7);
assert_eq!(
delivery.boundary_submission(),
DiagnosticSubmission::Enqueued
);
assert_eq!(fail().submission(), DiagnosticSubmission::Full);
handle.status.closed.store(true, Ordering::Release);
assert_eq!(fail().submission(), DiagnosticSubmission::Closed);
let source: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let boundary: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(source["context"], boundary["context"]);
assert_eq!(
source["diagnostic_reference"],
boundary["diagnostic_reference"]
);
assert_eq!(source["context"]["rpc_id"]["value"], "0.1");
assert_eq!(source["context"]["span_id"]["value"], "0000000000000002");
assert!(source["diagnostic"].is_object());
assert!(boundary["diagnostic"].is_null());
assert_eq!(handle.snapshot().written, 0);
}
#[test]
fn early_request_source_boundary_admitted_full_closed_and_missing_output() {
use crate::{DiagnosticRequestPhase, EarlyRequestContext, RequestDiagnosticScope};
use saddle_core::{BoundedDiagnosticCause, DiagnosticOutcomeAxes, OperationOutcome};
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let known_call = saddle_core::CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"operation".into(),
saddle_core::TraceId::from_u128(1),
saddle_core::SpanId::from_u64(2),
);
let known_event = crate::EventContext::new(
crate::RequestIdentity::new("real-request").unwrap(),
crate::RouteIdentity::new("/route").unwrap(),
1,
)
.unwrap();
let ledger = saddle_admission::ProcessLedger::new(saddle_admission::ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
let axes = DiagnosticOutcomeAxes {
operation: OperationOutcome::Failed,
..Default::default()
};
let cause = || {
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDecode,
DiagnosticCode::new("request.read_failed").unwrap(),
)
};
let enriched = RequestDiagnosticScope::early(
None,
EarlyRequestContext::socket_accepted("app"),
)
.with_task(crate::DiagnosticTaskId::from_runtime_id("42").unwrap())
.unwrap_or_else(|_| panic!("task conflict"))
.bind_established(&known_call, &known_event)
.unwrap_or_else(|_| panic!("context conflict"));
let unchanged =
enriched
.reborrow()
.fail(22u32, DiagnosticCategory::UnexpectedError, cause());
assert_eq!(
unchanged.submission(),
DiagnosticSubmission::OutputUnavailable
);
let scope = RequestDiagnosticScope::early(
Some(&handle),
EarlyRequestContext::socket_accepted("app"),
)
.with_phase(DiagnosticRequestPhase::ReadingHead);
let failure = scope.fail(23u32, DiagnosticCategory::UnexpectedError, cause());
let id = failure.source_diagnostic().id();
assert_eq!(failure.submission(), DiagnosticSubmission::Enqueued);
let (failure, delivered) = failure.finish_boundary_retained(Some(&handle), &axes);
assert_eq!(
delivered.boundary_submission(),
DiagnosticSubmission::Enqueued
);
assert_eq!(*failure.error(), 23);
let full = scope.fail(24u32, DiagnosticCategory::UnexpectedError, cause());
assert_eq!(full.submission(), DiagnosticSubmission::Full);
let (full, delivered) = full.finish_boundary_retained(None, &axes);
assert_eq!(*full.error(), 24);
assert_eq!(delivered.source_submission(), DiagnosticSubmission::Full);
assert_eq!(
delivered.boundary_submission(),
DiagnosticSubmission::OutputUnavailable
);
handle.status.closed.store(true, Ordering::Release);
let closed = scope.fail(25u32, DiagnosticCategory::UnexpectedError, cause());
assert_eq!(closed.submission(), DiagnosticSubmission::Closed);
let no_output = scope.with_output(None).fail(
26u32,
DiagnosticCategory::UnexpectedError,
cause(),
);
assert_eq!(
no_output.submission(),
DiagnosticSubmission::OutputUnavailable
);
assert_eq!(*no_output.error(), 26);
id
})
.unwrap();
let mut future = std::pin::pin!(future);
let result = future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop()));
assert!(
matches!(result, Poll::Ready(Ok(_))),
"early path escaped: {result:?}"
);
let source: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let boundary: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(source["context"], boundary["context"]);
assert_eq!(
source["diagnostic_reference"],
boundary["diagnostic_reference"]
);
assert_eq!(source["context"]["trace_id"]["state"], "not_established");
assert_eq!(
source["context"]["lifecycle"]["value"]["value"],
"reading_head"
);
assert!(source["diagnostic"].is_object());
assert!(boundary["diagnostic"].is_null());
assert_eq!(boundary["axes"]["operation"], "failed");
assert_eq!(handle.snapshot().written, 0); }
use std::time::Instant;
#[test]
fn early_retained_reference_late_output_keeps_original_submission_and_context() {
use crate::{EarlyRequestContext, RequestDiagnosticScope};
use saddle_core::{BoundedDiagnostic, BoundedDiagnosticCause, DiagnosticOutcomeAxes};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let scope =
RequestDiagnosticScope::early(None, EarlyRequestContext::socket_accepted("app"));
let diagnostic = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDecode,
DiagnosticCode::new("request.read_failed").unwrap(),
),
);
let source = serde_json::to_value(&diagnostic).unwrap();
let reference = scope.capture_required(diagnostic).into_reference();
let (retained, delivery) =
reference.finish_retained(Some(&handle), &DiagnosticOutcomeAxes::default());
assert_eq!(
delivery.source_submission(),
DiagnosticSubmission::OutputUnavailable
);
assert_eq!(
delivery.boundary_submission(),
DiagnosticSubmission::Enqueued
);
assert_eq!(
serde_json::to_value(retained.source_diagnostic()).unwrap(),
source
);
let packet = receiver.try_recv().unwrap();
let row: serde_json::Value = serde_json::from_slice(packet.bytes()).unwrap();
assert_eq!(row["source_submission"], "output_unavailable");
assert_eq!(row["context"]["trace_id"]["state"], "not_established");
assert!(row["diagnostic"].is_null()); assert!(receiver.try_recv().is_err());
assert_eq!(handle.snapshot().written, 0);
}
static FILE_TEST: std::sync::Mutex<()> = std::sync::Mutex::new(());
#[test]
fn zone_projection_exact_missing_and_unsafe_are_explicit() {
use crate::{
DiagnosticContextMissing as Missing, DiagnosticZone, EventContext,
RequestDiagnosticScope, RequestIdentity, RouteIdentity,
};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let call = CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"query".into(),
saddle_core::TraceId::from_u128(1),
saddle_core::SpanId::from_u64(2),
);
let event = EventContext::new(
RequestIdentity::new("r1").unwrap(),
RouteIdentity::new("/query").unwrap(),
1,
)
.unwrap();
let cause = saddle_core::BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
);
for zone in ["zone-a", "杭州-zone/蓝"] {
let scope = RequestDiagnosticScope::established(&handle, &call, &event)
.with_zone(DiagnosticZone::from_validated_ingress(zone).unwrap())
.with_zone_missing(Missing::Unavailable);
let source = scope
.fail((), DiagnosticCategory::UnexpectedError, cause)
.into_reference();
source.record(&handle, &saddle_core::DiagnosticOutcomeAxes::default());
for _ in 0..2 {
let row: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert_eq!(
row["context"]["zone"],
serde_json::json!({"state":"present","value":zone})
);
}
}
for (reason, expected) in [
(Missing::NotApplicable, "not_applicable"),
(Missing::NotEstablished, "not_established"),
(Missing::Unavailable, "unavailable"),
] {
RequestDiagnosticScope::established(&handle, &call, &event)
.with_zone_missing(reason)
.fail((), DiagnosticCategory::UnexpectedError, cause);
let row: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert_eq!(
row["context"]["zone"],
serde_json::json!({"state":expected})
);
}
for invalid in [
"",
" ",
"https://user:SECRET@host",
"zone\nSECRET",
"zone?token=SECRET",
] {
assert!(DiagnosticZone::from_validated_ingress(invalid).is_err());
}
assert!(DiagnosticZone::from_validated_ingress(&"a".repeat(257)).is_err());
assert!(DiagnosticZone::from_validated_ingress(&"a".repeat(256)).is_ok());
}
#[test]
fn existing_source_receipt_keeps_stack_and_submits_once() {
use crate::{EventContext, RequestDiagnosticScope, RequestIdentity, RouteIdentity};
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let call = CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"query".into(),
saddle_core::TraceId::from_u128(1),
saddle_core::SpanId::from_u64(2),
);
let event = EventContext::new(
RequestIdentity::new("request-1").unwrap(),
RouteIdentity::new("/query").unwrap(),
1,
)
.unwrap();
let diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
),
);
let before = serde_json::to_value(&diagnostic).unwrap();
let receipt = RequestDiagnosticScope::established(&handle, &call, &event)
.capture_existing(diagnostic, None);
assert_eq!(receipt.submission(), DiagnosticSubmission::Enqueued);
let reference = receipt.into_reference();
assert_eq!(
serde_json::to_value(reference.source_diagnostic()).unwrap(),
before
);
let source: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert_eq!(source["diagnostic"], before);
assert!(receiver.try_recv().is_err());
assert_eq!(handle.snapshot().enqueued, 1);
assert_eq!(
reference.record(&handle, &saddle_core::DiagnosticOutcomeAxes::default()),
DiagnosticSubmission::Enqueued
);
let boundary: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert!(boundary["diagnostic"].is_null());
assert_eq!(
source["diagnostic_reference"],
boundary["diagnostic_reference"]
);
assert_eq!(source["context"], boundary["context"]);
assert_eq!(handle.snapshot().enqueued, 2);
assert_eq!(handle.snapshot().written, 0);
}
#[test]
fn controlled_request_source_preserves_error_projection_and_degradation() {
use crate::{EventContext, RequestDiagnosticScope, RequestIdentity, RouteIdentity};
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let call = CallContext::new(
"app".into(),
"module".into(),
"service".into(),
"query".into(),
TraceId::from_u128(1),
SpanId::from_u64(2),
);
let event = EventContext::new(
RequestIdentity::new("req-1").unwrap(),
RouteIdentity::new("/agent/{code}").unwrap(),
2,
)
.unwrap();
let scope = RequestDiagnosticScope::established(&handle, &call, &event);
let cause = BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
);
let failure = scope.fail(17u32, DiagnosticCategory::UnexpectedError, cause);
assert_eq!(failure.submission(), DiagnosticSubmission::Enqueued);
let failure = failure.map_error(|error| (error, false));
assert_eq!(failure.error(), &(17, false));
assert_eq!(
failure.record_boundary(&handle, &DiagnosticOutcomeAxes::default()),
DiagnosticSubmission::Enqueued
);
let full = scope.fail(19, DiagnosticCategory::UnexpectedError, cause);
assert_eq!(full.submission(), DiagnosticSubmission::Full);
assert_eq!(full.error(), &19);
let source: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
let boundary: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert_eq!(source["context"], boundary["context"]);
assert_eq!(
source["diagnostic_reference"],
boundary["diagnostic_reference"]
);
assert!(boundary["diagnostic"].is_null());
assert!(source["diagnostic"].is_object());
let c = &source["context"];
for (key, expected) in [
("application", "app"),
("module", "module"),
("service", "service"),
("operation", "query"),
("request", "req-1"),
("route", "/agent/{code}"),
] {
assert_eq!(c[key]["value"]["value"], expected);
}
assert_eq!(c["attempt"]["value"], 2);
assert_eq!(c["span_id"]["value"], "0000000000000002");
assert_eq!(c["rpc_id"]["state"], "unavailable");
for field in ["scope", "task", "lifecycle", "zone", "target"] {
assert_eq!(c[field]["state"], "unavailable");
}
handle.status.closed.store(true, Ordering::Release);
let closed = scope.fail(21, DiagnosticCategory::UnexpectedError, cause);
assert_eq!(closed.submission(), DiagnosticSubmission::Closed);
assert_eq!(closed.error(), &21);
assert_eq!(handle.snapshot().written, 0);
let unavailable = RequestDiagnosticScope::output_unavailable(&call, &event);
let original = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
);
let original_id = original.id();
let reference = unavailable.capture_required(original).into_reference();
assert_eq!(
reference.source_submission(),
DiagnosticSubmission::OutputUnavailable
);
assert_eq!(reference.source_diagnostic().id(), original_id);
assert_eq!(
reference.record(&handle, &DiagnosticOutcomeAxes::default()),
DiagnosticSubmission::Closed
);
}
#[test]
fn source_detail_and_boundary_reference_are_separate() {
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(4);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let cause = BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
);
let primary = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
);
let cleanup = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
)
.during_cleanup_of(&primary);
let reference = cleanup.occurrence();
let axes = DiagnosticOutcomeAxes::default();
assert_eq!(
handle.submit_bounded(Some(&cleanup), &axes, None),
DiagnosticSubmission::Enqueued
);
drop(cleanup);
assert_eq!(
handle.submit_boundary(Some(reference), &axes, None),
DiagnosticSubmission::Enqueued
);
assert_eq!(
handle.submit_boundary_parts(Some(reference), &axes, None, Some(("database", 1))),
DiagnosticSubmission::Enqueued
);
let source: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert!(source["diagnostic"].is_object());
for event in ["framework.boundary.outcome", "framework.stage.finished"] {
let boundary: serde_json::Value =
serde_json::from_slice(receiver.recv().unwrap().bytes()).unwrap();
assert_eq!(boundary["event"], event);
assert!(boundary["diagnostic"].is_null());
assert_eq!(
boundary["diagnostic_reference"],
source["diagnostic_reference"]
);
assert_eq!(
boundary["diagnostic_reference"]["primary_diagnostic_id"],
primary.id()
);
assert!(!boundary.to_string().contains("db.decode.type"));
}
assert_eq!(handle.snapshot().written, 0);
}
#[test]
fn deferred_unavailable_preserves_source_through_output_worker() {
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(1);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
),
);
let original = serde_json::to_value(&diagnostic).unwrap();
assert_eq!(original["stack_status"], "unavailable_deferred");
assert!(diagnostic.deferred_stack().is_none());
assert!(original["stack"].as_str().unwrap().is_empty());
assert_eq!(handle.submit(&diagnostic), DiagnosticSubmission::Enqueued);
let packet = receiver.try_recv().unwrap();
let queued: serde_json::Value = serde_json::from_slice(packet.bytes()).unwrap();
assert_eq!(queued["diagnostic"], original);
let status = handle.status.clone();
let row = std::thread::spawn(move || {
let mut packet = packet;
packet.resolve_on_worker(&status);
serde_json::from_slice::<serde_json::Value>(packet.bytes()).unwrap()
})
.join()
.unwrap();
for key in ["diagnostic_id", "primary_diagnostic_id", "origin", "causes"] {
assert_eq!(row["diagnostic"][key], original[key]);
}
assert_eq!(row["diagnostic"]["stack_status"], "unavailable_deferred");
assert!(row["diagnostic"]["stack"].as_str().unwrap().is_empty());
assert_eq!(serde_json::to_value(&diagnostic).unwrap(), original);
}
#[test]
fn request_stage_finishes_once_with_original_receipt_context() {
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(4);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let observer =
crate::Observer::with_writer(crate::ObserverConfig::default(), io::sink()).unwrap();
let (call, _) = observer
.start_external_call_checked("app", "module", "service", "query", None)
.unwrap();
let event = crate::EventContext::new(
crate::RequestIdentity::new("original-request").unwrap(),
crate::RouteIdentity::new("/query").unwrap(),
1,
)
.unwrap();
let mut scope = crate::RequestDiagnosticScope::established(&handle, call.context(), &event)
.with_zone(crate::DiagnosticZone::from_validated_ingress("zone-a").unwrap())
.with_db_operation(
crate::DiagnosticDbOperation::from_registered("AgentByCode").unwrap(),
)
.with_db_operation_missing(crate::DiagnosticContextMissing::Unavailable);
let reference = scope
.capture_required(BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
),
))
.into_reference();
scope.set_db_operation(
crate::DiagnosticDbOperation::from_registered("NextOperation").unwrap(),
);
let stage = observer.start_stage(call.context(), crate::Stage::Database, event.clone());
let axes = DiagnosticOutcomeAxes {
operation: OperationOutcome::Failed,
physical: PhysicalDispositionFact::Returned,
..Default::default()
};
let delivery = stage.finish_request_failure(Some(&handle), &axes, &reference);
assert_eq!(
delivery.boundary_submission(),
DiagnosticSubmission::Enqueued
);
let source: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
let boundary: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(source["context"], boundary["context"]);
assert_eq!(source["context"]["db_operation"]["value"], "AgentByCode");
assert_eq!(source["context"]["operation"]["value"]["value"], "query");
let next = scope.fail(
(),
DiagnosticCategory::ExpectedRejection,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.next").unwrap(),
),
);
assert_eq!(next.submission(), DiagnosticSubmission::Enqueued);
let next: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert_eq!(next["context"]["db_operation"]["value"], "NextOperation");
for field in ["operation", "request", "scope", "zone"] {
assert_eq!(next["context"][field], source["context"][field]);
}
assert!(crate::DiagnosticDbOperation::from_registered("SELECT * FROM private").is_none());
assert!(crate::DiagnosticDbOperation::from_registered("https://secret@host").is_none());
assert_eq!(
source["diagnostic_reference"],
boundary["diagnostic_reference"]
);
assert_eq!(boundary["source_submission"], "enqueued");
assert_eq!(boundary["axes"]["operation"], "failed");
assert!(boundary["diagnostic"].is_null());
assert!(boundary["elapsed_ms"].is_u64());
assert!(receiver.try_recv().is_err());
let stage = observer.start_stage(call.context(), crate::Stage::Database, event);
assert_eq!(
stage
.finish_request_failure(None, &axes, &reference)
.boundary_submission(),
DiagnosticSubmission::OutputUnavailable
);
assert!(receiver.try_recv().is_err());
assert_eq!(
observer
.metrics_snapshot()
.stage_latency(crate::Stage::Database)
.iter()
.sum::<u64>(),
2
);
call.succeed();
let _ = observer.shutdown();
}
#[test]
fn nonfailure_finish_optional_output_and_no_failed_bypass() {
use saddle_core::*;
let (sender, receiver) = mpsc::sync_channel(2);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let observer =
crate::Observer::with_writer(crate::ObserverConfig::default(), io::sink()).unwrap();
let (call, _) = observer
.start_external_call_checked("app", "module", "service", "query", None)
.unwrap();
let event = crate::EventContext::new(
crate::RequestIdentity::new("r").unwrap(),
crate::RouteIdentity::new("/query").unwrap(),
1,
)
.unwrap();
for operation in [OperationOutcome::Succeeded, OperationOutcome::Rejected] {
let stage = observer.start_stage(call.context(), crate::Stage::Database, event.clone());
let axes = DiagnosticOutcomeAxes {
operation,
..Default::default()
};
assert!(matches!(
stage.finish_nonfailure_bounded(Some(&handle), &axes),
Ok(DiagnosticSubmission::Enqueued)
));
let row: serde_json::Value =
serde_json::from_slice(receiver.try_recv().unwrap().bytes()).unwrap();
assert!(row["diagnostic_reference"].is_null());
assert!(row["diagnostic"].is_null());
assert!(receiver.try_recv().is_err());
let stage = observer.start_stage(call.context(), crate::Stage::Database, event.clone());
assert!(matches!(
stage.finish_nonfailure_bounded(None, &axes),
Ok(DiagnosticSubmission::OutputUnavailable)
));
}
for operation in [
OperationOutcome::Failed,
OperationOutcome::Panicked,
OperationOutcome::Unknown,
OperationOutcome::TimedOut,
OperationOutcome::Cancelled,
] {
let stage = observer.start_stage(call.context(), crate::Stage::Database, event.clone());
let axes = DiagnosticOutcomeAxes {
operation,
..Default::default()
};
let stage = match stage.finish_nonfailure_bounded(None, &axes) {
Err(stage) => stage,
Ok(_) => panic!("technical axis bypass"),
};
assert!(receiver.try_recv().is_err());
let reference =
crate::RequestDiagnosticScope::output_unavailable(call.context(), &event)
.capture_required(BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
),
))
.into_reference();
assert_eq!(
stage
.finish_request_failure(None, &axes, &reference)
.boundary_submission(),
DiagnosticSubmission::OutputUnavailable
);
}
assert_eq!(
observer
.metrics_snapshot()
.stage_latency(crate::Stage::Database)
.iter()
.sum::<u64>(),
8
);
call.succeed();
let _ = observer.shutdown();
}
#[test]
fn bounded_diagnostic_admitted_no_escape_full_closed() {
use saddle_core::*;
use std::{
future::Future,
task::{Context, Poll, Waker},
};
let (sender, receiver) = mpsc::sync_channel(1);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let observer =
crate::Observer::with_writer(crate::ObserverConfig::default(), io::sink()).unwrap();
let (call, _) = observer
.start_external_call_checked("app", "module", "service", "query", None)
.unwrap();
let stage = observer.start_stage(
call.context(),
crate::Stage::Database,
crate::EventContext::new(
crate::RequestIdentity::new("r1").unwrap(),
crate::RouteIdentity::new("query").unwrap(),
1,
)
.unwrap(),
);
let diagnostic_event = crate::EventContext::new(
crate::RequestIdentity::new("r1").unwrap(),
crate::RouteIdentity::new("/query").unwrap(),
1,
)
.unwrap();
let diagnostic_call = call.context().clone();
let success_stage = observer.start_stage(
call.context(),
crate::Stage::Database,
diagnostic_event.clone(),
);
let receipt_stage = observer.start_stage(
call.context(),
crate::Stage::Database,
diagnostic_event.clone(),
);
let unavailable_stage = observer.start_stage(
call.context(),
crate::Stage::Database,
diagnostic_event.clone(),
);
let ledger = saddle_admission::ProcessLedger::new(saddle_admission::ResourceConfig {
managed_capacity: 1024 * 1024,
entry_reserve: 1024 * 1024,
framework_reserve: 1024 * 1024,
task_reserve: 1024 * 1024,
process_state_reserve: 1024 * 1024,
system_estimate: 1024 * 1024,
safety_margin: 1024 * 1024,
process_limit: 8 * 1024 * 1024,
max_active_requests: 1,
})
.unwrap();
let future = ledger
.try_envelope(65536, 65536, |_| async {
assert!(matches!(
success_stage.finish_nonfailure_bounded(
None,
&DiagnosticOutcomeAxes {
operation: OperationOutcome::Succeeded,
..Default::default()
}
),
Ok(DiagnosticSubmission::OutputUnavailable)
));
observer.observe_database_disposition(crate::DatabaseDisposition::Returned);
let cause = BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.type").unwrap(),
)
.with_column(Some(0), Some(3))
.with_types(
Some(InlineDiagnosticText::metadata("Option<i64>")),
Some(InlineDiagnosticText::metadata("BIGINT")),
None,
);
let primary = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
cause,
);
let axes = DiagnosticOutcomeAxes {
operation: OperationOutcome::Failed,
physical: PhysicalDispositionFact::Returned,
..Default::default()
};
assert_eq!(
handle.submit_bounded(Some(&primary), &axes, None),
DiagnosticSubmission::Enqueued
);
assert_eq!(
handle.submit_bounded(Some(&primary), &axes, None),
DiagnosticSubmission::Full
);
assert_eq!(
stage.finish_bounded(&handle, &axes, Some(&primary)),
DiagnosticSubmission::Full
);
let controlled = crate::RequestDiagnosticScope::established(
&handle,
&diagnostic_call,
&diagnostic_event,
)
.with_zone(crate::DiagnosticZone::from_validated_ingress("zone-a").unwrap())
.with_db_operation(
crate::DiagnosticDbOperation::from_registered("AgentByCode").unwrap(),
);
let failure = controlled.fail(17u32, DiagnosticCategory::UnexpectedError, cause);
assert_eq!(failure.submission(), DiagnosticSubmission::Full);
let (original, delivery) = failure.finish_boundary(&handle, &axes);
assert_eq!(original, 17);
assert_eq!(delivery.boundary_submission(), DiagnosticSubmission::Full);
let reference = controlled.capture_required(primary).into_reference();
let delivery =
receipt_stage.finish_request_failure(Some(&handle), &axes, &reference);
assert_eq!(delivery.source_submission(), DiagnosticSubmission::Full);
assert_eq!(delivery.boundary_submission(), DiagnosticSubmission::Full);
let delivery = unavailable_stage.finish_request_failure(None, &axes, &reference);
assert_eq!(
delivery.boundary_submission(),
DiagnosticSubmission::OutputUnavailable
);
handle.status.closed.store(true, Ordering::Release);
assert_eq!(
handle.submit_bounded(Some(reference.source_diagnostic()), &axes, None),
DiagnosticSubmission::Closed
);
reference.source_diagnostic().id()
})
.unwrap();
let mut future = std::pin::pin!(future);
let result = future
.as_mut()
.poll(&mut Context::from_waker(Waker::noop()));
assert!(
matches!(result, Poll::Ready(Ok(_))),
"admitted capture/submit must not escape: {result:?}"
);
let packet = receiver.try_recv().unwrap();
let row: serde_json::Value = serde_json::from_slice(packet.bytes()).unwrap();
assert_eq!(row["axes"]["operation"], "failed");
assert_eq!(row["axes"]["physical"], "returned");
assert_eq!(handle.snapshot().dropped, 7);
assert_eq!(
observer
.metrics_snapshot()
.database(crate::DatabaseDisposition::Returned),
1
);
assert_eq!(
observer
.metrics_snapshot()
.database(crate::DatabaseDisposition::Discarded),
0
);
assert_eq!(
observer
.metrics_snapshot()
.database(crate::DatabaseDisposition::NotUsed),
0
);
assert_eq!(
observer
.metrics_snapshot()
.stage_latency(crate::Stage::Database)
.iter()
.sum::<u64>(),
4
);
call.succeed();
let _ = observer.shutdown();
}
#[test]
fn bounded_disconnected_and_oversized_context_are_explicit() {
let (sender, receiver) = mpsc::sync_channel(1);
drop(receiver);
let handle = EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
};
let axes = saddle_core::DiagnosticOutcomeAxes::default();
assert_eq!(
handle.submit_bounded(None, &axes, None),
DiagnosticSubmission::Closed
);
let observer =
crate::Observer::with_writer(crate::ObserverConfig::default(), io::sink()).unwrap();
let (call, _) = observer
.start_external_call_checked("app", "module", "service", "x".repeat(257), None)
.unwrap();
assert_eq!(
handle.submit_bounded_parts(
None,
&axes,
Some((call.context(), "request", "route")),
None
),
DiagnosticSubmission::EncodingFailed
);
assert_eq!(handle.snapshot().written, 0);
assert_eq!(handle.snapshot().dropped, 2);
call.succeed();
let _ = observer.shutdown();
}
#[test]
fn diagnostic_protocol_rpc_is_not_internal_span() {
let observer =
crate::Observer::with_writer(crate::ObserverConfig::default(), io::sink()).unwrap();
let (call, _) = observer
.start_external_call_with_rpc(
"test",
"test",
"test",
"query",
Some("opaque-trace"),
saddle_core::RpcCorrelationId::new("0.7").unwrap(),
)
.unwrap();
let event = crate::EventContext::new(
crate::RequestIdentity::new("request-7").unwrap(),
crate::RouteIdentity::new("bill.query").unwrap(),
1,
)
.unwrap();
let diagnostic = sample();
let row = serde_json::to_value(Envelope::new(&diagnostic, Some((call.context(), &event))))
.unwrap();
assert_eq!(row["rpc_id"], "0.7");
assert_eq!(row["span_id"], call.context().span_id().to_string());
assert_ne!(row["rpc_id"], row["span_id"]);
assert_eq!(row["trace_id"], "opaque-trace");
assert_eq!(row["operation"], "query");
assert_eq!(row["request"], "request-7");
assert!(!row.to_string().contains("DO_NOT_LOG_SECRET"));
let missing = serde_json::to_value(Envelope::new(&diagnostic, None)).unwrap();
for key in [
"trace_id",
"rpc_id",
"span_id",
"operation",
"request",
"route",
] {
assert!(missing[key].is_null());
}
call.succeed();
let _ = observer.shutdown();
}
fn sample() -> Diagnostic {
Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.connection_failed").unwrap(),
)
.with_io(&io::Error::other("DO_NOT_LOG_SECRET"))
.with_driver_details(saddle_core::DiagnosticDriverDetails::new(
Some(0),
Some(4),
saddle_core::DiagnosticTypeName::from_metadata("core::option::Option<i64>"),
saddle_core::DiagnosticTypeName::from_metadata("BIGINT"),
))
.with_input_location(
saddle_core::DiagnosticInputLocation::new(Some(17), Some(3))
.with_file(
saddle_core::DiagnosticLocator::from_projection(
"mapping.json",
false,
false,
)
.unwrap(),
)
.with_column(
saddle_core::DiagnosticLocator::from_projection(
"DO_NOT_LOG_SECRET",
true,
true,
)
.unwrap(),
),
),
)
}
struct Failed;
impl Write for Failed {
fn write(&mut self, _: &[u8]) -> io::Result<usize> {
Err(io::Error::from_raw_os_error(28))
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
#[test]
fn diagnostic_independent_file_after_main_failure() {
let _serial = FILE_TEST.lock().unwrap();
let directory = std::env::temp_dir().join(format!(
"saddle-emergency-{}-{}",
std::process::id(),
sample().id()
));
let config = FileLoggingConfig::new(&directory, crate::Rotation::Daily);
let mut emergency = EmergencyDiagnostics::start(&config).unwrap();
struct MarkFailure(Arc<AtomicU64>);
impl Write for MarkFailure {
fn write(&mut self, _: &[u8]) -> io::Result<usize> {
self.0.fetch_add(1, Ordering::Release);
Err(io::Error::from_raw_os_error(28))
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
let main_failed = Arc::new(AtomicU64::new(0));
let observer = crate::Observer::with_writer(
crate::ObserverConfig::default(),
MarkFailure(main_failed.clone()),
)
.unwrap();
let diagnostic = sample();
let (call, _) = observer
.start_external_call_checked("test", "test", "test", "operation", Some("safe-trace"))
.unwrap();
let event = crate::EventContext::new(
crate::RequestIdentity::new("request-1").unwrap(),
crate::RouteIdentity::new("route.test").unwrap(),
1,
)
.unwrap();
let end = Instant::now() + Duration::from_secs(3);
while main_failed.load(Ordering::Acquire) == 0 && Instant::now() < end {
thread::sleep(Duration::from_millis(1));
}
assert!(
main_failed.load(Ordering::Acquire) > 0,
"primary writer failed before diagnostic submission"
);
assert_eq!(
observer.record_diagnostic(
&diagnostic,
&emergency.handle(),
Some((call.context(), &event))
),
DiagnosticSubmission::Enqueued
);
let end = Instant::now() + Duration::from_secs(3);
while emergency.snapshot().written != 1 && Instant::now() < end {
thread::sleep(Duration::from_millis(1));
}
assert_eq!(emergency.snapshot().written, 1);
let text = fs::read_to_string(emergency.target()).unwrap();
let row: serde_json::Value = serde_json::from_str(text.trim()).unwrap();
assert_eq!(row["trace_id"], "safe-trace");
assert_eq!(row["request"], "request-1");
assert!(row["rpc_id"].is_null());
assert_eq!(row["span_id"], call.context().span_id().to_string());
assert_eq!(row["operation"], "operation");
assert_eq!(row["diagnostic"]["diagnostic_id"], diagnostic.id());
let driver = &row["diagnostic"]["causes"][0]["driver_details"];
assert_eq!(driver["column_index"], 0);
assert_eq!(driver["column_count"], 4);
assert_eq!(
driver["target_rust_type"]["value"],
"core::option::Option<i64>"
);
assert_eq!(driver["actual_db_type"]["value"], "BIGINT");
let input = &row["diagnostic"]["causes"][0]["input_location"];
assert_eq!(input["json_line"], 17);
assert_eq!(input["json_column"], 3);
assert_eq!(input["file"]["value"], "mapping.json");
assert_eq!(input["column"]["value"], "[redacted]");
assert_eq!(input["locator_truncated"], true);
assert_eq!(input["locator_redacted"], true);
assert!(!text.contains("DO_NOT_LOG_SECRET"));
let bounded = saddle_core::BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
DiagnosticStage::RequestDb,
DiagnosticCode::new("db.decode.failed").unwrap(),
),
);
let axes = saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Failed,
physical: saddle_core::PhysicalDispositionFact::Returned,
business: saddle_core::BusinessOutcome::Failure,
business_code: saddle_core::RegisteredDiagnosticCode::from_registered(
"INVALID_ROW",
&["INVALID_ROW"],
),
delivery: saddle_core::ResponseDelivery::LocalWriteComplete,
cleanup: saddle_core::CleanupOutcome::Succeeded,
bytes_written: Some(100),
};
assert_eq!(
emergency.handle().submit_bounded(
Some(&bounded),
&axes,
Some((call.context(), &event))
),
DiagnosticSubmission::Enqueued
);
let until = Instant::now() + Duration::from_secs(3);
while emergency.snapshot().written != 2 && Instant::now() < until {
thread::sleep(Duration::from_millis(1));
}
assert_eq!(emergency.snapshot().written, 2);
let content = fs::read_to_string(emergency.target()).unwrap();
let second: serde_json::Value =
serde_json::from_str(content.lines().nth(1).unwrap()).unwrap();
assert_eq!(second["axes"]["operation"], "failed");
assert_eq!(second["axes"]["physical"], "returned");
assert_eq!(second["axes"]["business_code"], "INVALID_ROW");
assert_eq!(second["axes"]["delivery"], "local_write_complete");
assert_eq!(second["axes"]["cleanup"], "succeeded");
assert_eq!(second["trace_id"], "safe-trace");
let end = Instant::now() + Duration::from_secs(3);
while emergency.shutdown() == DiagnosticShutdown::Pending && Instant::now() < end {
thread::sleep(Duration::from_millis(1));
}
assert_eq!(emergency.shutdown(), DiagnosticShutdown::Finished);
assert_eq!(
emergency.handle().submit(&diagnostic),
DiagnosticSubmission::Closed
);
call.succeed();
let _shutdown = observer.shutdown();
fs::remove_file(emergency.target()).unwrap();
fs::remove_dir(directory).unwrap();
}
#[test]
fn diagnostic_directory_failure_has_real_os_cause() {
let _serial = FILE_TEST.lock().unwrap();
let path = std::env::temp_dir().join(format!(
"saddle-emergency-file-{}-{}",
std::process::id(),
sample().id()
));
fs::write(&path, b"not a directory").unwrap();
let mut owner =
EmergencyDiagnostics::start(&FileLoggingConfig::new(&path, crate::Rotation::Daily))
.unwrap();
let end = Instant::now() + Duration::from_secs(3);
while !owner.snapshot().initialized && Instant::now() < end {
thread::sleep(Duration::from_millis(1));
}
let failure = owner.snapshot().first_failure.unwrap();
assert_eq!(failure.stage, "initialize");
assert!(failure.os_code.is_some());
assert_eq!(owner.target(), path.join(EMERGENCY_FILE_NAME));
while owner.shutdown() == DiagnosticShutdown::Pending && Instant::now() < end {
thread::sleep(Duration::from_millis(1));
}
assert_eq!(owner.shutdown(), DiagnosticShutdown::Finished);
fs::remove_file(path).unwrap();
}
#[test]
fn diagnostic_stuck_worker_shutdown_never_joins_early() {
let (sender, _receiver) = mpsc::sync_channel(1);
let (release, wait) = mpsc::channel();
let (started, reached) = mpsc::channel();
let worker = thread::spawn(move || {
started.send(()).unwrap();
wait.recv().unwrap();
});
reached.recv_timeout(Duration::from_secs(2)).unwrap();
let mut owner = EmergencyDiagnostics {
handle: EmergencyDiagnosticHandle {
sender,
status: Arc::new(Status::default()),
},
worker: Some(worker),
target: PathBuf::from("unused"),
};
assert_eq!(owner.shutdown(), DiagnosticShutdown::Pending);
release.send(()).unwrap();
let end = Instant::now() + Duration::from_secs(3);
while owner.shutdown() == DiagnosticShutdown::Pending && Instant::now() < end {
thread::yield_now();
}
assert_eq!(owner.shutdown(), DiagnosticShutdown::Finished);
}
#[test]
fn diagnostic_double_failure_is_counted_without_recursion() {
let status = Status::default();
let mut output: io::Result<Failed> = Err(io::Error::from_raw_os_error(13));
let mut stderr = Vec::new();
let mut attempted = false;
deliver(
&mut output,
&mut stderr,
std::path::Path::new("/safe/logs/saddle.emergency.log"),
b"{\"original\":true}\n",
&status,
&mut attempted,
);
let text = String::from_utf8(stderr).unwrap();
assert!(text.contains("\"os_code\":13") && text.contains("\"original\":true"));
attempted = false;
deliver(
&mut output,
&mut Failed,
std::path::Path::new("./logs"),
b"{}\n",
&status,
&mut attempted,
);
for _ in 0..10 {
deliver(
&mut output,
&mut Failed,
std::path::Path::new("./logs"),
b"{}\n",
&status,
&mut attempted,
);
}
assert_eq!(status.snapshot().stderr_failed, 1);
assert_eq!(status.snapshot().stderr_suppressed, 10);
assert_eq!(status.snapshot().output_failed, 12);
}
#[test]
fn diagnostic_full_queue_submission_does_not_wait() {
let (sender, _receiver) = mpsc::sync_channel(1);
let status = Arc::new(Status::default());
let handle = EmergencyDiagnosticHandle { sender, status };
let sample = sample();
assert_eq!(handle.submit(&sample), DiagnosticSubmission::Enqueued);
assert_eq!(handle.submit(&sample), DiagnosticSubmission::Full);
assert_eq!(handle.snapshot().dropped, 1);
assert_eq!(handle.snapshot().written, 0);
}
}