use saddle_core::{
CaptureSite, Diagnostic, DiagnosticCategory, DiagnosticCause, DiagnosticCode, DiagnosticStage,
SaddleError,
};
use saddle_observability::{CallContext, EmergencyDiagnosticHandle, EventContext, Observer};
use std::{
cell::RefCell,
future::Future,
panic::{AssertUnwindSafe, PanicHookInfo, catch_unwind, resume_unwind},
sync::OnceLock,
};
static OUTPUT: OnceLock<EmergencyDiagnosticHandle> = OnceLock::new();
#[derive(Debug)]
pub struct RuntimeDiagnosticExit {
pub shutdown: saddle_observability::DiagnosticShutdown,
pub snapshot: saddle_observability::DiagnosticOutputSnapshot,
}
pub(crate) fn close_output(
mut owner: saddle_observability::EmergencyDiagnostics,
deadline: Option<std::time::Instant>,
) -> RuntimeDiagnosticExit {
let shutdown = loop {
let status = owner.shutdown();
if status != saddle_observability::DiagnosticShutdown::Pending {
break status;
}
let remaining = deadline
.map(|d| d.saturating_duration_since(std::time::Instant::now()))
.unwrap_or_default();
if remaining.is_zero() {
break status;
}
std::thread::sleep(remaining.min(std::time::Duration::from_millis(1)));
};
RuntimeDiagnosticExit {
shutdown,
snapshot: owner.snapshot(),
}
}
type RequestContext = (Observer, CallContext, EventContext);
#[cfg(test)]
pub(crate) fn frame_layout() -> std::alloc::Layout { std::alloc::Layout::new::<Frame>() }
struct Frame {
root_view: Option<crate::request_task::reserved::ReservedRequestView>,
root_panic: Option<crate::request_task::reserved::ReservedRequestFailure>,
stage: DiagnosticStage,
task: &'static str,
context: Option<RequestContext>,
panic: Option<Diagnostic>,
primary: Option<Diagnostic>,
request_mode: bool,
request_panic: Option<RequestCaptured>,
request_primary: Option<RequestCaptured>,
cleanup_occurrence: Option<saddle_core::DiagnosticOccurrence>,
request_scope: Option<saddle_observability::RequestDiagnosticScope<'static>>,
request_output: Option<EmergencyDiagnosticHandle>,
}
pub(crate) enum RequestCaptured {
Established(saddle_observability::RequestBoundaryReference<Diagnostic>),
MissingContext(Diagnostic),
}
impl RequestCaptured {
pub(crate) fn diagnostic(&self) -> &Diagnostic {
match self {
Self::Established(reference) => reference.source_diagnostic(),
Self::MissingContext(diagnostic) => diagnostic,
}
}
pub(crate) fn occurrence(&self) -> saddle_core::DiagnosticOccurrence {
self.diagnostic().occurrence()
}
pub(crate) fn record(&self, axes: &saddle_core::DiagnosticOutcomeAxes) {
self.record_with_output(OUTPUT.get(), axes)
}
pub(crate) fn record_with_output(&self, output: Option<&EmergencyDiagnosticHandle>, axes: &saddle_core::DiagnosticOutcomeAxes) {
match self {
Self::Established(reference) => {
let _ = reference.record_optional(output, axes);
}
Self::MissingContext(_) => boundary(Some(self.occurrence()), axes, None),
}
}
}
pub(crate) fn live_request_scope(
projection: &saddle_core::DbScopeDiagnosticContext<(&CallContext, &EventContext)>,
) -> saddle_observability::RequestDiagnosticScope<'static> {
saddle_observability::RequestDiagnosticScope::live_db_scope(OUTPUT.get(), projection)
}
fn capture_request(
diagnostic: Diagnostic,
context: Option<&RequestContext>,
checked: Option<&saddle_observability::RequestDiagnosticScope<'static>>,
output: Option<&EmergencyDiagnosticHandle>,
) -> RequestCaptured {
if let Some(scope) = checked {
return RequestCaptured::Established(
scope.reborrow().with_output(output).capture_existing(diagnostic, context.map(|(observer, _, _)| observer)).into_reference(),
);
}
if let Some((observer, call, event)) = context {
if let Some(scope) = checked {
return RequestCaptured::Established(
scope
.capture_existing(diagnostic, Some(observer))
.into_reference(),
);
}
let scope = match OUTPUT.get() {
Some(output) => {
saddle_observability::RequestDiagnosticScope::established(output, call, event)
}
None => saddle_observability::RequestDiagnosticScope::output_unavailable(call, event),
};
RequestCaptured::Established(
scope
.capture_existing(diagnostic, Some(observer))
.into_reference(),
)
} else {
emit(&diagnostic, None);
RequestCaptured::MissingContext(diagnostic)
}
}
pub(crate) fn capture_task_source(
diagnostic: Diagnostic,
scope: &saddle_observability::RequestDiagnosticScope<'static>,
output: Option<&EmergencyDiagnosticHandle>,
) -> RequestCaptured {
capture_request(diagnostic, None, Some(scope), output)
}
thread_local! { static CURRENT: RefCell<Option<Frame>> = const { RefCell::new(None) }; }
pub(crate) fn update_root_view(view:&crate::request_task::reserved::ReservedRequestView) {
CURRENT.with(|slot| {
if let Some(frame)=slot.borrow_mut().as_mut() {
if frame.root_view.as_ref().is_some_and(|current|current.same_request(view)) {
frame.root_view=Some(view.clone());
}
}
});
}
pub(crate) fn refresh_task_scope(scope: saddle_observability::RequestDiagnosticScope<'static>) {
CURRENT.with(|slot| {
if let Ok(mut slot) = slot.try_borrow_mut() {
if let Some(frame) = slot.as_mut() {
if frame.request_mode && frame.task == "runtime.formal_request_task" {
frame.request_scope = Some(scope);
}
}
}
});
}
pub fn install_output(handle: EmergencyDiagnosticHandle) -> Result<(), EmergencyDiagnosticHandle> {
OUTPUT.set(handle)
}
pub fn capture_current_panic(info: &PanicHookInfo<'_>) -> bool {
CURRENT.with(|slot| {
let Ok(mut slot) = slot.try_borrow_mut() else {
return false;
};
let Some(frame) = slot.as_mut() else {
return false;
};
let mut diagnostic =
Diagnostic::capture_panic(info, frame.stage).with_task(code(frame.task));
if let Some(primary) = frame.primary.as_ref() {
diagnostic = diagnostic.during_cleanup_of(primary);
}
if let Some(primary) = frame.request_primary.as_ref() {
diagnostic = diagnostic.during_cleanup_of(primary.diagnostic());
}
if let Some(primary) = frame.cleanup_occurrence.as_ref() {
diagnostic = diagnostic.during_cleanup_of_occurrence(primary);
}
if let Some(view) = &frame.root_view {
frame.root_panic = Some(view.panic_source(info, diagnostic, frame.request_output.as_ref()));
return true;
}
if frame.request_mode {
frame.request_panic = Some(capture_request(
diagnostic,
frame.context.as_ref(),
frame.request_scope.as_ref(),
frame.request_output.as_ref(),
));
return true;
}
emit(&diagnostic, frame.context.as_ref());
frame.panic = Some(diagnostic);
true
})
}
fn code(value: &'static str) -> DiagnosticCode {
DiagnosticCode::new(value).expect("Runtime diagnostic codes are static schema identifiers")
}
#[track_caller]
pub(crate) fn bounded_source(
stage: DiagnosticStage,
name: &'static str,
axes: &saddle_core::DiagnosticOutcomeAxes,
) -> (
saddle_core::BoundedDiagnostic,
Option<saddle_observability::DiagnosticSubmission>,
) {
let diagnostic = saddle_core::BoundedDiagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(stage, code(name)),
);
let submission = OUTPUT
.get()
.map(|output| output.submit_bounded(Some(&diagnostic), axes, None));
#[cfg(test)]
if name == "runtime.driver_ledger_shutdown_failed"
&& std::env::var_os("RUNTIME_LEDGER_SOURCE_TEST").is_some()
{
eprintln!("RUNTIME_LEDGER_SOURCE_SUBMISSION={submission:?}");
}
(diagnostic, submission)
}
#[track_caller]
pub(crate) fn request_stop_source(
timed_out: bool,
physical_return: bool,
stage: DiagnosticStage,
context: Option<(&CallContext, &EventContext)>,
) -> saddle_core::DiagnosticOccurrence {
let diagnostic = saddle_core::BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
stage,
code(if physical_return && timed_out {
"runtime.physical_return_deadline"
} else if physical_return {
"runtime.physical_return_cancelled"
} else if timed_out {
"runtime.scope_deadline"
} else {
"runtime.scope_cancelled"
}),
),
);
if let Some(output) = OUTPUT.get() {
let _submission = output.submit_bounded(
Some(&diagnostic),
&saddle_core::DiagnosticOutcomeAxes {
operation: if physical_return {
saddle_core::OperationOutcome::Unknown
} else if timed_out {
saddle_core::OperationOutcome::TimedOut
} else {
saddle_core::OperationOutcome::Cancelled
},
..Default::default()
},
context,
);
}
diagnostic.occurrence()
}
pub(crate) fn boundary(
occurrence: Option<saddle_core::DiagnosticOccurrence>,
axes: &saddle_core::DiagnosticOutcomeAxes,
context: Option<(&CallContext, &EventContext)>,
) {
if let Some(output) = OUTPUT.get() {
let _submission = output.submit_boundary(occurrence, axes, context);
}
}
#[track_caller]
pub(crate) fn finalizer_contract_abort(name: &'static str) -> ! {
let _source = bounded_source(
DiagnosticStage::FinalizerResource,
name,
&saddle_core::DiagnosticOutcomeAxes {
cleanup: saddle_core::CleanupOutcome::Failed,
..Default::default()
},
);
std::process::abort()
}
pub(crate) fn emit(diagnostic: &Diagnostic, context: Option<&RequestContext>) {
if let Some(output) = OUTPUT.get() {
if let Some((observer, call, event)) = context {
let _ = observer.record_diagnostic(diagnostic, output, Some((call, event)));
} else {
let _ = output.submit(diagnostic);
}
}
}
#[track_caller]
pub(crate) fn failure(
stage: DiagnosticStage,
category: DiagnosticCategory,
name: &'static str,
) -> Diagnostic {
Diagnostic::capture(
category,
CaptureSite::FirstObserved,
DiagnosticCause::new(stage, code(name)),
)
}
pub(crate) fn catching_root<T>(
view: crate::request_task::reserved::ReservedRequestView,
output: Option<EmergencyDiagnosticHandle>,
primary: Option<saddle_core::DiagnosticOccurrence>,
cleanup: bool,
f: impl FnOnce() -> T,
) -> Result<T, crate::request_task::reserved::ReservedRequestFailure> {
let stage = if cleanup { DiagnosticStage::FinalizerResource } else { DiagnosticStage::BackgroundTask };
let previous = CURRENT.with(|slot| slot.replace(Some(Frame {
root_view: Some(view), root_panic: None,
stage, task: "runtime.reserved_request_task", context: None,
panic: None, primary: None, request_mode: true,
request_panic: None, request_primary: None, cleanup_occurrence: primary,
request_scope: None, request_output: output,
})));
let result = catch_unwind(AssertUnwindSafe(f));
let mut frame = CURRENT.with(|slot| slot.replace(previous)).expect("synchronous root frame");
match result {
Ok(value) => Ok(value),
Err(payload) => {
if let Some(failure) = frame.root_panic.take() { return Err(failure); }
let mut diagnostic = failure(stage, DiagnosticCategory::Panic, "runtime.panic_without_hook");
if let Some(primary) = primary.as_ref() { diagnostic = diagnostic.during_cleanup_of_occurrence(primary); }
let description = payload.downcast_ref::<&str>().copied()
.or_else(|| payload.downcast_ref::<String>().map(String::as_str))
.unwrap_or("panic payload does not expose a string description");
Err(frame.root_view.as_ref().expect("root frame").existing_description(
description, diagnostic, frame.request_output.as_ref()))
}
}
}
pub(crate) fn catching<T>(
stage: DiagnosticStage,
task: &'static str,
context: Option<RequestContext>,
f: impl FnOnce() -> T,
) -> Result<T, Diagnostic> {
catching_cleanup(stage, task, context, None, f)
}
pub(crate) fn catching_cleanup<T>(
stage: DiagnosticStage,
task: &'static str,
context: Option<RequestContext>,
primary: Option<Diagnostic>,
f: impl FnOnce() -> T,
) -> Result<T, Diagnostic> {
let previous = CURRENT.with(|slot| {
slot.replace(Some(Frame {
root_view: None,
root_panic: None,
stage,
task,
context,
panic: None,
primary,
request_mode: false,
request_panic: None,
request_primary: None,
cleanup_occurrence: None,
request_scope: None,
request_output: None,
}))
});
let outcome = catch_unwind(AssertUnwindSafe(f));
let frame = CURRENT
.with(|slot| slot.replace(previous))
.expect("matching synchronous diagnostic frame");
match outcome {
Ok(value) => Ok(value),
Err(payload) => {
if let Ok(diagnostic) = payload.downcast::<Diagnostic>() {
return Err(*diagnostic);
}
let diagnostic = frame.panic.unwrap_or_else(|| {
let mut d = failure(
stage,
DiagnosticCategory::Panic,
"runtime.panic_without_hook",
)
.with_task(code(task));
if let Some(primary) = frame.primary.as_ref() {
d = d.during_cleanup_of(primary);
}
emit(&d, frame.context.as_ref());
d
});
Err(diagnostic)
}
}
}
pub(crate) fn catching_request<T>(
stage: DiagnosticStage,
task: &'static str,
context: Option<RequestContext>,
request_scope: Option<saddle_observability::RequestDiagnosticScope<'static>>,
primary: Option<RequestCaptured>,
f: impl FnOnce() -> T,
) -> (Result<T, RequestCaptured>, Option<RequestCaptured>) {
catching_request_with_output(stage, task, context, request_scope, OUTPUT.get().cloned(), primary, f)
}
pub(crate) fn catching_request_with_output<T>(
stage: DiagnosticStage,
task: &'static str,
context: Option<RequestContext>,
request_scope: Option<saddle_observability::RequestDiagnosticScope<'static>>,
request_output: Option<EmergencyDiagnosticHandle>,
primary: Option<RequestCaptured>,
f: impl FnOnce() -> T,
) -> (Result<T, RequestCaptured>, Option<RequestCaptured>) {
let previous = CURRENT.with(|slot| {
slot.replace(Some(Frame {
root_view: None,
root_panic: None,
stage,
task,
context,
panic: None,
primary: None,
request_mode: true,
request_panic: None,
request_primary: primary,
cleanup_occurrence: None,
request_scope,
request_output,
}))
});
let outcome = catch_unwind(AssertUnwindSafe(f));
let mut frame = CURRENT
.with(|slot| slot.replace(previous))
.expect("synchronous request frame");
let outcome = match outcome {
Ok(value) => Ok(value),
Err(payload) => match payload.downcast::<RequestCaptured>() {
Ok(captured) => Err(*captured),
Err(payload) => match payload.downcast::<Diagnostic>() {
Ok(diagnostic) => Err(RequestCaptured::MissingContext(*diagnostic)),
Err(_) => Err(frame.request_panic.take().unwrap_or_else(|| {
let mut diagnostic = failure(
stage,
DiagnosticCategory::Panic,
"runtime.panic_without_hook",
)
.with_task(code(task));
if let Some(primary) = frame.request_primary.as_ref() {
diagnostic = diagnostic.during_cleanup_of(primary.diagnostic());
}
if let Some(primary) = frame.cleanup_occurrence.as_ref() {
diagnostic = diagnostic.during_cleanup_of_occurrence(primary);
}
capture_request(
diagnostic,
frame.context.as_ref(),
frame.request_scope.as_ref(),
frame.request_output.as_ref(),
)
})),
},
},
};
(outcome, frame.request_primary)
}
pub(crate) fn link_task_cleanup(primary: Option<saddle_core::DiagnosticOccurrence>) {
CURRENT.with(|slot| {
if let Some(frame) = slot.borrow_mut().as_mut() {
if frame.request_mode && frame.task == "runtime.formal_request_task_drop" {
frame.cleanup_occurrence = primary;
}
}
});
}
pub(crate) async fn task<F: Future>(
future: F,
stage: DiagnosticStage,
name: &'static str,
) -> F::Output {
let mut future = std::pin::pin!(future);
std::future::poll_fn(|cx| {
match catching(stage, name, None, || {
let poll = future.as_mut().poll(cx);
#[cfg(test)]
if poll.is_ready()
&& name == "runtime.managed_request_task"
&& std::env::var_os("RUNTIME_DIAGNOSTIC_MANAGER_FAULT").is_some()
{
panic!("DIAGNOSTIC_PRIVATE_SENTINEL");
}
poll
}) {
Ok(poll) => poll,
Err(diagnostic) => resume_unwind(Box::new(diagnostic)),
}
})
.await
}
pub(crate) fn joined(error: tokio::task::JoinError) {
if error.is_cancelled() {
boundary(
None,
&saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Cancelled,
..Default::default()
},
None,
);
return;
}
if error.is_panic() {
let payload = error.into_panic();
if let Ok(diagnostic) = payload.downcast::<Diagnostic>() {
boundary(
Some(diagnostic.occurrence()),
&saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Panicked,
..Default::default()
},
None,
);
return;
} }
let d = failure(
DiagnosticStage::BackgroundTask,
DiagnosticCategory::UnexpectedError,
"runtime.task_join_failed",
);
emit(&d, None);
boundary(
Some(d.occurrence()),
&saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Panicked,
..Default::default()
},
None,
);
}
#[track_caller]
pub(crate) fn attach(
error: SaddleError,
stage: DiagnosticStage,
name: &'static str,
) -> SaddleError {
if error.diagnostic().is_some() {
error
} else {
error.with_diagnostic(failure(stage, DiagnosticCategory::UnexpectedError, name))
}
}
pub(crate) fn report(error: &SaddleError) {
if let Some(d) = error.diagnostic() {
emit(d, None);
}
}
fn lifecycle_application(
application: Option<&str>,
) -> saddle_core::ContextFact<saddle_core::ContextLabel> {
application
.and_then(|name| saddle_core::ContextLabel::checked(name).ok())
.map_or(saddle_core::ContextFact::Unavailable, saddle_core::ContextFact::Present)
}
fn component_original_confirmed(
error: &SaddleError,
expected: saddle_observability::root_diagnostic::ComponentSourceKind,
) -> bool {
error.diagnostic().zip(error.source_receipt::<
saddle_observability::root_diagnostic::ComponentSourceReceipt>())
.is_some_and(|(diagnostic, receipt)|
receipt.kind() == expected
&& receipt.original_capture()
== saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
&& receipt.occurrence().matches_diagnostic(diagnostic))
}
pub(crate) fn component_start_source(
error: SaddleError,
application: Option<&str>,
) -> SaddleError {
if error.diagnostic().is_some() {
let error = if component_original_confirmed(&error,
saddle_observability::root_diagnostic::ComponentSourceKind::Start) { error }
else { error.with_unconfirmed_source() };
report(&error);
return error;
}
let raw = SaddleError::new(error.kind(), error.code(), error.message());
let mut error = attach(error, DiagnosticStage::StartupListener, "runtime.component_start_failed");
let receipt = saddle_observability::root_diagnostic::process_component_start_error(
OUTPUT.get(), lifecycle_application(application),
error.diagnostic().expect("component startup diagnostic"), &raw,
);
if receipt.original_capture()
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
{
error = error.with_unconfirmed_source();
}
report(&error);
error
}
pub(crate) fn shutdown_signal_source(
error: SaddleError,
application: Option<&str>,
) -> SaddleError {
if error.code() == "runtime.request_audit_failed" && error.diagnostic().is_some() {
report(&error);
return error;
}
let raw = SaddleError::new(error.kind(), error.code(), error.message());
let mut error = attach(error, DiagnosticStage::ShutdownComponent, "runtime.shutdown_signal_failed");
let receipt = saddle_observability::root_diagnostic::process_shutdown_signal_error(
OUTPUT.get(), lifecycle_application(application),
error.diagnostic().expect("shutdown signal diagnostic"), &raw,
);
if receipt.original_capture()
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
{
error = error.with_unconfirmed_source();
}
report(&error);
error
}
pub(crate) fn request_drain_source(
error: SaddleError,
primary: Option<&SaddleError>,
application: Option<&str>,
) -> SaddleError {
let raw = SaddleError::new(error.kind(), error.code(), error.message());
let mut error = attach(error, DiagnosticStage::ShutdownComponent, "runtime.request_drain_failed");
if let Some(primary) = primary {
error = error.during_cleanup_of(primary);
}
let receipt = saddle_observability::root_diagnostic::process_request_drain_error(
OUTPUT.get(), lifecycle_application(application),
error.diagnostic().expect("request drain diagnostic"), &raw,
);
if receipt.original_capture()
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
{
error = if primary.is_some() {
error.with_unconfirmed_cleanup_source()
} else {
error.with_unconfirmed_source()
};
}
report(&error);
error
}
pub(crate) fn post_driver_source(
error: SaddleError,
primary: Option<&SaddleError>,
application: Option<&str>,
) -> SaddleError {
let raw = SaddleError::new(error.kind(), error.code(), error.message());
post_driver_source_with_raw(error, primary, application, &raw)
}
pub(crate) fn post_driver_source_with_raw(
error: SaddleError,
primary: Option<&SaddleError>,
application: Option<&str>,
raw: &(dyn std::error::Error + 'static),
) -> SaddleError {
post_driver_source_with_occurrence(error, primary.and_then(|error| error.diagnostic().map(|d| d.occurrence())), application, raw)
}
pub(crate) fn post_driver_source_with_occurrence(
error: SaddleError,
primary: Option<saddle_core::DiagnosticOccurrence>,
application: Option<&str>,
raw: &(dyn std::error::Error + 'static),
) -> SaddleError {
let mut error = attach(error, DiagnosticStage::FinalizerResource, "runtime.post_driver_failed");
if let Some(primary) = primary {
error = error.during_cleanup_of_occurrence(&primary);
}
let receipt = saddle_observability::root_diagnostic::process_post_driver_error(
OUTPUT.get(), lifecycle_application(application),
error.diagnostic().expect("post-driver diagnostic"), raw,
);
if receipt.original_capture()
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
{
error = if primary.is_some() {
error.with_unconfirmed_cleanup_source()
} else {
error.with_unconfirmed_source()
};
}
report(&error);
error
}
pub(crate) fn component_cleanup_source(
error: SaddleError,
primary: Option<&SaddleError>,
application: Option<&str>,
) -> SaddleError {
if error.diagnostic().is_some() {
let mut error = error;
if let Some(primary) = primary { error = error.during_cleanup_of(primary); }
if !component_original_confirmed(&error,
saddle_observability::root_diagnostic::ComponentSourceKind::Cleanup) {
error = error.with_unconfirmed_cleanup_source();
}
report(&error);
return error;
}
let raw = SaddleError::new(error.kind(), error.code(), error.message());
let mut error = attach(error, DiagnosticStage::ShutdownComponent, "runtime.cleanup_failed");
if let Some(primary) = primary {
error = error.during_cleanup_of(primary);
}
let receipt = saddle_observability::root_diagnostic::process_component_cleanup_error(
OUTPUT.get(), lifecycle_application(application),
error.diagnostic().expect("component cleanup diagnostic"), &raw,
);
if receipt.original_capture()
!= saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten
{
error = error.with_unconfirmed_cleanup_source();
}
report(&error);
error
}
#[cfg(test)]
mod component_receipt_tests {
use super::*;
fn recorded(code: &'static str, output: &EmergencyDiagnosticHandle, cleanup: bool) -> SaddleError {
let diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(if cleanup { DiagnosticStage::ShutdownComponent }
else { DiagnosticStage::StartupListener },
DiagnosticCode::new(code).unwrap()),
);
let raw = SaddleError::new(saddle_core::ErrorKind::Infrastructure, code, "raw source");
let receipt = if cleanup {
saddle_observability::root_diagnostic::process_component_cleanup_error(
Some(output), saddle_core::ContextFact::Unavailable, &diagnostic, &raw)
} else {
saddle_observability::root_diagnostic::process_component_start_error(
Some(output), saddle_core::ContextFact::Unavailable, &diagnostic, &raw)
};
assert_eq!(receipt.original_capture(),
saddle_observability::root_diagnostic::OriginalCaptureState::CompleteWritten);
SaddleError::new(saddle_core::ErrorKind::Infrastructure, code, "safe")
.with_diagnostic(diagnostic).with_source_receipt(receipt)
}
#[test]
fn application_reuses_component_source_receipt() {
let root = std::env::temp_dir().join(format!("saddle-runtime-component-receipt-{}",
std::process::id()));
std::fs::create_dir(&root).unwrap();
let mut output = saddle_observability::EmergencyDiagnostics::start(
&saddle_observability::FileLoggingConfig::new(root.clone(),
saddle_observability::Rotation::Daily)).unwrap();
let handle = output.handle();
let startup = recorded("runtime.component_start_failed", &handle, false);
let id = startup.diagnostic().unwrap().id();
let received = component_start_source(startup, None);
assert_eq!(received.diagnostic().unwrap().id(), id);
assert!(component_original_confirmed(&received,
saddle_observability::root_diagnostic::ComponentSourceKind::Start));
assert!(!received.source_unavailable());
let primary = recorded("runtime.primary_failure", &handle, false);
let cleanup = recorded("runtime.component_cleanup_failed", &handle, true);
let cleanup_id = cleanup.diagnostic().unwrap().id();
let received = component_cleanup_source(cleanup, Some(&primary), None);
assert_eq!(received.diagnostic().unwrap().id(), cleanup_id);
assert_eq!(serde_json::to_value(received.diagnostic().unwrap()).unwrap()
["primary_diagnostic_id"], primary.diagnostic().unwrap().id());
assert!(component_original_confirmed(&received,
saddle_observability::root_diagnostic::ComponentSourceKind::Cleanup));
assert!(!received.cleanup_source_unavailable());
let unconfirmed = SaddleError::new(saddle_core::ErrorKind::Infrastructure,
"runtime.unconfirmed", "safe").with_diagnostic(Diagnostic::capture(
DiagnosticCategory::UnexpectedError, CaptureSite::FirstObserved,
DiagnosticCause::new(DiagnosticStage::StartupListener,
DiagnosticCode::new("runtime.unconfirmed").unwrap())));
assert!(component_start_source(unconfirmed, None).source_unavailable());
let unrelated = Diagnostic::capture(DiagnosticCategory::UnexpectedError,
CaptureSite::FirstObserved,
DiagnosticCause::new(DiagnosticStage::StartupListener,
DiagnosticCode::new("runtime.wrong_occurrence").unwrap()));
let wrong_occurrence = recorded("runtime.wrong_occurrence", &handle, false)
.with_diagnostic(unrelated);
assert!(!component_original_confirmed(&wrong_occurrence,
saddle_observability::root_diagnostic::ComponentSourceKind::Start));
assert!(component_start_source(wrong_occurrence, None).source_unavailable());
let wrong_stage = recorded("runtime.wrong_stage", &handle, false);
assert!(component_cleanup_source(wrong_stage, None, None).cleanup_source_unavailable());
let rows: Vec<serde_json::Value> = std::fs::read_to_string(output.target()).unwrap()
.lines().map(|line| serde_json::from_str(line).unwrap()).collect();
for id in [id, cleanup_id] {
assert_eq!(rows.iter().filter(|row| row["event"] == "request_error_original"
&& row["occurrence"]["diagnostic_id"] == id
&& row["channel"] == "context").count(), 1);
}
drop(handle);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while output.shutdown() == saddle_observability::DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline { std::thread::yield_now(); }
let target = output.target().to_owned();
drop(output);
std::fs::remove_file(target).unwrap();
std::fs::remove_dir(root).unwrap();
}
}