commit 78e2c5ac569b4ac57da6cb78b522bc82ec3d5c12
parent 6f274725752514996c9fbbbb534410c77e7dbc0e
Author: triesap <tyson@radroots.org>
Date: Mon, 3 Aug 2026 08:47:30 +0000
sync: implement delivery and satisfaction evaluation
- claim bounded outbox work and deliver through the injected sink
- persist exact per-target evidence with request-bound atomic digests
- evaluate accumulated Any, All, quorum, and required-target outcomes
- verify partial success, terminal failure, malformed receipts, and retries
Diffstat:
4 files changed, 649 insertions(+), 51 deletions(-)
diff --git a/crates/storage/src/outbox.rs b/crates/storage/src/outbox.rs
@@ -492,15 +492,6 @@ impl OutboxRecord {
.receipt
.validate_for_request(&self.request)
.map_err(|_| Error::InvalidDeliveryEvidence)?;
- let satisfied = value
- .receipt
- .is_satisfied(&self.request)
- .map_err(|_| Error::InvalidDeliveryEvidence)?;
- let retryable = value
- .receipt
- .target_receipts()
- .iter()
- .any(|receipt| matches!(receipt.outcome().retryability(), Retryability::Retryable));
self.evidence
.extend(
value
@@ -516,13 +507,7 @@ impl OutboxRecord {
}),
);
self.last_attempt = Some(value.attempt);
- self.satisfaction = if satisfied {
- SatisfactionResult::Satisfied
- } else if retryable {
- SatisfactionResult::Pending
- } else {
- SatisfactionResult::Exhausted
- };
+ self.satisfaction = evaluate_satisfaction(&self.request, &self.evidence);
self.stage = match self.satisfaction {
SatisfactionResult::Pending => OutboxStage::Retryable,
SatisfactionResult::Satisfied => OutboxStage::Satisfied,
@@ -874,7 +859,6 @@ fn validate_evidence(
}
let mut previous_recorded_at = created_at_unix_ms;
- let mut latest_receipt = None;
for attempt in 1..=last_attempt.get() {
let mut recorded_at = None;
let receipts = request
@@ -911,29 +895,68 @@ fn validate_evidence(
return Err(Error::CorruptOutboxRecord);
}
previous_recorded_at = recorded_at;
- latest_receipt = Some(
- DeliveryReceipt::for_request(request, receipts)
- .map_err(|_| Error::CorruptOutboxRecord)?,
- );
- }
- let receipt = latest_receipt.ok_or(Error::CorruptOutboxRecord)?;
- let expected = if receipt
- .is_satisfied(request)
- .map_err(|_| Error::CorruptOutboxRecord)?
- {
- SatisfactionResult::Satisfied
- } else if receipt
- .target_receipts()
- .iter()
- .any(|receipt| receipt.outcome().is_retryable())
- {
- SatisfactionResult::Pending
- } else {
- SatisfactionResult::Exhausted
- };
+ DeliveryReceipt::for_request(request, receipts).map_err(|_| Error::CorruptOutboxRecord)?;
+ }
+ let expected = evaluate_satisfaction(request, evidence);
if expected == satisfaction {
Ok(())
} else {
Err(Error::CorruptOutboxRecord)
}
}
+
+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| {
+ 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 if let Some(required) = policy.required_targets() {
+ (
+ required.iter().all(&is_successful),
+ required
+ .iter()
+ .all(|target| is_successful(target) || is_retryable(target)),
+ )
+ } else {
+ (false, false)
+ };
+ if satisfied {
+ SatisfactionResult::Satisfied
+ } else if possible {
+ SatisfactionResult::Pending
+ } else {
+ SatisfactionResult::Exhausted
+ }
+}
diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs
@@ -207,6 +207,8 @@ pub enum Error {
SignerFailed,
SignerDeadlineExceeded,
InvalidSignerOutput,
+ InvalidDeliveryRequest,
+ MissingSink,
}
impl core::fmt::Display for Error {
@@ -234,6 +236,8 @@ impl core::fmt::Display for Error {
Self::SignerFailed => "sync signer did not produce an event",
Self::SignerDeadlineExceeded => "sync signer exceeded its deadline",
Self::InvalidSignerOutput => "sync signer output failed canonical verification",
+ Self::InvalidDeliveryRequest => "sync delivery request is invalid",
+ Self::MissingSink => "sync engine has no event sink",
})
}
}
diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs
@@ -18,12 +18,16 @@ use radroots_storage::{
IdempotencyDigest, IdempotencyKey, JournalStage, JournalState, OperationInstanceId,
PrepareOperation,
},
- outbox::{DeliveryPlanDigest, EnqueueOutboxItem, OutboxItemId, OutboxRecord},
+ outbox::{
+ ClaimOutboxItems, DeliveryAttempt, DeliveryAttemptEvidence, DeliveryPlanDigest,
+ EnqueueOutboxItem, LeaseId, LeaseOwner, OUTBOX_CLAIM_LIMIT_MAX, OutboxItemId, OutboxRecord,
+ },
};
use radroots_transport::{
- DeliveryRequest, Target, TransportId,
+ DeliveryReceipt, DeliveryRequest, Target, TransportId,
+ outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability},
policy::SatisfactionPolicy,
- sink::DeliveryPayload,
+ sink::{DeliveryPayload, DeliveryTargetReceipt},
source::{EventProvenance, ObservedEvent},
target::TargetSet,
};
@@ -35,6 +39,8 @@ use crate::{
policy::{Error, OperationKind, SyncId},
};
+const MAX_DELIVERY_LEASE_MS: u64 = 86_400_000;
+
/// Caller-owned, replay-stable inputs for one outbound operation.
#[derive(Clone)]
pub struct PushRequest {
@@ -131,6 +137,59 @@ impl PushReceipt {
}
}
+/// Bounds and lease authority for one explicit outbox delivery pass.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryRunRequest {
+ owner: LeaseOwner,
+ lease_seed: SyncId,
+ lease_duration_ms: u64,
+ limit: u16,
+}
+
+impl DeliveryRunRequest {
+ pub fn new(
+ owner: LeaseOwner,
+ lease_seed: SyncId,
+ lease_duration_ms: u64,
+ limit: u16,
+ ) -> Result<Self, Error> {
+ if lease_duration_ms == 0
+ || lease_duration_ms > MAX_DELIVERY_LEASE_MS
+ || limit == 0
+ || limit > OUTBOX_CLAIM_LIMIT_MAX
+ {
+ return Err(Error::InvalidDeliveryRequest);
+ }
+ Ok(Self {
+ owner,
+ lease_seed,
+ lease_duration_ms,
+ limit,
+ })
+ }
+}
+
+/// Independent durable outcomes from one bounded delivery pass.
+#[derive(Clone, Debug, Eq, PartialEq)]
+pub struct DeliveryRunReceipt {
+ outcomes: Vec<Result<OutboxRecord, Error>>,
+}
+
+impl DeliveryRunReceipt {
+ pub fn outcomes(&self) -> &[Result<OutboxRecord, Error>] {
+ self.outcomes.as_slice()
+ }
+ pub fn succeeded(&self) -> usize {
+ self.outcomes
+ .iter()
+ .filter(|outcome| outcome.is_ok())
+ .count()
+ }
+ pub fn failed(&self) -> usize {
+ self.outcomes.len() - self.succeeded()
+ }
+}
+
impl Engine {
/// Authorizes, signs, verifies, and durably enqueues one outbound event.
///
@@ -295,6 +354,157 @@ impl Engine {
replay: false,
})
}
+
+ /// Claims and delivers at most the caller-bounded number of outbox items.
+ pub async fn deliver_pending(
+ &self,
+ request: DeliveryRunRequest,
+ ) -> Result<DeliveryRunReceipt, Error> {
+ if request.lease_duration_ms > self.deadlines.timeout_ms(OperationKind::Deliver) {
+ return Err(Error::InvalidDeliveryRequest);
+ }
+ let sink = self.sink.as_deref().ok_or(Error::MissingSink)?;
+ let now = self.clock.now_unix_ms()?;
+ let expires = now
+ .checked_add(request.lease_duration_ms)
+ .ok_or(Error::DeadlineOverflow)?;
+ let claimed = Outbox::claim(
+ self.storage.as_ref(),
+ ClaimOutboxItems::new(
+ request.owner,
+ LeaseId::new(*request.lease_seed.as_bytes()).map_err(map_storage_error)?,
+ now,
+ expires,
+ request.limit,
+ )
+ .map_err(map_storage_error)?,
+ )
+ .await
+ .map_err(map_storage_error)?;
+ let mut outcomes = Vec::with_capacity(claimed.len());
+ for item in claimed {
+ let delivery_request = item.record().request().clone();
+ let attempted_at = self.clock.now_unix_ms()?;
+ let receipt = if attempted_at >= delivery_request.deadline_unix_ms() {
+ synthetic_receipt(&delivery_request, false)?
+ } else {
+ match sink.deliver(delivery_request.clone()).await {
+ Ok(receipt) => receipt,
+ Err(_) => synthetic_receipt(&delivery_request, true)?,
+ }
+ };
+ if receipt.validate_for_request(&delivery_request).is_err() {
+ let released_at = self.clock.now_unix_ms()?;
+ let released = Outbox::release(
+ self.storage.as_ref(),
+ item.record().item_id(),
+ item.lease().id(),
+ item.record().revision(),
+ released_at,
+ None,
+ )
+ .await
+ .map_err(map_storage_error);
+ outcomes.push(match released {
+ Ok(_) => Err(Error::InvalidDeliveryRequest),
+ Err(error) => Err(error),
+ });
+ continue;
+ }
+ let attempt = DeliveryAttempt::new(
+ item.record()
+ .last_attempt()
+ .map_or(1, |attempt| attempt.get().saturating_add(1)),
+ )
+ .map_err(map_storage_error)?;
+ let evidence = DeliveryAttemptEvidence::new(
+ item.record().item_id(),
+ item.lease().id(),
+ item.record().revision(),
+ attempt,
+ receipt,
+ attempted_at,
+ )
+ .map_err(map_storage_error)?;
+ let digest = delivery_evidence_digest(&evidence);
+ let commit = AtomicCommit::new(
+ next_commit_id(self, OperationKind::Deliver)?,
+ digest,
+ attempted_at,
+ AtomicWorkflow::Delivered(Box::new(evidence)),
+ )
+ .map_err(map_storage_error)?;
+ let outcome = match self.storage.commit(commit).await {
+ Ok(receipt) => match receipt.outcome() {
+ AtomicCommitOutcome::Delivered { outbox } => Ok((**outbox).clone()),
+ _ => Err(Error::StorageFailed),
+ },
+ Err(error) => Err(map_storage_error(error)),
+ };
+ outcomes.push(outcome);
+ }
+ Ok(DeliveryRunReceipt { outcomes })
+ }
+}
+
+fn synthetic_receipt(request: &DeliveryRequest, retryable: bool) -> Result<DeliveryReceipt, Error> {
+ let outcome = if retryable {
+ DeliveryOutcome::unavailable()
+ } else {
+ DeliveryOutcome::failed(Retryability::Terminal)
+ .map_err(|_| Error::InvalidDeliveryRequest)?
+ };
+ let targets = request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .map(|target| DeliveryTargetReceipt::skipped(target, outcome.clone()))
+ .collect::<Result<Vec<_>, _>>()
+ .map_err(|_| Error::InvalidDeliveryRequest)?;
+ DeliveryReceipt::for_request(request, targets).map_err(|_| Error::InvalidDeliveryRequest)
+}
+
+fn delivery_evidence_digest(evidence: &DeliveryAttemptEvidence) -> AtomicCommitDigest {
+ let mut hasher = Sha256::new();
+ hash_field(&mut hasher, b"radroots.sync.delivery-evidence.v1");
+ hash_field(&mut hasher, evidence.item_id().as_bytes());
+ hash_field(&mut hasher, evidence.lease_id().as_bytes());
+ hasher.update(evidence.expected_revision().get().to_be_bytes());
+ hasher.update(evidence.attempt().get().to_be_bytes());
+ hasher.update(evidence.recorded_at_unix_ms().to_be_bytes());
+ hash_field(
+ &mut hasher,
+ evidence.receipt().request_id().as_str().as_bytes(),
+ );
+ for target in evidence.receipt().target_receipts() {
+ hash_field(
+ &mut hasher,
+ target.target().fingerprint().as_str().as_bytes(),
+ );
+ hasher.update([u8::from(target.was_attempted())]);
+ hasher.update([match target.outcome().kind() {
+ DeliveryOutcomeKind::Accepted => 0,
+ DeliveryOutcomeKind::Delivered => 1,
+ DeliveryOutcomeKind::Rejected => 2,
+ DeliveryOutcomeKind::Unavailable => 3,
+ DeliveryOutcomeKind::Failed => 4,
+ }]);
+ hasher.update([match target.outcome().retryability() {
+ Retryability::NotApplicable => 0,
+ Retryability::Retryable => 1,
+ Retryability::Terminal => 2,
+ }]);
+ hash_field(
+ &mut hasher,
+ target.outcome().code().unwrap_or_default().as_bytes(),
+ );
+ hash_field(
+ &mut hasher,
+ target.outcome().message().unwrap_or_default().as_bytes(),
+ );
+ }
+ AtomicCommitDigest::new(hasher.finalize().into())
}
fn outbound_admission(
diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs
@@ -1,6 +1,9 @@
-use std::sync::{
- Arc,
- atomic::{AtomicU64, AtomicUsize, Ordering},
+use std::{
+ collections::VecDeque,
+ sync::{
+ Arc, Mutex,
+ atomic::{AtomicU64, AtomicUsize, Ordering},
+ },
};
use futures::{FutureExt, task::noop_waker_ref};
@@ -15,16 +18,19 @@ use radroots_storage::{
event::{EventQuery, EventQueryBounds, SourceGeneration},
journal::{IdempotencyKey, JournalStage, OperationInstanceId},
memory::MemoryStorage,
- outbox::OutboxStage,
+ outbox::{LeaseOwner, OutboxStage, SatisfactionResult},
};
use radroots_sync::{
Engine, PushRequest,
policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId},
+ push::DeliveryRunRequest,
};
use radroots_transport::{
DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus, Target,
TargetSet, TransportId,
+ outcome::DeliveryOutcome,
policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
+ sink::DeliveryTargetReceipt,
};
use secp256k1::{Keypair, Message, Secp256k1, SecretKey};
@@ -44,6 +50,67 @@ impl EventSink for MockSink {
}
}
+enum DeliveryBehavior {
+ Outcomes(Vec<DeliveryOutcome>),
+ AdapterError,
+ MismatchedRequest,
+}
+
+struct ScriptedSink {
+ behaviors: Mutex<VecDeque<DeliveryBehavior>>,
+ requests: Mutex<Vec<DeliveryRequest>>,
+}
+
+impl ScriptedSink {
+ fn new(behaviors: impl IntoIterator<Item = DeliveryBehavior>) -> Self {
+ Self {
+ behaviors: Mutex::new(behaviors.into_iter().collect()),
+ requests: Mutex::new(Vec::new()),
+ }
+ }
+}
+
+impl EventSink for ScriptedSink {
+ fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
+ Box::pin(async { unreachable!("delivery does not inspect sink status") })
+ }
+
+ fn deliver(
+ &self,
+ request: DeliveryRequest,
+ ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> {
+ self.requests
+ .lock()
+ .expect("scripted request lock")
+ .push(request.clone());
+ let behavior = self
+ .behaviors
+ .lock()
+ .expect("scripted behavior lock")
+ .pop_front()
+ .expect("scripted delivery behavior");
+ Box::pin(async move {
+ match behavior {
+ DeliveryBehavior::Outcomes(outcomes) => receipt(&request, outcomes),
+ DeliveryBehavior::AdapterError => Err(TransportError::UnsupportedOperation),
+ DeliveryBehavior::MismatchedRequest => {
+ let mismatched = DeliveryRequest::new(
+ "mismatched-request",
+ request.payload().clone(),
+ request.target_set().clone(),
+ request.satisfaction().clone(),
+ request.deadline_unix_ms(),
+ )?;
+ receipt(
+ &mismatched,
+ vec![DeliveryOutcome::accepted(); mismatched.target_set().len()],
+ )
+ }
+ }
+ })
+ }
+}
+
struct TestClock(AtomicU64);
impl Clock for TestClock {
@@ -149,6 +216,20 @@ fn signed_event(request: &SignRequest) -> SignedEvent {
}
fn request(operation_byte: u8, relay: &str) -> PushRequest {
+ request_with_policy(
+ operation_byte,
+ &[relay],
+ SatisfactionClass::Accepted,
+ TargetPolicy::any(),
+ )
+}
+
+fn request_with_policy(
+ operation_byte: u8,
+ relays: &[&str],
+ class: SatisfactionClass,
+ target_policy: TargetPolicy,
+) -> PushRequest {
let pubkey = public_key_hex();
let draft = EventDraft::new(
"radroots.social.geochat.v1",
@@ -170,32 +251,68 @@ fn request(operation_byte: u8, relay: &str) -> PushRequest {
IdempotencyKey::parse(format!("push-{operation_byte}")).expect("idempotency key"),
actor,
draft,
- TargetSet::new(vec![
- Target::new(TransportId::NOSTR, relay).expect("target"),
- ])
+ TargetSet::new(
+ relays
+ .iter()
+ .map(|relay| Target::new(TransportId::NOSTR, *relay).expect("target"))
+ .collect(),
+ )
.expect("targets"),
- SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::any()),
+ SatisfactionPolicy::new(class, target_policy),
CancellationPolicy::PreservePublishedRequest,
)
.expect("push request")
}
+fn receipt(
+ request: &DeliveryRequest,
+ outcomes: Vec<DeliveryOutcome>,
+) -> Result<DeliveryReceipt, TransportError> {
+ let targets = request
+ .target_set()
+ .targets()
+ .iter()
+ .cloned()
+ .zip(outcomes)
+ .map(|(target, outcome)| DeliveryTargetReceipt::attempted(target, outcome))
+ .collect();
+ DeliveryReceipt::for_request(request, targets)
+}
+
fn setup_engine(signer: Arc<MockSigner>) -> (Engine, Arc<MemoryStorage>) {
+ setup_engine_with_sink(signer, Arc::new(MockSink)).0
+}
+
+fn setup_engine_with_sink(
+ signer: Arc<MockSigner>,
+ sink: Arc<dyn EventSink>,
+) -> ((Engine, Arc<MemoryStorage>), Arc<TestClock>) {
let storage = Arc::new(MemoryStorage::new(
SourceGeneration::new([6; 32]).expect("generation"),
));
let capability: Arc<dyn Storage> = storage.clone();
+ let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
let engine = Engine::builder(
capability,
- Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))),
+ clock.clone(),
Arc::new(TestIds(AtomicU64::new(10))),
DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"),
)
- .sink(Arc::new(MockSink))
+ .sink(sink)
.signer(signer)
.build()
.expect("engine");
- (engine, storage)
+ ((engine, storage), clock)
+}
+
+fn delivery_run(seed: u8, limit: u16) -> DeliveryRunRequest {
+ DeliveryRunRequest::new(
+ LeaseOwner::parse("sync-delivery-test").expect("lease owner"),
+ SyncId::new([seed; 16]).expect("lease seed"),
+ 1_000,
+ limit,
+ )
+ .expect("delivery run")
}
#[test]
@@ -298,3 +415,247 @@ fn idempotency_conflict_and_cancellation_before_commit_fail_closed() {
.is_none()
);
}
+
+#[test]
+fn delivery_evaluates_any_all_quorum_required_and_partial_outcomes() {
+ let sink = Arc::new(ScriptedSink::new([
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::accepted(),
+ DeliveryOutcome::rejected(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::accepted(),
+ DeliveryOutcome::delivered(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::delivered(),
+ DeliveryOutcome::accepted(),
+ DeliveryOutcome::rejected(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::rejected(),
+ DeliveryOutcome::accepted(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::accepted(),
+ DeliveryOutcome::unavailable(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::rejected(),
+ DeliveryOutcome::unavailable(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::rejected(),
+ DeliveryOutcome::unavailable(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::unavailable(),
+ DeliveryOutcome::rejected(),
+ ]),
+ ]));
+ let signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix: 1_800_000_200,
+ }));
+ let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone());
+ let two = ["wss://one.example", "wss://two.example"];
+ let three = [
+ "wss://one.example",
+ "wss://two.example",
+ "wss://three.example",
+ ];
+ let required = Target::new(TransportId::NOSTR, two[1])
+ .expect("required target")
+ .fingerprint()
+ .clone();
+ let plans = [
+ request_with_policy(11, &two, SatisfactionClass::Accepted, TargetPolicy::any()),
+ request_with_policy(12, &two, SatisfactionClass::Accepted, TargetPolicy::all()),
+ request_with_policy(
+ 13,
+ &three,
+ SatisfactionClass::Accepted,
+ TargetPolicy::quorum(2).expect("quorum"),
+ ),
+ request_with_policy(
+ 14,
+ &two,
+ SatisfactionClass::Accepted,
+ TargetPolicy::required(vec![required]).expect("required policy"),
+ ),
+ request_with_policy(15, &two, SatisfactionClass::Accepted, TargetPolicy::all()),
+ request_with_policy(16, &two, SatisfactionClass::Accepted, TargetPolicy::all()),
+ request_with_policy(17, &two, SatisfactionClass::Accepted, TargetPolicy::any()),
+ request_with_policy(
+ 18,
+ &two,
+ SatisfactionClass::Accepted,
+ TargetPolicy::required(vec![
+ Target::new(TransportId::NOSTR, two[1])
+ .expect("terminal required target")
+ .fingerprint()
+ .clone(),
+ ])
+ .expect("terminal required policy"),
+ ),
+ ];
+ for plan in plans {
+ block_on(engine.sign_and_enqueue(plan)).expect("enqueue delivery plan");
+ }
+
+ let delivered = block_on(engine.deliver_pending(delivery_run(31, 8))).expect("deliver batch");
+ assert_eq!(delivered.succeeded(), 8);
+ assert_eq!(delivered.failed(), 0);
+ let records: Vec<_> = delivered
+ .outcomes()
+ .iter()
+ .map(|outcome| outcome.as_ref().expect("durable outcome"))
+ .collect();
+ assert_eq!(records[0].stage(), OutboxStage::Satisfied);
+ assert_eq!(records[1].stage(), OutboxStage::Satisfied);
+ assert_eq!(records[2].stage(), OutboxStage::Satisfied);
+ assert_eq!(records[3].stage(), OutboxStage::Satisfied);
+ assert_eq!(records[4].stage(), OutboxStage::Retryable);
+ assert_eq!(records[4].satisfaction(), SatisfactionResult::Pending);
+ assert_eq!(records[5].stage(), OutboxStage::Exhausted);
+ assert_eq!(records[6].stage(), OutboxStage::Retryable);
+ assert_eq!(records[7].stage(), OutboxStage::Exhausted);
+ for record in records {
+ assert_eq!(record.evidence().len(), record.request().target_set().len());
+ let expected: Vec<_> = record
+ .request()
+ .target_set()
+ .targets()
+ .iter()
+ .map(|target| target.fingerprint())
+ .collect();
+ let actual: Vec<_> = record
+ .evidence()
+ .iter()
+ .map(|evidence| evidence.target())
+ .collect();
+ assert_eq!(actual, expected);
+ }
+ assert_eq!(sink.requests.lock().expect("request log").len(), 8);
+}
+
+#[test]
+fn transport_failure_is_durable_and_retry_preserves_the_exact_plan() {
+ let sink = Arc::new(ScriptedSink::new([
+ DeliveryBehavior::AdapterError,
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::accepted(),
+ DeliveryOutcome::unavailable(),
+ ]),
+ DeliveryBehavior::Outcomes(vec![
+ DeliveryOutcome::unavailable(),
+ DeliveryOutcome::accepted(),
+ ]),
+ ]));
+ let signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix: 1_800_000_200,
+ }));
+ let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone());
+ block_on(engine.sign_and_enqueue(request_with_policy(
+ 21,
+ &["wss://one.example", "wss://two.example"],
+ SatisfactionClass::Accepted,
+ TargetPolicy::all(),
+ )))
+ .expect("enqueue retry plan");
+
+ let first = block_on(engine.deliver_pending(delivery_run(41, 1))).expect("first delivery");
+ let first = first.outcomes()[0].as_ref().expect("durable failure");
+ assert_eq!(first.stage(), OutboxStage::Retryable);
+ assert_eq!(first.evidence().len(), 2);
+ assert!(first.evidence().iter().all(|evidence| {
+ !evidence.was_attempted() && evidence.outcome() == &DeliveryOutcome::unavailable()
+ }));
+
+ let second = block_on(engine.deliver_pending(delivery_run(42, 1))).expect("partial retry");
+ let second = second.outcomes()[0].as_ref().expect("durable retry");
+ assert_eq!(second.stage(), OutboxStage::Retryable);
+ assert_eq!(second.last_attempt().expect("attempt").get(), 2);
+
+ let third = block_on(engine.deliver_pending(delivery_run(43, 1))).expect("completed retry");
+ let third = third.outcomes()[0].as_ref().expect("durable retry");
+ assert_eq!(third.stage(), OutboxStage::Satisfied);
+ assert_eq!(third.last_attempt().expect("attempt").get(), 3);
+ assert_eq!(third.evidence().len(), 6);
+ assert_eq!(
+ third
+ .latest_target_evidence(third.request().target_set().targets()[0].fingerprint())
+ .expect("latest first target evidence")
+ .outcome(),
+ &DeliveryOutcome::unavailable()
+ );
+ assert!(
+ third.evidence()[2]
+ .outcome()
+ .satisfies(SatisfactionClass::Accepted)
+ );
+ let requests = sink.requests.lock().expect("request log");
+ assert_eq!(requests.len(), 3);
+ assert_eq!(requests[0], requests[1]);
+ assert_eq!(requests[1], requests[2]);
+}
+
+#[test]
+fn malformed_receipts_release_work_and_expired_plans_terminalize() {
+ let malformed_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::MismatchedRequest]));
+ let signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix: 1_800_000_200,
+ }));
+ let ((engine, storage), _) = setup_engine_with_sink(signer, malformed_sink);
+ let enqueued = block_on(engine.sign_and_enqueue(request(31, "wss://one.example")))
+ .expect("enqueue malformed receipt plan");
+ let malformed =
+ block_on(engine.deliver_pending(delivery_run(51, 1))).expect("malformed delivery pass");
+ assert_eq!(malformed.outcomes(), &[Err(Error::InvalidDeliveryRequest)]);
+ let released = block_on(Outbox::item(&*storage, enqueued.outbox().item_id()))
+ .expect("released lookup")
+ .expect("released record");
+ assert_eq!(released.stage(), OutboxStage::Pending);
+ assert!(released.lease().is_none());
+ assert!(released.evidence().is_empty());
+
+ let expired_sink = Arc::new(ScriptedSink::new([]));
+ let signer = Arc::new(MockSigner::new(SignBehavior::Success {
+ completed_at_unix: 1_800_000_200,
+ }));
+ let ((engine, _), clock) = setup_engine_with_sink(signer, expired_sink.clone());
+ let enqueued = block_on(engine.sign_and_enqueue(request(32, "wss://one.example")))
+ .expect("enqueue expiring plan");
+ clock.0.store(
+ enqueued.outbox().request().deadline_unix_ms() + 1,
+ Ordering::Relaxed,
+ );
+ let expired =
+ block_on(engine.deliver_pending(delivery_run(52, 1))).expect("expired delivery pass");
+ let expired = expired.outcomes()[0]
+ .as_ref()
+ .expect("durable terminal state");
+ assert_eq!(expired.stage(), OutboxStage::Exhausted);
+ assert_eq!(expired.satisfaction(), SatisfactionResult::Exhausted);
+ assert!(expired.evidence()[0].outcome().is_terminal());
+ assert!(
+ expired_sink
+ .requests
+ .lock()
+ .expect("request log")
+ .is_empty()
+ );
+}
+
+#[test]
+fn delivery_run_rejects_unbounded_claims() {
+ let owner = LeaseOwner::parse("sync-delivery-test").expect("lease owner");
+ let seed = SyncId::new([71; 16]).expect("lease seed");
+ assert_eq!(
+ DeliveryRunRequest::new(owner.clone(), seed, 0, 1),
+ Err(Error::InvalidDeliveryRequest)
+ );
+ assert_eq!(
+ DeliveryRunRequest::new(owner, seed, 1_000, 0),
+ Err(Error::InvalidDeliveryRequest)
+ );
+}