use std::collections::VecDeque;
use std::fmt;
use std::future::Future;
use std::sync::{Arc, Mutex};
use std::time::{Duration, SystemTime, UNIX_EPOCH};
use tonic::transport::{Channel, Endpoint};
pub mod ingress;
pub mod proto {
tonic::include_proto!("saddle.profusecontract.v1");
}
pub const FILE_DESCRIPTOR_SET: &[u8] = include_bytes!(concat!(
env!("CARGO_MANIFEST_DIR"),
"/proto/saddle-boundary-descriptor.bin"
));
pub const FILE_DESCRIPTOR_SET_SHA256: &str =
"7ab1e94ea05523935159dd55b36a62e37796d26c782b01ce69a40a6887a7e6ab";
const MAX_TRACE_ID_BYTES: usize = 256;
pub const MAX_UNARY_PAYLOAD_BYTES: usize = 1_048_576;
pub use proto::invoke_response::Outcome;
pub use proto::{
CallerContext, Completed, ExecutionCertainty, InvokeRequest, InvokeResponse, TechnicalCode,
TechnicalFailure,
};
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct InvocationTarget {
pub business_unit: String,
pub function: String,
}
impl InvocationTarget {
pub fn new(
business_unit: impl Into<String>,
function: impl Into<String>,
) -> Result<Self, BoundaryError> {
let target = Self {
business_unit: business_unit.into(),
function: function.into(),
};
if target.business_unit.trim().is_empty() {
return Err(BoundaryError::invalid_request("business_unit is empty"));
}
if target.function.trim().is_empty() {
return Err(BoundaryError::invalid_request("function is empty"));
}
Ok(target)
}
}
impl InvokeRequest {
pub fn unary(
request_id: impl Into<String>,
call_id: impl Into<String>,
target: InvocationTarget,
deadline_unix_ms: i64,
caller_context: CallerContext,
payload: impl Into<Vec<u8>>,
) -> Result<Self, BoundaryError> {
let request = Self {
request_id: request_id.into(),
call_id: call_id.into(),
business_unit: target.business_unit,
function: target.function,
deadline_unix_ms,
caller_context: Some(caller_context),
payload: payload.into(),
};
request.validate()?;
Ok(request)
}
pub fn validate(&self) -> Result<(), BoundaryError> {
for (name, value) in [
("request_id", self.request_id.as_str()),
("call_id", self.call_id.as_str()),
("business_unit", self.business_unit.as_str()),
("function", self.function.as_str()),
] {
if value.trim().is_empty() {
return Err(BoundaryError::invalid_request(format!("{name} is empty")));
}
}
if self.deadline_unix_ms <= 0 {
return Err(BoundaryError::invalid_request(
"deadline_unix_ms must be positive",
));
}
if self.caller_context.is_none() {
return Err(BoundaryError::invalid_request("caller_context is missing"));
}
if self.payload.len() > MAX_UNARY_PAYLOAD_BYTES {
return Err(BoundaryError::invalid_request(
"payload exceeds the frozen unary request cap",
));
}
let caller = self
.caller_context
.as_ref()
.expect("checked caller context");
let trace = caller
.trace_info
.as_ref()
.ok_or_else(|| BoundaryError::invalid_request("trace_info is missing"))?;
let ldc = caller
.ldc_info
.as_ref()
.ok_or_else(|| BoundaryError::invalid_request("ldc_info is missing"))?;
for (name, value) in [
("rpc_id", trace.rpc_id.as_str()),
("zone", ldc.zone.as_str()),
("idc", ldc.idc.as_str()),
("env", ldc.env.as_str()),
] {
if value.trim().is_empty() {
return Err(BoundaryError::invalid_request(format!("{name} is empty")));
}
}
if trace.trace_id.is_empty()
|| trace.trace_id.len() > MAX_TRACE_ID_BYTES
|| trace.trace_id.chars().any(char::is_control)
{
return Err(BoundaryError::invalid_request("trace_id is invalid"));
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct BoundaryError {
pub code: TechnicalCode,
pub certainty: ExecutionCertainty,
pub message: String,
}
impl BoundaryError {
fn invalid_request(message: impl Into<String>) -> Self {
Self {
code: TechnicalCode::FunctionRequestInvalid,
certainty: ExecutionCertainty::NotExecuted,
message: message.into(),
}
}
fn transport(message: impl Into<String>, certainty: ExecutionCertainty) -> Self {
Self {
code: TechnicalCode::TransportFailure,
certainty,
message: message.into(),
}
}
fn invalid_response(message: impl Into<String>) -> Self {
Self {
code: TechnicalCode::ContractResultInvalid,
certainty: ExecutionCertainty::MayHaveExecuted,
message: message.into(),
}
}
}
impl fmt::Display for BoundaryError {
fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(
formatter,
"{:?}/{:?}: {}",
self.code, self.certainty, self.message
)
}
}
impl std::error::Error for BoundaryError {}
pub trait BoundaryTransport: Send + Sync {
fn invoke(
&self,
request: InvokeRequest,
) -> impl Future<Output = Result<InvokeResponse, BoundaryError>> + Send;
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ProfuseContractEndpoint(String);
impl ProfuseContractEndpoint {
pub fn new(endpoint: impl Into<String>) -> Result<Self, BoundaryError> {
let endpoint = endpoint.into();
let authority = endpoint.strip_prefix("http://").ok_or_else(|| {
BoundaryError::invalid_request("profusecontract endpoint must use http://")
})?;
if authority.is_empty()
|| authority.contains('/')
|| authority.contains('?')
|| authority.contains('#')
|| authority.contains('@')
{
return Err(BoundaryError::invalid_request(
"profusecontract endpoint must contain only an authority",
));
}
Endpoint::from_shared(endpoint.clone()).map_err(|error| {
BoundaryError::invalid_request(format!("invalid profusecontract endpoint: {error}"))
})?;
Ok(Self(endpoint))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct ProfuseContractAuthorityTemplate(String);
impl ProfuseContractAuthorityTemplate {
pub fn new(template: impl Into<String>) -> Result<Self, BoundaryError> {
let template = template.into();
if template.matches("{zone}").count() > 1
|| template.replace("{zone}", "").contains(['{', '}'])
{
return Err(BoundaryError::invalid_request(
"profusecontract authority may contain at most one {zone}",
));
}
ProfuseContractEndpoint::new(template.replace("{zone}", "valid-zone"))?;
Ok(Self(template))
}
pub fn resolve(&self, zone: &str) -> Result<ProfuseContractEndpoint, BoundaryError> {
if zone.is_empty()
|| zone.len() > ingress::MAX_ID_BYTES
|| !zone
.bytes()
.all(|byte| byte.is_ascii_alphanumeric() || matches!(byte, b'-' | b'_' | b'.'))
{
return Err(BoundaryError::invalid_request("invalid ldc zone"));
}
ProfuseContractEndpoint::new(self.0.replace("{zone}", zone))
}
pub fn as_str(&self) -> &str {
&self.0
}
}
#[derive(Clone)]
pub struct TonicBoundary {
channel: Channel,
}
impl TonicBoundary {
pub async fn connect(endpoint: ProfuseContractEndpoint) -> Result<Self, BoundaryError> {
let endpoint = Endpoint::from_shared(endpoint.0).map_err(|error| {
BoundaryError::transport(error.to_string(), ExecutionCertainty::NotExecuted)
})?;
let channel = endpoint.connect().await.map_err(|error| {
BoundaryError::transport(error.to_string(), ExecutionCertainty::NotExecuted)
})?;
Ok(Self { channel })
}
async fn connect_until(
endpoint: ProfuseContractEndpoint,
deadline_unix_ms: i64,
) -> Result<Self, BoundaryError> {
let remaining = deadline_timeout(deadline_unix_ms)?;
let endpoint = Endpoint::from_shared(endpoint.0)
.map_err(|error| {
BoundaryError::transport(error.to_string(), ExecutionCertainty::NotExecuted)
})?
.connect_timeout(remaining)
.timeout(remaining)
.concurrency_limit(1)
.buffer_size(1);
let channel = endpoint.connect().await.map_err(|error| {
BoundaryError::transport(error.to_string(), ExecutionCertainty::NotExecuted)
})?;
Ok(Self { channel })
}
pub async fn invoke_routed(
template: &ProfuseContractAuthorityTemplate,
request: InvokeRequest,
) -> Result<InvokeResponse, BoundaryError> {
request.validate()?;
let caller = request
.caller_context
.as_ref()
.expect("validated caller context");
let trace = caller.trace_info.as_ref().expect("validated trace info");
let ldc = caller.ldc_info.as_ref().expect("validated ldc info");
let endpoint = template.resolve(&ldc.zone)?;
let event_context = saddle_observability::RequestIdentity::new(request.request_id.clone())
.ok()
.and_then(|request_identity| {
saddle_observability::RouteIdentity::new(request.function.clone())
.ok()
.and_then(|route| {
saddle_observability::EventContext::new(request_identity, route, 1).ok()
})
});
let observed_zone = ldc.zone.clone();
let deadline_unix_ms = request.deadline_unix_ms;
let root = saddle_observability::global().and_then(|observer| {
observer
.start_external_call_checked(
"saddle",
ldc.zone.clone(),
"profusecontract",
"invoke",
Some(&trace.trace_id),
)
.ok()
.map(|started| started.0)
});
let connect = root.as_ref().and_then(|root| {
saddle_observability::global().map(|observer| {
observer.start_child_call(
root.context(),
saddle_observability::CallKind::ExternalConnect,
ldc.zone.clone(),
"profusecontract",
"connect",
)
})
});
let connect_stage = root.as_ref().and_then(|root| {
event_context.as_ref().map(|event| {
root_observer().start_stage(
root.context(),
saddle_observability::Stage::ProfuseContract,
event.clone(),
)
})
});
let boundary = match Self::connect_until(endpoint, request.deadline_unix_ms).await {
Ok(boundary) => {
if let Some(stage) = connect_stage {
stage.succeed();
}
if let Some(connect) = connect {
connect.succeed();
}
boundary
}
Err(error) => {
if let Some(stage) = connect_stage {
stage.fail(&observation_error(
"boundary.profusecontract_connect_failed",
));
}
if let Some(connect) = connect {
connect.fail(&saddle_core::SaddleError::new(
saddle_core::ErrorKind::Unavailable,
"boundary.profusecontract_connect_failed",
"profusecontract connect failed",
));
}
let outcome = if deadline_timeout(deadline_unix_ms).is_err() {
saddle_observability::OutboundResult::Timeout
} else {
saddle_observability::OutboundResult::Failure
};
record_outbound(
root.as_ref(),
event_context.as_ref(),
&observed_zone,
outcome,
);
if let Some(root) = root {
root.fail(&saddle_core::SaddleError::new(
saddle_core::ErrorKind::Unavailable,
"boundary.profusecontract_connect_failed",
"profusecontract connect failed",
));
}
return Err(error);
}
};
let unary = root.as_ref().and_then(|root| {
saddle_observability::global().map(|observer| {
observer.start_child_call(
root.context(),
saddle_observability::CallKind::ExternalUnary,
ldc.zone.clone(),
"profusecontract",
"unary",
)
})
});
let unary_stage = root.as_ref().and_then(|root| {
event_context.as_ref().map(|event| {
root_observer().start_stage(
root.context(),
saddle_observability::Stage::ProfuseContract,
event.clone(),
)
})
});
let result = boundary.invoke(request).await;
match &result {
Ok(_) => {
if let Some(stage) = unary_stage {
stage.succeed();
}
record_outbound(
root.as_ref(),
event_context.as_ref(),
&observed_zone,
saddle_observability::OutboundResult::Success,
);
if let Some(unary) = unary {
unary.succeed();
}
if let Some(root) = root {
root.succeed();
}
}
Err(_) => {
if let Some(stage) = unary_stage {
stage.fail(&observation_error("boundary.profusecontract_unary_failed"));
}
let outcome = if deadline_timeout(deadline_unix_ms).is_err() {
saddle_observability::OutboundResult::Timeout
} else {
saddle_observability::OutboundResult::Failure
};
record_outbound(
root.as_ref(),
event_context.as_ref(),
&observed_zone,
outcome,
);
let failure = saddle_core::SaddleError::new(
saddle_core::ErrorKind::Unavailable,
"boundary.profusecontract_unary_failed",
"profusecontract unary failed",
);
if let Some(unary) = unary {
unary.fail(&failure);
}
if let Some(root) = root {
root.fail(&failure);
}
}
}
result
}
}
fn root_observer() -> saddle_observability::Observer {
saddle_observability::global()
.expect("an observed call has a global observer")
.clone()
}
fn observation_error(code: &'static str) -> saddle_core::SaddleError {
saddle_core::SaddleError::new(
saddle_core::ErrorKind::Unavailable,
code,
"typed transport observation",
)
}
fn record_outbound(
root: Option<&saddle_observability::ActiveCall>,
event: Option<&saddle_observability::EventContext>,
zone: &str,
result: saddle_observability::OutboundResult,
) {
let (Some(observer), Some(root), Some(event)) = (saddle_observability::global(), root, event)
else {
return;
};
let (Ok(zone), Ok(authority)) = (
saddle_observability::RouteIdentity::new(zone),
saddle_observability::OutboundAuthority::new("profusecontract"),
) else {
return;
};
observer.record_outbound(
root.context(),
event,
saddle_observability::OutboundObservation::new(zone, authority, result),
);
}
impl BoundaryTransport for TonicBoundary {
async fn invoke(&self, request: InvokeRequest) -> Result<InvokeResponse, BoundaryError> {
request.validate()?;
let timeout = deadline_timeout(request.deadline_unix_ms)?;
let identity = (request.request_id.clone(), request.call_id.clone());
let mut client =
proto::profuse_contract_boundary_client::ProfuseContractBoundaryClient::new(
self.channel.clone(),
);
let mut tonic_request = tonic::Request::new(request);
tonic_request.set_timeout(timeout);
let response = client
.invoke(tonic_request)
.await
.map(tonic::Response::into_inner)
.map_err(|status| {
BoundaryError::transport(status.to_string(), ExecutionCertainty::MayHaveExecuted)
})?;
validate_response(&identity, &response)?;
Ok(response)
}
}
fn deadline_timeout(deadline_unix_ms: i64) -> Result<Duration, BoundaryError> {
let deadline = u64::try_from(deadline_unix_ms)
.ok()
.and_then(|millis| UNIX_EPOCH.checked_add(Duration::from_millis(millis)))
.ok_or_else(|| BoundaryError::invalid_request("deadline_unix_ms is invalid"))?;
deadline
.duration_since(SystemTime::now())
.map_err(|_| BoundaryError {
code: TechnicalCode::DeadlineExceeded,
certainty: ExecutionCertainty::NotExecuted,
message: "absolute deadline has elapsed".into(),
})
}
#[derive(Clone, Default)]
pub struct FakeProfuseContractBoundary {
state: Arc<Mutex<FakeState>>,
}
#[derive(Default)]
struct FakeState {
scripted: VecDeque<FakeStep>,
observed: Vec<FakeAttempt>,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum FakeTechnicalCode {
FunctionNotFound,
FunctionRequestInvalid,
CapacityRejected,
DeadlineExceeded,
DependencyUnavailable,
ContractResultInvalid,
InternalFailure,
TransportFailure,
}
impl From<FakeTechnicalCode> for TechnicalCode {
fn from(code: FakeTechnicalCode) -> Self {
match code {
FakeTechnicalCode::FunctionNotFound => Self::FunctionNotFound,
FakeTechnicalCode::FunctionRequestInvalid => Self::FunctionRequestInvalid,
FakeTechnicalCode::CapacityRejected => Self::CapacityRejected,
FakeTechnicalCode::DeadlineExceeded => Self::DeadlineExceeded,
FakeTechnicalCode::DependencyUnavailable => Self::DependencyUnavailable,
FakeTechnicalCode::ContractResultInvalid => Self::ContractResultInvalid,
FakeTechnicalCode::InternalFailure => Self::InternalFailure,
FakeTechnicalCode::TransportFailure => Self::TransportFailure,
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub enum FakeExecutionCertainty {
NotExecuted,
Executed,
MayHaveExecuted,
}
impl From<FakeExecutionCertainty> for ExecutionCertainty {
fn from(certainty: FakeExecutionCertainty) -> Self {
match certainty {
FakeExecutionCertainty::NotExecuted => Self::NotExecuted,
FakeExecutionCertainty::Executed => Self::Executed,
FakeExecutionCertainty::MayHaveExecuted => Self::MayHaveExecuted,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct FakeStep {
outcome: FakeStepOutcome,
}
#[derive(Clone, Debug, Eq, PartialEq)]
enum FakeStepOutcome {
Completed(Vec<u8>),
TechnicalFailure(FakeTechnicalCode, FakeExecutionCertainty),
}
impl FakeStep {
pub fn completed<M: prost::Message>(result: &M) -> Self {
Self {
outcome: FakeStepOutcome::Completed(result.encode_to_vec()),
}
}
pub const fn technical_failure(
code: FakeTechnicalCode,
certainty: FakeExecutionCertainty,
) -> Self {
Self {
outcome: FakeStepOutcome::TechnicalFailure(code, certainty),
}
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub struct FakeAttempt {
pub count: usize,
pub request_id: String,
pub call_id: String,
pub trace_id: String,
pub rpc_id: String,
pub zone: String,
pub idc: String,
pub env: String,
pub function: String,
pub deadline_unix_ms: i64,
}
impl FakeProfuseContractBoundary {
pub fn scripted(steps: impl IntoIterator<Item = FakeStep>) -> Self {
Self {
state: Arc::new(Mutex::new(FakeState {
scripted: steps.into_iter().collect(),
observed: Vec::new(),
})),
}
}
pub fn attempts(&self) -> Vec<FakeAttempt> {
self.state
.lock()
.expect("fake boundary mutex poisoned")
.observed
.clone()
}
pub fn attempt_count(&self) -> usize {
self.state
.lock()
.expect("fake boundary mutex poisoned")
.observed
.len()
}
}
impl BoundaryTransport for FakeProfuseContractBoundary {
async fn invoke(&self, request: InvokeRequest) -> Result<InvokeResponse, BoundaryError> {
request.validate()?;
let identity = (request.request_id.clone(), request.call_id.clone());
let mut state = self.state.lock().expect("fake boundary mutex poisoned");
let count = state.observed.len() + 1;
let caller = request
.caller_context
.as_ref()
.expect("validated caller context");
let trace = caller.trace_info.as_ref().expect("validated trace info");
let ldc = caller.ldc_info.as_ref().expect("validated ldc info");
state.observed.push(FakeAttempt {
count,
request_id: request.request_id,
call_id: request.call_id,
trace_id: trace.trace_id.clone(),
rpc_id: trace.rpc_id.clone(),
zone: ldc.zone.clone(),
idc: ldc.idc.clone(),
env: ldc.env.clone(),
function: request.function,
deadline_unix_ms: request.deadline_unix_ms,
});
let step = state.scripted.pop_front().ok_or_else(|| {
BoundaryError::transport(
"fake response script exhausted",
ExecutionCertainty::NotExecuted,
)
})?;
let outcome = match step.outcome {
FakeStepOutcome::Completed(result) => Outcome::Completed(Completed { result }),
FakeStepOutcome::TechnicalFailure(code, certainty) => {
Outcome::TechnicalFailure(TechnicalFailure {
code: TechnicalCode::from(code) as i32,
certainty: ExecutionCertainty::from(certainty) as i32,
})
}
};
let response = InvokeResponse {
request_id: identity.0.clone(),
call_id: identity.1.clone(),
outcome: Some(outcome),
};
validate_response(&identity, &response)?;
Ok(response)
}
}
fn validate_response(
request_identity: &(String, String),
response: &InvokeResponse,
) -> Result<(), BoundaryError> {
if response.request_id != request_identity.0 || response.call_id != request_identity.1 {
return Err(BoundaryError::invalid_response(
"response identity mismatch",
));
}
match response.outcome.as_ref() {
Some(Outcome::Completed(completed)) => {
if completed.result.len() > MAX_UNARY_PAYLOAD_BYTES {
return Err(BoundaryError::invalid_response(
"result exceeds the frozen unary response cap",
));
}
Ok(())
}
Some(Outcome::TechnicalFailure(failure)) => {
if TechnicalCode::try_from(failure.code)
.ok()
.filter(|code| *code != TechnicalCode::Unspecified)
.is_none()
{
return Err(BoundaryError::invalid_response(
"technical failure code is unspecified or unknown",
));
}
match ExecutionCertainty::try_from(failure.certainty).ok() {
Some(
ExecutionCertainty::NotExecuted
| ExecutionCertainty::Executed
| ExecutionCertainty::MayHaveExecuted,
) => Ok(()),
_ => Err(BoundaryError::invalid_response(
"execution certainty is unspecified or unknown",
)),
}
}
None => Err(BoundaryError::invalid_response(
"response outcome is missing",
)),
}
}
#[cfg(test)]
mod tests {
use super::*;
use prost::Message;
use sha2::{Digest, Sha256};
#[derive(Clone, PartialEq, prost::Message)]
struct TypedResult {
#[prost(string, tag = "1")]
value: String,
}
fn request() -> InvokeRequest {
InvokeRequest::unary(
"request-1",
"call-1",
InvocationTarget::new("交易", "创建订单").unwrap(),
1_800_000_000_000,
CallerContext {
trace_info: Some(proto::TraceInfo {
trace_id: "trace-1".into(),
rpc_id: "0.1".into(),
}),
ldc_info: Some(proto::LdcInfo {
zone: "z1".into(),
idc: "i1".into(),
env: "test".into(),
}),
},
br#"{"amount":1}"#.to_vec(),
)
.unwrap()
}
#[tokio::test]
async fn fake_is_scripted_and_records_unicode_target() {
let fake = FakeProfuseContractBoundary::scripted([FakeStep::completed(&TypedResult {
value: "完成".into(),
})]);
let response = fake.invoke(request()).await.unwrap();
let Some(Outcome::Completed(completed)) = response.outcome else {
panic!("expected completed outcome")
};
assert_eq!(
TypedResult::decode(completed.result.as_slice())
.unwrap()
.value,
"完成"
);
let observed = fake.attempts();
assert_eq!(fake.attempt_count(), 1);
assert_eq!(observed[0].count, 1);
assert_eq!(observed[0].request_id, "request-1");
assert_eq!(observed[0].call_id, "call-1");
assert_eq!(observed[0].function, "创建订单");
assert_eq!(observed[0].deadline_unix_ms, 1_800_000_000_000);
}
#[test]
fn invocation_envelope_rejects_noncanonical_trace_storage() {
let mut request = request();
request
.caller_context
.as_mut()
.unwrap()
.trace_info
.as_mut()
.unwrap()
.trace_id = "bad\ntrace".into();
assert_eq!(
request.validate().unwrap_err().code,
TechnicalCode::FunctionRequestInvalid
);
}
#[test]
fn unary_payload_caps_are_enforced_in_both_directions() {
let mut oversized = request();
oversized.payload = vec![0; MAX_UNARY_PAYLOAD_BYTES + 1];
assert_eq!(
oversized.validate().unwrap_err().code,
TechnicalCode::FunctionRequestInvalid
);
let identity = ("request-1".to_owned(), "call-1".to_owned());
let response = InvokeResponse {
request_id: identity.0.clone(),
call_id: identity.1.clone(),
outcome: Some(Outcome::Completed(Completed {
result: vec![0; MAX_UNARY_PAYLOAD_BYTES + 1],
})),
};
assert_eq!(
validate_response(&identity, &response).unwrap_err().code,
TechnicalCode::ContractResultInvalid
);
}
#[tokio::test]
async fn routed_connect_and_unary_share_one_absolute_deadline() {
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap();
let address = listener.local_addr().unwrap();
let blackhole = tokio::spawn(async move {
let (_socket, _) = listener.accept().await.unwrap();
tokio::time::sleep(Duration::from_secs(2)).await;
});
let template = ProfuseContractAuthorityTemplate::new(format!(
"http://127.0.0.{{zone}}:{}",
address.port()
))
.unwrap();
let mut request = request();
request
.caller_context
.as_mut()
.unwrap()
.ldc_info
.as_mut()
.unwrap()
.zone = "1".into();
request.deadline_unix_ms = i64::try_from(
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap()
.as_millis()
+ 150,
)
.unwrap();
let started = std::time::Instant::now();
let error = TonicBoundary::invoke_routed(&template, request)
.await
.unwrap_err();
assert_eq!(error.code, TechnicalCode::TransportFailure);
assert!(started.elapsed() < Duration::from_secs(1));
blackhole.abort();
}
#[tokio::test]
async fn fake_exhaustion_is_not_executed() {
let error = FakeProfuseContractBoundary::default()
.invoke(request())
.await
.unwrap_err();
assert_eq!(error.code, TechnicalCode::TransportFailure);
assert_eq!(error.certainty, ExecutionCertainty::NotExecuted);
}
#[tokio::test]
async fn fake_scripts_all_failure_codes_and_certainties() {
let codes = [
FakeTechnicalCode::FunctionNotFound,
FakeTechnicalCode::FunctionRequestInvalid,
FakeTechnicalCode::CapacityRejected,
FakeTechnicalCode::DeadlineExceeded,
FakeTechnicalCode::DependencyUnavailable,
FakeTechnicalCode::ContractResultInvalid,
FakeTechnicalCode::InternalFailure,
FakeTechnicalCode::TransportFailure,
];
let certainties = [
FakeExecutionCertainty::NotExecuted,
FakeExecutionCertainty::Executed,
FakeExecutionCertainty::MayHaveExecuted,
];
for code in codes {
for certainty in certainties {
let fake = FakeProfuseContractBoundary::scripted([FakeStep::technical_failure(
code, certainty,
)]);
let response = fake.invoke(request()).await.unwrap();
let Some(Outcome::TechnicalFailure(failure)) = response.outcome else {
panic!("expected technical failure")
};
assert_eq!(failure.code, TechnicalCode::from(code) as i32);
assert_eq!(
failure.certainty,
ExecutionCertainty::from(certainty) as i32
);
assert_eq!(fake.attempt_count(), 1);
}
}
}
#[test]
fn conditional_branch_can_assert_not_called() {
let fake = FakeProfuseContractBoundary::scripted([FakeStep::completed(&TypedResult {
value: "unused".into(),
})]);
assert_eq!(fake.attempt_count(), 0);
assert!(fake.attempts().is_empty());
}
#[test]
fn certainty_has_exactly_three_executable_states() {
assert_eq!(ExecutionCertainty::NotExecuted as i32, 1);
assert_eq!(ExecutionCertainty::Executed as i32, 2);
assert_eq!(ExecutionCertainty::MayHaveExecuted as i32, 3);
}
#[test]
fn endpoint_is_configurable_but_has_no_tls_path_or_credentials() {
let endpoint = ProfuseContractEndpoint::new("http://127.0.0.1:50051").unwrap();
assert_eq!(endpoint.as_str(), "http://127.0.0.1:50051");
for invalid in [
"https://contracts.example",
"http://user@contracts.example",
"http://contracts.example/service",
"dns:///contracts.example",
] {
let error = ProfuseContractEndpoint::new(invalid).unwrap_err();
assert_eq!(error.code, TechnicalCode::FunctionRequestInvalid);
assert_eq!(error.certainty, ExecutionCertainty::NotExecuted);
}
}
#[test]
fn authority_accepts_fixed_or_one_validated_zone_substitution() {
let template =
ProfuseContractAuthorityTemplate::new("http://profusecontract-{zone}.internal:50051")
.unwrap();
assert_eq!(
template.resolve("cn-shanghai-1").unwrap().as_str(),
"http://profusecontract-cn-shanghai-1.internal:50051"
);
let fixed =
ProfuseContractAuthorityTemplate::new("http://profusecontract.internal:50051").unwrap();
assert_eq!(
fixed.resolve("cn-shanghai-1").unwrap().as_str(),
"http://profusecontract.internal:50051"
);
for invalid in [
"http://{zone}.{zone}.internal:50051",
"http://{region}.internal:50051",
] {
assert!(ProfuseContractAuthorityTemplate::new(invalid).is_err());
}
for invalid_zone in ["", "../admin", "zone:50052", "zone/name"] {
assert!(template.resolve(invalid_zone).is_err());
}
}
#[test]
fn elapsed_absolute_deadline_fails_before_transport() {
let error = deadline_timeout(1).unwrap_err();
assert_eq!(error.code, TechnicalCode::DeadlineExceeded);
assert_eq!(error.certainty, ExecutionCertainty::NotExecuted);
}
#[test]
fn empty_dynamic_target_is_rejected() {
let error = InvocationTarget::new("交易", " ").unwrap_err();
assert_eq!(error.code, TechnicalCode::FunctionRequestInvalid);
assert_eq!(error.certainty, ExecutionCertainty::NotExecuted);
}
#[test]
fn frozen_descriptor_has_one_generic_unary_boundary() {
let actual = format!("{:x}", Sha256::digest(FILE_DESCRIPTOR_SET));
assert_eq!(actual, FILE_DESCRIPTOR_SET_SHA256);
let descriptor = prost_types::FileDescriptorSet::decode(FILE_DESCRIPTOR_SET).unwrap();
let boundary = descriptor
.file
.iter()
.find(|file| file.package.as_deref() == Some("saddle.profusecontract.v1"))
.unwrap();
let service = &boundary.service[0];
assert_eq!(service.name.as_deref(), Some("ProfuseContractBoundary"));
assert_eq!(service.method.len(), 1);
assert_eq!(service.method[0].name.as_deref(), Some("Invoke"));
assert!(!service.method[0].client_streaming.unwrap_or(false));
assert!(!service.method[0].server_streaming.unwrap_or(false));
}
}