use super::{
ExecutionCertainty, TechnicalClassification, TechnicalFailure, TechnicalFailureCode,
TechnicalRetention,
};
use crate::{
database_capability::DatabaseRequest,
grpc::{Method, Response},
};
use saddle_boundary::{
ProfuseContractAuthorityTemplate,
standard_rpc::{Attempt, Failure, FailureKind},
};
use saddle_observability::root_diagnostic::RootOutcomeFacts;
use saddle_runtime::request_task::reserved::{ReservedActiveStage, ReservedObservationStage};
struct Owner<'a> {
attempt: Option<Attempt<'a>>,
stage: Option<ReservedActiveStage<'a>>,
}
fn classification(
kind: FailureKind,
certainty: saddle_boundary::ExecutionCertainty,
) -> TechnicalClassification {
TechnicalClassification {
code: match kind {
FailureKind::Configuration => TechnicalFailureCode::FunctionRequestInvalid,
FailureKind::Connection => TechnicalFailureCode::DependencyUnavailable,
FailureKind::Remote | FailureKind::Cancelled => TechnicalFailureCode::TransportFailure,
FailureKind::Protocol => TechnicalFailureCode::ContractResultInvalid,
FailureKind::Resource => TechnicalFailureCode::CapacityRejected,
FailureKind::Deadline => TechnicalFailureCode::DeadlineExceeded,
FailureKind::Supervision => TechnicalFailureCode::InternalFailure,
},
certainty: match certainty {
saddle_boundary::ExecutionCertainty::NotExecuted => ExecutionCertainty::NotExecuted,
saddle_boundary::ExecutionCertainty::Executed => ExecutionCertainty::Executed,
saddle_boundary::ExecutionCertainty::MayHaveExecuted
| saddle_boundary::ExecutionCertainty::Unspecified => {
ExecutionCertainty::MayHaveExecuted
}
},
}
}
fn facts(operation: saddle_core::OperationOutcome) -> RootOutcomeFacts {
RootOutcomeFacts {
axes: saddle_core::DiagnosticOutcomeAxes {
operation,
..Default::default()
},
..Default::default()
}
}
impl Owner<'_> {
fn failure(&mut self, failure: Failure) -> TechnicalFailure {
let (kind, certainty, source) = failure.into_parts();
let operation = match kind {
FailureKind::Deadline => saddle_core::OperationOutcome::TimedOut,
FailureKind::Cancelled => saddle_core::OperationOutcome::Cancelled,
_ => saddle_core::OperationOutcome::Failed,
};
let receipt = self
.stage
.take()
.unwrap()
.finish_failure_public(source, facts(operation))
.unwrap_or_else(|_| unreachable!("standard call retains its original root"));
TechnicalFailure {
retained: TechnicalRetention::Reserved {
classification: classification(kind, certainty),
_receipt: receipt,
},
}
}
}
impl Drop for Owner<'_> {
fn drop(&mut self) {
let Some(attempt) = self.attempt.take() else {
return;
};
match attempt.finish() {
Some(failure) if self.stage.is_some() => {
let _retained = self.failure(failure);
}
None if self.stage.is_some() => {
let _ = self
.stage
.take()
.unwrap()
.finish_nonfailure(facts(saddle_core::OperationOutcome::Succeeded));
}
_ => {}
}
}
}
struct Preparation(saddle_runtime::profusegw::ProfuseGwReservedScopeFailure);
impl std::fmt::Debug for Preparation {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
use saddle_runtime::profusegw::ProfuseGwReservedScopeFailure as E;
match &self.0 {
E::Preparation(reason) => f.debug_tuple("Preparation").field(reason).finish(),
E::Execution(reason) => f.debug_tuple("Execution").field(reason).finish(),
}
}
}
impl std::fmt::Display for Preparation {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self)
}
}
impl std::error::Error for Preparation {}
fn child_rpc(
parent: &saddle_core::CallContext,
number: u64,
) -> Option<saddle_core::RpcCorrelationId> {
let parent = parent.rpc_correlation_id()?.as_str();
let mut bytes = [0u8; saddle_core::MAX_TRACE_CORRELATION_ID_BYTES];
let mut digits = [0u8; 20];
let mut at = digits.len();
let mut value = number;
loop {
at -= 1;
digits[at] = b'0' + (value % 10) as u8;
value /= 10;
if value == 0 {
break;
}
}
let len = parent
.len()
.checked_add(1)?
.checked_add(digits.len() - at)?;
if len > bytes.len() {
return None;
}
bytes[..parent.len()].copy_from_slice(parent.as_bytes());
bytes[parent.len()] = b'.';
bytes[parent.len() + 1..len].copy_from_slice(&digits[at..]);
saddle_core::RpcCorrelationId::new(std::str::from_utf8(&bytes[..len]).ok()?)
}
pub(crate) async fn call<M: Method>(
database: &mut DatabaseRequest,
template: Option<&ProfuseContractAuthorityTemplate>,
zone: &str,
business_unit: &str,
number: u64,
registry: Option<&crate::grpc::FrozenRegistry>,
request_id: &str,
trace: saddle_boundary::standard_rpc::TraceContext<'_>,
fields: &M::Fields<'_>,
) -> Result<Response<M>, TechnicalFailure> {
let context = database
.reserved_context()
.expect("standard calls require the formal request task");
let root_view = context.view();
let observer = database
.reserved_observer()
.expect("formal request observer")
.clone();
let output = database.diagnostic_output();
let child = registry
.and_then(|registry| registry.child_context::<M>(&observer, database.grpc_parent_context()))
.and_then(|child| {
child_rpc(database.grpc_parent_context(), number)
.map(|rpc| child.with_rpc_correlation_id(Some(rpc)))
});
let (view, child_error) = match child {
Some(child) if number <= u32::MAX as u64 => {
match root_view.child(&child, request_id, M::SPEC.path, number as u32) {
Ok(view) => (view, None),
Err(error) => (root_view, Some(error)),
}
}
_ if registry.is_some_and(|registry| crate::grpc::registered::<M>(registry.methods())) => (
root_view,
Some(
saddle_runtime::request_task::reserved::ReservedContextError::Context(
saddle_core::request_context::ContextConflict::ChildRelation,
),
),
),
_ => (root_view, None),
};
let methods = registry
.map(crate::grpc::FrozenRegistry::methods)
.unwrap_or(&[]);
let stage = view.start_stage(
&observer,
output.as_ref(),
ReservedObservationStage::Outbound,
);
let deadline = match database.grpc_deadline() {
Ok(deadline) => deadline,
Err(error) => {
let code = saddle_core::DiagnosticCode::new("grpc.preparation").unwrap();
let source = 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,
code,
),
),
code,
output.as_ref(),
saddle_observability::root_diagnostic::RootRequestEvent::Outbound,
facts(saddle_core::OperationOutcome::Rejected),
);
let receipt = stage
.finish_failure_public(source, facts(saddle_core::OperationOutcome::Rejected))
.unwrap_or_else(|_| unreachable!("same preparation root"));
return Err(TechnicalFailure {
retained: TechnicalRetention::Reserved {
classification: classification(
FailureKind::Resource,
saddle_boundary::ExecutionCertainty::NotExecuted,
),
_receipt: receipt,
},
});
}
};
let mut owner = Owner {
attempt: Some(Attempt::new(
view.clone(),
output.as_ref(),
&M::SPEC,
template
.map(ProfuseContractAuthorityTemplate::as_str)
.unwrap_or("NotProvided"),
business_unit,
number,
deadline,
)),
stage: Some(stage),
};
if !crate::grpc::registered::<M>(methods) || template.is_none() {
let failure = owner.attempt.as_mut().unwrap().capture(
&crate::grpc::BuildError::Unregistered,
FailureKind::Configuration,
None,
);
return Err(owner.failure(failure));
}
if let Some(error) = child_error {
#[derive(Debug)]
struct Child(saddle_runtime::request_task::reserved::ReservedContextError);
impl std::fmt::Display for Child {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
write!(f, "{:?}", self.0)
}
}
impl std::error::Error for Child {}
let failure =
owner
.attempt
.as_mut()
.unwrap()
.capture(&Child(error), FailureKind::Resource, None);
return Err(owner.failure(failure));
}
let template = template.expect("selected registered endpoint");
let request =
match crate::grpc::prepare::<M>(&database.framework_future_memory(), methods, fields) {
Ok(request) => request,
Err(error) => {
let kind = match &error {
crate::grpc::BuildError::Resource(_) => FailureKind::Resource,
crate::grpc::BuildError::Wire(_) => FailureKind::Protocol,
_ => FailureKind::Configuration,
};
let failure = owner.attempt.as_mut().unwrap().capture(&error, kind, None);
return Err(owner.failure(failure));
}
};
let memory = database.framework_future_memory();
let outcome = database
.supervise_external_concrete(|stage| {
saddle_boundary::standard_rpc::unary(
template,
zone,
&M::SPEC,
request.wire,
memory,
owner.attempt.as_mut().unwrap(),
stage,
registry
.map(|registry| registry.metadata::<M>())
.unwrap_or(&[]),
trace,
)
})
.await;
match outcome {
Ok(outcome) => {
if outcome.completion.supervision().panicked
|| outcome.completion.supervision().stop.is_some()
{
let timed_out = matches!(
outcome.result,
Err(
saddle_runtime::profusegw::ProfuseGwReservedScopeFailure::Execution(
saddle_runtime::profusegw::ProfuseGwScopeFailure::Stopped(
saddle_runtime::profusegw::ProfuseGwScopeStop::TimedOut
)
)
)
);
let kind = if timed_out {
FailureKind::Deadline
} else if outcome.completion.supervision().panicked {
FailureKind::Supervision
} else {
FailureKind::Cancelled
};
let control = outcome.result.err().unwrap_or_else(|| {
unreachable!("stopped supervision carries its control outcome")
});
let progress =
owner
.attempt
.as_mut()
.unwrap()
.capture(&Preparation(control), kind, None);
let retained = owner.failure(progress);
let TechnicalRetention::Reserved {
classification,
_receipt: outbound,
} = retained.retained
else {
unreachable!("outbound stage has reserved retention")
};
let receipt = outcome.completion.into_failure().unwrap_or_else(|_| {
unreachable!("failed supervision retains its original source")
});
return Err(TechnicalFailure {
retained: TechnicalRetention::StandardSupervised {
classification,
_receipt: receipt,
_outbound: outbound,
},
});
}
match outcome.result {
Ok(Ok(wire)) => Ok(Response::from_wire(wire)),
Ok(Err(failure)) => Err(owner.failure(failure)),
Err(error) => {
let failure = owner.attempt.as_mut().unwrap().capture(
&Preparation(error),
FailureKind::Configuration,
None,
);
Err(owner.failure(failure))
}
}
}
Err(saddle_runtime::profusegw::ProfuseGwReservedScopeFailure::Execution(
saddle_runtime::profusegw::ProfuseGwScopeFailure::Admission(error),
)) => {
let failure =
owner
.attempt
.as_mut()
.unwrap()
.capture(&error, FailureKind::Resource, None);
Err(owner.failure(failure))
}
Err(error) => {
let failure = owner.attempt.as_mut().unwrap().capture(
&Preparation(error),
FailureKind::Configuration,
None,
);
Err(owner.failure(failure))
}
}
}
#[doc(hidden)]
pub async fn declared_standard_call<M: Method>(
seal: &super::ApplicationContractSeal,
database: &mut DatabaseRequest,
business_unit: &'static str,
fields: &M::Fields<'_>,
) -> Result<Response<M>, TechnicalFailure> {
let registry = match &seal.adapter.boundary {
super::ContractBoundary::Standard(registry) => Some(registry),
_ => None,
};
let endpoint = registry.and_then(|registry| registry.endpoint::<M>());
let number = seal
.adapter
.next_call
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
call::<M>(
database,
endpoint,
&seal.adapter.zone,
business_unit,
number,
registry,
&seal.adapter.request_id,
saddle_boundary::standard_rpc::TraceContext {
trace_id: &seal.adapter.trace_id,
rpc_id: &seal.adapter.rpc_id_prefix,
},
fields,
)
.await
}
#[doc(hidden)]
pub fn unsupported_standard_adapter(seal: &super::ApplicationContractSeal) -> TechnicalFailure {
let scope = seal.adapter.scope.as_ref().map(|scope| scope.reborrow()).unwrap_or_else(||
saddle_observability::RequestDiagnosticScope::early(None, saddle_observability::EarlyRequestContext::unavailable()));
TechnicalFailure::capture(scope, TechnicalFailureCode::FunctionRequestInvalid,
ExecutionCertainty::NotExecuted, "service.outbound.standard_adapter_unavailable")
}
#[cfg(test)]
mod adapter_tests {
use super::*;
#[test]
fn fake_adapter_rejects_standard_wire_without_a_peer_attempt() {
let boundary = super::super::FakeProfuseContractBoundary::scripted([]);
let observed = boundary.boundary.clone();
let seal = super::super::ApplicationContractSeal::from_fake(boundary.boundary, 1_800_000_000_000);
let failure = unsupported_standard_adapter(&seal);
assert_eq!(failure.code(), TechnicalFailureCode::FunctionRequestInvalid);
assert_eq!(failure.certainty(), ExecutionCertainty::NotExecuted);
assert!(observed.attempts().is_empty());
}
}