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)]
pub enum DiagnosticSubmission {
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<Vec<u8>>,
status: Arc<Status>,
}
#[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>,
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.map(|(call, _)| call.span_id().to_string()),
request: context.map(|(_, event)| event.diagnostic_request()),
route: context.map(|(_, event)| event.diagnostic_route()),
}
}
}
impl EmergencyDiagnosticHandle {
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 {
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, &Envelope::new(diagnostic, context)).is_err() {
bytes.0.clear();
let mut reduced = match serde_json::to_value(Envelope::new(diagnostic, context)) {
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(bytes.0) {
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);
if let Ok(serde_json::Value::Object(fields)) =
serde_json::to_value(Envelope::new(diagnostic, context))
{
let level = if matches!(
diagnostic.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);
}
result
}
}
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::<Vec<u8>>(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(bytes) => deliver(
&mut output,
&mut io::stderr(),
&path,
&bytes,
&worker_status,
&mut stderr_attempted,
),
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);
}
}
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,
) {
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()
|| 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::*;
use saddle_core::{
CaptureSite, DiagnosticCategory, DiagnosticCause, DiagnosticCode, DiagnosticStage,
};
use std::time::Instant;
static FILE_TEST: std::sync::Mutex<()> = std::sync::Mutex::new(());
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_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 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);
}
}