use super::*;
use std::{error::Error, fmt, io};
const PAYLOAD_BYTES: usize = 6144;
#[derive(Clone, Copy, Debug, Eq, PartialEq, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum OriginalCaptureState {
CompleteEnqueued,
OutputUnavailable,
OutputRejected,
FormattingFailed,
CauseCycle,
LegacyProjectionOnly,
}
#[derive(Serialize)]
struct Segment<'a> {
schema_version: u8,
event: &'static str,
occurrence: DiagnosticOccurrence,
sequence: u64,
channel: &'static str,
cause_depth: usize,
payload: &'a str,
state: &'static str,
}
trait OriginalOutput {
fn submit(&self, segment: &Segment<'_>) -> DiagnosticSubmission;
}
impl OriginalOutput for EmergencyDiagnosticHandle {
fn submit(&self, segment: &Segment<'_>) -> DiagnosticSubmission {
self.submit_original_record(segment)
}
}
struct Stream<'a, O: OriginalOutput> {
output: &'a O,
occurrence: DiagnosticOccurrence,
sequence: u64,
depth: usize,
channel: &'static str,
bytes: [u8; PAYLOAD_BYTES],
len: usize,
escaped: usize,
failure: Option<OriginalCaptureState>,
submission: DiagnosticSubmission,
}
impl<O: OriginalOutput> Stream<'_, O> {
fn emit(&mut self, state: &'static str) -> fmt::Result {
let payload = std::str::from_utf8(&self.bytes[..self.len]).map_err(|_| fmt::Error)?;
let submission = self.output.submit(&Segment {
schema_version: 1,
event: "request_error_original",
occurrence: self.occurrence,
sequence: self.sequence,
channel: self.channel,
cause_depth: self.depth,
payload,
state,
});
if submission != DiagnosticSubmission::Enqueued {
self.submission = submission;
self.failure = Some(OriginalCaptureState::OutputRejected);
return Err(fmt::Error);
}
self.sequence += 1;
self.len = 0;
self.escaped = 0;
Ok(())
}
fn push(&mut self, text: &str) -> fmt::Result {
if self.failure.is_some() {
return Err(fmt::Error);
}
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 json<T: Serialize>(&mut self, value: &T) -> fmt::Result {
serde_json::to_writer(&mut *self, value).map_err(|_| fmt::Error)?;
self.emit("field_end")
}
fn report_failure(&mut self, state: OriginalCaptureState) {
if self.failure == Some(OriginalCaptureState::OutputRejected) {
return;
}
self.failure = Some(state);
let marker = match state {
OriginalCaptureState::FormattingFailed => "formatting_failed",
OriginalCaptureState::CauseCycle => "cause_cycle",
_ => "capture_failed",
};
let _ = self.emit(marker);
}
}
impl<O: OriginalOutput> fmt::Write for Stream<'_, O> {
fn write_str(&mut self, text: &str) -> fmt::Result {
self.push(text)
}
}
impl<O: OriginalOutput> io::Write for Stream<'_, O> {
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(())
}
}
#[derive(Serialize)]
struct Header<'a> {
timestamp_unix_ms: u128,
context: crate::AdmissionEventContext<'a>,
stage: RootRequestEvent,
facts: SourceDetail<'a>,
outcome: RootOutcomeFacts,
original_contract: &'static str,
source_interface: &'static str,
}
enum RawSource<'a> {
Error(&'a (dyn Error + 'static)),
Description(&'a dyn fmt::Display, &'a dyn fmt::Debug),
}
#[derive(Serialize)]
struct IoFacts {
kind: Kind,
os_code: Option<i32>,
}
struct Kind(io::ErrorKind);
impl Serialize for Kind {
fn serialize<S: serde::Serializer>(&self, serializer: S) -> Result<S::Ok, S::Error> {
serializer.collect_str(&format_args!("{:?}", self.0))
}
}
impl RootDiagnosticScope<'_> {
#[track_caller]
pub fn source_error(
&self,
error: &(dyn Error + 'static),
source_stage: saddle_core::DiagnosticStage,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
let code = DiagnosticCode::new("framework.error").expect("static code");
let facts = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(source_stage, code),
);
self.source_error_with_facts(error, facts, code, stage, axes)
}
pub fn source_error_with_facts(
&self,
error: &(dyn Error + 'static),
facts: BoundedDiagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.original_record(
RawSource::Error(error),
SourceDetail::Bounded(&facts),
classification,
stage,
axes,
)
}
pub fn source_description<D: fmt::Display + fmt::Debug>(
&self,
description: &D,
facts: BoundedDiagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.original_record(
RawSource::Description(description, description),
SourceDetail::Bounded(&facts),
classification,
stage,
axes,
)
}
pub fn source_existing_description<D: fmt::Display + fmt::Debug>(
&self,
description: &D,
diagnostic: Diagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.original_record(
RawSource::Description(description, description),
SourceDetail::Existing(&diagnostic),
classification,
stage,
axes,
)
}
pub fn source_existing_error(
&self,
error: &(dyn Error + 'static),
diagnostic: Diagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.original_record(
RawSource::Error(error),
SourceDetail::Existing(&diagnostic),
classification,
stage,
axes,
)
}
fn original_record(
&self,
raw: RawSource<'_>,
facts: SourceDetail<'_>,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
let (occurrence, category) = match &facts {
SourceDetail::Bounded(d) => (d.occurrence(), d.category()),
SourceDetail::Existing(d) => (d.occurrence(), d.category()),
};
let source_interface = match raw {
RawSource::Error(_) => "error_source",
RawSource::Description(_, _) => "unavailable_description_debug_only",
};
let (submission, original_capture) = match self.output {
None => (
DiagnosticSubmission::OutputUnavailable,
OriginalCaptureState::OutputUnavailable,
),
Some(output) => capture(
output,
occurrence,
raw,
&Header {
timestamp_unix_ms: timestamp(),
context: crate::AdmissionEventContext::Rooted(self.view),
stage,
facts,
outcome: axes,
original_contract: "display-debug-source/v1",
source_interface,
},
),
};
RootRequestFailure {
source: self.view.clone(),
occurrence,
classification,
source_submission: submission,
category,
outcome: axes,
original_capture,
}
}
}
#[must_use = "capture facts describe submission, not durable output or resource completion"]
pub struct UnrootedCaptureFacts {
occurrence: DiagnosticOccurrence,
submission: DiagnosticSubmission,
original_capture: OriginalCaptureState,
}
impl UnrootedCaptureFacts {
pub fn occurrence(&self) -> DiagnosticOccurrence {
self.occurrence
}
pub fn submission(&self) -> DiagnosticSubmission {
self.submission
}
pub fn original_capture(&self) -> OriginalCaptureState {
self.original_capture
}
}
pub struct UnrootedDiagnosticScope<'a> {
application: &'a saddle_core::ContextLabel,
lifecycle: saddle_core::RequestViewPhase,
output: Option<&'a EmergencyDiagnosticHandle>,
}
impl<'a> UnrootedDiagnosticScope<'a> {
pub fn new(
application: &'a saddle_core::ContextLabel,
lifecycle: saddle_core::RequestViewPhase,
output: Option<&'a EmergencyDiagnosticHandle>,
) -> Self {
Self {
application,
lifecycle,
output,
}
}
#[track_caller]
pub fn source_error(
&self,
error: &(dyn Error + 'static),
source_stage: saddle_core::DiagnosticStage,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> UnrootedCaptureFacts {
let facts = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
source_stage,
DiagnosticCode::new("framework.error").unwrap(),
),
);
self.original_record(RawSource::Error(error), &facts, stage, axes)
}
#[track_caller]
pub fn source_description<D: fmt::Display + fmt::Debug>(
&self,
description: &D,
source_stage: saddle_core::DiagnosticStage,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> UnrootedCaptureFacts {
let facts = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
source_stage,
DiagnosticCode::new("framework.error").unwrap(),
),
);
self.original_record(
RawSource::Description(description, description),
&facts,
stage,
axes,
)
}
fn original_record(
&self,
raw: RawSource<'_>,
facts: &BoundedDiagnostic,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> UnrootedCaptureFacts {
let occurrence = facts.occurrence();
let source_interface = match raw {
RawSource::Error(_) => "error_source",
RawSource::Description(_, _) => "unavailable_description_debug_only",
};
let (submission, original_capture) = match self.output {
None => (
DiagnosticSubmission::OutputUnavailable,
OriginalCaptureState::OutputUnavailable,
),
Some(output) => capture(
output,
occurrence,
raw,
&Header {
timestamp_unix_ms: timestamp(),
context: crate::AdmissionEventContext::Unrooted {
application: self.application,
lifecycle: self.lifecycle,
},
stage,
facts: SourceDetail::Bounded(facts),
outcome: axes,
original_contract: "display-debug-source/v1",
source_interface,
},
),
};
UnrootedCaptureFacts {
occurrence,
submission,
original_capture,
}
}
}
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
}
}
fn capture<O: OriginalOutput>(
output: &O,
occurrence: DiagnosticOccurrence,
raw: RawSource<'_>,
header: &impl Serialize,
) -> (DiagnosticSubmission, OriginalCaptureState) {
let mut stream = Stream {
output,
occurrence,
sequence: 0,
depth: 0,
channel: "context",
bytes: [0; PAYLOAD_BYTES],
len: 0,
escaped: 0,
failure: None,
submission: DiagnosticSubmission::Enqueued,
};
let result = (|| -> fmt::Result {
stream.json(header)?;
let error = match raw {
RawSource::Error(error) => error,
RawSource::Description(description, debug) => {
stream.channel = "description";
if fmt::write(&mut stream, format_args!("{description}")).is_err() {
let state = stream
.failure
.unwrap_or(OriginalCaptureState::FormattingFailed);
stream.report_failure(state);
return Err(fmt::Error);
}
stream.emit("field_end")?;
stream.channel = "debug";
if fmt::write(&mut stream, format_args!("{debug:?}")).is_err() {
let state = stream
.failure
.unwrap_or(OriginalCaptureState::FormattingFailed);
stream.report_failure(state);
return Err(fmt::Error);
}
stream.emit("field_end")?;
stream.depth = 1;
stream.channel = "terminal";
return stream.emit("description_complete_source_unavailable");
}
};
let mut cycle = CycleDetector::new(error);
let mut next = Some(error);
while let Some(current) = next {
stream.channel = "description";
if fmt::write(&mut stream, format_args!("{current}")).is_err() {
let failure = stream
.failure
.unwrap_or(OriginalCaptureState::FormattingFailed);
stream.report_failure(failure);
return Err(fmt::Error);
}
stream.emit("field_end")?;
stream.channel = "debug";
if fmt::write(&mut stream, format_args!("{current:?}")).is_err() {
let state = stream
.failure
.unwrap_or(OriginalCaptureState::FormattingFailed);
stream.report_failure(state);
return Err(fmt::Error);
}
stream.emit("field_end")?;
if let Some(io) = current.downcast_ref::<io::Error>() {
stream.channel = "io_facts";
stream.json(&IoFacts {
kind: Kind(io.kind()),
os_code: io.raw_os_error(),
})?;
}
next = current.source();
stream.depth += 1;
if next.is_some_and(|next| cycle.repeats(next)) {
stream.channel = "description";
stream.report_failure(OriginalCaptureState::CauseCycle);
return Err(fmt::Error);
}
}
stream.channel = "terminal";
stream.emit("exposed_chain_complete")
})();
if result.is_ok() {
(stream.submission, OriginalCaptureState::CompleteEnqueued)
} else {
let state = stream
.failure
.unwrap_or(OriginalCaptureState::FormattingFailed);
if stream.submission == DiagnosticSubmission::Enqueued {
(DiagnosticSubmission::EncodingFailed, state)
} else {
(stream.submission, state)
}
}
}
pub fn original_capture_layout() -> (std::alloc::Layout, std::alloc::Layout, usize) {
(
std::alloc::Layout::new::<Stream<'static, EmergencyDiagnosticHandle>>(),
std::alloc::Layout::new::<CycleDetector<'static>>(),
PAYLOAD_BYTES,
)
}
#[cfg(test)]
mod content_tests {
use super::*;
use std::cell::RefCell;
struct Node {
text: String,
next: Option<Box<Node>>,
}
impl fmt::Display for Node {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.text)
}
}
impl fmt::Debug for Node {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.write_str(&self.text)
}
}
impl Error for Node {
fn source(&self) -> Option<&(dyn Error + 'static)> {
self.next.as_deref().map(|v| v as &dyn Error)
}
}
struct Readback<'a> {
expected: &'a [String],
state: RefCell<(u64, usize, usize, bool)>,
}
impl OriginalOutput for Readback<'_> {
fn submit(&self, segment: &Segment<'_>) -> DiagnosticSubmission {
let encoded = serde_json::to_vec(segment).unwrap();
assert!(encoded.len() < 8192);
let row: serde_json::Value = serde_json::from_slice(&encoded).unwrap();
let mut s = self.state.borrow_mut();
assert_eq!(row["sequence"].as_u64().unwrap(), s.0);
s.0 += 1;
if matches!(segment.channel, "description" | "debug") {
let expected = &self.expected[segment.cause_depth];
let payload = row["payload"].as_str().unwrap();
assert_eq!(payload, &expected[s.1..s.1 + payload.len()]);
s.1 += payload.len();
if segment.state == "field_end" {
assert_eq!(s.1, expected.len(), "original tail missing");
s.1 = 0;
s.2 += 1;
}
}
if segment.channel == "terminal" {
assert_eq!(segment.state, "exposed_chain_complete");
assert_eq!(segment.cause_depth, self.expected.len());
assert_eq!(s.2, self.expected.len() * 2);
s.3 = true;
}
DiagnosticSubmission::Enqueued
}
}
fn assert_content(expected: Vec<String>) -> u64 {
assert_content_context(expected, false)
}
fn assert_content_context(expected: Vec<String>, unrooted: bool) -> u64 {
let mut head = None;
for text in expected.iter().rev() {
head = Some(Box::new(Node {
text: text.clone(),
next: head,
}));
}
let publisher = saddle_core::RequestRootPublisher::create(
saddle_core::ContextLabel::checked("content-test").unwrap(),
saddle_core::ContextFact::NotEstablished,
)
.unwrap();
let view = publisher
.reference()
.view(saddle_core::RequestLocalFacts::new(
saddle_core::RequestViewPhase::Reading,
));
let diagnostic = BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestDecode,
DiagnosticCode::new("content.original").unwrap(),
),
);
let output = Readback {
expected: &expected,
state: RefCell::new((0, 0, 0, false)),
};
let application = saddle_core::ContextLabel::checked("content-test").unwrap();
let result = capture(
&output,
diagnostic.occurrence(),
RawSource::Error(head.as_deref().unwrap()),
&Header {
timestamp_unix_ms: 0,
context: if unrooted {
crate::AdmissionEventContext::Unrooted {
application: &application,
lifecycle: saddle_core::RequestViewPhase::Reading,
}
} else {
crate::AdmissionEventContext::Rooted(&view)
},
stage: RootRequestEvent::Ingress,
facts: SourceDetail::Bounded(&diagnostic),
outcome: RootOutcomeFacts::default(),
original_contract: "display-debug-source/v1",
source_interface: "error_source",
},
);
assert_eq!(
result,
(
DiagnosticSubmission::Enqueued,
OriginalCaptureState::CompleteEnqueued
)
);
assert!(output.state.borrow().3);
output.state.into_inner().0
}
#[test]
fn rooted_and_unrooted_share_lossless_long_content() {
let expected = vec![
format!("HEAD-中文\n\"{}\"-TAIL", "原文".repeat(2200)),
"nested-cause\n完整".into(),
];
let rooted = assert_content_context(expected.clone(), false);
let unrooted = assert_content_context(expected, true);
assert_eq!(rooted, unrooted);
}
#[test]
fn finite_chain_beyond_old_256_keeps_every_cause() {
let segments = assert_content(
(0..300)
.map(|n| format!("unknown-{n}-原文\n\"tail\""))
.collect(),
);
assert!(segments > 600);
println!("content readback: 300 original causes and Debug fields complete");
}
#[test]
fn cycle_detector_handles_prefix_and_nontrivial_cycle() {
let nodes: Vec<Node> = (0..4)
.map(|n| Node {
text: n.to_string(),
next: None,
})
.collect();
let mut cycle = CycleDetector::new(&nodes[0]);
for node in [&nodes[1], &nodes[2], &nodes[3]] {
assert!(!cycle.repeats(node));
}
let found = (0..32).any(|step| cycle.repeats(&nodes[1 + step % 3]));
assert!(found);
}
#[test]
fn finite_text_beyond_old_4096_segments_keeps_tail() {
let text = format!("{}END-原文", "x".repeat(6144 * 4097));
let segments = assert_content(vec![text]);
assert!(segments > 8194);
println!("content readback: {segments} segments, complete Display/Debug tails");
}
}
#[doc(hidden)]
pub fn database_background_error(
output: Option<&EmergencyDiagnosticHandle>,
datasource: u64,
connection: Option<u64>,
sweep: u64,
phase: &'static str,
maintenance: bool,
error: &(dyn Error + 'static),
facts: BoundedDiagnostic,
) -> UnrootedCaptureFacts {
#[derive(Serialize)]
struct Context {
schema_version: u8,
datasource: u64,
connection: ContextFact<u64>,
maintenance_sweep: u64,
maintenance_phase: &'static str,
request: ContextFact<()>,
trace_id: ContextFact<()>,
application: ContextFact<()>,
}
#[derive(Serialize)]
struct BackgroundHeader<'a> {
timestamp_unix_ms: u128,
context: Context,
stage: &'static str,
facts: &'a BoundedDiagnostic,
original_contract: &'static str,
source_interface: &'static str,
}
let occurrence = facts.occurrence();
let missing = if maintenance { ContextFact::NotApplicable } else { ContextFact::Unavailable };
let header = BackgroundHeader {
timestamp_unix_ms: timestamp(),
context: Context { schema_version: 2, datasource,
connection: connection.map(ContextFact::Present).unwrap_or(ContextFact::NotEstablished),
maintenance_sweep: sweep, maintenance_phase: phase,
request: missing, trace_id: missing, application: ContextFact::Unavailable },
stage: if maintenance { "database_maintenance" } else { "database_connection" }, facts: &facts,
original_contract: "display-debug-source/v1", source_interface: "error_source",
};
let (submission, original_capture) = match output {
Some(output) => capture(output, occurrence, RawSource::Error(error), &header),
None => (DiagnosticSubmission::OutputUnavailable, OriginalCaptureState::OutputUnavailable),
};
UnrootedCaptureFacts { occurrence, submission, original_capture }
}