saddle-observability 0.3.29

Saddle structured logging and trace correlation
Documentation
//! Request-scoped original capture through the existing synchronous file writer.
use super::{Projection, timestamp};
use crate::{DiagnosticSubmission, EmergencyDiagnosticHandle};
use saddle_core::BoundedDiagnostic;
use serde::Serialize;
use std::{error::Error, fmt, io};

const PAYLOAD_BYTES: usize = 6144;

#[derive(Serialize)]
struct Segment<'a> {
    schema_version: u8,
    event: &'static str,
    occurrence: saddle_core::DiagnosticOccurrence,
    sequence: u64,
    channel: &'static str,
    cause_depth: usize,
    payload: &'a str,
    state: &'static str,
}

#[derive(Serialize)]
struct IoFacts {
    kind: IoKind,
    os_code: Option<i32>,
}
struct IoKind(io::ErrorKind);
impl Serialize for IoKind {
    fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
        serializer.collect_str(&format_args!("{:?}", self.0))
    }
}

struct Stream<'a> {
    output: &'a EmergencyDiagnosticHandle,
    occurrence: saddle_core::DiagnosticOccurrence,
    sequence: u64,
    channel: &'static str,
    depth: usize,
    bytes: [u8; PAYLOAD_BYTES],
    len: usize,
    escaped: usize,
}

struct CycleDetector<'a> {
    anchor: &'a (dyn Error + 'static),
    power: usize,
    distance: usize,
}
impl<'a> CycleDetector<'a> {
    fn new(anchor: &'a (dyn Error + 'static)) -> Self {
        Self {
            anchor,
            power: 1,
            distance: 0,
        }
    }
    fn repeats(&mut self, next: &'a (dyn Error + 'static)) -> bool {
        if std::ptr::eq(self.anchor, next) {
            return true;
        }
        self.distance += 1;
        if self.distance == self.power {
            self.anchor = next;
            self.power = self.power.saturating_mul(2);
            self.distance = 0;
        }
        false
    }
}
impl Stream<'_> {
    fn emit(&mut self, state: &'static str) -> fmt::Result {
        let payload = std::str::from_utf8(&self.bytes[..self.len]).map_err(|_| fmt::Error)?;
        self.output
            .write_original_record(&Segment {
                schema_version: 1,
                event: "request_error_original",
                occurrence: self.occurrence,
                sequence: self.sequence,
                channel: self.channel,
                cause_depth: self.depth,
                payload,
                state,
            })
            .map_err(|_| fmt::Error)?;
        self.sequence += 1;
        self.len = 0;
        self.escaped = 0;
        Ok(())
    }
    fn push(&mut self, text: &str) -> fmt::Result {
        for c in text.chars() {
            let cost = match c {
                '"' | '\\' | '\n' | '\r' | '\t' | '\u{8}' | '\u{c}' => 2,
                '\u{0}'..='\u{1f}' => 6,
                _ => c.len_utf8(),
            };
            if self.escaped + cost > PAYLOAD_BYTES || self.len + c.len_utf8() > PAYLOAD_BYTES {
                self.emit("continuation")?;
            }
            let size = c.len_utf8();
            c.encode_utf8(&mut self.bytes[self.len..self.len + size]);
            self.len += size;
            self.escaped += cost;
        }
        Ok(())
    }
    fn field(
        &mut self,
        channel: &'static str,
        depth: usize,
        value: fmt::Arguments<'_>,
    ) -> fmt::Result {
        self.channel = channel;
        self.depth = depth;
        fmt::write(self, value)?;
        self.emit("field_end")
    }
}
impl fmt::Write for Stream<'_> {
    fn write_str(&mut self, text: &str) -> fmt::Result {
        self.push(text)
    }
}
impl io::Write for Stream<'_> {
    fn write(&mut self, bytes: &[u8]) -> io::Result<usize> {
        let text = std::str::from_utf8(bytes).map_err(|_| io::ErrorKind::InvalidData)?;
        self.push(text).map_err(|_| io::ErrorKind::Other)?;
        Ok(bytes.len())
    }
    fn flush(&mut self) -> io::Result<()> {
        Ok(())
    }
}

/// No payload is retained by this scope after return. Partial files have no
/// successful terminal segment, and therefore cannot issue a Written receipt.
pub(super) fn capture<T: Serialize>(
    output: Option<&EmergencyDiagnosticHandle>,
    context: &Projection,
    diagnostic: &BoundedDiagnostic,
    error: &(dyn Error + 'static),
    target: &T,
) -> DiagnosticSubmission {
    capture_inner(output, context, diagnostic, None, Some(error), target)
}

pub(super) fn capture_borrowed<T: Serialize, E: Error>(
    output: Option<&EmergencyDiagnosticHandle>,
    context: &Projection,
    diagnostic: &BoundedDiagnostic,
    error: &E,
    target: &T,
) -> DiagnosticSubmission {
    capture_inner(output, context, diagnostic, Some(error), None, target)
}

fn capture_inner<T: Serialize>(
    output: Option<&EmergencyDiagnosticHandle>,
    context: &Projection,
    diagnostic: &BoundedDiagnostic,
    borrowed_head: Option<&dyn Error>,
    static_head: Option<&(dyn Error + 'static)>,
    target: &T,
) -> DiagnosticSubmission {
    let Some(output) = output else {
        return DiagnosticSubmission::OutputUnavailable;
    };
    let mut stream = Stream {
        output,
        occurrence: diagnostic.occurrence(),
        sequence: 0,
        channel: "context",
        depth: 0,
        bytes: [0; PAYLOAD_BYTES],
        len: 0,
        escaped: 0,
    };
    #[derive(Serialize)]
    struct Header<'a, T> {
        timestamp_unix_ms: u128,
        context: &'a Projection,
        diagnostic: &'a BoundedDiagnostic,
        target: &'a T,
        original_contract: &'static str,
    }
    let header = Header {
        timestamp_unix_ms: timestamp(),
        context,
        diagnostic,
        target,
        original_contract: "display-debug-source/v1",
    };
    let result = (|| -> fmt::Result {
        serde_json::to_writer(&mut stream, &header).map_err(|_| fmt::Error)?;
        stream.emit("field_end")?;
        if let Some(error) = borrowed_head {
            stream.field("description", 0, format_args!("{error}"))?;
            stream.field("debug", 0, format_args!("{error:?}"))?;
        }
        let mut current = static_head.or_else(|| borrowed_head.and_then(Error::source));
        let mut cycle = current.map(CycleDetector::new);
        let mut depth = usize::from(borrowed_head.is_some());
        while let Some(value) = current {
            stream.field("description", depth, format_args!("{value}"))?;
            stream.channel = "debug";
            stream.depth = depth;
            if let Some(json) = value.downcast_ref::<serde_json::Error>() {
                crate::root_diagnostic::write_json_error_debug(&mut stream, json)?;
            } else if let Some(saddle_admission::InputDecodeFailure::Syntax(json)) =
                value.downcast_ref::<saddle_admission::InputDecodeFailure>() {
                stream.push("Syntax(")?;
                crate::root_diagnostic::write_json_error_debug(&mut stream, json)?;
                stream.push(")")?;
            } else {
                fmt::write(&mut stream, format_args!("{value:?}"))?;
            }
            stream.emit("field_end")?;
            if let Some(io) = value.downcast_ref::<io::Error>() {
                stream.channel = "io_facts";
                stream.depth = depth;
                serde_json::to_writer(
                    &mut stream,
                    &IoFacts {
                        kind: IoKind(io.kind()),
                        os_code: io.raw_os_error(),
                    },
                )
                .map_err(|_| fmt::Error)?;
                stream.emit("field_end")?;
            }
            current = value.source();
            depth += 1;
            if current.is_some_and(|next| cycle.as_mut().is_some_and(|cycle| cycle.repeats(next))) {
                stream.channel = "terminal";
                stream.depth = depth;
                stream.emit("cause_cycle")?;
                return Err(fmt::Error);
            }
        }
        stream.channel = "terminal";
        stream.depth = depth;
        stream.emit("exposed_chain_complete")
    })();
    if result.is_ok() {
        DiagnosticSubmission::Written
    } else {
        // Keep any successfully formatted prefix while making incompleteness
        // machine visible. A failed file write may also reject this marker.
        let _ = stream.emit("capture_failed");
        DiagnosticSubmission::EncodingFailed
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use crate::{EmergencyDiagnostics, FileLoggingConfig, Rotation};
    use saddle_core::{
        BoundedDiagnosticCause, DiagnosticCategory, DiagnosticCode, DiagnosticStage,
    };

    #[test]
    fn borrowed_first_error_keeps_its_source_and_io_facts() {
        #[derive(Debug)]
        struct Borrowed<'a> { input: &'a str, cause: io::Error }
        impl fmt::Display for Borrowed<'_> {
            fn fmt(&self, out: &mut fmt::Formatter<'_>) -> fmt::Result {
                write!(out, "raw input {:?}: {}", self.input, self.cause)
            }
        }
        impl Error for Borrowed<'_> {
            fn source(&self) -> Option<&(dyn Error + 'static)> { Some(&self.cause) }
        }
        let raw_input = String::from("borrowed-json-value");
        let original = Borrowed {
            input: &raw_input,
            cause: io::Error::from_raw_os_error(5),
        };
        let directory = std::env::temp_dir().join(format!(
            "request-borrowed-original-{}-{}", std::process::id(),
            std::time::SystemTime::now().duration_since(std::time::UNIX_EPOCH)
                .unwrap().as_nanos(),
        ));
        std::fs::create_dir_all(&directory).unwrap();
        let mut writer = EmergencyDiagnostics::start(
            &FileLoggingConfig::new(&directory, Rotation::Daily),
        ).unwrap();
        let handle = writer.handle();
        let receipt = super::super::RequestDiagnosticScope::unassociated_transport(&handle)
            .source_borrowed_error_with_target(
                &original, &"row 0", DiagnosticCategory::UnexpectedError,
                BoundedDiagnosticCause::new(
                    DiagnosticStage::RequestDb,
                    DiagnosticCode::new("db.decode.invalid_json").unwrap(),
                ),
            );
        assert_eq!(receipt.submission(), DiagnosticSubmission::Written);
        while matches!(writer.shutdown(), crate::DiagnosticShutdown::Pending) {
            std::thread::yield_now();
        }
        let lines: Vec<serde_json::Value> = std::fs::read_to_string(writer.target())
            .unwrap().lines().map(|line| serde_json::from_str(line).unwrap()).collect();
        assert!(lines.iter().all(|line| line["occurrence"] == lines[0]["occurrence"]));
        assert!(lines.iter().any(|line| line["payload"].as_str().unwrap().contains(&raw_input)));
        assert!(lines.iter().any(|line| line["channel"] == "io_facts"));
        assert_eq!(lines.last().unwrap()["state"], "exposed_chain_complete");
        std::fs::remove_dir_all(directory).unwrap();
    }

    #[test]
    fn failed_original_format_cannot_confirm_original() {
        struct FailingError;
        impl std::fmt::Display for FailingError {
            fn fmt(&self, out: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
                out.write_str("preserved-prefix")?;
                Err(std::fmt::Error)
            }
        }
        impl std::fmt::Debug for FailingError {
            fn fmt(&self, out: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
                out.write_str("FailingError")
            }
        }
        impl std::error::Error for FailingError {}
        let directory = std::env::temp_dir().join(format!(
            "request-original-closed-{}-{}",
            std::process::id(),
            std::time::SystemTime::now()
                .duration_since(std::time::UNIX_EPOCH)
                .unwrap()
                .as_nanos(),
        ));
        std::fs::create_dir_all(&directory).unwrap();
        let mut writer =
            EmergencyDiagnostics::start(&FileLoggingConfig::new(&directory, Rotation::Daily))
                .unwrap();
        let handle = writer.handle();
        let scope = super::super::RequestDiagnosticScope::unassociated_transport(&handle);
        let receipt = scope.source_error_with_target(
            &FailingError,
            &"validated target",
            DiagnosticCategory::UnexpectedError,
            BoundedDiagnosticCause::new(
                DiagnosticStage::RequestOutbound,
                DiagnosticCode::new("transport.connect_failed").unwrap(),
            ),
        );
        assert_ne!(receipt.submission(), DiagnosticSubmission::Written);
        while matches!(writer.shutdown(), crate::DiagnosticShutdown::Pending) {
            std::thread::yield_now();
        }
        let raw = std::fs::read_to_string(writer.target()).unwrap();
        assert!(raw.contains("preserved-prefix"));
        assert!(raw.contains("capture_failed"));
        assert!(!raw.contains("exposed_chain_complete"));
        std::fs::remove_dir_all(directory).unwrap();
    }
}