commit d12f7985b3c2e097e6ae01ae392c5d0e4394a502
parent 035c3b551a4a8f4b407873710b80487207b94cff
Author: triesap <tyson@radroots.org>
Date: Wed, 5 Aug 2026 01:23:44 +0000
refactor(transport): converge delivery policy
- centralize satisfied pending and exhausted policy evaluation
- retain typed sink retry timing and partial target evidence
- migrate storage sync and preview adapters to canonical contracts
- verify transport adapters storage features clippy and packaging
Diffstat:
14 files changed, 503 insertions(+), 137 deletions(-)
diff --git a/crates/storage/src/outbox.rs b/crates/storage/src/outbox.rs
@@ -7,7 +7,10 @@ use core::fmt;
pub use radroots_transport::{
BoxFuture, DeliveryReceipt, DeliveryRequest, TransportId,
outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability},
- policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ policy::{
+ SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy,
+ evaluate_satisfaction as evaluate_transport_satisfaction,
+ },
sink::{DeliveryPayload, DeliveryTargetReceipt},
target::{
TARGET_SET_MAX_ITEMS, Target, TargetFingerprint, TargetLabel, TargetScope, TargetSet,
@@ -921,56 +924,16 @@ fn evaluate_satisfaction(
request: &DeliveryRequest,
evidence: &[TargetDeliveryEvidence],
) -> SatisfactionResult {
- let class = request.satisfaction().class();
- let targets = request.target_set().targets();
- let is_successful = |target: &TargetFingerprint| {
- evidence
- .iter()
- .any(|entry| entry.target() == target && entry.outcome().satisfies(class))
- };
- let is_retryable = |target: &TargetFingerprint| {
+ match evaluate_transport_satisfaction(
+ request.satisfaction(),
+ request.target_set(),
evidence
.iter()
- .rev()
- .find(|entry| entry.target() == target)
- .is_some_and(|entry| entry.outcome().is_retryable())
- };
- let successful = targets
- .iter()
- .filter(|target| is_successful(target.fingerprint()))
- .count();
- let retryable = targets
- .iter()
- .filter(|target| !is_successful(target.fingerprint()) && is_retryable(target.fingerprint()))
- .count();
- let policy = request.satisfaction().targets();
- let (satisfied, possible) = if policy.is_any() {
- (successful != 0, successful + retryable != 0)
- } else if policy.is_all() {
- (
- successful == targets.len(),
- successful + retryable == targets.len(),
- )
- } else if let Some(threshold) = policy.quorum_threshold() {
- let threshold = usize::from(threshold);
- (successful >= threshold, successful + retryable >= threshold)
- } else {
- // `TargetPolicy` is closed over any/all/quorum/required; after the
- // preceding branches, the required-target slice is necessarily set.
- let required = policy.required_targets().unwrap_or_default();
- (
- required.iter().all(&is_successful),
- required
- .iter()
- .all(|target| is_successful(target) || is_retryable(target)),
- )
- };
- if satisfied {
- SatisfactionResult::Satisfied
- } else if possible {
- SatisfactionResult::Pending
- } else {
- SatisfactionResult::Exhausted
+ .map(|entry| (entry.target(), entry.outcome())),
+ ) {
+ Ok(SatisfactionState::Satisfied) => SatisfactionResult::Satisfied,
+ Ok(SatisfactionState::Pending) => SatisfactionResult::Pending,
+ Ok(SatisfactionState::Exhausted) | Err(_) => SatisfactionResult::Exhausted,
}
}
@@ -1787,7 +1750,7 @@ mod tests {
let any_request = request_with(TargetPolicy::any());
assert_eq!(
evaluate_satisfaction(&any_request, &[]),
- SatisfactionResult::Exhausted
+ SatisfactionResult::Pending
);
assert_eq!(
evaluate_satisfaction(&any_request, &retryable),
@@ -1800,7 +1763,7 @@ mod tests {
let quorum_request = request_with(TargetPolicy::quorum(2).unwrap());
assert_eq!(
evaluate_satisfaction(&quorum_request, &one_accepted),
- SatisfactionResult::Exhausted
+ SatisfactionResult::Pending
);
let required_request =
request_with(TargetPolicy::required(vec![targets[0].fingerprint().clone()]).unwrap());
diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs
@@ -64,7 +64,7 @@ impl PushRequest {
satisfaction: SatisfactionPolicy,
cancellation: CancellationPolicy,
) -> Result<Self, Error> {
- if !valid_satisfaction(&satisfaction, &targets) {
+ if satisfaction.validate_for(&targets).is_err() {
return Err(Error::InvalidPushRequest);
}
Ok(Self {
@@ -609,22 +609,6 @@ fn hash_satisfaction(hasher: &mut Sha256, policy: &SatisfactionPolicy) {
}
}
-fn valid_satisfaction(policy: &SatisfactionPolicy, targets: &TargetSet) -> bool {
- let selection = policy.targets();
- if let Some(threshold) = selection.quorum_threshold() {
- return usize::from(threshold) <= targets.len();
- }
- if let Some(required) = selection.required_targets() {
- return required.iter().all(|required| {
- targets
- .targets()
- .iter()
- .any(|target| target.fingerprint() == required)
- });
- }
- selection.is_any() || selection.is_all()
-}
-
fn atomic_digest(domain: &[u8], input: &[u8]) -> AtomicCommitDigest {
let mut hasher = Sha256::new();
hash_field(&mut hasher, domain);
diff --git a/crates/transport/examples/host_transport.rs b/crates/transport/examples/host_transport.rs
@@ -41,8 +41,11 @@ impl EventSink for HostTransport {
})
}
- fn deliver(&self, _request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
- Box::pin(async { Err(Error::UnsupportedOperation) })
+ fn deliver(
+ &self,
+ _request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>> {
+ Box::pin(async move { Err(radroots_transport::SinkFailure::invalid_contract(&_request)) })
}
}
diff --git a/crates/transport/src/lib.rs b/crates/transport/src/lib.rs
@@ -32,7 +32,7 @@ pub use kind::{
RadrootsTransportImplementationState,
};
pub use payload::RadrootsTransportPayload;
-pub use sink::{DeliveryReceipt, DeliveryRequest, EventSink, SinkStatus};
+pub use sink::{DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure, SinkStatus};
pub use source::{BoxFuture, EventSource, FetchPage, FetchRequest, SourceStatus};
pub use status::{
RadrootsTransportCapabilities, RadrootsTransportDeliveryTargetStatus, RadrootsTransportOutcome,
diff --git a/crates/transport/src/outcome.rs b/crates/transport/src/outcome.rs
@@ -253,23 +253,36 @@ impl DeliveryOutcome {
}
}
-fn validate_delivery_detail(code: &str, message: &str) -> Result<(), crate::Error> {
- let valid_code = !code.is_empty()
+pub(crate) fn validate_delivery_code(code: &str) -> Result<(), crate::Error> {
+ let valid = !code.is_empty()
&& code.len() <= DELIVERY_OUTCOME_CODE_MAX_BYTES
&& code.bytes().all(|byte| {
byte.is_ascii_lowercase() || byte.is_ascii_digit() || matches!(byte, b'_' | b'-' | b'.')
});
+ if valid {
+ Ok(())
+ } else {
+ Err(crate::Error::InvalidDeliveryOutcome)
+ }
+}
+
+pub(crate) fn validate_delivery_message(message: &str) -> Result<(), crate::Error> {
let valid_message = !message.is_empty()
&& message.len() <= DELIVERY_OUTCOME_MESSAGE_MAX_BYTES
&& message == message.trim()
&& !message.chars().any(char::is_control);
- if valid_code && valid_message {
+ if valid_message {
Ok(())
} else {
Err(crate::Error::InvalidDeliveryOutcome)
}
}
+fn validate_delivery_detail(code: &str, message: &str) -> Result<(), crate::Error> {
+ validate_delivery_code(code)?;
+ validate_delivery_message(message)
+}
+
#[cfg(feature = "serde")]
impl<'de> serde::Deserialize<'de> for DeliveryOutcome {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
diff --git a/crates/transport/src/policy.rs b/crates/transport/src/policy.rs
@@ -2,9 +2,23 @@
use crate::{
Error,
+ outcome::DeliveryOutcome,
target::{TargetFingerprint, TargetSet},
};
-use alloc::{collections::BTreeSet, vec::Vec};
+use alloc::{collections::BTreeMap, vec::Vec};
+
+/// Current result of evaluating delivery evidence against one exact policy.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Copy, Debug, Eq, PartialEq)]
+pub enum SatisfactionState {
+ /// The evidence already satisfies the policy.
+ Satisfied,
+ /// The policy is not satisfied, but unattempted or retryable work can satisfy it.
+ Pending,
+ /// The available evidence proves that the policy can no longer be satisfied.
+ Exhausted,
+}
/// Success level a caller requires from selected targets.
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
@@ -111,14 +125,9 @@ impl TargetPolicy {
}
}
TargetPolicyKind::Required(required) => {
- let requested: BTreeSet<&str> = targets
- .targets()
- .iter()
- .map(|target| target.fingerprint().as_str())
- .collect();
if required
.iter()
- .any(|fingerprint| !requested.contains(fingerprint.as_str()))
+ .any(|fingerprint| !targets.contains(fingerprint))
{
Err(Error::RequiredTargetNotRequested)
} else {
@@ -128,14 +137,14 @@ impl TargetPolicy {
}
}
- pub(crate) fn is_satisfied(&self, total_targets: usize, satisfied: &BTreeSet<&str>) -> bool {
+ fn is_satisfied(&self, total_targets: usize, satisfied: usize, fingerprints: &[&str]) -> bool {
match &self.kind {
- TargetPolicyKind::Any => !satisfied.is_empty(),
- TargetPolicyKind::All => satisfied.len() == total_targets,
- TargetPolicyKind::Quorum(threshold) => satisfied.len() >= usize::from(*threshold),
+ TargetPolicyKind::Any => satisfied != 0,
+ TargetPolicyKind::All => satisfied == total_targets,
+ TargetPolicyKind::Quorum(threshold) => satisfied >= usize::from(*threshold),
TargetPolicyKind::Required(required) => required
.iter()
- .all(|fingerprint| satisfied.contains(fingerprint.as_str())),
+ .all(|fingerprint| fingerprints.contains(&fingerprint.as_str())),
}
}
}
@@ -175,6 +184,69 @@ impl SatisfactionPolicy {
}
}
+/// Evaluates target evidence using the transport-owned satisfaction law.
+///
+/// Evidence is ordered from oldest to newest when a target occurs more than
+/// once. A prior success remains authoritative; otherwise the newest outcome
+/// determines whether the target can be retried. Targets without evidence are
+/// pending. Evidence for a target outside `targets` is rejected.
+pub fn evaluate_satisfaction<'a, I>(
+ policy: &SatisfactionPolicy,
+ targets: &TargetSet,
+ evidence: I,
+) -> Result<SatisfactionState, Error>
+where
+ I: IntoIterator<Item = (&'a TargetFingerprint, &'a DeliveryOutcome)>,
+{
+ policy.validate_for(targets)?;
+ let mut states: BTreeMap<&str, (bool, bool)> = targets
+ .targets()
+ .iter()
+ .map(|target| (target.fingerprint().as_str(), (false, true)))
+ .collect();
+
+ for (target, outcome) in evidence {
+ outcome.validate()?;
+ let Some((satisfied, retryable)) = states.get_mut(target.as_str()) else {
+ return Err(Error::UnexpectedDeliveryTargetReceipt);
+ };
+ if outcome.satisfies(policy.class()) {
+ *satisfied = true;
+ *retryable = false;
+ } else if !*satisfied {
+ *retryable = outcome.is_retryable();
+ }
+ }
+
+ let satisfied_targets: Vec<&str> = states
+ .iter()
+ .filter_map(|(target, (satisfied, _))| satisfied.then_some(*target))
+ .collect();
+ if policy.targets.is_satisfied(
+ targets.len(),
+ satisfied_targets.len(),
+ satisfied_targets.as_slice(),
+ ) {
+ return Ok(SatisfactionState::Satisfied);
+ }
+
+ let possible_targets: Vec<&str> = states
+ .iter()
+ .filter_map(|(target, (satisfied, retryable))| {
+ (*satisfied || *retryable).then_some(*target)
+ })
+ .collect();
+ if policy.targets.is_satisfied(
+ targets.len(),
+ possible_targets.len(),
+ possible_targets.as_slice(),
+ ) {
+ Ok(SatisfactionState::Pending)
+ } else {
+ Ok(SatisfactionState::Exhausted)
+ }
+}
+
#[cfg(feature = "serde")]
impl<'de> serde::Deserialize<'de> for TargetPolicy {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
diff --git a/crates/transport/src/sink.rs b/crates/transport/src/sink.rs
@@ -2,8 +2,8 @@
use crate::{
Error,
- outcome::DeliveryOutcome,
- policy::SatisfactionPolicy,
+ outcome::{DeliveryOutcome, Retryability, validate_delivery_code, validate_delivery_message},
+ policy::{SatisfactionPolicy, SatisfactionState, evaluate_satisfaction},
source::BoxFuture,
target::{Target, TargetSet},
};
@@ -19,6 +19,19 @@ pub use crate::status::SinkStatus;
/// Maximum encoded delivery request identity length.
pub const DELIVERY_REQUEST_ID_MAX_BYTES: usize = 256;
+/// Sink-wide typed failure retaining safe retry and partial-target evidence.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SinkFailure {
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ code: String,
+ retryability: Retryability,
+ retry_after_unix_ms: Option<u64>,
+ message: Option<String>,
+ partial_evidence: Vec<DeliveryTargetReceipt>,
+}
+
/// Validated caller identity for one delivery operation.
#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
pub struct DeliveryRequestId(String);
@@ -257,17 +270,22 @@ impl DeliveryReceipt {
/// Returns whether the receipt satisfies the request's exact policy.
pub fn is_satisfied(&self, request: &DeliveryRequest) -> Result<bool, Error> {
+ Ok(matches!(
+ self.satisfaction(request)?,
+ SatisfactionState::Satisfied
+ ))
+ }
+
+ /// Evaluates this receipt as satisfied, pending, or exhausted.
+ pub fn satisfaction(&self, request: &DeliveryRequest) -> Result<SatisfactionState, Error> {
self.validate_for_request(request)?;
- let satisfied: BTreeSet<&str> = self
- .target_receipts
- .iter()
- .filter(|receipt| receipt.outcome().satisfies(request.satisfaction().class()))
- .map(|receipt| receipt.target().fingerprint().as_str())
- .collect();
- Ok(request
- .satisfaction()
- .targets()
- .is_satisfied(self.target_set.len(), &satisfied))
+ evaluate_satisfaction(
+ request.satisfaction(),
+ request.target_set(),
+ self.target_receipts
+ .iter()
+ .map(|receipt| (receipt.target().fingerprint(), receipt.outcome())),
+ )
}
/// Returns the request identity.
@@ -281,6 +299,103 @@ impl DeliveryReceipt {
}
}
+impl SinkFailure {
+ /// Creates a request-bound sink-wide failure with validated partial evidence.
+ pub fn for_request(
+ request: &DeliveryRequest,
+ code: impl Into<String>,
+ retryability: Retryability,
+ retry_after_unix_ms: Option<u64>,
+ message: Option<String>,
+ partial_evidence: Vec<DeliveryTargetReceipt>,
+ ) -> Result<Self, Error> {
+ let failure = Self {
+ request_id: request.request_id.clone(),
+ target_set: request.target_set.clone(),
+ code: code.into(),
+ retryability,
+ retry_after_unix_ms,
+ message,
+ partial_evidence,
+ };
+ failure.validate_for_request(request)?;
+ Ok(failure)
+ }
+
+ /// Returns a terminal adapter-contract failure for an exact request.
+ pub fn invalid_contract(request: &DeliveryRequest) -> Self {
+ Self::for_request(
+ request,
+ "invalid_transport_contract",
+ Retryability::Terminal,
+ None,
+ Some("transport adapter returned invalid evidence".to_string()),
+ Vec::new(),
+ )
+ .expect("static sink failure is valid")
+ }
+
+ /// Validates identity, retry timing, and bounded partial evidence.
+ pub fn validate_for_request(&self, request: &DeliveryRequest) -> Result<(), Error> {
+ if self.request_id != *request.request_id() {
+ return Err(Error::DeliveryReceiptRequestIdMismatch);
+ }
+ if self.target_set != *request.target_set() {
+ return Err(Error::DeliveryReceiptTargetSetMismatch);
+ }
+ validate_delivery_code(self.code.as_str())?;
+ if matches!(self.retryability, Retryability::NotApplicable)
+ || matches!(self.retry_after_unix_ms, Some(0))
+ || (self.retry_after_unix_ms.is_some()
+ && !matches!(self.retryability, Retryability::Retryable))
+ {
+ return Err(Error::InvalidDeliveryOutcome);
+ }
+ if let Some(message) = &self.message {
+ validate_delivery_message(message)?;
+ }
+ let mut observed = BTreeSet::new();
+ for receipt in &self.partial_evidence {
+ receipt.validate()?;
+ if !request
+ .target_set()
+ .contains(receipt.target().fingerprint())
+ {
+ return Err(Error::UnexpectedDeliveryTargetReceipt);
+ }
+ if !observed.insert(receipt.target().fingerprint().as_str()) {
+ return Err(Error::DuplicateDeliveryTargetReceipt);
+ }
+ }
+ Ok(())
+ }
+
+ /// Returns the stable normalized failure code.
+ pub fn code(&self) -> &str {
+ self.code.as_str()
+ }
+
+ /// Returns whether retrying the same request may be useful.
+ pub const fn retryability(&self) -> Retryability {
+ self.retryability
+ }
+
+ /// Returns the earliest absolute Unix millisecond retry time, when supplied.
+ pub const fn retry_after_unix_ms(&self) -> Option<u64> {
+ self.retry_after_unix_ms
+ }
+
+ /// Returns bounded caller-safe diagnostic detail.
+ pub fn message(&self) -> Option<&str> {
+ self.message.as_deref()
+ }
+
+ /// Returns safe target evidence collected before the sink-wide failure.
+ pub fn partial_evidence(&self) -> &[DeliveryTargetReceipt] {
+ self.partial_evidence.as_slice()
+ }
+}
+
/// Host SPI for outbound event delivery.
///
/// This trait supports external implementations and is dyn-compatible. Its
@@ -299,7 +414,10 @@ pub trait EventSink: Send + Sync {
fn status(&self) -> BoxFuture<'_, Result<SinkStatus, Error>>;
/// Delivers an event according to the request's bounded target policy.
- fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>>;
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>>;
}
#[cfg(feature = "serde")]
@@ -394,4 +512,60 @@ mod serde_impl {
.map_err(serde::de::Error::custom)
}
}
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct SinkFailureWire {
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ code: String,
+ retryability: Retryability,
+ retry_after_unix_ms: Option<u64>,
+ message: Option<String>,
+ partial_evidence: Vec<DeliveryTargetReceipt>,
+ }
+
+ impl<'de> serde::Deserialize<'de> for SinkFailure {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = SinkFailureWire::deserialize(deserializer)?;
+ let failure = Self {
+ request_id: wire.request_id,
+ target_set: wire.target_set,
+ code: wire.code,
+ retryability: wire.retryability,
+ retry_after_unix_ms: wire.retry_after_unix_ms,
+ message: wire.message,
+ partial_evidence: wire.partial_evidence,
+ };
+ validate_delivery_code(failure.code.as_str()).map_err(serde::de::Error::custom)?;
+ if matches!(failure.retryability, Retryability::NotApplicable)
+ || matches!(failure.retry_after_unix_ms, Some(0))
+ || (failure.retry_after_unix_ms.is_some()
+ && !matches!(failure.retryability, Retryability::Retryable))
+ {
+ return Err(serde::de::Error::custom(Error::InvalidDeliveryOutcome));
+ }
+ if let Some(message) = &failure.message {
+ validate_delivery_message(message).map_err(serde::de::Error::custom)?;
+ }
+ let mut observed = BTreeSet::new();
+ for receipt in &failure.partial_evidence {
+ receipt.validate().map_err(serde::de::Error::custom)?;
+ if !failure.target_set.contains(receipt.target().fingerprint()) {
+ return Err(serde::de::Error::custom(
+ Error::UnexpectedDeliveryTargetReceipt,
+ ));
+ }
+ if !observed.insert(receipt.target().fingerprint().as_str()) {
+ return Err(serde::de::Error::custom(
+ Error::DuplicateDeliveryTargetReceipt,
+ ));
+ }
+ }
+ Ok(failure)
+ }
+ }
}
diff --git a/crates/transport/src/target.rs b/crates/transport/src/target.rs
@@ -469,6 +469,13 @@ impl TargetSet {
&self.targets
}
+ /// Returns whether this set contains the exact target fingerprint.
+ pub fn contains(&self, fingerprint: &TargetFingerprint) -> bool {
+ self.targets
+ .iter()
+ .any(|target| target.fingerprint() == fingerprint)
+ }
+
pub fn len(&self) -> usize {
self.targets.len()
}
diff --git a/crates/transport/tests/conformance/suite.rs b/crates/transport/tests/conformance/suite.rs
@@ -134,10 +134,8 @@ pub(crate) fn assert_sink_conformance(harness: &impl SinkConformanceHarness) {
assert_eq!(harness.captured_request().as_ref(), Some(&request));
let expired = delivery_request("sink-expired", harness.target_set(), harness.now_unix_ms());
- assert_eq!(
- block_on(sink.deliver(expired)).expect_err("expired delivery"),
- Error::InvalidDeliveryDeadline
- );
+ let failure = block_on(sink.deliver(expired)).expect_err("expired delivery");
+ assert_eq!(failure.code(), "invalid_transport_contract");
}
pub(crate) fn assert_request_boundaries() {
@@ -183,16 +181,18 @@ pub(crate) fn assert_source_error(harness: &impl SourceConformanceHarness, expec
);
}
-pub(crate) fn assert_sink_error(harness: &impl SinkConformanceHarness, expected: Error) {
+pub(crate) fn assert_sink_error(harness: &impl SinkConformanceHarness, _expected: Error) {
let request = delivery_request(
"sink-error",
harness.target_set(),
harness.now_unix_ms() + 100,
);
- assert_eq!(
- block_on(harness.sink().deliver(request)).expect_err("sink error"),
- expected
- );
+ let failure = block_on(harness.sink().deliver(request)).expect_err("sink error");
+ assert_eq!(failure.code(), "invalid_transport_contract");
+ assert!(matches!(
+ failure.retryability(),
+ radroots_transport::outcome::Retryability::Terminal
+ ));
}
pub(crate) fn assert_source_cancellation(harness: &impl SourceConformanceHarness) {
diff --git a/crates/transport/tests/conformance/support.rs b/crates/transport/tests/conformance/support.rs
@@ -6,7 +6,7 @@ use std::sync::{
use futures::future;
use radroots_transport::{
BoxFuture, DeliveryReceipt, DeliveryRequest, Error, EventSink, EventSource, FetchPage,
- FetchRequest, SinkStatus, SourceStatus, TransportId,
+ FetchRequest, SinkFailure, SinkStatus, SourceStatus, TransportId,
capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
outcome::{DeliveryOutcome, FetchTargetOutcome, FetchTargetState},
sink::DeliveryTargetReceipt,
@@ -184,16 +184,19 @@ impl EventSink for MockSink {
Box::pin(async { Ok(sink_status()) })
}
- fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
let state = Arc::clone(&self.state);
let mode = self.mode.clone();
Box::pin(async move {
*state.request.lock().expect("sink request lock") = Some(request.clone());
if request.deadline_unix_ms() <= NOW_UNIX_MS {
- return Err(Error::InvalidDeliveryDeadline);
+ return Err(SinkFailure::invalid_contract(&request));
}
match mode {
- Mode::Fail(error) => Err(error),
+ Mode::Fail(_) => Err(SinkFailure::invalid_contract(&request)),
Mode::Pending => {
state.published.store(true, Ordering::SeqCst);
let state_for_drop = Arc::clone(&state);
@@ -222,6 +225,7 @@ impl EventSink for MockSink {
})
.collect();
DeliveryReceipt::for_request(&request, receipts)
+ .map_err(|_| SinkFailure::invalid_contract(&request))
}
}
})
@@ -303,7 +307,10 @@ impl EventSink for CombinedAdapter {
EventSink::status(&self.sink)
}
- fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>> {
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
self.sink.deliver(request)
}
}
diff --git a/crates/transport/tests/delivery_contract.rs b/crates/transport/tests/delivery_contract.rs
@@ -1,8 +1,11 @@
use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
use radroots_transport::{
- DeliveryReceipt, DeliveryRequest, Error, Target, TargetSet,
+ DeliveryReceipt, DeliveryRequest, Error, SinkFailure, Target, TargetSet,
outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability},
- policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ policy::{
+ SatisfactionClass, SatisfactionPolicy, SatisfactionState, TargetPolicy,
+ evaluate_satisfaction,
+ },
sink::{DeliveryPayload, DeliveryTargetReceipt},
};
@@ -255,6 +258,115 @@ fn retryability_and_terminality_are_explicit_normalized_data() {
}
#[test]
+fn canonical_evaluator_covers_pending_exhausted_and_historical_evidence() {
+ let set = targets();
+ let policy = SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all());
+ assert_eq!(
+ evaluate_satisfaction(&policy, &set, core::iter::empty()).expect("empty evidence"),
+ SatisfactionState::Pending
+ );
+
+ let first = set.targets()[0].fingerprint();
+ let second = set.targets()[1].fingerprint();
+ let third = set.targets()[2].fingerprint();
+ let accepted = DeliveryOutcome::accepted();
+ let terminal = DeliveryOutcome::rejected();
+ let retryable = DeliveryOutcome::unavailable();
+ assert_eq!(
+ evaluate_satisfaction(
+ &policy,
+ &set,
+ [(first, &accepted), (second, &terminal), (third, &retryable),],
+ )
+ .expect("mixed evidence"),
+ SatisfactionState::Exhausted
+ );
+ assert_eq!(
+ evaluate_satisfaction(
+ &policy,
+ &set,
+ [
+ (first, &accepted),
+ (second, &accepted),
+ (third, &accepted),
+ (first, &terminal),
+ ],
+ )
+ .expect("historical success"),
+ SatisfactionState::Satisfied
+ );
+
+ let foreign = Target::nostr_relay("wss://foreign.example").expect("foreign");
+ assert_eq!(
+ evaluate_satisfaction(&policy, &set, [(foreign.fingerprint(), &accepted)])
+ .expect_err("foreign evidence"),
+ Error::UnexpectedDeliveryTargetReceipt
+ );
+}
+
+#[test]
+fn sink_failures_retain_validated_retry_timing_and_partial_evidence() {
+ let request = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::all(),
+ ));
+ let partial = DeliveryTargetReceipt::attempted(
+ request.target_set().targets()[0].clone(),
+ DeliveryOutcome::accepted(),
+ );
+ let failure = SinkFailure::for_request(
+ &request,
+ "relay_batch_unavailable",
+ Retryability::Retryable,
+ Some(1_700_000_200_000),
+ Some("relay batch unavailable".to_owned()),
+ vec![partial.clone()],
+ )
+ .expect("sink failure");
+ assert_eq!(failure.code(), "relay_batch_unavailable");
+ assert_eq!(failure.retryability(), Retryability::Retryable);
+ assert_eq!(failure.retry_after_unix_ms(), Some(1_700_000_200_000));
+ assert_eq!(failure.message(), Some("relay batch unavailable"));
+ assert_eq!(failure.partial_evidence(), core::slice::from_ref(&partial));
+ assert_eq!(
+ SinkFailure::for_request(
+ &request,
+ "terminal_failure",
+ Retryability::Terminal,
+ Some(1),
+ None,
+ Vec::new(),
+ )
+ .expect_err("terminal retry timing"),
+ Error::InvalidDeliveryOutcome
+ );
+ assert_eq!(
+ SinkFailure::for_request(
+ &request,
+ "invalid_retry_time",
+ Retryability::Retryable,
+ Some(0),
+ None,
+ Vec::new(),
+ )
+ .expect_err("zero retry timing"),
+ Error::InvalidDeliveryOutcome
+ );
+ assert_eq!(
+ SinkFailure::for_request(
+ &request,
+ "duplicate_evidence",
+ Retryability::Retryable,
+ None,
+ None,
+ vec![partial.clone(), partial],
+ )
+ .expect_err("duplicate evidence"),
+ Error::DuplicateDeliveryTargetReceipt
+ );
+}
+
+#[test]
#[cfg(feature = "serde")]
fn serde_revalidates_policy_outcome_and_receipt_invariants() {
let request = request(SatisfactionPolicy::new(
@@ -288,4 +400,19 @@ fn serde_revalidates_policy_outcome_and_receipt_invariants() {
.expect("receipts array")
.pop();
assert!(serde_json::from_value::<DeliveryReceipt>(missing).is_err());
+
+ let failure = SinkFailure::for_request(
+ &request,
+ "relay_unavailable",
+ Retryability::Retryable,
+ Some(1_700_000_200_000),
+ None,
+ vec![receipt.target_receipts()[0].clone()],
+ )
+ .expect("sink failure");
+ let encoded_failure = serde_json::to_string(&failure).expect("failure json");
+ assert_eq!(
+ serde_json::from_str::<SinkFailure>(&encoded_failure).expect("failure round trip"),
+ failure
+ );
}
diff --git a/crates/transport/tests/spi.rs b/crates/transport/tests/spi.rs
@@ -37,7 +37,7 @@ impl EventSink for SinkOnly {
fn deliver(
&self,
request: DeliveryRequest,
- ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>> {
Box::pin(async move {
let receipts = request
.target_set()
@@ -49,6 +49,7 @@ impl EventSink for SinkOnly {
})
.collect();
DeliveryReceipt::for_request(&request, receipts)
+ .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))
})
}
}
@@ -79,7 +80,7 @@ impl EventSink for Bidirectional {
fn deliver(
&self,
request: DeliveryRequest,
- ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>> {
self.sink.deliver(request)
}
}
diff --git a/crates/transport_nostr/src/sink.rs b/crates/transport_nostr/src/sink.rs
@@ -4,7 +4,7 @@ use crate::{NostrTransport, RelayUrl, status};
use core::time::Duration;
use radroots_nostr::event::Event;
use radroots_transport::{
- BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink,
+ BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, SinkFailure,
outcome::DeliveryOutcome,
sink::{DeliveryTargetReceipt, SinkStatus},
};
@@ -98,7 +98,7 @@ impl EventSink for NostrTransport {
fn deliver(
&self,
request: DeliveryRequest,
- ) -> BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::Error>> {
+ ) -> BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
Box::pin(async move {
let mut requested = Vec::new();
let mut skipped = Vec::new();
@@ -107,20 +107,25 @@ impl EventSink for NostrTransport {
Ok(relay) if self.config().relays().contains(&relay) => {
requested.push((relay, target.clone()));
}
- _ => skipped.push(DeliveryTargetReceipt::skipped(
- target.clone(),
- DeliveryOutcome::rejected().with_detail(
- "target_denied",
- "target is not configured for this sink",
- )?,
- )?),
+ _ => skipped.push(
+ DeliveryTargetReceipt::skipped(
+ target.clone(),
+ DeliveryOutcome::rejected()
+ .with_detail(
+ "target_denied",
+ "target is not configured for this sink",
+ )
+ .map_err(|_| SinkFailure::invalid_contract(&request))?,
+ )
+ .map_err(|_| SinkFailure::invalid_contract(&request))?,
+ ),
}
}
let event = match radroots_nostr::event::to_nostr(request.payload().event().envelope())
{
Ok(event) => event,
- Err(_) => return Err(radroots_transport::Error::InvalidDeliveryOutcome),
+ Err(_) => return Err(SinkFailure::invalid_contract(&request)),
};
let remaining_ms = request.deadline_unix_ms().saturating_sub(unix_time_ms());
let operation_timeout_ms = remaining_ms.min(self.config().request_timeout_ms());
@@ -131,7 +136,8 @@ impl EventSink for NostrTransport {
.expect("normalized timeout cannot satisfy delivery")
}));
self.status.record_sink(0, skipped.len(), Some("timeout"));
- return DeliveryReceipt::for_request(&request, skipped);
+ return DeliveryReceipt::for_request(&request, skipped)
+ .map_err(|_| SinkFailure::invalid_contract(&request));
}
let expected: BTreeSet<_> = requested.iter().map(|(relay, _)| relay.clone()).collect();
let results = self
@@ -150,7 +156,7 @@ impl EventSink for NostrTransport {
{
continue;
}
- return Err(radroots_transport::Error::InvalidDeliveryOutcome);
+ return Err(SinkFailure::invalid_contract(&request));
}
let mut receipts = skipped;
@@ -172,6 +178,7 @@ impl EventSink for NostrTransport {
.find_map(|receipt| receipt.outcome().message());
self.status.record_sink(accepted, failed, diagnostic);
DeliveryReceipt::for_request(&request, receipts)
+ .map_err(|_| SinkFailure::invalid_contract(&request))
})
}
}
@@ -411,8 +418,10 @@ mod tests {
},
]);
assert_eq!(
- futures::executor::block_on(duplicate.deliver(request())),
- Err(radroots_transport::Error::InvalidDeliveryOutcome)
+ futures::executor::block_on(duplicate.deliver(request()))
+ .expect_err("duplicate relay evidence")
+ .code(),
+ "invalid_transport_contract"
);
let other = RelayUrl::parse("wss://other.example", RelayUrlPolicy::Public).expect("other");
@@ -421,8 +430,10 @@ mod tests {
outcome: DeliveryOutcome::accepted(),
}]);
assert_eq!(
- futures::executor::block_on(unexpected.deliver(request())),
- Err(radroots_transport::Error::InvalidDeliveryOutcome)
+ futures::executor::block_on(unexpected.deliver(request()))
+ .expect_err("unexpected relay evidence")
+ .code(),
+ "invalid_transport_contract"
);
let denied_request = DeliveryRequest::new(
diff --git a/crates/transport_reticulum/src/lib.rs b/crates/transport_reticulum/src/lib.rs
@@ -343,20 +343,24 @@ impl EventSink for RadrootsReticulumTransport {
fn deliver(
&self,
request: DeliveryRequest,
- ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, RadrootsTransportError>> {
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, radroots_transport::SinkFailure>>
+ {
Box::pin(async move {
ensure_reticulum_targets(request.target_set().targets())
- .map_err(reticulum_error_to_transport_error)?;
+ .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))?;
let outcome = DeliveryOutcome::unavailable()
- .with_detail(UNAVAILABLE_CODE, RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE)?;
+ .with_detail(UNAVAILABLE_CODE, RADROOTS_RETICULUM_UNAVAILABLE_MESSAGE)
+ .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))?;
let receipts = request
.target_set()
.targets()
.iter()
.cloned()
.map(|target| DeliveryTargetReceipt::skipped(target, outcome.clone()))
- .collect::<Result<Vec<_>, _>>()?;
+ .collect::<Result<Vec<_>, _>>()
+ .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))?;
DeliveryReceipt::for_request(&request, receipts)
+ .map_err(|_| radroots_transport::SinkFailure::invalid_contract(&request))
})
}
}