use std::future::Future;
use std::marker::PhantomData;
use std::pin::Pin;
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use std::task::{Context, Poll};
use prost::Message;
use saddle_boundary::{
BoundaryTransport, FakeProfuseContractBoundary as TransportFakeBoundary, InvocationTarget,
ProfuseContractEndpoint, TonicBoundary,
};
use serde::de::DeserializeOwned;
pub use saddle_boundary::{FakeAttempt, FakeExecutionCertainty, FakeStep, FakeTechnicalCode};
pub const MAX_PROFUSE_GW_USER_ID_BYTES: usize = 64;
#[derive(Clone)]
pub struct TraceInfo {
trace_id: String,
rpc_id: String,
}
impl TraceInfo {
pub fn trace_id(&self) -> &str {
&self.trace_id
}
pub fn rpc_id(&self) -> &str {
&self.rpc_id
}
}
#[derive(Clone)]
pub struct LdcInfo {
zone: String,
idc: String,
env: String,
}
impl LdcInfo {
pub fn zone(&self) -> &str {
&self.zone
}
pub fn idc(&self) -> &str {
&self.idc
}
pub fn env(&self) -> &str {
&self.env
}
}
#[derive(Clone)]
pub struct ProfuseGwContext {
user_id: [u8; MAX_PROFUSE_GW_USER_ID_BYTES],
user_id_len: u8,
trace_info: TraceInfo,
ldc_info: LdcInfo,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[allow(dead_code)]
pub(crate) enum ProfuseGwContextError {
MissingUserId,
UserIdTooLong,
}
impl ProfuseGwContext {
#[allow(dead_code)]
pub(crate) fn from_framework(
user_id: &str,
trace_id: &str,
rpc_id: &str,
zone: &str,
idc: &str,
env: &str,
) -> Result<Self, ProfuseGwContextError> {
if user_id.is_empty() {
return Err(ProfuseGwContextError::MissingUserId);
}
if user_id.len() > MAX_PROFUSE_GW_USER_ID_BYTES {
return Err(ProfuseGwContextError::UserIdTooLong);
}
let mut stored = [0; MAX_PROFUSE_GW_USER_ID_BYTES];
stored[..user_id.len()].copy_from_slice(user_id.as_bytes());
Ok(Self {
user_id: stored,
user_id_len: user_id.len() as u8,
trace_info: TraceInfo {
trace_id: trace_id.to_owned(),
rpc_id: rpc_id.to_owned(),
},
ldc_info: LdcInfo {
zone: zone.to_owned(),
idc: idc.to_owned(),
env: env.to_owned(),
},
})
}
pub fn user_id(&self) -> &str {
std::str::from_utf8(&self.user_id[..usize::from(self.user_id_len)])
.expect("ProfuseGwContext is constructed from validated UTF-8")
}
pub fn trace_info(&self) -> &TraceInfo {
&self.trace_info
}
pub fn ldc_info(&self) -> &LdcInfo {
&self.ldc_info
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ExecutionCertainty {
NotExecuted,
Executed,
MayHaveExecuted,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum TechnicalFailureCode {
FunctionNotFound,
FunctionRequestInvalid,
CapacityRejected,
DeadlineExceeded,
DependencyUnavailable,
ContractResultInvalid,
InternalFailure,
TransportFailure,
}
pub struct TechnicalFailure {
retained: saddle_observability::FrameworkRequestFailure<TechnicalClassification>,
}
#[derive(Debug)]
struct TechnicalClassification {
code: TechnicalFailureCode,
certainty: ExecutionCertainty,
}
impl std::fmt::Debug for TechnicalFailure {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TechnicalFailure")
.field("code", &self.code())
.field("certainty", &self.certainty())
.finish()
}
}
impl TechnicalFailure {
pub(crate) fn capture(
scope: saddle_observability::RequestDiagnosticScope<'_>,
code: TechnicalFailureCode,
certainty: ExecutionCertainty,
reason: &'static str,
) -> Self {
let retained = scope.fail(
TechnicalClassification { code, certainty },
saddle_core::DiagnosticCategory::UnexpectedError,
saddle_core::BoundedDiagnosticCause::new(
saddle_core::DiagnosticStage::RequestOutbound,
saddle_core::DiagnosticCode::new(reason).expect("closed source code"),
),
);
Self { retained }
}
#[cfg(test)]
fn from_framework(code: TechnicalFailureCode, certainty: ExecutionCertainty) -> Self {
Self::capture(
saddle_observability::RequestDiagnosticScope::early(
None,
saddle_observability::EarlyRequestContext::unavailable(),
),
code,
certainty,
"framework.test_failure",
)
}
pub fn code(&self) -> TechnicalFailureCode {
self.retained.error().code
}
pub fn certainty(&self) -> ExecutionCertainty {
self.retained.error().certainty
}
pub(crate) fn finish(
self,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) -> Self {
let (retained, _) = self.retained.finish_boundary_retained(
output,
&saddle_core::DiagnosticOutcomeAxes {
operation: saddle_core::OperationOutcome::Failed,
..Default::default()
},
);
Self { retained }
}
}
pub enum ExternalFunctionResult<T> {
Completed(T),
TechnicalFailure(TechnicalFailure),
}
#[doc(hidden)]
pub struct ApplicationContractSeal {
adapter: ContractAdapter,
}
impl ApplicationContractSeal {
#[cfg(test)]
fn from_fake(boundary: TransportFakeBoundary, deadline_unix_ms: i64) -> Self {
Self {
adapter: ContractAdapter {
boundary: ContractBoundary::Fake(boundary),
scope: None,
cancellation: None,
request_id: "alpha1-application".into(),
call_id_prefix: "call".into(),
trace_id: "test-trace".into(),
rpc_id_prefix: "test-rpc".into(),
zone: "test-zone".into(),
idc: "test-idc".into(),
env: "test".into(),
deadline_unix_ms,
next_call: Arc::new(AtomicU64::new(1)),
},
}
}
}
#[derive(Clone)]
struct ContractAdapter {
scope: Option<Arc<saddle_observability::RequestDiagnosticScope<'static>>>,
cancellation: Option<AttemptRetentions>,
boundary: ContractBoundary,
request_id: String,
call_id_prefix: String,
trace_id: String,
rpc_id_prefix: String,
zone: String,
idc: String,
env: String,
deadline_unix_ms: i64,
next_call: Arc<AtomicU64>,
}
#[derive(Clone)]
enum ContractBoundary {
Fake(TransportFakeBoundary),
Tonic(TonicBoundary),
Routed(
saddle_boundary::ProfuseContractAuthorityTemplate,
Option<saddle_observability::EmergencyDiagnosticHandle>,
),
}
fn accepted_diagnostic_scope(
accepted: &saddle_boundary::ingress::AcceptedIngress,
) -> Arc<saddle_observability::RequestDiagnosticScope<'static>> {
use saddle_observability::{EarlyRequestContext, RequestDiagnosticScope};
let mut early = EarlyRequestContext::unavailable();
if let Ok(trace) = saddle_core::TraceCorrelationId::new(accepted.trace_id.as_str()) {
early = early
.with_trace(&trace)
.unwrap_or_else(|e| e.into_original());
}
if let Some(rpc) = saddle_core::RpcCorrelationId::new(&accepted.rpc_id) {
early = early.with_rpc(&rpc).unwrap_or_else(|e| e.into_original());
}
if let Ok(request) = saddle_observability::RequestIdentity::new(&accepted.identity.request_id) {
early = early
.with_request(&request)
.unwrap_or_else(|e| e.into_original());
}
let scope = RequestDiagnosticScope::early(None, early);
Arc::new(
match saddle_observability::DiagnosticZone::from_validated_ingress(&accepted.zone) {
Ok(zone) => scope.with_zone(zone),
Err(_) => {
scope.with_zone_missing(saddle_observability::DiagnosticContextMissing::Unavailable)
}
},
)
}
impl ContractBoundary {
async fn invoke(
&self,
request: saddle_boundary::InvokeRequest,
attempt: &mut saddle_boundary::request_diagnostics::RequestAttempt,
) -> Result<
saddle_boundary::InvokeResponse,
saddle_boundary::request_diagnostics::RequiredBoundaryFailure,
> {
match self {
Self::Fake(boundary) => {
let result = boundary.invoke(request).await;
attempt.complete_adapter();
result.map_err(|error| attempt.capture_boundary_failure(error))
}
Self::Tonic(boundary) => boundary.invoke_required(request, attempt).await,
Self::Routed(template, _) => {
TonicBoundary::invoke_routed_required(template, request, attempt).await
}
}
}
}
#[doc(hidden)]
#[derive(Clone)]
pub struct ProfuseContractDeployment {
boundary: ContractBoundary,
}
impl ProfuseContractDeployment {
#[doc(hidden)]
pub fn bind_accepted_request(
&self,
accepted: &saddle_boundary::ingress::AcceptedIngress,
database: &crate::database_capability::DatabaseRequest,
) -> ApplicationContractSeal {
let mut seal = self
.clone()
.with_diagnostic_output(database.diagnostic_output())
.bind_accepted(accepted);
seal.adapter.scope = Some(Arc::new(
database.request_diagnostic_scope().with_output(None),
));
seal.adapter.cancellation = Some(database.outbound_retentions());
seal
}
#[doc(hidden)]
pub fn routed(template: saddle_boundary::ProfuseContractAuthorityTemplate) -> Self {
Self {
boundary: ContractBoundary::Routed(template, None),
}
}
pub(crate) fn with_diagnostic_output(
mut self,
output: Option<saddle_observability::EmergencyDiagnosticHandle>,
) -> Self {
if let ContractBoundary::Routed(_, slot) = &mut self.boundary {
*slot = output;
}
self
}
#[doc(hidden)]
pub async fn connect(endpoint: ProfuseContractEndpoint) -> Result<Self, TechnicalFailure> {
TonicBoundary::connect(endpoint)
.await
.map(|boundary| Self {
boundary: ContractBoundary::Tonic(boundary),
})
.map_err(|error| match boundary_error::<()>(error) {
ExternalFunctionResult::TechnicalFailure(failure) => failure,
ExternalFunctionResult::Completed(()) => unreachable!(),
})
}
#[doc(hidden)]
pub fn bind_accepted(
&self,
accepted: &saddle_boundary::ingress::AcceptedIngress,
) -> ApplicationContractSeal {
ApplicationContractSeal {
adapter: ContractAdapter {
boundary: self.boundary.clone(),
scope: Some(accepted_diagnostic_scope(accepted)),
cancellation: None,
request_id: accepted.identity.request_id.clone(),
call_id_prefix: accepted.identity.call_id.clone(),
trace_id: accepted.trace_id.clone(),
rpc_id_prefix: accepted.rpc_id.clone(),
zone: accepted.zone.clone(),
idc: accepted.idc.clone(),
env: accepted.env.clone(),
deadline_unix_ms: accepted.identity.deadline_unix_ms,
next_call: Arc::new(AtomicU64::new(1)),
},
}
}
}
#[must_use = "a declared external-function call must be explicitly handled"]
pub struct DeclaredExternalFunctionCall<Application, Function, Request, Response> {
future: Pin<Box<dyn Future<Output = ExternalFunctionResult<Response>> + Send>>,
_type: PhantomData<fn(Application, Function) -> Response>,
_request: PhantomData<fn(Request)>,
}
struct AttemptFuture<T> {
future: Option<Pin<Box<dyn Future<Output = ExternalFunctionResult<T>> + Send>>>,
attempt: Arc<tokio::sync::Mutex<Option<saddle_boundary::request_diagnostics::RequestAttempt>>>,
output: Option<saddle_observability::EmergencyDiagnosticHandle>,
cancellation: Option<AttemptRetentions>,
}
pub(crate) type AttemptRetentions =
Arc<std::sync::Mutex<Vec<saddle_boundary::request_diagnostics::RequiredBoundaryFailure>>>;
impl<T> Drop for AttemptFuture<T> {
fn drop(&mut self) {
drop(self.future.take());
let mut slot = self
.attempt
.try_lock()
.expect("invoke dropped before attempt consumption");
if let Some(attempt) = slot.take() {
if let saddle_boundary::request_diagnostics::AttemptCompletion::Interrupted(failure) =
attempt.finish()
{
let (retained, _) = failure.finish_boundary_retained(
self.output.as_ref(),
&saddle_core::DiagnosticOutcomeAxes {
operation: if std::thread::panicking() {
saddle_core::OperationOutcome::Panicked
} else {
saddle_core::OperationOutcome::Cancelled
},
..Default::default()
},
);
if let Some(owner) = &self.cancellation {
owner
.lock()
.unwrap_or_else(|p| p.into_inner())
.push(retained);
}
}
}
}
}
impl<T> Future for AttemptFuture<T> {
type Output = ExternalFunctionResult<T>;
fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll<Self::Output> {
self.future
.as_mut()
.expect("attempt polled after completion")
.as_mut()
.poll(cx)
}
}
impl<Application, Function, Request, Response>
DeclaredExternalFunctionCall<Application, Function, Request, Response>
where
Request: Message + Send + 'static,
Response: Message + Default + Send + 'static,
{
#[doc(hidden)]
pub fn with_database_request(
mut self,
database: Option<&crate::database_capability::DatabaseRequest>,
) -> Self {
if let Some(database) = database {
self.future = database.supervise_external(self.future);
}
self
}
#[doc(hidden)]
pub fn from_declared(
request: Request,
seal: &ApplicationContractSeal,
business_unit: &'static str,
function: &'static str,
) -> Self {
let adapter = seal.adapter.clone();
let call_number = adapter.next_call.fetch_add(1, Ordering::Relaxed);
let request_id = adapter.request_id.clone();
let call_id = format!("{}-{call_number}", adapter.call_id_prefix);
let child_rpc_id = format!("{}.{}", adapter.rpc_id_prefix, call_number);
Self::from_declared_identity(
request,
adapter,
request_id,
call_id,
child_rpc_id,
business_unit,
function,
)
}
#[doc(hidden)]
pub fn from_declared_with_identity(
request: Request,
seal: &ApplicationContractSeal,
request_id: impl Into<String>,
call_id: impl Into<String>,
business_unit: &'static str,
function: &'static str,
) -> Self {
Self::from_declared_identity(
request,
seal.adapter.clone(),
request_id.into(),
call_id.into(),
format!("{}.1", seal.adapter.rpc_id_prefix),
business_unit,
function,
)
}
fn from_declared_identity(
request: Request,
adapter: ContractAdapter,
request_id: String,
call_id: String,
child_rpc_id: String,
business_unit: &'static str,
function: &'static str,
) -> Self {
let output = match &adapter.boundary {
ContractBoundary::Routed(_, output) => output.clone(),
_ => None,
};
let scope = adapter
.scope
.as_ref()
.map(|scope| scope.reborrow())
.unwrap_or_else(|| {
saddle_observability::RequestDiagnosticScope::early(
None,
saddle_observability::EarlyRequestContext::unavailable(),
)
});
let attempt = Arc::new(tokio::sync::Mutex::new(Some(
saddle_boundary::request_diagnostics::RequestAttempt::new(scope, output.as_ref()),
)));
let invoke_attempt = Arc::clone(&attempt);
let invoke_output = output.clone();
let cancellation = adapter.cancellation.clone();
let future = Box::pin(async move {
let mut slot = invoke_attempt.lock().await;
let attempt = slot.as_mut().expect("one invoke per attempt");
let request = match saddle_boundary::InvokeRequest::unary(
request_id,
call_id,
match InvocationTarget::new(business_unit, function) {
Ok(target) => target,
Err(error) => {
return required_boundary_error(
attempt.capture_boundary_failure(error),
invoke_output.as_ref(),
);
}
},
adapter.deadline_unix_ms,
saddle_boundary::CallerContext {
trace_info: Some(saddle_boundary::proto::TraceInfo {
trace_id: adapter.trace_id,
rpc_id: child_rpc_id,
}),
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) => {
return required_boundary_error(
attempt.capture_boundary_failure(error),
invoke_output.as_ref(),
);
}
};
match adapter.boundary.invoke(request, attempt).await {
Err(error) => required_boundary_error(error, invoke_output.as_ref()),
Ok(response) => match response.outcome {
Some(saddle_boundary::Outcome::Completed(completed)) => {
match Response::decode(completed.result.as_slice()) {
Ok(result) => ExternalFunctionResult::Completed(result),
Err(_) => ExternalFunctionResult::TechnicalFailure(
TechnicalFailure::capture(
attempt.diagnostic_scope(),
TechnicalFailureCode::ContractResultInvalid,
ExecutionCertainty::Executed,
"framework.response_decode_failed",
)
.finish(invoke_output.as_ref()),
),
}
}
Some(saddle_boundary::Outcome::TechnicalFailure(failure)) => {
ExternalFunctionResult::TechnicalFailure(
TechnicalFailure::capture(
attempt.diagnostic_scope(),
map_code(failure.code),
map_certainty(failure.certainty),
"framework.adapter_technical_failure",
)
.finish(invoke_output.as_ref()),
)
}
None => ExternalFunctionResult::TechnicalFailure(
TechnicalFailure::capture(
attempt.diagnostic_scope(),
TechnicalFailureCode::ContractResultInvalid,
ExecutionCertainty::MayHaveExecuted,
"framework.response_outcome_missing",
)
.finish(invoke_output.as_ref()),
),
},
}
});
Self {
future: Box::pin(AttemptFuture {
future: Some(future),
attempt,
output,
cancellation,
}),
_type: PhantomData,
_request: PhantomData,
}
}
}
impl<Application, Function, Request, Response> Future
for DeclaredExternalFunctionCall<Application, Function, Request, Response>
{
type Output = ExternalFunctionResult<Response>;
fn poll(mut self: Pin<&mut Self>, context: &mut Context<'_>) -> Poll<Self::Output> {
self.future.as_mut().poll(context)
}
}
fn boundary_error<T>(error: saddle_boundary::BoundaryError) -> ExternalFunctionResult<T> {
if let Some(diagnostic) = error.diagnostic() {
let _submission =
crate::process::diagnostics::record(diagnostic, error.diagnostic_context());
}
ExternalFunctionResult::TechnicalFailure(TechnicalFailure::capture(
saddle_observability::RequestDiagnosticScope::early(
None,
saddle_observability::EarlyRequestContext::unavailable(),
),
map_boundary_code(error.code),
map_boundary_certainty(error.certainty),
"framework.connection_failed",
))
}
fn required_boundary_error<T>(
failure: saddle_boundary::request_diagnostics::RequiredBoundaryFailure,
output: Option<&saddle_observability::EmergencyDiagnosticHandle>,
) -> ExternalFunctionResult<T> {
ExternalFunctionResult::TechnicalFailure(
TechnicalFailure {
retained: failure.map_error(|error| TechnicalClassification {
code: map_boundary_code(error.code),
certainty: map_boundary_certainty(error.certainty),
}),
}
.finish(output),
)
}
fn map_code(code: i32) -> TechnicalFailureCode {
saddle_boundary::TechnicalCode::try_from(code)
.map(map_boundary_code)
.unwrap_or(TechnicalFailureCode::ContractResultInvalid)
}
fn map_certainty(certainty: i32) -> ExecutionCertainty {
saddle_boundary::ExecutionCertainty::try_from(certainty)
.map(map_boundary_certainty)
.unwrap_or(ExecutionCertainty::MayHaveExecuted)
}
fn map_boundary_code(code: saddle_boundary::TechnicalCode) -> TechnicalFailureCode {
match code {
saddle_boundary::TechnicalCode::FunctionNotFound => TechnicalFailureCode::FunctionNotFound,
saddle_boundary::TechnicalCode::FunctionRequestInvalid => {
TechnicalFailureCode::FunctionRequestInvalid
}
saddle_boundary::TechnicalCode::CapacityRejected => TechnicalFailureCode::CapacityRejected,
saddle_boundary::TechnicalCode::DeadlineExceeded => TechnicalFailureCode::DeadlineExceeded,
saddle_boundary::TechnicalCode::DependencyUnavailable => {
TechnicalFailureCode::DependencyUnavailable
}
saddle_boundary::TechnicalCode::ContractResultInvalid => {
TechnicalFailureCode::ContractResultInvalid
}
saddle_boundary::TechnicalCode::InternalFailure => TechnicalFailureCode::InternalFailure,
saddle_boundary::TechnicalCode::TransportFailure => TechnicalFailureCode::TransportFailure,
saddle_boundary::TechnicalCode::Unspecified => TechnicalFailureCode::ContractResultInvalid,
}
}
fn map_boundary_certainty(certainty: saddle_boundary::ExecutionCertainty) -> ExecutionCertainty {
match certainty {
saddle_boundary::ExecutionCertainty::NotExecuted => ExecutionCertainty::NotExecuted,
saddle_boundary::ExecutionCertainty::Executed => ExecutionCertainty::Executed,
saddle_boundary::ExecutionCertainty::MayHaveExecuted
| saddle_boundary::ExecutionCertainty::Unspecified => ExecutionCertainty::MayHaveExecuted,
}
}
#[derive(Clone)]
pub struct FakeProfuseContractBoundary {
boundary: TransportFakeBoundary,
}
pub struct FakeApplicationBinding {
context: ProfuseGwContext,
seal: ApplicationContractSeal,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum ProfuseGwDispatchError {
InterfaceNotFound,
RequestDataInvalid,
ContextInvalid,
IdentityMismatch,
}
#[doc(hidden)]
pub fn ingress_matches_contract(
accepted: &saddle_boundary::ingress::AcceptedIngress,
seal: &ApplicationContractSeal,
) -> bool {
seal.adapter.request_id == accepted.identity.request_id
&& seal.adapter.call_id_prefix == accepted.identity.call_id
&& seal.adapter.deadline_unix_ms == accepted.identity.deadline_unix_ms
&& seal.adapter.trace_id == accepted.trace_id
&& seal.adapter.rpc_id_prefix == accepted.rpc_id
&& seal.adapter.zone == accepted.zone
&& seal.adapter.idc == accepted.idc
&& seal.adapter.env == accepted.env
}
#[doc(hidden)]
pub fn decode_accepted_profusegw<Request: DeserializeOwned>(
accepted: saddle_boundary::ingress::AcceptedIngress,
) -> Result<(Request, ProfuseGwContext), ProfuseGwDispatchError> {
let context = ProfuseGwContext::from_framework(
&accepted.user_id,
&accepted.trace_id,
&accepted.rpc_id,
&accepted.zone,
&accepted.idc,
&accepted.env,
)
.map_err(|_| ProfuseGwDispatchError::ContextInvalid)?;
let request = serde_json::from_value(accepted.request_data)
.map_err(|_| ProfuseGwDispatchError::RequestDataInvalid)?;
Ok((request, context))
}
impl FakeApplicationBinding {
pub fn profusegw_context(&self) -> ProfuseGwContext {
self.context.clone()
}
pub fn into_contract_seal(self) -> ApplicationContractSeal {
self.seal
}
}
impl FakeProfuseContractBoundary {
pub fn completed<Response: Message>(result: Response) -> Self {
Self {
boundary: TransportFakeBoundary::scripted([FakeStep::completed(&result)]),
}
}
pub fn scripted(steps: impl IntoIterator<Item = FakeStep>) -> Self {
Self {
boundary: TransportFakeBoundary::scripted(steps),
}
}
pub fn attempt_count(&self) -> usize {
self.boundary.attempt_count()
}
pub fn attempts(&self) -> Vec<FakeAttempt> {
self.boundary.attempts()
}
#[doc(hidden)]
pub fn bind_accepted(
self,
accepted: &saddle_boundary::ingress::AcceptedIngress,
) -> Result<FakeApplicationBinding, ProfuseGwDispatchError> {
let context = ProfuseGwContext::from_framework(
&accepted.user_id,
&accepted.trace_id,
&accepted.rpc_id,
&accepted.zone,
&accepted.idc,
&accepted.env,
)
.map_err(|_| ProfuseGwDispatchError::ContextInvalid)?;
let seal = ApplicationContractSeal {
adapter: ContractAdapter {
boundary: ContractBoundary::Fake(self.boundary),
scope: Some(accepted_diagnostic_scope(accepted)),
cancellation: None,
request_id: accepted.identity.request_id.clone(),
call_id_prefix: accepted.identity.call_id.clone(),
trace_id: accepted.trace_id.clone(),
rpc_id_prefix: accepted.rpc_id.clone(),
zone: accepted.zone.clone(),
idc: accepted.idc.clone(),
env: accepted.env.clone(),
deadline_unix_ms: accepted.identity.deadline_unix_ms,
next_call: Arc::new(AtomicU64::new(1)),
},
};
Ok(FakeApplicationBinding { context, seal })
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn required_attempt_cancel_releases_borrow_before_retaining_failure() {
use saddle_boundary::request_diagnostics::RequestAttempt;
use saddle_observability::{EarlyRequestContext, RequestDiagnosticScope};
let attempt = Arc::new(tokio::sync::Mutex::new(Some(RequestAttempt::new(
RequestDiagnosticScope::early(
None,
EarlyRequestContext::socket_accepted("cancel-test"),
),
None,
))));
let owner = Arc::new(std::sync::Mutex::new(Vec::new()));
let inner_attempt = Arc::clone(&attempt);
let mut future = Box::pin(AttemptFuture::<()> {
future: Some(Box::pin(async move {
let _borrow = inner_attempt.lock().await;
std::future::pending::<ExternalFunctionResult<()>>().await
})),
attempt: Arc::clone(&attempt),
output: None,
cancellation: Some(Arc::clone(&owner)),
});
assert!(matches!(
std::future::poll_fn(|cx| Poll::Ready(future.as_mut().poll(cx))).await,
Poll::Pending
));
assert!(attempt.try_lock().is_err());
drop(future);
assert!(attempt.try_lock().unwrap().is_none());
let retained = owner.lock().unwrap();
assert_eq!(retained.len(), 1);
assert_eq!(
retained[0].error().certainty,
saddle_boundary::ExecutionCertainty::NotExecuted
);
assert_eq!(
retained[0].submission(),
saddle_observability::DiagnosticSubmission::OutputUnavailable
);
}
#[test]
fn startup_diagnostic_boundary_mapping_submits_before_terminal() {
use saddle_observability::{
DiagnosticShutdown, EmergencyDiagnostics, FileLoggingConfig, Rotation,
};
let root = std::env::temp_dir().join(format!(
"saddle-boundary-mapping-diagnostic-{}",
std::process::id()
));
std::fs::create_dir(&root).unwrap();
let mut output =
EmergencyDiagnostics::start(&FileLoggingConfig::new(root.clone(), Rotation::Daily))
.unwrap();
let registration = crate::process::diagnostics::OutputRegistration::install(&output);
let socket = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
let address = socket.local_addr().unwrap();
drop(socket);
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
let error = match runtime.block_on(TonicBoundary::connect(
ProfuseContractEndpoint::new(format!("http://{address}")).unwrap(),
)) {
Ok(_) => panic!("closed endpoint unexpectedly connected"),
Err(error) => error,
};
assert!(error.diagnostic().is_some());
let expected_code = map_boundary_code(error.code);
let expected_certainty = map_boundary_certainty(error.certainty);
let terminal = boundary_error::<()>(error);
match terminal {
ExternalFunctionResult::TechnicalFailure(failure) => {
assert_eq!(failure.code(), expected_code);
assert_eq!(failure.certainty(), expected_certainty);
}
_ => panic!("error mapping changed terminal"),
}
drop(registration);
let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5);
while output.shutdown() == DiagnosticShutdown::Pending
&& std::time::Instant::now() < deadline
{
std::thread::sleep(std::time::Duration::from_millis(1));
}
assert_eq!(output.shutdown(), DiagnosticShutdown::Finished);
assert_eq!(output.snapshot().written, 1);
assert!(output.snapshot().initialized && output.snapshot().first_failure.is_none());
let value: serde_json::Value =
serde_json::from_slice(&std::fs::read(output.target()).unwrap()).unwrap();
assert_eq!(
value["diagnostic"]["causes"][0]["code"],
"transport.connect_failed"
);
assert!(value["request"].is_null());
std::fs::remove_file(output.target()).unwrap();
std::fs::remove_dir(root).unwrap();
}
#[test]
fn profusegw_context_is_fixed_bounded_and_read_only() {
assert!(matches!(
ProfuseGwContext::from_framework("", "trace", "rpc", "zone", "idc", "env"),
Err(ProfuseGwContextError::MissingUserId)
));
assert!(matches!(
ProfuseGwContext::from_framework(
&"x".repeat(MAX_PROFUSE_GW_USER_ID_BYTES + 1),
"trace",
"rpc",
"zone",
"idc",
"env",
),
Err(ProfuseGwContextError::UserIdTooLong)
));
let context =
ProfuseGwContext::from_framework("2088用户", "trace", "rpc", "zone", "idc", "env")
.unwrap();
assert_eq!(context.user_id(), "2088用户");
}
#[test]
fn technical_failure_keeps_code_and_execution_certainty_separate() {
let failure = TechnicalFailure::from_framework(
TechnicalFailureCode::DependencyUnavailable,
ExecutionCertainty::MayHaveExecuted,
);
assert_eq!(failure.code(), TechnicalFailureCode::DependencyUnavailable);
assert_eq!(failure.certainty(), ExecutionCertainty::MayHaveExecuted);
}
#[tokio::test]
async fn declared_call_executes_transport_fake_and_decodes_typed_result() {
#[derive(Clone, PartialEq, Message)]
struct Request {
#[prost(uint64, tag = "1")]
value: u64,
}
#[derive(Clone, PartialEq, Message)]
struct Response {
#[prost(bool, tag = "1")]
accepted: bool,
}
struct Application;
struct Function;
let boundary = FakeProfuseContractBoundary::scripted([
FakeStep::completed(&Response { accepted: true }),
FakeStep::completed(&Response { accepted: true }),
]);
let seal = ApplicationContractSeal::from_fake(boundary.boundary, 1_800_000_000_000);
let observed = match &seal.adapter.boundary {
ContractBoundary::Fake(boundary) => boundary.clone(),
ContractBoundary::Tonic(_) | ContractBoundary::Routed(..) => {
panic!("unit test binds the fake transport")
}
};
let call =
DeclaredExternalFunctionCall::<Application, Function, Request, Response>::from_declared(
Request { value: 7 },
&seal,
"puc",
"query",
);
match call.await {
ExternalFunctionResult::Completed(response) => assert!(response.accepted),
ExternalFunctionResult::TechnicalFailure(_) => panic!("fake should complete"),
}
let second =
DeclaredExternalFunctionCall::<Application, Function, Request, Response>::from_declared(
Request { value: 8 },
&seal,
"puc",
"query",
);
assert!(matches!(second.await, ExternalFunctionResult::Completed(_)));
let requests = observed.attempts();
assert_eq!(requests.len(), 2);
assert_eq!(requests[0].function, "query");
assert_eq!(requests[0].request_id, "alpha1-application");
assert_eq!(requests[0].call_id, "call-1");
assert_eq!(requests[0].trace_id, "test-trace");
assert_eq!(requests[0].rpc_id, "test-rpc.1");
assert_eq!(requests[1].rpc_id, "test-rpc.2");
}
#[test]
fn transport_codes_map_to_the_closed_eight_code_set() {
let cases = [
(
saddle_boundary::TechnicalCode::FunctionNotFound,
TechnicalFailureCode::FunctionNotFound,
),
(
saddle_boundary::TechnicalCode::FunctionRequestInvalid,
TechnicalFailureCode::FunctionRequestInvalid,
),
(
saddle_boundary::TechnicalCode::CapacityRejected,
TechnicalFailureCode::CapacityRejected,
),
(
saddle_boundary::TechnicalCode::DeadlineExceeded,
TechnicalFailureCode::DeadlineExceeded,
),
(
saddle_boundary::TechnicalCode::DependencyUnavailable,
TechnicalFailureCode::DependencyUnavailable,
),
(
saddle_boundary::TechnicalCode::ContractResultInvalid,
TechnicalFailureCode::ContractResultInvalid,
),
(
saddle_boundary::TechnicalCode::InternalFailure,
TechnicalFailureCode::InternalFailure,
),
(
saddle_boundary::TechnicalCode::TransportFailure,
TechnicalFailureCode::TransportFailure,
),
];
for (boundary, facade) in cases {
assert_eq!(map_boundary_code(boundary), facade);
}
}
}