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, UnrootedCaptureFacts, UnrootedDiagnosticScope, original_capture_layout,
database_background_error,
};
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,
}
#[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: RootRequestFailure) -> Self {
Self {
result: Err(failure),
}
}
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 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>>(),
}
}
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,
}
}