commit f5777f5b1864240f6e353496f90738d1153dab0d
parent be39b6fbcace4970f50cebbe6b9f816c3c7288da
Author: triesap <tyson@radroots.org>
Date: Thu, 30 Jul 2026 17:40:25 +0000
transport: define delivery requests, receipts, and satisfaction policy
- bind signed event payloads to explicit targets deadlines and request identities
- model accepted and delivered satisfaction across any all quorum and required targets
- normalize retryable and terminal per-target outcomes with bounded diagnostics
- verify partial success receipt invariants serde no-std wasm and workspace gates
Diffstat:
8 files changed, 1074 insertions(+), 20 deletions(-)
diff --git a/crates/transport/src/error.rs b/crates/transport/src/error.rs
@@ -37,6 +37,8 @@ pub enum RadrootsTransportError {
EmptyDeliveryRequestId,
InvalidDeliveryRequestId,
InvalidDeliveryTimestamp,
+ InvalidDeliveryDeadline,
+ InvalidDeliveryOutcome,
UnexpectedDeliveryTargetReceipt,
DuplicateDeliveryTargetReceipt,
MissingDeliveryTargetReceipt,
@@ -114,6 +116,8 @@ impl fmt::Display for RadrootsTransportError {
Self::InvalidDeliveryTimestamp => {
f.write_str("transport delivery timestamp is invalid")
}
+ Self::InvalidDeliveryDeadline => f.write_str("transport delivery deadline is invalid"),
+ Self::InvalidDeliveryOutcome => f.write_str("transport delivery outcome is invalid"),
Self::UnexpectedDeliveryTargetReceipt => {
f.write_str("transport delivery receipt contains an unexpected target")
}
diff --git a/crates/transport/src/outcome.rs b/crates/transport/src/outcome.rs
@@ -3,6 +3,11 @@
use crate::target::TargetFingerprint;
use alloc::string::String;
+/// Maximum encoded normalized outcome code length.
+pub const DELIVERY_OUTCOME_CODE_MAX_BYTES: usize = 64;
+/// Maximum encoded normalized outcome message length.
+pub const DELIVERY_OUTCOME_MESSAGE_MAX_BYTES: usize = 1_024;
+
/// Target-local result of one bounded fetch attempt.
#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
@@ -74,8 +79,220 @@ impl FetchTargetOutcome {
self.state
}
+ /// Returns bounded adapter-normalized diagnostic detail.
+ pub fn message(&self) -> Option<&str> {
+ self.message.as_deref()
+ }
+}
+
+/// Normalized result class for one delivery target.
+#[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 DeliveryOutcomeKind {
+ /// The target accepted responsibility for the event.
+ Accepted,
+ /// The target confirmed final delivery.
+ Delivered,
+ /// The target rejected the event permanently.
+ Rejected,
+ /// The target was temporarily unavailable.
+ Unavailable,
+ /// The adapter reported another normalized failure.
+ Failed,
+}
+
+/// Whether a failed outcome can be retried without changing the request.
+#[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 Retryability {
+ /// Outcome is successful and retry classification does not apply.
+ NotApplicable,
+ /// A caller may decide to retry the same target.
+ Retryable,
+ /// Retrying the same target and payload is not useful.
+ Terminal,
+}
+
+/// Validated normalized outcome for one delivery target.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryOutcome {
+ kind: DeliveryOutcomeKind,
+ retryability: Retryability,
+ code: Option<String>,
+ message: Option<String>,
+}
+
+impl DeliveryOutcome {
+ /// The target accepted responsibility for the event.
+ pub const fn accepted() -> Self {
+ Self::new_success(DeliveryOutcomeKind::Accepted)
+ }
+
+ /// The target confirmed final delivery.
+ pub const fn delivered() -> Self {
+ Self::new_success(DeliveryOutcomeKind::Delivered)
+ }
+
+ /// The target rejected the event permanently.
+ pub const fn rejected() -> Self {
+ Self::new_failure(DeliveryOutcomeKind::Rejected, Retryability::Terminal)
+ }
+
+ /// The target was temporarily unavailable.
+ pub const fn unavailable() -> Self {
+ Self::new_failure(DeliveryOutcomeKind::Unavailable, Retryability::Retryable)
+ }
+
+ /// Creates another normalized failure with explicit retry classification.
+ pub const fn failed(retryability: Retryability) -> Result<Self, crate::Error> {
+ if matches!(retryability, Retryability::NotApplicable) {
+ return Err(crate::Error::InvalidDeliveryOutcome);
+ }
+ Ok(Self::new_failure(DeliveryOutcomeKind::Failed, retryability))
+ }
+
+ const fn new_success(kind: DeliveryOutcomeKind) -> Self {
+ Self {
+ kind,
+ retryability: Retryability::NotApplicable,
+ code: None,
+ message: None,
+ }
+ }
+
+ const fn new_failure(kind: DeliveryOutcomeKind, retryability: Retryability) -> Self {
+ Self {
+ kind,
+ retryability,
+ code: None,
+ message: None,
+ }
+ }
+
+ /// Attaches adapter-normalized diagnostic fields.
+ pub fn with_detail(
+ mut self,
+ code: impl Into<String>,
+ message: impl Into<String>,
+ ) -> Result<Self, crate::Error> {
+ let code = code.into();
+ let message = message.into();
+ validate_delivery_detail(code.as_str(), message.as_str())?;
+ self.code = Some(code);
+ self.message = Some(message);
+ Ok(self)
+ }
+
+ /// Returns the normalized result kind.
+ pub const fn kind(&self) -> DeliveryOutcomeKind {
+ self.kind
+ }
+
+ /// Returns the explicit retry classification.
+ pub const fn retryability(&self) -> Retryability {
+ self.retryability
+ }
+
+ /// Whether this outcome satisfies the requested success class.
+ pub const fn satisfies(&self, class: crate::policy::SatisfactionClass) -> bool {
+ match class {
+ crate::policy::SatisfactionClass::Accepted => matches!(
+ self.kind,
+ DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered
+ ),
+ crate::policy::SatisfactionClass::Delivered => {
+ matches!(self.kind, DeliveryOutcomeKind::Delivered)
+ }
+ }
+ }
+
+ /// Whether the same target and payload may be retried.
+ pub const fn is_retryable(&self) -> bool {
+ matches!(self.retryability, Retryability::Retryable)
+ }
+
+ /// Whether the failure is terminal for the same target and payload.
+ pub const fn is_terminal(&self) -> bool {
+ matches!(self.retryability, Retryability::Terminal)
+ }
+
+ /// Returns the adapter-normalized code.
+ pub fn code(&self) -> Option<&str> {
+ self.code.as_deref()
+ }
+
/// Returns caller-safe diagnostic detail.
pub fn message(&self) -> Option<&str> {
self.message.as_deref()
}
+
+ pub(crate) fn validate(&self) -> Result<(), crate::Error> {
+ let valid = match self.kind {
+ DeliveryOutcomeKind::Accepted | DeliveryOutcomeKind::Delivered => {
+ matches!(self.retryability, Retryability::NotApplicable)
+ }
+ DeliveryOutcomeKind::Rejected => matches!(self.retryability, Retryability::Terminal),
+ DeliveryOutcomeKind::Unavailable => {
+ matches!(self.retryability, Retryability::Retryable)
+ }
+ DeliveryOutcomeKind::Failed => {
+ !matches!(self.retryability, Retryability::NotApplicable)
+ }
+ };
+ if !valid {
+ return Err(crate::Error::InvalidDeliveryOutcome);
+ }
+ match (&self.code, &self.message) {
+ (None, None) => Ok(()),
+ (Some(code), Some(message)) => validate_delivery_detail(code, message),
+ (None, Some(_)) | (Some(_), None) => Err(crate::Error::InvalidDeliveryOutcome),
+ }
+ }
+}
+
+fn validate_delivery_detail(code: &str, message: &str) -> Result<(), crate::Error> {
+ let valid_code = !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'.')
+ });
+ 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 {
+ Ok(())
+ } else {
+ Err(crate::Error::InvalidDeliveryOutcome)
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for DeliveryOutcome {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct Wire {
+ kind: DeliveryOutcomeKind,
+ retryability: Retryability,
+ code: Option<String>,
+ message: Option<String>,
+ }
+
+ let wire = Wire::deserialize(deserializer)?;
+ let outcome = Self {
+ kind: wire.kind,
+ retryability: wire.retryability,
+ code: wire.code,
+ message: wire.message,
+ };
+ outcome.validate().map_err(serde::de::Error::custom)?;
+ Ok(outcome)
+ }
}
diff --git a/crates/transport/src/policy.rs b/crates/transport/src/policy.rs
@@ -1 +1,172 @@
-//! Transport-neutral delivery and fetch policy.
+//! Transport-neutral delivery satisfaction policy.
+
+use crate::{
+ Error,
+ target::{TargetFingerprint, TargetSet},
+};
+use alloc::{collections::BTreeSet, vec::Vec};
+
+/// Success level a caller requires from selected targets.
+#[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 SatisfactionClass {
+ /// The target accepted responsibility for the event.
+ Accepted,
+ /// The target confirmed final delivery.
+ Delivered,
+}
+
+/// Which requested targets must reach the satisfaction class.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[cfg_attr(feature = "serde", serde(transparent))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct TargetPolicy {
+ kind: TargetPolicyKind,
+}
+
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+enum TargetPolicyKind {
+ Any,
+ All,
+ Quorum(u16),
+ Required(Vec<TargetFingerprint>),
+}
+
+impl TargetPolicy {
+ /// Any one requested target must satisfy the class.
+ pub const fn any() -> Self {
+ Self {
+ kind: TargetPolicyKind::Any,
+ }
+ }
+
+ /// Every requested target must satisfy the class.
+ pub const fn all() -> Self {
+ Self {
+ kind: TargetPolicyKind::All,
+ }
+ }
+
+ /// At least this non-zero count of requested targets must satisfy the class.
+ pub const fn quorum(threshold: u16) -> Result<Self, Error> {
+ if threshold == 0 {
+ return Err(Error::InvalidSatisfactionPolicy);
+ }
+ Ok(Self {
+ kind: TargetPolicyKind::Quorum(threshold),
+ })
+ }
+
+ /// These exact, unique target fingerprints must satisfy the class.
+ pub fn required(mut targets: Vec<TargetFingerprint>) -> Result<Self, Error> {
+ if targets.is_empty() {
+ return Err(Error::EmptyRequiredTargetSet);
+ }
+ targets.sort();
+ if targets.windows(2).any(|pair| pair[0] == pair[1]) {
+ return Err(Error::DuplicateRequiredTargetFingerprint);
+ }
+ Ok(Self {
+ kind: TargetPolicyKind::Required(targets),
+ })
+ }
+
+ /// Returns exact required fingerprints for a required-target policy.
+ pub fn required_targets(&self) -> Option<&[TargetFingerprint]> {
+ match &self.kind {
+ TargetPolicyKind::Required(targets) => Some(targets.as_slice()),
+ TargetPolicyKind::Any | TargetPolicyKind::All | TargetPolicyKind::Quorum(_) => None,
+ }
+ }
+
+ pub(crate) fn validate_for(&self, targets: &TargetSet) -> Result<(), Error> {
+ match &self.kind {
+ TargetPolicyKind::Any | TargetPolicyKind::All => Ok(()),
+ TargetPolicyKind::Quorum(threshold) => {
+ if usize::from(*threshold) > targets.len() {
+ Err(Error::InvalidSatisfactionPolicy)
+ } else {
+ Ok(())
+ }
+ }
+ 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()))
+ {
+ Err(Error::RequiredTargetNotRequested)
+ } else {
+ Ok(())
+ }
+ }
+ }
+ }
+
+ pub(crate) fn is_satisfied(&self, total_targets: usize, satisfied: &BTreeSet<&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::Required(required) => required
+ .iter()
+ .all(|fingerprint| satisfied.contains(fingerprint.as_str())),
+ }
+ }
+}
+
+/// Required success level and target selection for one delivery request.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct SatisfactionPolicy {
+ class: SatisfactionClass,
+ targets: TargetPolicy,
+}
+
+impl SatisfactionPolicy {
+ /// Creates an explicit satisfaction policy.
+ pub const fn new(class: SatisfactionClass, targets: TargetPolicy) -> Self {
+ Self { class, targets }
+ }
+
+ /// Returns the accepted or delivered success level.
+ pub const fn class(&self) -> SatisfactionClass {
+ self.class
+ }
+
+ /// Returns the selected target policy.
+ pub const fn targets(&self) -> &TargetPolicy {
+ &self.targets
+ }
+
+ pub(crate) fn validate_for(&self, targets: &TargetSet) -> Result<(), Error> {
+ self.targets.validate_for(targets)
+ }
+}
+
+#[cfg(feature = "serde")]
+impl<'de> serde::Deserialize<'de> for TargetPolicy {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ match TargetPolicyKind::deserialize(deserializer)? {
+ TargetPolicyKind::Any => Ok(Self::any()),
+ TargetPolicyKind::All => Ok(Self::all()),
+ TargetPolicyKind::Quorum(threshold) => {
+ Self::quorum(threshold).map_err(serde::de::Error::custom)
+ }
+ TargetPolicyKind::Required(targets) => {
+ Self::required(targets).map_err(serde::de::Error::custom)
+ }
+ }
+ }
+}
diff --git a/crates/transport/src/sink.rs b/crates/transport/src/sink.rs
@@ -1,19 +1,285 @@
-//! Outbound event delivery SPI and request models.
+//! Outbound event delivery SPI and bounded request models.
use crate::{
- Error, RadrootsTransportDeliveryReceipt, RadrootsTransportDeliveryRequest, source::BoxFuture,
+ Error,
+ outcome::DeliveryOutcome,
+ policy::SatisfactionPolicy,
+ source::BoxFuture,
+ target::{Target, TargetSet},
};
+use alloc::{
+ collections::{BTreeMap, BTreeSet},
+ string::{String, ToString},
+ vec::Vec,
+};
+use radroots_event::SignedEvent;
pub use crate::status::SinkStatus;
-/// Bounded delivery request.
-///
-/// The dedicated delivery checkpoint replaces this compatibility alias with
-/// the final request model.
-pub type DeliveryRequest = RadrootsTransportDeliveryRequest;
+/// Maximum encoded delivery request identity length.
+pub const DELIVERY_REQUEST_ID_MAX_BYTES: usize = 256;
+
+/// Validated caller identity for one delivery operation.
+#[derive(Clone, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)]
+pub struct DeliveryRequestId(String);
+
+impl DeliveryRequestId {
+ /// Parses a non-empty, bounded, printable request identity.
+ pub fn parse(value: impl Into<String>) -> Result<Self, Error> {
+ let value = value.into();
+ if value.is_empty() {
+ return Err(Error::EmptyDeliveryRequestId);
+ }
+ if value.len() > DELIVERY_REQUEST_ID_MAX_BYTES
+ || value != value.trim()
+ || value.chars().any(char::is_control)
+ {
+ return Err(Error::InvalidDeliveryRequestId);
+ }
+ Ok(Self(value))
+ }
+
+ /// Returns the validated request identity.
+ pub fn as_str(&self) -> &str {
+ self.0.as_str()
+ }
+}
+
+/// Transport-neutral outbound event payload.
+#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))]
+#[cfg_attr(feature = "serde", serde(deny_unknown_fields))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryPayload {
+ event: SignedEvent,
+}
+
+impl DeliveryPayload {
+ /// Wraps an ID-checked signed event for delivery.
+ pub const fn new(event: SignedEvent) -> Self {
+ Self { event }
+ }
+
+ /// Returns the signed event. Signature verification remains a caller concern.
+ pub const fn event(&self) -> &SignedEvent {
+ &self.event
+ }
+}
+
+/// Bounded multi-target delivery request.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryRequest {
+ request_id: DeliveryRequestId,
+ payload: DeliveryPayload,
+ target_set: TargetSet,
+ satisfaction: SatisfactionPolicy,
+ deadline_unix_ms: u64,
+}
+
+impl DeliveryRequest {
+ /// Creates and validates one explicit delivery request.
+ pub fn new(
+ request_id: impl Into<String>,
+ payload: DeliveryPayload,
+ target_set: TargetSet,
+ satisfaction: SatisfactionPolicy,
+ deadline_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if deadline_unix_ms == 0 {
+ return Err(Error::InvalidDeliveryDeadline);
+ }
+ satisfaction.validate_for(&target_set)?;
+ Ok(Self {
+ request_id: DeliveryRequestId::parse(request_id)?,
+ payload,
+ target_set,
+ satisfaction,
+ deadline_unix_ms,
+ })
+ }
+
+ /// Returns the request identity.
+ pub const fn request_id(&self) -> &DeliveryRequestId {
+ &self.request_id
+ }
+
+ /// Returns the signed event payload.
+ pub const fn payload(&self) -> &DeliveryPayload {
+ &self.payload
+ }
+
+ /// Returns the exact non-empty target set.
+ pub const fn target_set(&self) -> &TargetSet {
+ &self.target_set
+ }
+
+ /// Returns the requested success and target policy.
+ pub const fn satisfaction(&self) -> &SatisfactionPolicy {
+ &self.satisfaction
+ }
+
+ /// Returns the absolute Unix deadline in milliseconds.
+ pub const fn deadline_unix_ms(&self) -> u64 {
+ self.deadline_unix_ms
+ }
+}
+
+/// Normalized result for one requested target.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryTargetReceipt {
+ target: Target,
+ attempted: bool,
+ outcome: DeliveryOutcome,
+}
+
+impl DeliveryTargetReceipt {
+ /// Records the outcome of an attempted target.
+ pub const fn attempted(target: Target, outcome: DeliveryOutcome) -> Self {
+ Self {
+ target,
+ attempted: true,
+ outcome,
+ }
+ }
+
+ /// Records an unattempted target and its normalized failure reason.
+ pub fn skipped(target: Target, outcome: DeliveryOutcome) -> Result<Self, Error> {
+ if outcome.satisfies(crate::policy::SatisfactionClass::Accepted) {
+ return Err(Error::DeliveryTargetReceiptAttemptMismatch);
+ }
+ Ok(Self {
+ target,
+ attempted: false,
+ outcome,
+ })
+ }
-/// Per-target delivery result.
-pub type DeliveryReceipt = RadrootsTransportDeliveryReceipt;
+ /// Returns the exact target.
+ pub const fn target(&self) -> &Target {
+ &self.target
+ }
+
+ /// Whether the adapter attempted remote publication.
+ pub const fn was_attempted(&self) -> bool {
+ self.attempted
+ }
+
+ /// Returns normalized target outcome data.
+ pub const fn outcome(&self) -> &DeliveryOutcome {
+ &self.outcome
+ }
+
+ fn validate(&self) -> Result<(), Error> {
+ self.outcome.validate()?;
+ if !self.attempted
+ && self
+ .outcome
+ .satisfies(crate::policy::SatisfactionClass::Accepted)
+ {
+ return Err(Error::DeliveryTargetReceiptAttemptMismatch);
+ }
+ Ok(())
+ }
+}
+
+/// Request-bound per-target delivery receipt.
+#[cfg_attr(feature = "serde", derive(serde::Serialize))]
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryReceipt {
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ target_receipts: Vec<DeliveryTargetReceipt>,
+}
+
+impl DeliveryReceipt {
+ /// Creates a complete result set in the original request target order.
+ pub fn for_request(
+ request: &DeliveryRequest,
+ target_receipts: Vec<DeliveryTargetReceipt>,
+ ) -> Result<Self, Error> {
+ Self::new(
+ request.request_id.clone(),
+ request.target_set.clone(),
+ target_receipts,
+ )
+ }
+
+ fn new(
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ target_receipts: Vec<DeliveryTargetReceipt>,
+ ) -> Result<Self, Error> {
+ let mut by_fingerprint = BTreeMap::new();
+ for receipt in target_receipts {
+ receipt.validate()?;
+ let fingerprint = receipt.target().fingerprint().as_str().to_string();
+ if !target_set
+ .targets()
+ .iter()
+ .any(|target| target.fingerprint().as_str() == fingerprint)
+ {
+ return Err(Error::UnexpectedDeliveryTargetReceipt);
+ }
+ if by_fingerprint.insert(fingerprint, receipt).is_some() {
+ return Err(Error::DuplicateDeliveryTargetReceipt);
+ }
+ }
+
+ let mut ordered = Vec::with_capacity(target_set.len());
+ for target in target_set.targets() {
+ let Some(receipt) = by_fingerprint.remove(target.fingerprint().as_str()) else {
+ return Err(Error::MissingDeliveryTargetReceipt);
+ };
+ ordered.push(receipt);
+ }
+ Ok(Self {
+ request_id,
+ target_set,
+ target_receipts: ordered,
+ })
+ }
+
+ /// Validates the receipt against the exact request identity and targets.
+ 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);
+ }
+ let rebuilt = Self::for_request(request, self.target_receipts.clone())?;
+ if rebuilt != *self {
+ return Err(Error::DeliveryReceiptTargetSetMismatch);
+ }
+ Ok(())
+ }
+
+ /// Returns whether the receipt satisfies the request's exact policy.
+ pub fn is_satisfied(&self, request: &DeliveryRequest) -> Result<bool, 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))
+ }
+
+ /// Returns the request identity.
+ pub const fn request_id(&self) -> &DeliveryRequestId {
+ &self.request_id
+ }
+
+ /// Returns per-target results in request order.
+ pub fn target_receipts(&self) -> &[DeliveryTargetReceipt] {
+ self.target_receipts.as_slice()
+ }
+}
/// Host SPI for outbound event delivery.
///
@@ -35,3 +301,97 @@ pub trait EventSink: Send + Sync {
/// Delivers an event according to the request's bounded target policy.
fn deliver(&self, request: DeliveryRequest) -> BoxFuture<'_, Result<DeliveryReceipt, Error>>;
}
+
+#[cfg(feature = "serde")]
+mod serde_impl {
+ use super::*;
+
+ impl serde::Serialize for DeliveryRequestId {
+ fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
+ where
+ S: serde::Serializer,
+ {
+ serializer.serialize_str(self.as_str())
+ }
+ }
+
+ impl<'de> serde::Deserialize<'de> for DeliveryRequestId {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let value = <String as serde::Deserialize>::deserialize(deserializer)?;
+ Self::parse(value).map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct DeliveryRequestWire {
+ request_id: String,
+ payload: DeliveryPayload,
+ target_set: TargetSet,
+ satisfaction: SatisfactionPolicy,
+ deadline_unix_ms: u64,
+ }
+
+ impl<'de> serde::Deserialize<'de> for DeliveryRequest {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = DeliveryRequestWire::deserialize(deserializer)?;
+ Self::new(
+ wire.request_id,
+ wire.payload,
+ wire.target_set,
+ wire.satisfaction,
+ wire.deadline_unix_ms,
+ )
+ .map_err(serde::de::Error::custom)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct DeliveryTargetReceiptWire {
+ target: Target,
+ attempted: bool,
+ outcome: DeliveryOutcome,
+ }
+
+ impl<'de> serde::Deserialize<'de> for DeliveryTargetReceipt {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = DeliveryTargetReceiptWire::deserialize(deserializer)?;
+ let receipt = Self {
+ target: wire.target,
+ attempted: wire.attempted,
+ outcome: wire.outcome,
+ };
+ receipt.validate().map_err(serde::de::Error::custom)?;
+ Ok(receipt)
+ }
+ }
+
+ #[derive(serde::Deserialize)]
+ #[serde(deny_unknown_fields)]
+ struct DeliveryReceiptWire {
+ request_id: DeliveryRequestId,
+ target_set: TargetSet,
+ target_receipts: Vec<DeliveryTargetReceipt>,
+ }
+
+ impl<'de> serde::Deserialize<'de> for DeliveryReceipt {
+ fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
+ where
+ D: serde::Deserializer<'de>,
+ {
+ let wire = DeliveryReceiptWire::deserialize(deserializer)?;
+ Self::new(wire.request_id, wire.target_set, wire.target_receipts)
+ .map_err(serde::de::Error::custom)
+ }
+ }
+}
diff --git a/crates/transport/tests/delivery_contract.rs b/crates/transport/tests/delivery_contract.rs
@@ -0,0 +1,290 @@
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
+use radroots_transport::{
+ DeliveryReceipt, DeliveryRequest, Error, Target, TargetSet,
+ outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability},
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::{DeliveryPayload, DeliveryTargetReceipt},
+};
+
+fn payload() -> DeliveryPayload {
+ let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
+ let wire = Nip01EventWire::parse_json(raw).expect("wire event");
+ DeliveryPayload::new(
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed delivery event"),
+ )
+}
+
+fn targets() -> TargetSet {
+ TargetSet::new(vec![
+ Target::nostr_relay("wss://one.example").expect("first"),
+ Target::nostr_relay("wss://two.example").expect("second"),
+ Target::nostr_relay("wss://three.example").expect("third"),
+ ])
+ .expect("targets")
+}
+
+fn request(policy: SatisfactionPolicy) -> DeliveryRequest {
+ DeliveryRequest::new(
+ "delivery-request",
+ payload(),
+ targets(),
+ policy,
+ 1_700_000_100_000,
+ )
+ .expect("delivery request")
+}
+
+fn mixed_receipt(request: &DeliveryRequest) -> DeliveryReceipt {
+ let targets = request.target_set().targets();
+ DeliveryReceipt::for_request(
+ request,
+ vec![
+ DeliveryTargetReceipt::attempted(
+ targets[2].clone(),
+ DeliveryOutcome::unavailable()
+ .with_detail("offline", "relay unavailable")
+ .expect("normalized detail"),
+ ),
+ DeliveryTargetReceipt::attempted(targets[0].clone(), DeliveryOutcome::delivered()),
+ DeliveryTargetReceipt::attempted(targets[1].clone(), DeliveryOutcome::accepted()),
+ ],
+ )
+ .expect("mixed receipt")
+}
+
+#[test]
+fn any_all_quorum_and_required_targets_are_exact() {
+ let any = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::any(),
+ ));
+ assert!(mixed_receipt(&any).is_satisfied(&any).expect("any"));
+
+ let all = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::all(),
+ ));
+ assert!(!mixed_receipt(&all).is_satisfied(&all).expect("all"));
+
+ let quorum = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::quorum(2).expect("quorum"),
+ ));
+ assert!(
+ mixed_receipt(&quorum)
+ .is_satisfied(&quorum)
+ .expect("quorum")
+ );
+
+ let delivered = request(SatisfactionPolicy::new(
+ SatisfactionClass::Delivered,
+ TargetPolicy::quorum(2).expect("quorum"),
+ ));
+ assert!(
+ !mixed_receipt(&delivered)
+ .is_satisfied(&delivered)
+ .expect("delivered quorum")
+ );
+
+ let selected = targets();
+ let required = TargetPolicy::required(vec![
+ selected.targets()[0].fingerprint().clone(),
+ selected.targets()[1].fingerprint().clone(),
+ ])
+ .expect("required targets");
+ let required_request = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ required,
+ ));
+ assert!(
+ mixed_receipt(&required_request)
+ .is_satisfied(&required_request)
+ .expect("required")
+ );
+
+ let unsatisfied = TargetPolicy::required(vec![
+ selected.targets()[0].fingerprint().clone(),
+ selected.targets()[2].fingerprint().clone(),
+ ])
+ .expect("required targets");
+ let unsatisfied_request = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ unsatisfied,
+ ));
+ assert!(
+ !mixed_receipt(&unsatisfied_request)
+ .is_satisfied(&unsatisfied_request)
+ .expect("required unsatisfied")
+ );
+}
+
+#[test]
+fn policies_and_requests_reject_empty_duplicate_and_impossible_inputs() {
+ assert_eq!(
+ TargetPolicy::quorum(0).expect_err("zero quorum"),
+ Error::InvalidSatisfactionPolicy
+ );
+ assert_eq!(
+ TargetPolicy::required(Vec::new()).expect_err("empty required"),
+ Error::EmptyRequiredTargetSet
+ );
+ let set = targets();
+ let duplicate = set.targets()[0].fingerprint().clone();
+ assert_eq!(
+ TargetPolicy::required(vec![duplicate.clone(), duplicate]).expect_err("duplicate required"),
+ Error::DuplicateRequiredTargetFingerprint
+ );
+ assert_eq!(
+ DeliveryRequest::new(
+ "request",
+ payload(),
+ set.clone(),
+ SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::quorum(4).expect("nonzero quorum"),
+ ),
+ 1,
+ )
+ .expect_err("impossible quorum"),
+ Error::InvalidSatisfactionPolicy
+ );
+ assert_eq!(
+ DeliveryRequest::new(
+ "",
+ payload(),
+ set.clone(),
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ 1,
+ )
+ .expect_err("empty id"),
+ Error::EmptyDeliveryRequestId
+ );
+ assert_eq!(
+ DeliveryRequest::new(
+ "request",
+ payload(),
+ set,
+ SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ 0,
+ )
+ .expect_err("zero deadline"),
+ Error::InvalidDeliveryDeadline
+ );
+ assert_eq!(
+ TargetSet::new(Vec::new()).expect_err("empty target set"),
+ Error::EmptyTargetSet
+ );
+}
+
+#[test]
+fn receipts_reject_duplicate_missing_unexpected_and_false_attempts() {
+ let request = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::all(),
+ ));
+ let first = request.target_set().targets()[0].clone();
+ let second = request.target_set().targets()[1].clone();
+ let third = request.target_set().targets()[2].clone();
+ let accepted = DeliveryTargetReceipt::attempted(first.clone(), DeliveryOutcome::accepted());
+ assert_eq!(
+ DeliveryTargetReceipt::skipped(first.clone(), DeliveryOutcome::accepted())
+ .expect_err("unattempted success"),
+ Error::DeliveryTargetReceiptAttemptMismatch
+ );
+ assert_eq!(
+ DeliveryReceipt::for_request(&request, vec![accepted.clone(), accepted])
+ .expect_err("duplicate receipt"),
+ Error::DuplicateDeliveryTargetReceipt
+ );
+ assert_eq!(
+ DeliveryReceipt::for_request(
+ &request,
+ vec![
+ DeliveryTargetReceipt::attempted(first, DeliveryOutcome::accepted()),
+ DeliveryTargetReceipt::attempted(second, DeliveryOutcome::accepted()),
+ ],
+ )
+ .expect_err("missing receipt"),
+ Error::MissingDeliveryTargetReceipt
+ );
+ let foreign = Target::nostr_relay("wss://foreign.example").expect("foreign");
+ assert_eq!(
+ DeliveryReceipt::for_request(
+ &request,
+ vec![
+ DeliveryTargetReceipt::attempted(foreign, DeliveryOutcome::accepted()),
+ DeliveryTargetReceipt::attempted(third, DeliveryOutcome::accepted()),
+ ],
+ )
+ .expect_err("unexpected receipt"),
+ Error::UnexpectedDeliveryTargetReceipt
+ );
+}
+
+#[test]
+fn retryability_and_terminality_are_explicit_normalized_data() {
+ let unavailable = DeliveryOutcome::unavailable();
+ assert_eq!(unavailable.kind(), DeliveryOutcomeKind::Unavailable);
+ assert_eq!(unavailable.retryability(), Retryability::Retryable);
+ assert!(unavailable.is_retryable());
+ assert!(!unavailable.is_terminal());
+
+ let rejected = DeliveryOutcome::rejected();
+ assert!(rejected.is_terminal());
+ assert!(!rejected.is_retryable());
+ assert_eq!(
+ DeliveryOutcome::failed(Retryability::NotApplicable).expect_err("unclassified failure"),
+ Error::InvalidDeliveryOutcome
+ );
+ assert!(
+ DeliveryOutcome::failed(Retryability::Retryable)
+ .expect("retryable failure")
+ .is_retryable()
+ );
+ assert!(
+ DeliveryOutcome::failed(Retryability::Terminal)
+ .expect("terminal failure")
+ .is_terminal()
+ );
+ assert_eq!(
+ DeliveryOutcome::unavailable()
+ .with_detail("INVALID", "relay unavailable")
+ .expect_err("invalid code"),
+ Error::InvalidDeliveryOutcome
+ );
+}
+
+#[test]
+fn serde_revalidates_policy_outcome_and_receipt_invariants() {
+ let request = request(SatisfactionPolicy::new(
+ SatisfactionClass::Accepted,
+ TargetPolicy::all(),
+ ));
+ let receipt = mixed_receipt(&request);
+ let encoded_request = serde_json::to_string(&request).expect("request json");
+ assert_eq!(
+ serde_json::from_str::<DeliveryRequest>(&encoded_request).expect("request round trip"),
+ request
+ );
+ let encoded_receipt = serde_json::to_string(&receipt).expect("receipt json");
+ assert_eq!(
+ serde_json::from_str::<DeliveryReceipt>(&encoded_receipt).expect("receipt round trip"),
+ receipt
+ );
+
+ let mut forged_outcome =
+ serde_json::to_value(DeliveryOutcome::accepted()).expect("outcome value");
+ forged_outcome["retryability"] = serde_json::json!("retryable");
+ assert!(serde_json::from_value::<DeliveryOutcome>(forged_outcome).is_err());
+
+ let mut forged_attempt = serde_json::to_value(&receipt).expect("receipt value");
+ forged_attempt["target_receipts"][0]["attempted"] = false.into();
+ assert!(serde_json::from_value::<DeliveryReceipt>(forged_attempt).is_err());
+
+ let mut missing = serde_json::to_value(&receipt).expect("receipt value");
+ missing["target_receipts"]
+ .as_array_mut()
+ .expect("receipts array")
+ .pop();
+ assert!(serde_json::from_value::<DeliveryReceipt>(missing).is_err());
+}
diff --git a/crates/transport/tests/spi.rs b/crates/transport/tests/spi.rs
@@ -1,10 +1,12 @@
use futures::executor::block_on;
+use radroots_event::{SignedEvent, wire::v1::Nip01EventWire};
use radroots_transport::{
BoxFuture, DeliveryReceipt, DeliveryRequest, EventSink, EventSource, FetchPage, FetchRequest,
- RadrootsTransportOutcome, RadrootsTransportOutcomeKind, RadrootsTransportPayload,
- RadrootsTransportSatisfactionPolicy, RadrootsTransportTarget, RadrootsTransportTargetReceipt,
- RadrootsTransportTargetSet, SinkStatus, SourceStatus, TransportId,
+ RadrootsTransportTarget, RadrootsTransportTargetSet, SinkStatus, SourceStatus, TransportId,
capability::{Availability, Maturity, SinkCapabilities, SourceCapabilities},
+ outcome::DeliveryOutcome,
+ policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::{DeliveryPayload, DeliveryTargetReceipt},
source::{FetchBounds, NextPage},
};
@@ -43,10 +45,7 @@ impl EventSink for SinkOnly {
.iter()
.cloned()
.map(|target| {
- RadrootsTransportTargetReceipt::new(
- target,
- RadrootsTransportOutcome::new(RadrootsTransportOutcomeKind::Delivered),
- )
+ DeliveryTargetReceipt::attempted(target, DeliveryOutcome::delivered())
})
.collect();
DeliveryReceipt::for_request(&request, receipts)
@@ -114,6 +113,14 @@ fn target_set() -> RadrootsTransportTargetSet {
.expect("target set")
}
+fn delivery_payload() -> DeliveryPayload {
+ let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
+ let wire = Nip01EventWire::parse_json(raw).expect("wire event");
+ DeliveryPayload::new(
+ SignedEvent::from_wire_verified_id(wire, raw).expect("signed delivery event"),
+ )
+}
+
fn assert_source_dyn_compatible(_: &dyn EventSource) {}
fn assert_sink_dyn_compatible(_: &dyn EventSink) {}
@@ -141,15 +148,16 @@ fn source_only_and_sink_only_implementations_are_independently_dispatchable() {
sink.deliver(
DeliveryRequest::new(
"deliver-1",
- RadrootsTransportPayload::opaque_bytes("spi", [1]).expect("payload"),
+ delivery_payload(),
target_set(),
- RadrootsTransportSatisfactionPolicy::all_delivered(),
+ SatisfactionPolicy::new(SatisfactionClass::Delivered, TargetPolicy::all()),
+ 1_700_000_100_000,
)
.expect("delivery request"),
),
)
.expect("delivery receipt");
- assert_eq!(receipt.request_id(), "deliver-1");
+ assert_eq!(receipt.request_id().as_str(), "deliver-1");
}
#[test]
diff --git a/crates/transport/tests/transport.rs b/crates/transport/tests/transport.rs
@@ -2101,6 +2101,8 @@ fn every_transport_error_has_a_stable_display_message() {
RadrootsTransportError::EmptyDeliveryRequestId,
RadrootsTransportError::InvalidDeliveryRequestId,
RadrootsTransportError::InvalidDeliveryTimestamp,
+ RadrootsTransportError::InvalidDeliveryDeadline,
+ RadrootsTransportError::InvalidDeliveryOutcome,
RadrootsTransportError::UnexpectedDeliveryTargetReceipt,
RadrootsTransportError::DuplicateDeliveryTargetReceipt,
RadrootsTransportError::MissingDeliveryTargetReceipt,
diff --git a/crates/transport_nostr/src/outbox.rs b/crates/transport_nostr/src/outbox.rs
@@ -802,6 +802,8 @@ fn transport_error_to_relay_error(error: RadrootsTransportError) -> RadrootsRela
| RadrootsTransportError::DuplicateFetchTargetOutcome
| RadrootsTransportError::FetchPageLimitExceeded
| RadrootsTransportError::FetchPageRequestMismatch
+ | RadrootsTransportError::InvalidDeliveryDeadline
+ | RadrootsTransportError::InvalidDeliveryOutcome
| RadrootsTransportError::UnexpectedDeliveryTargetReceipt
| RadrootsTransportError::DuplicateDeliveryTargetReceipt
| RadrootsTransportError::MissingDeliveryTargetReceipt