use crate::request_diagnostics::BoundaryFailure;
use crate::source_diagnostics::{SourceCapture, SourceReceipt, code};
use crate::{
BoundaryError, InvokeRequest, InvokeResponse, ProfuseContractAuthorityTemplate, TonicBoundary,
};
use saddle_core::{
BoundedDiagnostic, BoundedDiagnosticCause, CaptureSite, DiagnosticCategory,
DiagnosticOutcomeAxes, DiagnosticStage, OperationOutcome,
};
use saddle_observability::root_diagnostic::{RootOutcomeFacts, RootRequestEvent};
use saddle_observability::{DiagnosticSubmission, EmergencyDiagnosticHandle, EventContext};
use saddle_runtime::request_task::reserved::{ReservedRequestFailure, ReservedRequestView};
use std::sync::Mutex;
#[must_use]
pub struct ReservedFailure<E> {
error: E,
receipt: ReservedRequestFailure,
}
impl<E> ReservedFailure<E> {
pub fn error(&self) -> &E {
&self.error
}
pub fn occurrence(&self) -> saddle_core::DiagnosticOccurrence {
self.receipt.occurrence()
}
pub fn map_error<F>(self, map: impl FnOnce(E) -> F) -> ReservedFailure<F> {
ReservedFailure {
error: map(self.error),
receipt: self.receipt,
}
}
pub fn into_parts(self) -> (E, ReservedRequestFailure) {
(self.error, self.receipt)
}
}
struct State {
view: ReservedRequestView,
failure: Option<ReservedRequestFailure>,
started: bool,
completed: bool,
}
pub struct ReservedAttempt<'a> {
state: Mutex<State>,
output: Option<&'a EmergencyDiagnosticHandle>,
}
#[must_use]
pub enum ReservedAttemptCompletion {
Returned,
Interrupted(ReservedFailure<BoundaryFailure>),
}
fn facts(operation: OperationOutcome) -> RootOutcomeFacts {
RootOutcomeFacts {
axes: DiagnosticOutcomeAxes {
operation,
..Default::default()
},
..Default::default()
}
}
#[track_caller]
pub fn capture_io(
view: &ReservedRequestView,
output: Option<&EmergencyDiagnosticHandle>,
operation: crate::source_diagnostics::IoOperation,
error: &std::io::Error,
) -> ReservedFailure<std::io::ErrorKind> {
struct Capture<'a> {
view: &'a ReservedRequestView,
output: Option<&'a EmergencyDiagnosticHandle>,
event: RootRequestEvent,
receipt: Mutex<Option<ReservedRequestFailure>>,
}
impl SourceCapture for Capture<'_> {
fn submit_cause(
&self,
_: DiagnosticCategory,
_: BoundedDiagnosticCause,
) -> Option<SourceReceipt> {
unreachable!("IO always supplies its actual error")
}
fn submit_error(
&self,
error: &(dyn std::error::Error + 'static),
category: DiagnosticCategory,
cause: BoundedDiagnosticCause,
) -> Option<SourceReceipt> {
let diagnostic =
BoundedDiagnostic::capture(category, CaptureSite::FirstObserved, cause);
*self.receipt.lock().unwrap_or_else(|p| p.into_inner()) =
Some(self.view.source_error_with_facts(
error,
diagnostic,
code("transport.io"),
self.output,
self.event,
facts(OperationOutcome::Failed),
));
None
}
fn boundary(
&self,
_: Option<SourceReceipt>,
_: &DiagnosticOutcomeAxes,
) -> DiagnosticSubmission {
unreachable!("source only")
}
}
use crate::source_diagnostics::IoOperation;
let event = match operation {
IoOperation::Accept | IoOperation::ReadHead | IoOperation::ReadBody => {
RootRequestEvent::Ingress
}
_ => RootRequestEvent::Response,
};
let capture = Capture {
view,
output,
event,
receipt: Mutex::new(None),
};
let _ = (&capture as &dyn SourceCapture).io_failure(operation, error);
ReservedFailure {
error: error.kind(),
receipt: capture
.receipt
.into_inner()
.unwrap_or_else(|p| p.into_inner())
.expect("actual IO captured before mapping"),
}
}
impl<'a> ReservedAttempt<'a> {
pub fn capture_boundary_failure(
&mut self,
error: BoundaryError,
) -> ReservedFailure<BoundaryFailure> {
self.capture(
Some(&error),
DiagnosticCategory::ExpectedRejection,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestOutbound,
code("transport.adapter_failure"),
),
);
let mut state = self.state.lock().unwrap_or_else(|p| p.into_inner());
state.completed = true;
ReservedFailure {
error: BoundaryFailure {
code: error.code,
certainty: error.certainty,
},
receipt: state.failure.take().expect("source before mapping"),
}
}
pub fn complete_adapter(&mut self) {
self.state
.lock()
.unwrap_or_else(|p| p.into_inner())
.completed = true;
}
pub fn timed_out(&self) {
self.local_failure(crate::source_diagnostics::LocalFailure::Deadline);
}
pub fn new(view: ReservedRequestView, output: Option<&'a EmergencyDiagnosticHandle>) -> Self {
Self {
state: Mutex::new(State {
view,
failure: None,
started: false,
completed: false,
}),
output,
}
}
pub fn view(&self) -> ReservedRequestView {
self.state
.lock()
.unwrap_or_else(|p| p.into_inner())
.view
.clone()
}
#[track_caller]
fn capture(
&self,
original: Option<&(dyn std::error::Error + 'static)>,
category: DiagnosticCategory,
cause: BoundedDiagnosticCause,
) {
self.capture_with_outcome(
original,
category,
cause,
"transport source has no Error object",
OperationOutcome::Failed,
);
}
#[track_caller]
fn capture_with_outcome(
&self,
original: Option<&(dyn std::error::Error + 'static)>,
category: DiagnosticCategory,
cause: BoundedDiagnosticCause,
description: &'static str,
operation: OperationOutcome,
) {
let mut state = self.state.lock().unwrap_or_else(|p| p.into_inner());
if state.failure.is_some() {
return;
}
let diagnostic = BoundedDiagnostic::capture(category, CaptureSite::FirstObserved, cause);
state.failure = Some(match original {
Some(error) => state.view.source_error_with_facts(
error,
diagnostic,
code("transport.source"),
self.output,
RootRequestEvent::Outbound,
facts(operation),
),
None => state.view.source_description(
&description,
diagnostic,
code("transport.source"),
self.output,
RootRequestEvent::Outbound,
facts(operation),
),
});
}
pub fn finish(self) -> ReservedAttemptCompletion {
if !self
.state
.lock()
.unwrap_or_else(|p| p.into_inner())
.completed
{
self.local_failure(crate::source_diagnostics::LocalFailure::Cancelled);
}
let state = self.state.into_inner().unwrap_or_else(|p| p.into_inner());
if state.completed {
return ReservedAttemptCompletion::Returned;
}
ReservedAttemptCompletion::Interrupted(ReservedFailure {
error: BoundaryFailure {
code: crate::TechnicalCode::TransportFailure,
certainty: if state.started {
crate::ExecutionCertainty::MayHaveExecuted
} else {
crate::ExecutionCertainty::NotExecuted
},
},
receipt: state.failure.expect("cancellation captured"),
})
}
}
impl SourceCapture for ReservedAttempt<'_> {
fn protocol_failure(&self, description: &'static str) -> Option<SourceReceipt> {
self.capture_with_outcome(
None,
DiagnosticCategory::UnexpectedError,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestOutbound,
code("transport.response_invalid"),
),
description,
OperationOutcome::Failed,
);
None
}
fn guarded(&self) -> bool {
true
}
fn local_error(
&self,
error: &BoundaryError,
failure: crate::source_diagnostics::LocalFailure,
) -> Option<SourceReceipt> {
use crate::source_diagnostics::LocalFailure;
let (reason, operation) = match failure {
LocalFailure::Deadline => ("transport.deadline_elapsed", OperationOutcome::TimedOut),
LocalFailure::Cancelled => ("transport.cancelled", OperationOutcome::Cancelled),
LocalFailure::InvalidAuthority => {
("transport.authority_invalid", OperationOutcome::Rejected)
}
LocalFailure::InvalidRequest => {
("transport.request_invalid", OperationOutcome::Rejected)
}
};
self.capture_with_outcome(
Some(error),
DiagnosticCategory::ExpectedRejection,
BoundedDiagnosticCause::new(DiagnosticStage::RequestOutbound, code(reason)),
reason,
operation,
);
None
}
fn local_failure(
&self,
failure: crate::source_diagnostics::LocalFailure,
) -> Option<SourceReceipt> {
use crate::source_diagnostics::LocalFailure;
let (reason, operation) = match failure {
LocalFailure::Cancelled => ("transport.cancelled", OperationOutcome::Cancelled),
LocalFailure::Deadline => ("transport.deadline_elapsed", OperationOutcome::TimedOut),
LocalFailure::InvalidRequest => {
("transport.request_invalid", OperationOutcome::Rejected)
}
LocalFailure::InvalidAuthority => {
("transport.authority_invalid", OperationOutcome::Rejected)
}
};
self.capture_with_outcome(
None,
DiagnosticCategory::ExpectedRejection,
BoundedDiagnosticCause::new(DiagnosticStage::RequestOutbound, code(reason)),
reason,
operation,
);
None
}
fn required(&self) -> bool {
true
}
fn submit_cause(
&self,
category: DiagnosticCategory,
cause: BoundedDiagnosticCause,
) -> Option<SourceReceipt> {
self.capture(None, category, cause);
None }
fn submit_error(
&self,
error: &(dyn std::error::Error + 'static),
category: DiagnosticCategory,
cause: BoundedDiagnosticCause,
) -> Option<SourceReceipt> {
self.capture(Some(error), category, cause);
None
}
fn boundary(
&self,
_: Option<SourceReceipt>,
_: &DiagnosticOutcomeAxes,
) -> DiagnosticSubmission {
DiagnosticSubmission::OutputUnavailable
}
fn bind_child(
&self,
call: &saddle_core::CallContext,
_: &EventContext,
request: &str,
route: &str,
) -> Result<(), BoundaryError> {
let child = self
.state
.lock()
.unwrap_or_else(|p| p.into_inner())
.view
.child(call, request, route, 1)
.and_then(|view| {
view.with_phase(saddle_core::request_context::RequestViewPhase::Outbound)
});
match child {
Ok(view) => {
self.state.lock().unwrap_or_else(|p| p.into_inner()).view = view;
Ok(())
}
Err(error) => {
#[derive(Debug)]
struct Description<'a>(
&'a saddle_runtime::request_task::reserved::ReservedContextError,
);
impl std::fmt::Display for Description<'_> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.0)
}
}
let mut state = self.state.lock().unwrap_or_else(|p| p.into_inner());
let diagnostic = BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestOutbound,
code("transport.child_context_rejected"),
),
);
state.failure = Some(state.view.source_description(
&Description(&error),
diagnostic,
code("transport.child_context_rejected"),
self.output,
RootRequestEvent::Outbound,
facts(OperationOutcome::Rejected),
));
Err(BoundaryError::invalid_request(
"outbound diagnostic context rejected",
))
}
}
}
}
impl TonicBoundary {
pub async fn invoke_routed_reserved(
template: &ProfuseContractAuthorityTemplate,
request: InvokeRequest,
attempt: &mut ReservedAttempt<'_>,
) -> Result<InvokeResponse, ReservedFailure<BoundaryFailure>> {
{
let mut state = attempt.state.lock().unwrap_or_else(|p| p.into_inner());
if state.started {
let receipt = state.view.source_description(
&"attempt owner reused",
BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestOutbound,
code("transport.attempt_reused"),
),
),
code("transport.attempt_reused"),
attempt.output,
RootRequestEvent::Outbound,
facts(OperationOutcome::Rejected),
);
return Err(ReservedFailure {
error: BoundaryFailure {
code: crate::TechnicalCode::FunctionRequestInvalid,
certainty: crate::ExecutionCertainty::NotExecuted,
},
receipt,
});
}
state.started = true;
}
let result =
Self::invoke_routed_inner(template, request, attempt.output, Some(attempt)).await;
let mut state = attempt.state.lock().unwrap_or_else(|p| p.into_inner());
state.completed = true;
result.map_err(|error| ReservedFailure {
error: BoundaryFailure {
code: error.code,
certainty: error.certainty,
},
receipt: state
.failure
.take()
.expect("every technical exit captures before mapping"),
})
}
}
pub struct ReservedDelivery<'a> {
view: ReservedRequestView,
output: Option<&'a EmergencyDiagnosticHandle>,
written: u64,
started: bool,
complete: bool,
failure: Option<ReservedFailure<crate::request_diagnostics::DeliveryReason>>,
}
#[must_use]
pub enum ReservedDeliveryOutcome {
Complete(DiagnosticOutcomeAxes),
Failed {
axes: DiagnosticOutcomeAxes,
failure: ReservedFailure<crate::request_diagnostics::DeliveryReason>,
},
}
impl<'a> ReservedDelivery<'a> {
pub fn new(view: ReservedRequestView, output: Option<&'a EmergencyDiagnosticHandle>) -> Self {
Self {
view,
output,
written: 0,
started: false,
complete: false,
failure: None,
}
}
pub fn begin_write(&mut self) -> bool {
if self.complete || self.failure.is_some() {
return false;
}
self.started = true;
true
}
pub fn record_write(
&mut self,
result: std::io::Result<usize>,
) -> crate::request_diagnostics::WriteStep {
use crate::request_diagnostics::WriteStep;
if self.complete || self.failure.is_some() {
return WriteStep::Stopped;
}
self.started = true;
match result {
Ok(n) if n > 0 => {
self.written = self.written.saturating_add(n as u64);
WriteStep::Progress
}
Ok(_) => {
self.io_failure(&std::io::Error::from(std::io::ErrorKind::WriteZero));
WriteStep::Stopped
}
Err(error) => {
self.io_failure(&error);
WriteStep::Stopped
}
}
}
pub fn local_write_complete(&mut self) {
if self.failure.is_none() {
self.complete = true;
}
}
pub fn bytes_written(&self) -> u64 {
self.written
}
#[track_caller]
pub fn io_failure(&mut self, error: &std::io::Error) {
if self.failure.is_some() {
return;
}
let expected = matches!(
error.kind(),
std::io::ErrorKind::BrokenPipe
| std::io::ErrorKind::ConnectionReset
| std::io::ErrorKind::ConnectionAborted
| std::io::ErrorKind::TimedOut
| std::io::ErrorKind::UnexpectedEof
| std::io::ErrorKind::Interrupted
| std::io::ErrorKind::WouldBlock
| std::io::ErrorKind::InvalidData
| std::io::ErrorKind::InvalidInput
);
let diagnostic = BoundedDiagnostic::capture(
if expected {
DiagnosticCategory::ExpectedRejection
} else {
DiagnosticCategory::UnexpectedError
},
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(
DiagnosticStage::RequestResponse,
code("transport.write_io"),
)
.with_system(
crate::source_diagnostics::io_kind(error),
error.raw_os_error(),
),
);
self.failure = Some(ReservedFailure {
error: crate::request_diagnostics::DeliveryReason::Io(error.kind()),
receipt: self.view.source_error_with_facts(
error,
diagnostic,
code("transport.write_io"),
self.output,
RootRequestEvent::Response,
facts(OperationOutcome::Failed),
),
});
}
pub fn timed_out(&mut self) {
self.stop(crate::request_diagnostics::DeliveryReason::Timeout);
}
pub fn cancelled(&mut self) {
self.stop(crate::request_diagnostics::DeliveryReason::Cancelled);
}
#[track_caller]
fn stop(&mut self, reason: crate::request_diagnostics::DeliveryReason) {
if self.failure.is_some() {
return;
}
let (text, operation) = match reason {
crate::request_diagnostics::DeliveryReason::Timeout => {
("transport.delivery_timeout", OperationOutcome::TimedOut)
}
_ => ("transport.delivery_cancelled", OperationOutcome::Cancelled),
};
let diagnostic = BoundedDiagnostic::capture(
DiagnosticCategory::ExpectedRejection,
CaptureSite::FirstObserved,
BoundedDiagnosticCause::new(DiagnosticStage::RequestResponse, code(text)),
);
self.failure = Some(ReservedFailure {
error: reason,
receipt: self.view.source_description(
&text,
diagnostic,
code(text),
self.output,
RootRequestEvent::Response,
facts(operation),
),
});
}
pub fn finish(mut self) -> ReservedDeliveryOutcome {
if !self.complete && self.failure.is_none() {
self.cancelled();
}
let delivery = if self.complete {
saddle_core::ResponseDelivery::LocalWriteComplete
} else if self.written > 0 {
saddle_core::ResponseDelivery::Partial
} else if self.started {
saddle_core::ResponseDelivery::Unknown
} else {
saddle_core::ResponseDelivery::NotStarted
};
let mut axes = DiagnosticOutcomeAxes {
delivery,
bytes_written: Some(self.written),
..Default::default()
};
match self.failure {
None => {
axes.operation = OperationOutcome::Succeeded;
ReservedDeliveryOutcome::Complete(axes)
}
Some(failure) => {
axes.operation = match failure.error {
crate::request_diagnostics::DeliveryReason::Timeout => {
OperationOutcome::TimedOut
}
crate::request_diagnostics::DeliveryReason::Cancelled => {
OperationOutcome::Cancelled
}
_ => OperationOutcome::Failed,
};
ReservedDeliveryOutcome::Failed { axes, failure }
}
}
}
}