use crate::{DiagnosticSubmission, EmergencyDiagnosticHandle, EventContext, Observer};
use saddle_core::{
BoundedDiagnostic, ContextFact, Diagnostic, DiagnosticCategory, DiagnosticCode,
DiagnosticOccurrence, DiagnosticOutcomeAxes, RequestExecutionView,
};
use serde::Serialize;
mod original;
pub use original::{
OriginalCaptureState, RemoteProtocolContext, RemoteTargetContext, UnrootedCaptureFacts, UnrootedDiagnosticScope,
WrittenComponentFailure, WrittenUnrootedFailure,
ComponentSourceKind, ComponentSourceReceipt, original_capture_layout,
database_background_error, database_transaction_error, database_operation_error,
database_operation_borrowed_error, database_operation_borrowed_input_bounded,
write_json_error_debug,
process_startup_error, process_config_error,
process_finalization_error, process_component_start_error, process_component_cleanup_error,
process_component_start_recorded, process_component_cleanup_recorded,
process_shutdown_signal_error, process_request_drain_error, process_post_driver_error,
process_bootstrap_error,
service_registry_source_error, service_registry_source_description, service_route_source_description,
service_entry_trace_source_error, ServiceRegistrySourceOutput,
service_production_fact_source_error, service_compiled_build_source_error,
service_request_source_error, service_request_source_description,
runtime_admission_source_description,
runtime_unrooted_admission_source_description,
};
#[must_use = "pass the recorded failure to its formal consumer"]
pub struct RecordedSaddleError {
safe: saddle_core::SaddleError,
written: WrittenSource,
}
enum WrittenSource {
Unrooted(WrittenUnrootedFailure),
Component(WrittenComponentFailure),
}
impl WrittenSource {
fn occurrence(&self) -> DiagnosticOccurrence {
match self {
Self::Unrooted(written) => written.occurrence(),
Self::Component(written) => written.occurrence(),
}
}
}
impl RecordedSaddleError {
pub(crate) fn previously_written_writer_error<E>(
source: E,
diagnostic: Diagnostic,
written: Option<WrittenComponentFailure>,
) -> LifecycleSourceFailure
where E: std::error::Error + Send + Sync + 'static {
let safe = saddle_core::SaddleError::new(
saddle_core::ErrorKind::Infrastructure,
"observability.shutdown_failed",
"structured log shutdown failed",
).with_diagnostic(diagnostic);
match written {
Some(written) => Self::from_component_written(safe, written)
.map_or_else(|safe| UnconfirmedSaddleSource { safe, source }.erase(),
LifecycleSourceFailure::from),
None => UnconfirmedSaddleSource {
safe: safe.with_unconfirmed_source(), source,
}.erase(),
}
}
pub fn capture_unrooted<E, P>(
source: E,
diagnostic: Diagnostic,
application: &saddle_core::ContextLabel,
lifecycle: saddle_core::RequestViewPhase,
output: &crate::SourceOutput,
stage: RootRequestEvent,
outcome: RootOutcomeFacts,
project: P,
) -> Result<Self, UnconfirmedSaddleSource<E>>
where
E: std::error::Error + Send + Sync + 'static,
P: FnOnce(&E, Diagnostic) -> saddle_core::SaddleError,
{
let scope = UnrootedDiagnosticScope::new(application, lifecycle, Some(output.handle()));
let written = scope.source_existing_error_recorded(&source, &diagnostic, stage, outcome);
let safe = project(&source, diagnostic);
match written {
Ok(written) => Self::from_written(safe, written)
.map_err(|safe| UnconfirmedSaddleSource { safe, source }),
Err(_) => Err(UnconfirmedSaddleSource {
safe: safe.with_unconfirmed_source(), source,
}),
}
}
pub(crate) fn capture_component<E, P>(
source: E,
diagnostic: Diagnostic,
application: ContextFact<saddle_core::ContextLabel>,
output: &crate::SourceOutput,
kind: ComponentSourceKind,
project: P,
) -> Result<Self, UnconfirmedSaddleSource<E>>
where
E: std::error::Error + Send + Sync + 'static,
P: FnOnce(&E, Diagnostic) -> saddle_core::SaddleError,
{
let written = match kind {
ComponentSourceKind::Start => process_component_start_recorded(
Some(output.handle()), application, &diagnostic, &source),
ComponentSourceKind::Cleanup => process_component_cleanup_recorded(
Some(output.handle()), application, &diagnostic, &source),
};
let safe = project(&source, diagnostic);
match written {
Ok(written) => Self::from_component_written(safe, written)
.map_err(|safe| UnconfirmedSaddleSource { safe, source }),
Err(_) => Err(UnconfirmedSaddleSource {
safe: safe.with_unconfirmed_source(), source,
}),
}
}
pub fn component_result<T, E>(
result: Result<T, E>,
application: &saddle_core::ContextLabel,
output: &crate::SourceOutput,
kind: ComponentSourceKind,
primary: Option<DiagnosticOccurrence>,
safe_kind: saddle_core::ErrorKind,
code: &'static str,
safe_message: &'static str,
) -> Result<T, LifecycleSourceFailure>
where E: std::error::Error + Send + Sync + 'static {
result.map_err(|raw| {
let source_stage = match kind {
ComponentSourceKind::Start => saddle_core::DiagnosticStage::StartupListener,
ComponentSourceKind::Cleanup => saddle_core::DiagnosticStage::ShutdownComponent,
};
let mut diagnostic = Diagnostic::capture(
DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::DiagnosticCause::new(source_stage,
DiagnosticCode::new(code).expect("static component code")),
);
if let Some(primary) = primary.as_ref() {
diagnostic = diagnostic.during_cleanup_of_occurrence(primary);
}
Self::capture_component(raw, diagnostic,
ContextFact::Present(application.clone()), output, kind,
|_, diagnostic| saddle_core::SaddleError::new(
safe_kind, code, safe_message).with_diagnostic(diagnostic),
).map_or_else(UnconfirmedSaddleSource::erase, LifecycleSourceFailure::from)
})
}
pub(crate) fn from_written(
safe: saddle_core::SaddleError,
written: WrittenUnrootedFailure,
) -> Result<Self, saddle_core::SaddleError> {
Self::pair(safe, WrittenSource::Unrooted(written))
}
pub(crate) fn from_component_written(
safe: saddle_core::SaddleError,
written: WrittenComponentFailure,
) -> Result<Self, saddle_core::SaddleError> {
Self::pair(safe, WrittenSource::Component(written))
}
fn pair(safe: saddle_core::SaddleError, written: WrittenSource)
-> Result<Self, saddle_core::SaddleError> {
if safe.diagnostic().is_some_and(|diagnostic|
written.occurrence().matches_diagnostic(diagnostic)) {
Ok(Self { safe, written })
} else {
Err(safe.with_unconfirmed_source())
}
}
pub fn safe(&self) -> &saddle_core::SaddleError { &self.safe }
pub fn component_kind(&self) -> Option<ComponentSourceKind> {
match &self.written {
WrittenSource::Component(written) => Some(written.kind()),
WrittenSource::Unrooted(_) => None,
}
}
}
#[must_use = "propagate the unconfirmed source to the formal failure consumer"]
pub struct UnconfirmedSaddleSource<E> {
safe: saddle_core::SaddleError,
source: E,
}
impl<E> UnconfirmedSaddleSource<E> {
pub fn safe(&self) -> &saddle_core::SaddleError { &self.safe }
pub fn source(&self) -> &E { &self.source }
pub fn erase(self) -> LifecycleSourceFailure
where E: std::error::Error + Send + Sync + 'static {
LifecycleSourceFailure(LifecycleSourceFailureState::Unconfirmed {
safe: self.safe,
source: Box::new(self.source),
})
}
}
#[must_use = "forward the lifecycle failure to Application"]
pub struct LifecycleSourceFailure(LifecycleSourceFailureState);
enum LifecycleSourceFailureState {
Written(RecordedSaddleError),
WrongPhase {
safe: saddle_core::SaddleError,
_recorded: RecordedSaddleError,
},
Unconfirmed {
safe: saddle_core::SaddleError,
source: Box<dyn std::error::Error + Send + Sync>,
},
}
impl LifecycleSourceFailure {
pub fn safe(&self) -> &saddle_core::SaddleError {
match &self.0 {
LifecycleSourceFailureState::Written(recorded) => recorded.safe(),
LifecycleSourceFailureState::WrongPhase { safe, .. } => safe,
LifecycleSourceFailureState::Unconfirmed { safe, .. } => safe,
}
}
pub fn into_safe(self) -> saddle_core::SaddleError {
match self.0 {
LifecycleSourceFailureState::Written(recorded) => recorded.safe,
LifecycleSourceFailureState::WrongPhase { safe, _recorded } =>
safe.with_unconfirmed_original(Box::new(_recorded)),
LifecycleSourceFailureState::Unconfirmed { safe, source } =>
safe.with_unconfirmed_original(source),
}
}
pub fn require_component_kind(self, expected: ComponentSourceKind) -> Self {
match self.0 {
LifecycleSourceFailureState::Written(recorded)
if recorded.component_kind() != Some(expected) => {
let safe = saddle_core::SaddleError::new(
saddle_core::ErrorKind::Infrastructure,
"component.source_phase_mismatch",
"component source was recorded for another lifecycle phase",
).with_unconfirmed_source();
Self(LifecycleSourceFailureState::WrongPhase { safe, _recorded: recorded })
}
state => Self(state),
}
}
pub fn original_if_unconfirmed(&self) -> Option<&(dyn std::error::Error + Send + Sync)> {
match &self.0 {
LifecycleSourceFailureState::Written(_) => None,
LifecycleSourceFailureState::WrongPhase { .. } => None,
LifecycleSourceFailureState::Unconfirmed { source, .. } => Some(source.as_ref()),
}
}
}
impl From<RecordedSaddleError> for LifecycleSourceFailure {
fn from(recorded: RecordedSaddleError) -> Self {
Self(LifecycleSourceFailureState::Written(recorded))
}
}
#[cfg(test)]
mod lifecycle_projection_tests {
use super::*;
#[test]
fn unconfirmed_application_projection_retains_owned_original_without_receipt() {
let safe = saddle_core::SaddleError::new(saddle_core::ErrorKind::Infrastructure,
"component.unconfirmed", "safe").with_unconfirmed_source();
let failure = LifecycleSourceFailure(LifecycleSourceFailureState::Unconfirmed {
safe, source: Box::new(std::io::Error::other("M_PRIVATE_ORIGINAL")),
});
let projected = failure.into_safe();
assert!(projected.source_unavailable());
assert!(projected.source_receipt::<ComponentSourceReceipt>().is_none());
assert_eq!(projected.unconfirmed_original().unwrap().to_string(), "M_PRIVATE_ORIGINAL");
assert!(!projected.to_string().contains("M_PRIVATE_ORIGINAL"));
}
}
pub type RecordedLifecycleFuture<'a> = std::pin::Pin<Box<dyn std::future::Future<
Output = Result<(), LifecycleSourceFailure>> + Send + 'a>>;
pub trait RecordedComponentLifecycle: Send + Sync {
fn name(&self) -> &'static str;
fn start<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput) -> RecordedLifecycleFuture<'a>;
fn shutdown<'a>(&'a self, application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput)
-> RecordedLifecycleFuture<'a>;
fn shutdown_with_primary<'a>(
&'a self,
application: &'a saddle_core::ContextLabel,
output: &'a crate::SourceOutput,
primary: Option<DiagnosticOccurrence>,
) -> RecordedLifecycleFuture<'a>;
}
impl std::fmt::Display for RecordedSaddleError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
self.safe.fmt(formatter)
}
}
impl std::fmt::Debug for RecordedSaddleError {
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
std::fmt::Debug::fmt(&self.safe, formatter)
}
}
impl std::error::Error for RecordedSaddleError {
fn source(&self) -> Option<&(dyn std::error::Error + 'static)> {
Some(&self.safe)
}
}
pub fn request_identity_group(
call: &saddle_core::CallContext,
event: &EventContext,
zone: ContextFact<saddle_core::ContextLabel>,
) -> Result<saddle_core::RequestIdentityGroup, saddle_core::ContextConflict> {
saddle_core::RequestIdentityGroup::from_validated(
call,
event.diagnostic_request(),
event.diagnostic_route(),
event.diagnostic_attempt(),
zone,
)
}
pub fn request_child_view(
parent: &RequestExecutionView,
call: &saddle_core::CallContext,
event: &EventContext,
) -> Result<RequestExecutionView, saddle_core::ContextConflict> {
parent.child(
call,
event.diagnostic_request(),
event.diagnostic_route(),
event.diagnostic_attempt(),
)
}
#[derive(Clone, Copy, Debug, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RootRequestEvent {
Ingress,
Admission,
Handler,
Database,
Outbound,
Response,
Finalization,
Supervision,
}
#[derive(Clone, Copy, Serialize)]
#[serde(rename_all = "snake_case")]
pub enum RequestTransactionFact {
Committed,
Rejected,
Unknown,
}
#[derive(Clone, Copy, Serialize)]
pub struct RootOutcomeFacts {
pub axes: DiagnosticOutcomeAxes,
pub transaction: ContextFact<RequestTransactionFact>,
}
impl Default for RootOutcomeFacts {
fn default() -> Self {
Self {
axes: DiagnosticOutcomeAxes::default(),
transaction: ContextFact::Unavailable,
}
}
}
#[derive(Serialize)]
#[serde(untagged)]
enum SourceDetail<'a> {
Bounded(&'a BoundedDiagnostic),
Existing(&'a Diagnostic),
}
fn timestamp() -> u128 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
}
#[derive(Serialize)]
struct RootRecord<'a> {
schema_version: u8,
timestamp_unix_ms: u128,
level: &'static str,
elapsed_ms: Option<u64>,
event: &'static str,
stage: RootRequestEvent,
context: &'a RequestExecutionView,
source_context: Option<&'a RequestExecutionView>,
occurrence: Option<DiagnosticOccurrence>,
classification: Option<DiagnosticCode>,
source_submission: Option<DiagnosticSubmission>,
original_capture: Option<OriginalCaptureState>,
axes: RootOutcomeFacts,
source_outcome: Option<RootOutcomeFacts>,
category: Option<DiagnosticCategory>,
detail_status: &'static str,
diagnostic: Option<SourceDetail<'a>>,
}
#[must_use = "carry the source receipt to its declared terminal boundary"]
pub struct RootRequestFailure {
source: RequestExecutionView,
occurrence: DiagnosticOccurrence,
classification: DiagnosticCode,
source_submission: DiagnosticSubmission,
category: DiagnosticCategory,
outcome: RootOutcomeFacts,
original_capture: OriginalCaptureState,
}
#[must_use = "carry the written source to its terminal boundary"]
pub struct WrittenRootFailure(RootRequestFailure);
impl WrittenRootFailure {
pub fn into_source(self) -> RootRequestFailure { self.0 }
pub fn occurrence(&self) -> DiagnosticOccurrence { self.0.occurrence }
}
#[derive(Clone, Copy, Serialize)]
pub struct PublicRequestFailure {
occurrence: DiagnosticOccurrence,
classification: DiagnosticCode,
source_submission: DiagnosticSubmission,
terminal_submission: DiagnosticSubmission,
category: DiagnosticCategory,
source_outcome: RootOutcomeFacts,
original_capture: OriginalCaptureState,
}
#[must_use = "the supervisor must consume the actual task result"]
pub struct RootSupervisionReturn<T> {
result: Result<T, RootRequestFailure>,
}
impl<T> RootSupervisionReturn<T> {
pub fn completed(value: T) -> Self {
Self { result: Ok(value) }
}
pub fn failed(failure: WrittenRootFailure) -> Self {
Self {
result: Err(failure.into_source()),
}
}
pub fn consume(self) -> Result<T, RootRequestFailure> {
self.result
}
}
pub struct RootDiagnosticScope<'a> {
view: &'a RequestExecutionView,
output: Option<&'a EmergencyDiagnosticHandle>,
}
impl<'a> RootDiagnosticScope<'a> {
pub fn new(
view: &'a RequestExecutionView,
output: Option<&'a EmergencyDiagnosticHandle>,
) -> Self {
Self { view, output }
}
pub fn source(
&self,
diagnostic: BoundedDiagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.submit_source(
diagnostic.occurrence(),
diagnostic.category(),
SourceDetail::Bounded(&diagnostic),
classification,
stage,
axes,
)
}
pub fn source_existing(
&self,
diagnostic: Diagnostic,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
self.submit_source(
diagnostic.occurrence(),
diagnostic.category(),
SourceDetail::Existing(&diagnostic),
classification,
stage,
axes,
)
}
fn submit_source(
&self,
occurrence: DiagnosticOccurrence,
category: DiagnosticCategory,
diagnostic: SourceDetail<'_>,
classification: DiagnosticCode,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> RootRequestFailure {
let record = RootRecord {
schema_version: 2,
timestamp_unix_ms: timestamp(),
elapsed_ms: None,
level: if matches!(category, DiagnosticCategory::ExpectedRejection) {
"warn"
} else {
"error"
},
event: "request_failure_source",
stage,
context: self.view,
source_context: None,
occurrence: Some(occurrence),
classification: Some(classification),
source_submission: None,
original_capture: Some(OriginalCaptureState::LegacyProjectionOnly),
axes,
diagnostic: Some(diagnostic),
source_outcome: None,
category: Some(category),
detail_status: "bounded",
};
let fallback = RootRecord {
diagnostic: None,
detail_status: "omitted_encoding_capacity",
..record
};
let submission = self
.output
.map_or(DiagnosticSubmission::OutputUnavailable, |output| {
output.submit_fixed_with_fallback(&record, &fallback)
});
RootRequestFailure {
source: self.view.clone(),
occurrence,
classification,
source_submission: submission,
category,
outcome: axes,
original_capture: OriginalCaptureState::LegacyProjectionOnly,
}
}
pub fn ordinary(
&self,
observer: &Observer,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> DiagnosticSubmission {
self.ordinary_at(observer, stage, axes, None, "request_stage")
}
fn ordinary_at(
&self,
observer: &Observer,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
elapsed_ms: Option<u64>,
event: &'static str,
) -> DiagnosticSubmission {
observer.emit_root_record(&RootRecord {
schema_version: 2,
timestamp_unix_ms: timestamp(),
level: "info",
elapsed_ms,
event,
stage,
context: self.view,
source_context: None,
occurrence: None,
classification: None,
source_submission: None,
original_capture: None,
axes,
diagnostic: None,
source_outcome: None,
category: None,
detail_status: "not_applicable",
})
}
pub fn start_stage(
self,
observer: &'a Observer,
stage: RootRequestEvent,
) -> RootActiveStage<'a> {
self.bounded_event(
stage,
RootOutcomeFacts::default(),
Some(0),
"request_stage_started",
);
RootActiveStage {
scope: self,
observer,
stage,
started: std::time::Instant::now(),
finished: false,
}
}
fn bounded_event(
&self,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
elapsed_ms: Option<u64>,
event: &'static str,
) -> DiagnosticSubmission {
self.output
.map_or(DiagnosticSubmission::OutputUnavailable, |output| {
output.submit_fixed_record(&RootRecord {
schema_version: 2,
timestamp_unix_ms: timestamp(),
level: "info",
elapsed_ms,
event,
stage,
context: self.view,
source_context: None,
occurrence: None,
classification: None,
source_submission: None,
original_capture: None,
axes,
source_outcome: None,
category: None,
detail_status: "not_applicable",
diagnostic: None,
})
})
}
}
impl RootRequestFailure {
fn into_public(self, terminal_submission: DiagnosticSubmission) -> PublicRequestFailure {
PublicRequestFailure {
occurrence: self.occurrence,
classification: self.classification,
source_submission: self.source_submission,
terminal_submission,
category: self.category,
source_outcome: self.outcome,
original_capture: self.original_capture,
}
}
pub fn original_capture(&self) -> OriginalCaptureState {
self.original_capture
}
pub fn occurrence(&self) -> DiagnosticOccurrence {
self.occurrence
}
pub fn submission(&self) -> DiagnosticSubmission {
self.source_submission
}
pub fn source_view(&self) -> &RequestExecutionView {
&self.source
}
pub fn map_classification(mut self, classification: DiagnosticCode) -> Self {
self.classification = classification;
self
}
pub fn boundary(
self,
current: &RequestExecutionView,
output: Option<&EmergencyDiagnosticHandle>,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
) -> Result<(Self, DiagnosticSubmission), Self> {
self.boundary_at(current, output, stage, axes, None, "request_failure_boundary")
}
fn boundary_at(
self,
current: &RequestExecutionView,
output: Option<&EmergencyDiagnosticHandle>,
stage: RootRequestEvent,
axes: RootOutcomeFacts,
elapsed_ms: Option<u64>,
event: &'static str,
) -> Result<(Self, DiagnosticSubmission), Self> {
if !self.source.same_request(current) {
return Err(self);
}
let submission = output.map_or(DiagnosticSubmission::OutputUnavailable, |output| {
output.submit_fixed_record(&RootRecord {
schema_version: 2,
timestamp_unix_ms: timestamp(),
level: "info",
elapsed_ms,
event,
stage,
context: current,
source_context: Some(&self.source),
occurrence: Some(self.occurrence),
classification: Some(self.classification),
source_submission: Some(self.source_submission),
original_capture: Some(self.original_capture),
axes,
diagnostic: None,
source_outcome: Some(self.outcome),
category: Some(self.category),
detail_status: "source_reference_only",
})
});
Ok((self, submission))
}
pub fn finish(
self,
current: &RequestExecutionView,
output: Option<&EmergencyDiagnosticHandle>,
axes: RootOutcomeFacts,
) -> Result<PublicRequestFailure, Self> {
let (receipt, terminal_submission) =
self.boundary(current, output, RootRequestEvent::Finalization, axes)?;
Ok(PublicRequestFailure {
occurrence: receipt.occurrence,
classification: receipt.classification,
source_submission: receipt.source_submission,
terminal_submission,
category: receipt.category,
source_outcome: receipt.outcome,
original_capture: receipt.original_capture,
})
}
}
pub struct RootActiveStage<'a> {
scope: RootDiagnosticScope<'a>,
observer: &'a Observer,
stage: RootRequestEvent,
started: std::time::Instant,
finished: bool,
}
impl RootActiveStage<'_> {
fn finish_metrics(&mut self, outcome: saddle_core::OperationOutcome) -> u64 {
self.finished = true;
let elapsed = u64::try_from(self.started.elapsed().as_millis()).unwrap_or(u64::MAX);
use crate::{Stage, StageOutcome};
use saddle_core::OperationOutcome as Outcome;
let stage = match self.stage {
RootRequestEvent::Ingress => Some(Stage::Ingress),
RootRequestEvent::Admission => Some(Stage::Admission),
RootRequestEvent::Handler => Some(Stage::Handler),
RootRequestEvent::Database => Some(Stage::Database),
RootRequestEvent::Outbound => Some(Stage::ProfuseContract),
RootRequestEvent::Response => Some(Stage::Response),
RootRequestEvent::Finalization => Some(Stage::ResourceFinalization),
RootRequestEvent::Supervision => None,
};
let outcome = match outcome {
Outcome::Succeeded => Some(StageOutcome::Success),
Outcome::Rejected => Some(StageOutcome::Rejected),
Outcome::Failed | Outcome::Panicked => Some(StageOutcome::Failure),
Outcome::Cancelled | Outcome::TimedOut => Some(StageOutcome::Cancelled),
Outcome::Unknown => None,
};
if let (Some(stage), Some(outcome)) = (stage, outcome) {
self.observer
.inner
.metrics
.stage_finished(stage, outcome, elapsed);
}
elapsed
}
pub fn finish_nonfailure(
mut self,
facts: RootOutcomeFacts,
) -> Result<DiagnosticSubmission, Self> {
if !matches!(
facts.axes.operation,
saddle_core::OperationOutcome::Succeeded | saddle_core::OperationOutcome::Rejected
) {
return Err(self);
}
let elapsed = self.finish_metrics(facts.axes.operation);
Ok(self
.scope
.bounded_event(self.stage, facts, Some(elapsed), "request_stage_finished"))
}
pub fn finish_failure_public(
self,
failure: RootRequestFailure,
facts: RootOutcomeFacts,
) -> Result<PublicRequestFailure, (Self, RootRequestFailure)> {
self.finish_failure(failure, facts)
.map(|(failure, submission)| failure.into_public(submission))
}
#[allow(clippy::result_large_err)] pub fn finish_failure(
self,
failure: RootRequestFailure,
facts: RootOutcomeFacts,
) -> Result<(RootRequestFailure, DiagnosticSubmission), (Self, RootRequestFailure)> {
self.finish_failure_record(failure, facts, "request_failure_boundary")
}
#[allow(clippy::result_large_err)] pub fn finish_failure_retained(
self,
failure: RootRequestFailure,
facts: RootOutcomeFacts,
) -> Result<(RootRequestFailure, DiagnosticSubmission), (Self, RootRequestFailure)> {
self.finish_failure_record(failure, facts, "request_stage_finished")
}
#[allow(clippy::result_large_err)] fn finish_failure_record(
mut self,
failure: RootRequestFailure,
facts: RootOutcomeFacts,
event: &'static str,
) -> Result<(RootRequestFailure, DiagnosticSubmission), (Self, RootRequestFailure)> {
if !failure.source.same_request(self.scope.view) {
return Err((self, failure));
}
let elapsed = self.finish_metrics(facts.axes.operation);
match failure.boundary_at(
self.scope.view,
self.scope.output,
self.stage,
facts,
Some(elapsed),
event,
) {
Ok(result) => Ok(result),
Err(failure) => Err((self, failure)),
}
}
}
impl Drop for RootActiveStage<'_> {
fn drop(&mut self) {
if !self.finished {
let elapsed = self.finish_metrics(saddle_core::OperationOutcome::Cancelled);
let facts = RootOutcomeFacts {
axes: DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Cancelled,
..Default::default()
},
..Default::default()
};
self.scope
.bounded_event(self.stage, facts, Some(elapsed), "request_stage_finished");
}
}
}
pub fn root_failure_layout() -> std::alloc::Layout {
std::alloc::Layout::new::<RootRequestFailure>()
}
pub struct RequestLoggingLayouts {
pub source_frame: std::alloc::Layout,
pub emergency_packet: std::alloc::Layout,
pub emergency_slots: usize,
pub ordinary_command: std::alloc::Layout,
pub ordinary_encoded_bytes_max: usize,
pub source_record: std::alloc::Layout,
pub confirmation_owner: std::alloc::Layout,
}
pub fn request_logging_layouts() -> RequestLoggingLayouts {
let (source_frame, emergency_packet, emergency_slots) = crate::diagnostic::root_frame_layouts();
RequestLoggingLayouts {
source_frame,
emergency_packet,
emergency_slots,
ordinary_command: crate::logger::root_queue_layout(),
ordinary_encoded_bytes_max: 8191,
source_record: std::alloc::Layout::new::<RootRecord<'static>>(),
confirmation_owner: crate::diagnostic::root_status_owner_layout(),
}
}
impl PublicRequestFailure {
pub fn occurrence(&self) -> DiagnosticOccurrence {
self.occurrence
}
pub fn source_submission(&self) -> DiagnosticSubmission {
self.source_submission
}
pub fn terminal_submission(&self) -> DiagnosticSubmission {
self.terminal_submission
}
}
impl std::fmt::Debug for PublicRequestFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PublicRequestFailure")
.field("classification", &self.classification)
.field("source_submission", &self.source_submission)
.finish_non_exhaustive()
}
}
impl std::fmt::Display for PublicRequestFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "framework request failure: {:?}", self.classification)
}
}
impl std::error::Error for PublicRequestFailure {}
#[doc(hidden)]
pub fn database_maintenance_success(
output: Option<&EmergencyDiagnosticHandle>, datasource: u64, connection: u64,
sweep: u64, recovery_acquisition: bool,
) -> DiagnosticSubmission {
#[derive(Serialize)]
struct CompletedProbe {
schema_version: u8,
event: &'static str,
datasource: u64,
connection: u64,
maintenance_sweep: u64,
recovery_acquisition: bool,
response: &'static str,
request: ContextFact<()>,
trace_id: ContextFact<()>,
}
match output {
Some(output) => output.submit_fixed_record(&CompletedProbe {
schema_version: 1, event: "database_maintenance_probe", datasource, connection,
maintenance_sweep: sweep, recovery_acquisition, response: "complete_valid_select1",
request: ContextFact::NotApplicable, trace_id: ContextFact::NotApplicable,
}),
None => DiagnosticSubmission::OutputUnavailable,
}
}