use super::*;
use saddle_boundary::reserved_diagnostics::{
ReservedAttempt, ReservedAttemptCompletion, ReservedFailure,
};
use saddle_observability::root_diagnostic::RootOutcomeFacts;
use saddle_runtime::request_task::reserved::{
ReservedActiveStage, ReservedObservationStage, ReservedRequestFailure,
};
struct AttemptOwner<'a> {
attempt: Option<ReservedAttempt<'a>>,
stage: Option<ReservedActiveStage<'a>>,
output: Option<&'a saddle_observability::EmergencyDiagnosticHandle>,
}
impl AttemptOwner<'_> {
fn capture_description<D: std::fmt::Display + std::fmt::Debug>(
&mut self,
description: &D,
diagnostic_code: &'static str,
code: TechnicalFailureCode,
certainty: ExecutionCertainty,
) -> TechnicalFailure {
let diagnostic_code = saddle_core::DiagnosticCode::new(diagnostic_code).unwrap();
let receipt = self.attempt.as_ref().unwrap().view().source_description(
description,
saddle_core::BoundedDiagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestOutbound,
diagnostic_code,
),
),
diagnostic_code,
self.output,
saddle_observability::root_diagnostic::RootRequestEvent::Outbound,
Default::default(),
);
self.attempt.as_mut().unwrap().complete_adapter();
self.finish_failure(receipt, TechnicalClassification { code, certainty })
}
fn finish_failure(
&mut self,
receipt: ReservedRequestFailure,
classification: TechnicalClassification,
) -> TechnicalFailure {
let facts = RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Failed,
..Default::default()
},
..Default::default()
};
let receipt = self
.stage
.take()
.unwrap()
.finish_failure_public(receipt, facts)
.unwrap_or_else(|_| unreachable!("same call root"));
TechnicalFailure {
retained: TechnicalRetention::Reserved {
classification,
_receipt: receipt,
},
}
}
fn returned_failure(
&mut self,
failure: ReservedFailure<saddle_boundary::request_diagnostics::BoundaryFailure>,
) -> TechnicalFailure {
let (error, receipt) = failure.into_parts();
self.finish_failure(
receipt,
TechnicalClassification {
code: map_boundary_code(error.code),
certainty: map_boundary_certainty(error.certainty),
},
)
}
fn capture_error(
&mut self,
error: &(dyn std::error::Error + 'static),
code: TechnicalFailureCode,
certainty: ExecutionCertainty,
) -> TechnicalFailure {
let diagnostic_code = saddle_core::DiagnosticCode::new("service.outbound.decode").unwrap();
let receipt = self
.attempt
.as_ref()
.unwrap()
.view()
.source_error_with_facts(
error,
saddle_core::BoundedDiagnostic::capture(
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::CaptureSite::FirstObserved,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestOutbound,
diagnostic_code,
),
),
diagnostic_code,
self.output,
saddle_observability::root_diagnostic::RootRequestEvent::Outbound,
Default::default(),
);
self.finish_failure(receipt, TechnicalClassification { code, certainty })
}
}
impl Drop for AttemptOwner<'_> {
fn drop(&mut self) {
let Some(attempt) = self.attempt.take() else {
return;
};
match attempt.finish() {
ReservedAttemptCompletion::Returned => {
if let Some(stage) = self.stage.take() {
let facts = RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Succeeded,
..Default::default()
},
..Default::default()
};
let _ = stage.finish_nonfailure(facts);
}
}
ReservedAttemptCompletion::Interrupted(failure) => {
let (_, receipt) = failure.into_parts();
let facts = RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Cancelled,
..Default::default()
},
..Default::default()
};
if let Some(stage) = self.stage.take() {
let _receipt = stage
.finish_failure_public(receipt, facts)
.unwrap_or_else(|_| unreachable!("original cancellation"));
}
}
}
}
}
#[doc(hidden)]
pub fn declared_call<A, F, Q, R>(
request: Q,
seal: &ApplicationContractSeal,
database: Option<&crate::database_capability::DatabaseRequest>,
business_unit: &'static str,
function: &'static str,
) -> impl Future<Output = ExternalFunctionResult<R>> + Send + 'static
where
Q: Message + Send + 'static,
R: Message + Default + Send + 'static,
A: 'static,
F: 'static,
{
let adapter = seal.adapter.clone();
let database = database.map(|database| database.shared());
async move {
let Some(mut database) = database.filter(|database| database.reserved_context().is_some())
else {
let seal = ApplicationContractSeal { adapter };
return DeclaredExternalFunctionCall::<A, F, Q, R>::from_declared(
request,
&seal,
business_unit,
function,
)
.await;
};
let context = database.reserved_context().unwrap();
let view = context.view();
let observer = database
.reserved_observer()
.expect("formal request observer")
.clone();
let output = database.diagnostic_output();
let mut owner = AttemptOwner {
attempt: Some(ReservedAttempt::new(view.clone(), output.as_ref())),
stage: Some(view.start_stage(
&observer,
output.as_ref(),
ReservedObservationStage::Outbound,
)),
output: output.as_ref(),
};
let number = adapter.next_call.fetch_add(1, Ordering::Relaxed);
let target = match InvocationTarget::new(business_unit, function) {
Ok(target) => target,
Err(error) => {
let failure = owner
.attempt
.as_mut()
.unwrap()
.capture_boundary_failure(error);
return ExternalFunctionResult::TechnicalFailure(owner.returned_failure(failure));
}
};
let request = match saddle_boundary::InvokeRequest::unary(
adapter.request_id,
format!("{}-{number}", adapter.call_id_prefix),
target,
adapter.deadline_unix_ms,
saddle_boundary::CallerContext {
trace_info: Some(saddle_boundary::proto::TraceInfo {
trace_id: adapter.trace_id,
rpc_id: format!("{}.{number}", adapter.rpc_id_prefix),
}),
ldc_info: Some(saddle_boundary::proto::LdcInfo {
zone: adapter.zone,
idc: adapter.idc,
env: adapter.env,
}),
},
request.encode_to_vec(),
) {
Ok(request) => request,
Err(error) => {
let failure = owner
.attempt
.as_mut()
.unwrap()
.capture_boundary_failure(error);
return ExternalFunctionResult::TechnicalFailure(owner.returned_failure(failure));
}
};
let result = {
let attempt = owner.attempt.as_mut().unwrap();
database
.supervise_external_concrete(async {
match adapter.boundary {
ContractBoundary::Routed(template, _) => {
TonicBoundary::invoke_routed_reserved(&template, request, attempt).await
}
ContractBoundary::Fake(boundary) => {
let result = boundary.invoke(request).await;
attempt.complete_adapter();
result.map_err(|error| attempt.capture_boundary_failure(error))
}
ContractBoundary::Tonic(boundary) => {
let result = boundary.invoke(request).await;
attempt.complete_adapter();
result.map_err(|error| attempt.capture_boundary_failure(error))
}
}
})
.await
};
let result = match result {
Ok(outcome) => {
let completion = outcome.completion;
if let Err(saddle_runtime::profusegw::ProfuseGwReservedScopeFailure::Execution(
ref reason,
)) = outcome.result
{
if completion.supervision().panicked || completion.supervision().stop.is_some()
{
let code = match reason {
saddle_runtime::profusegw::ProfuseGwScopeFailure::Stopped(
saddle_runtime::profusegw::ProfuseGwScopeStop::TimedOut,
) => TechnicalFailureCode::DeadlineExceeded,
saddle_runtime::profusegw::ProfuseGwScopeFailure::Stopped(
saddle_runtime::profusegw::ProfuseGwScopeStop::Cancelled,
) => TechnicalFailureCode::TransportFailure,
_ => TechnicalFailureCode::InternalFailure,
};
let receipt = completion.into_failure().unwrap_or_else(|_| {
unreachable!("finished failed supervision retains its source")
});
return ExternalFunctionResult::TechnicalFailure(TechnicalFailure {
retained: TechnicalRetention::Supervised {
classification: TechnicalClassification {
code,
certainty: ExecutionCertainty::MayHaveExecuted,
},
_receipt: receipt,
},
});
}
}
outcome.result
}
Err(error) => Err(error),
};
let response = match result {
Ok(Ok(response)) => response,
Ok(Err(failure)) => {
return ExternalFunctionResult::TechnicalFailure(owner.returned_failure(failure));
}
Err(saddle_runtime::profusegw::ProfuseGwReservedScopeFailure::Preparation(reason)) => {
#[derive(Debug)]
struct Description(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)
}
}
return ExternalFunctionResult::TechnicalFailure(owner.capture_description(
&Description(reason),
"service.outbound.preparation",
TechnicalFailureCode::InternalFailure,
ExecutionCertainty::NotExecuted,
));
}
Err(saddle_runtime::profusegw::ProfuseGwReservedScopeFailure::Execution(reason)) => {
use saddle_runtime::profusegw::{ProfuseGwScopeFailure, ProfuseGwScopeStop};
let timed_out = match reason {
ProfuseGwScopeFailure::Stopped(ProfuseGwScopeStop::TimedOut) => true,
ProfuseGwScopeFailure::Stopped(ProfuseGwScopeStop::Cancelled) => false,
ProfuseGwScopeFailure::Panicked
| ProfuseGwScopeFailure::AlreadySupervised
| ProfuseGwScopeFailure::ScopeAlreadyEntered => {
#[derive(Debug)]
struct Description(ProfuseGwScopeFailure);
impl std::fmt::Display for Description {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.0)
}
}
return ExternalFunctionResult::TechnicalFailure(
owner.capture_description(
&Description(reason),
"service.outbound.supervision_failed",
TechnicalFailureCode::InternalFailure,
ExecutionCertainty::MayHaveExecuted,
),
);
}
};
if timed_out {
owner.attempt.as_ref().unwrap().timed_out();
}
let attempt = owner.attempt.take().unwrap();
let ReservedAttemptCompletion::Interrupted(failure) = attempt.finish() else {
unreachable!("stopped outstanding call")
};
let (error, receipt) = failure.into_parts();
let facts = RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation: if timed_out {
saddle_core::OperationOutcome::TimedOut
} else {
saddle_core::OperationOutcome::Cancelled
},
..Default::default()
},
..Default::default()
};
let receipt = owner
.stage
.take()
.unwrap()
.finish_failure_public(receipt, facts)
.unwrap_or_else(|_| unreachable!("same call"));
return ExternalFunctionResult::TechnicalFailure(TechnicalFailure {
retained: TechnicalRetention::Reserved {
classification: TechnicalClassification {
code: if timed_out {
TechnicalFailureCode::DeadlineExceeded
} else {
TechnicalFailureCode::TransportFailure
},
certainty: map_boundary_certainty(error.certainty),
},
_receipt: receipt,
},
});
}
};
match response.outcome {
Some(saddle_boundary::Outcome::Completed(completed)) => {
match R::decode(completed.result.as_slice()) {
Ok(value) => ExternalFunctionResult::Completed(value),
Err(error) => ExternalFunctionResult::TechnicalFailure(owner.capture_error(
&error,
TechnicalFailureCode::ContractResultInvalid,
ExecutionCertainty::Executed,
)),
}
}
Some(saddle_boundary::Outcome::TechnicalFailure(failure)) => {
#[derive(Debug)]
struct RemoteClassification {
code: i32,
certainty: i32,
}
impl std::fmt::Display for RemoteClassification {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(
f,
"remote technical code={} certainty={}; remote cause unavailable",
self.code, self.certainty
)
}
}
ExternalFunctionResult::TechnicalFailure(owner.capture_description(
&RemoteClassification {
code: failure.code,
certainty: failure.certainty,
},
"service.outbound.remote_failure",
map_code(failure.code),
map_certainty(failure.certainty),
))
}
None => ExternalFunctionResult::TechnicalFailure(owner.capture_description(
&"response outcome absent",
"service.outbound.outcome_missing",
TechnicalFailureCode::ContractResultInvalid,
ExecutionCertainty::MayHaveExecuted,
)),
}
}
}