lib

Core libraries for Radroots
git clone https://radroots.dev/git/lib.git
Log | Files | Refs | README

commit 0aa05380aa18a504ee2758daac4853b872639541
parent 934f377e5dc22771d044a99699f7d3a6078b517a
Author: triesap <tyson@radroots.org>
Date:   Wed,  5 Aug 2026 04:08:15 +0000

refactor(sync): execute durable delivery plans

- Claim and fence authored delivery execution.
- Persist typed outcomes, evidence, and retry timing.
- Report complete operation settlement without hiding failures.
- Cover memory, SQLite, recovery, and malformed adapter paths.


Diffstat:
Mcrates/storage/src/authored.rs | 99+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
Mcrates/storage/src/authored_delivery.rs | 43+++++++++++++++++++++++++++++++++++++++++++
Mcrates/sync/src/lib.rs | 3++-
Mcrates/sync/src/policy.rs | 2++
Mcrates/sync/src/push.rs | 266++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++---
Mcrates/sync/tests/push_enqueue.rs | 250+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
6 files changed, 653 insertions(+), 10 deletions(-)

diff --git a/crates/storage/src/authored.rs b/crates/storage/src/authored.rs @@ -6,7 +6,11 @@ use radroots_event_codec::authoring::{AuthoredEventPlan, HistoricalPlanIntegrity use sha2::{Digest, Sha256}; use std::{collections::BTreeSet, string::String, vec::Vec}; -use crate::{Error, journal::OperationInstanceId}; +use crate::{ + Error, + authored_delivery::{AuthoredDeliveryPlan, AuthoredDeliveryState}, + journal::OperationInstanceId, +}; pub const AUTHORED_OPERATION_ARTIFACTS_MAX: usize = 256; pub const WORK_CLAIM_OWNER_MAX_BYTES: usize = 128; @@ -1214,6 +1218,13 @@ pub struct OperationSettlement { indeterminate: u16, failed_terminal: u16, cancelled: u16, + delivery_plans: u16, + delivery_satisfied: u16, + delivery_pending: u16, + delivery_retryable: u16, + delivery_exhausted: u16, + delivery_failed_terminal: u16, + delivery_cancelled: u16, } impl OperationSettlement { @@ -1221,6 +1232,14 @@ impl OperationSettlement { operation: &AuthoredOperation, artifacts: &[AuthoredArtifact], ) -> Result<Self, Error> { + Self::evaluate_complete(operation, artifacts, &[]) + } + + pub fn evaluate_complete( + operation: &AuthoredOperation, + artifacts: &[AuthoredArtifact], + delivery_plans: &[AuthoredDeliveryPlan], + ) -> Result<Self, Error> { if artifacts.len() != operation.artifact_ids.len() || artifacts.iter().enumerate().any(|(ordinal, artifact)| { artifact.operation_id != operation.operation_id @@ -1231,6 +1250,22 @@ impl OperationSettlement { { return Err(Error::InvalidAuthoredOperation); } + let artifact_ids = artifacts + .iter() + .map(AuthoredArtifact::artifact_id) + .collect::<BTreeSet<_>>(); + if delivery_plans + .iter() + .map(AuthoredDeliveryPlan::plan_id) + .collect::<BTreeSet<_>>() + .len() + != delivery_plans.len() + || delivery_plans + .iter() + .any(|plan| !artifact_ids.contains(&plan.artifact_id()) || plan.validate().is_err()) + { + return Err(Error::InvalidAuthoredOperation); + } let mut settlement = Self { artifacts: u16::try_from(artifacts.len()) .map_err(|_| Error::InvalidAuthoredOperation)?, @@ -1241,6 +1276,14 @@ impl OperationSettlement { indeterminate: 0, failed_terminal: 0, cancelled: 0, + delivery_plans: u16::try_from(delivery_plans.len()) + .map_err(|_| Error::InvalidAuthoredOperation)?, + delivery_satisfied: 0, + delivery_pending: 0, + delivery_retryable: 0, + delivery_exhausted: 0, + delivery_failed_terminal: 0, + delivery_cancelled: 0, }; for artifact in artifacts { if artifact.signed.is_some() { @@ -1264,6 +1307,18 @@ impl OperationSettlement { }, } } + for plan in delivery_plans { + match plan.state() { + AuthoredDeliveryState::Pending => settlement.delivery_pending += 1, + AuthoredDeliveryState::Retryable => settlement.delivery_retryable += 1, + AuthoredDeliveryState::Satisfied => settlement.delivery_satisfied += 1, + AuthoredDeliveryState::Exhausted => settlement.delivery_exhausted += 1, + AuthoredDeliveryState::FailedTerminal => { + settlement.delivery_failed_terminal += 1; + } + AuthoredDeliveryState::Cancelled => settlement.delivery_cancelled += 1, + } + } Ok(settlement) } @@ -1291,8 +1346,48 @@ impl OperationSettlement { pub const fn cancelled(self) -> u16 { self.cancelled } + pub const fn delivery_plans(self) -> u16 { + self.delivery_plans + } + pub const fn delivery_satisfied(self) -> u16 { + self.delivery_satisfied + } + pub const fn delivery_pending(self) -> u16 { + self.delivery_pending + } + pub const fn delivery_retryable(self) -> u16 { + self.delivery_retryable + } + pub const fn delivery_exhausted(self) -> u16 { + self.delivery_exhausted + } + pub const fn delivery_failed_terminal(self) -> u16 { + self.delivery_failed_terminal + } + pub const fn delivery_cancelled(self) -> u16 { + self.delivery_cancelled + } pub const fn is_settled(self) -> bool { - self.pending == 0 && self.retryable == 0 && self.indeterminate == 0 + self.pending == 0 + && self.retryable == 0 + && self.indeterminate == 0 + && self.delivery_pending == 0 + && self.delivery_retryable == 0 + } + pub const fn has_failures(self) -> bool { + self.indeterminate != 0 + || self.failed_terminal != 0 + || self.cancelled != 0 + || self.delivery_exhausted != 0 + || self.delivery_failed_terminal != 0 + || self.delivery_cancelled != 0 + } + pub const fn is_successful(self) -> bool { + self.is_settled() + && !self.has_failures() + && self.signed == self.artifacts + && self.admitted == self.artifacts + && self.delivery_satisfied == self.delivery_plans } } diff --git a/crates/storage/src/authored_delivery.rs b/crates/storage/src/authored_delivery.rs @@ -429,6 +429,10 @@ impl AuthoredDeliveryPlan { .validate_for_request(request) .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; let satisfaction = self.evaluate_with(DeliveryAttemptOutcome::Receipt(receipt.clone()))?; + let at_attempt_limit = self.attempt_count.saturating_add(1) == DELIVERY_PLAN_ATTEMPTS_MAX; + if satisfaction == SatisfactionState::Pending && at_attempt_limit && retry.is_some() { + return Err(Error::InvalidRetrySchedule); + } let (state, last_failure) = match satisfaction { SatisfactionState::Satisfied if retry.is_none() => { (AuthoredDeliveryState::Satisfied, None) @@ -436,6 +440,16 @@ impl AuthoredDeliveryPlan { SatisfactionState::Exhausted if retry.is_none() => { (AuthoredDeliveryState::Exhausted, None) } + SatisfactionState::Pending if at_attempt_limit && retry.is_none() => ( + AuthoredDeliveryState::Exhausted, + Some(WorkFailure::new( + "delivery_attempt_limit", + WorkPhase::Delivery, + FailureClass::Terminal, + None, + None, + )?), + ), SatisfactionState::Pending => { let schedule = retry.as_ref().ok_or(Error::InvalidRetrySchedule)?; if schedule.failure().phase() != WorkPhase::Delivery { @@ -479,6 +493,10 @@ impl AuthoredDeliveryPlan { .map_err(|_| Error::InvalidAuthoredDeliveryPlan)?; let outcome = DeliveryAttemptOutcome::SinkFailure(failure.clone()); let satisfaction = self.evaluate_with(outcome.clone())?; + let at_attempt_limit = self.attempt_count.saturating_add(1) == DELIVERY_PLAN_ATTEMPTS_MAX; + if satisfaction == SatisfactionState::Pending && at_attempt_limit && retry.is_some() { + return Err(Error::InvalidRetrySchedule); + } let typed = WorkFailure::new( failure.code(), WorkPhase::Delivery, @@ -499,6 +517,18 @@ impl AuthoredDeliveryPlan { return Err(Error::InvalidRetrySchedule); } (AuthoredDeliveryState::Exhausted, None, Some(typed)) + } else if at_attempt_limit && retry.is_none() { + ( + AuthoredDeliveryState::Exhausted, + None, + Some(WorkFailure::new( + "delivery_attempt_limit", + WorkPhase::Delivery, + FailureClass::Terminal, + None, + None, + )?), + ) } else if failure.retryability() == Retryability::Retryable { let schedule = retry.ok_or(Error::InvalidRetrySchedule)?; if schedule.failure() != &typed { @@ -698,6 +728,19 @@ impl AuthoredDeliveryPlan { pub const fn revision(&self) -> NonZeroU64 { self.revision } + + /// Evaluates one prospective attempt against all durable prior evidence. + pub fn evaluate_next_attempt( + &self, + outcome: &DeliveryAttemptOutcome, + ) -> Result<SatisfactionState, Error> { + let request = self + .request + .as_ref() + .ok_or(Error::InvalidAuthoredDeliveryPlan)?; + outcome.validate_for(request)?; + self.evaluate_with(outcome.clone()) + } } #[cfg(feature = "serde")] diff --git a/crates/sync/src/lib.rs b/crates/sync/src/lib.rs @@ -13,6 +13,7 @@ pub use engine::Engine; pub use policy::Error; pub use pull::{PullReceipt, PullRequest}; pub use push::{ - AdmissionRunReceipt, PushPreparation, PushReceipt, PushRequest, PushStatus, SigningRunReceipt, + AdmissionRunReceipt, DeliveryExecutionReceipt, PushPreparation, PushReceipt, PushRequest, + PushStatus, SigningRunReceipt, }; pub use status::SyncStatus; diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -239,6 +239,7 @@ pub enum Error { InvalidSignerOutput, AdmissionFailed, InvalidDeliveryRequest, + DeliveryDeferred, MissingSink, InvalidStatusRequest, } @@ -278,6 +279,7 @@ impl core::fmt::Display for Error { Self::InvalidSignerOutput => "sync signer output failed canonical verification", Self::AdmissionFailed => "sync local admission did not complete", Self::InvalidDeliveryRequest => "sync delivery request is invalid", + Self::DeliveryDeferred => "sync delivery retry is not yet eligible", Self::MissingSink => "sync engine has no event sink", Self::InvalidStatusRequest => "sync status request is invalid", }) diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -20,14 +20,18 @@ use radroots_storage::{ }, authored::{ AdmissionState, AuthoredArtifact, AuthoredArtifactId, AuthoredOperation, FailureClass, - RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, + OperationSettlement, RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, }, authored_atomic::{ - ApplyAdmissionResult, ApplySignedArtifact, ApplyWorkFailure, AuthoredAtomicCommand, - AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, - ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, WorkFence, + ApplyAdmissionResult, ApplyDeliveryAttempt, ApplySignedArtifact, ApplyWorkFailure, + AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, + CancelAuthoredWork, ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, + WorkFence, + }, + authored_delivery::{ + AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, + DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, }, - authored_delivery::{AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId}, event::{AdmissionDisposition, EventAdmission, EventStore}, journal::{ IdempotencyDigest, IdempotencyKey, JournalStage, JournalState, OperationInstanceId, @@ -39,9 +43,9 @@ use radroots_storage::{ }, }; use radroots_transport::{ - DeliveryReceipt, DeliveryRequest, Target, TransportId, + DeliveryReceipt, DeliveryRequest, SinkFailure, Target, TransportId, outcome::{DeliveryOutcome, DeliveryOutcomeKind, Retryability}, - policy::SatisfactionPolicy, + policy::{SatisfactionPolicy, SatisfactionState}, sink::{DeliveryPayload, DeliveryTargetReceipt}, source::{EventProvenance, ObservedEvent}, target::TargetSet, @@ -59,6 +63,7 @@ const SIGNING_CLAIM_OWNER_EXACT: &str = "radroots-sync-signing-exact"; const SIGNING_CLAIM_OWNER_LOCAL: &str = "radroots-sync-signing-local"; const SIGNING_CLAIM_OWNER_NON_REPLAYABLE: &str = "radroots-sync-signing-non-replayable"; const ADMISSION_CLAIM_OWNER: &str = "radroots-sync-admission"; +const DELIVERY_CLAIM_OWNER: &str = "radroots-sync-delivery"; const WORK_RETRY_DELAY_MS: u64 = 1_000; /// Caller-owned, replay-stable inputs for one outbound operation. @@ -176,6 +181,7 @@ pub struct PushStatus { operation: AuthoredOperation, artifact: AuthoredArtifact, delivery_plan: AuthoredDeliveryPlan, + settlement: OperationSettlement, } impl PushStatus { @@ -188,6 +194,9 @@ impl PushStatus { pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { &self.delivery_plan } + pub const fn settlement(&self) -> OperationSettlement { + self.settlement + } } /// Durable result of one bounded signing execution. @@ -213,6 +222,22 @@ pub struct AdmissionRunReceipt { replay: bool, } +/// Durable result of one bounded authored-delivery execution. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct DeliveryExecutionReceipt { + plan: AuthoredDeliveryPlan, + replay: bool, +} + +impl DeliveryExecutionReceipt { + pub const fn plan(&self) -> &AuthoredDeliveryPlan { + &self.plan + } + pub const fn is_replay(&self) -> bool { + self.replay + } +} + impl AdmissionRunReceipt { pub const fn artifact(&self) -> &AuthoredArtifact { &self.artifact @@ -394,10 +419,17 @@ impl Engine { { return Err(Error::StorageFailed); } + let settlement = OperationSettlement::evaluate_complete( + &operation, + core::slice::from_ref(&artifact), + core::slice::from_ref(&delivery_plan), + ) + .map_err(map_storage_error)?; Ok(Some(PushStatus { operation, artifact, delivery_plan, + settlement, })) } @@ -695,6 +727,141 @@ impl Engine { }) } + /// Claims and executes one durable authored delivery plan. + /// + /// Terminal plans replay without invoking the sink. Retry timing and the + /// request deadline are enforced from durable intent before any adapter + /// call, and every adapter result is fenced into the plan as evidence. + pub async fn deliver_push( + &self, + operation_id: SyncId, + ) -> Result<DeliveryExecutionReceipt, Error> { + let sink = self.sink.as_deref().ok_or(Error::MissingSink)?; + let status = self + .push_status(operation_id) + .await? + .ok_or(Error::StorageFailed)?; + if status.artifact.signing_state() != SigningState::Signed { + return Err(Error::InvalidSignerOutput); + } + if !status.artifact.admission_state().is_admitted() { + return Err(Error::AdmissionFailed); + } + let plan = status.delivery_plan; + if plan.state().is_terminal() { + return Ok(DeliveryExecutionReceipt { plan, replay: true }); + } + let request = plan.request().cloned().ok_or(Error::InvalidSignerOutput)?; + let now = self.clock.now_unix_ms()?.max(plan.updated_at_unix_ms()); + if plan + .claim_evidence() + .is_some_and(|claim| now < claim.expires_at_unix_ms()) + { + return Err(Error::WorkClaimConflict); + } + if plan + .retry() + .is_some_and(|retry| now < retry.not_before_unix_ms()) + { + return Err(Error::DeliveryDeferred); + } + let (claimed, fence) = self.claim_delivery_plan(plan, now).await?; + let execution_started_at = self.clock.now_unix_ms()?.max(claimed.updated_at_unix_ms()); + let result = if execution_started_at >= request.deadline_unix_ms() { + DeliveryAttemptOutcome::SinkFailure( + SinkFailure::for_request( + &request, + "delivery_deadline_exceeded", + Retryability::Terminal, + None, + None, + Vec::new(), + ) + .map_err(|_| Error::InvalidDeliveryRequest)?, + ) + } else { + match sink.deliver(request.clone()).await { + Ok(receipt) if receipt.validate_for_request(&request).is_ok() => { + DeliveryAttemptOutcome::Receipt(receipt) + } + Ok(_) => { + DeliveryAttemptOutcome::SinkFailure(SinkFailure::invalid_contract(&request)) + } + Err(failure) if failure.validate_for_request(&request).is_ok() => { + DeliveryAttemptOutcome::SinkFailure(failure) + } + Err(_) => { + DeliveryAttemptOutcome::SinkFailure(SinkFailure::invalid_contract(&request)) + } + } + }; + let attempted_at = self.clock.now_unix_ms()?.max(execution_started_at); + let outcome = match result { + DeliveryAttemptOutcome::SinkFailure(failure) => DeliveryAttemptOutcome::SinkFailure( + normalize_sink_failure(&request, failure, attempted_at)?, + ), + outcome => outcome, + }; + let satisfaction = claimed + .evaluate_next_attempt(&outcome) + .map_err(map_storage_error)?; + let retry = delivery_retry_schedule(&claimed, &outcome, satisfaction, attempted_at)?; + let command = AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new(claimed.plan_id(), fence, outcome, retry, attempted_at) + .map_err(map_storage_error)?, + ); + let receipt = self + .storage + .execute_authored(command) + .await + .map_err(map_claim_error)?; + let AuthoredAtomicOutcome::DeliveryPlan(plan) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + Ok(DeliveryExecutionReceipt { + plan: plan.clone(), + replay: false, + }) + } + + async fn claim_delivery_plan( + &self, + plan: AuthoredDeliveryPlan, + acquired_at: u64, + ) -> Result<(AuthoredDeliveryPlan, WorkFence), Error> { + let generation = plan + .claim_evidence() + .map_or(1, |claim| claim.generation().get().saturating_add(1)); + let generation = NonZeroU64::new(generation).ok_or(Error::StorageFailed)?; + let expires_at = acquired_at + .checked_add(self.deadlines.timeout_ms(OperationKind::Deliver)) + .ok_or(Error::DeadlineOverflow)?; + let claim = WorkClaim::new( + *self.ids.next_id(OperationKind::Deliver)?.as_bytes(), + DELIVERY_CLAIM_OWNER, + generation, + acquired_at, + expires_at, + plan.revision(), + ) + .map_err(map_storage_error)?; + let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) + .map_err(map_storage_error)?; + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(plan.plan_id()), + claim, + )); + let receipt = self + .storage + .execute_authored(command) + .await + .map_err(map_claim_error)?; + let AuthoredAtomicOutcome::DeliveryPlan(claimed) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + Ok((claimed.clone(), fence)) + } + async fn claim_artifact( &self, artifact: AuthoredArtifact, @@ -1065,6 +1232,91 @@ fn retry_schedule( Ok(Some(schedule)) } +fn normalize_sink_failure( + request: &DeliveryRequest, + failure: SinkFailure, + attempted_at_unix_ms: u64, +) -> Result<SinkFailure, Error> { + if failure.retryability() != Retryability::Retryable { + return Ok(failure); + } + let default_retry_at = attempted_at_unix_ms + .checked_add(WORK_RETRY_DELAY_MS) + .ok_or(Error::DeadlineOverflow)?; + let retry_at = failure + .retry_after_unix_ms() + .unwrap_or(default_retry_at) + .max( + attempted_at_unix_ms + .checked_add(1) + .ok_or(Error::DeadlineOverflow)?, + ); + if failure.retry_after_unix_ms() == Some(retry_at) { + return Ok(failure); + } + SinkFailure::for_request( + request, + failure.code(), + failure.retryability(), + Some(retry_at), + failure.message().map(str::to_owned), + failure.partial_evidence().to_vec(), + ) + .map_err(|_| Error::InvalidDeliveryRequest) +} + +fn delivery_retry_schedule( + plan: &AuthoredDeliveryPlan, + outcome: &DeliveryAttemptOutcome, + satisfaction: SatisfactionState, + attempted_at_unix_ms: u64, +) -> Result<Option<RetrySchedule>, Error> { + let attempt = plan + .attempt_count() + .checked_add(1) + .and_then(NonZeroU32::new) + .ok_or(Error::StorageFailed)?; + if satisfaction != SatisfactionState::Pending || attempt.get() >= DELIVERY_PLAN_ATTEMPTS_MAX { + return Ok(None); + } + let failure = match outcome { + DeliveryAttemptOutcome::Receipt(_) => { + let retry_at = attempted_at_unix_ms + .checked_add(WORK_RETRY_DELAY_MS) + .ok_or(Error::DeadlineOverflow)?; + WorkFailure::new( + "delivery_pending", + WorkPhase::Delivery, + FailureClass::Retryable, + Some(retry_at), + None, + ) + .map_err(map_storage_error)? + } + DeliveryAttemptOutcome::SinkFailure(failure) + if failure.retryability() == Retryability::Retryable => + { + WorkFailure::new( + failure.code(), + WorkPhase::Delivery, + FailureClass::Retryable, + failure.retry_after_unix_ms(), + failure.message().map(str::to_owned), + ) + .map_err(map_storage_error)? + } + DeliveryAttemptOutcome::SinkFailure(_) => return Ok(None), + }; + let retry_at = failure.retry_after_unix_ms().unwrap_or( + attempted_at_unix_ms + .checked_add(WORK_RETRY_DELAY_MS) + .ok_or(Error::DeadlineOverflow)?, + ); + RetrySchedule::new(attempt, retry_at, failure) + .map(Some) + .map_err(map_storage_error) +} + fn delivery_evidence_digest(evidence: &DeliveryAttemptEvidence) -> AtomicCommitDigest { let mut hasher = Sha256::new(); hash_field(&mut hasher, b"radroots.sync.delivery-evidence.v1"); diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -24,6 +24,7 @@ use radroots_signing::{ }; use radroots_storage::{ EventStore, Journal, Outbox, + authored_delivery::AuthoredDeliveryState, event::{EventQuery, EventQueryBounds, SourceGeneration}, journal::{IdempotencyKey, OperationInstanceId}, memory::MemoryStorage, @@ -64,6 +65,7 @@ enum DeliveryBehavior { Outcomes(Vec<DeliveryOutcome>), AdapterError, MismatchedRequest, + Pending, } struct ScriptedSink { @@ -128,6 +130,7 @@ impl EventSink for ScriptedSink { ) .expect("mismatched receipt")) } + DeliveryBehavior::Pending => std::future::pending().await, } }) } @@ -767,6 +770,10 @@ async fn sqlite_signing_and_admission_recover_across_every_reopen_boundary() { ); assert!(signed.artifact().signed().is_some()); assert_eq!(signer.calls.load(Ordering::Relaxed), 1); + assert_eq!( + engine.deliver_push(push.operation_id()).await, + Err(Error::AdmissionFailed) + ); } { @@ -836,6 +843,249 @@ async fn sqlite_signing_and_admission_recover_across_every_reopen_boundary() { } #[test] +fn authored_delivery_persists_retry_evidence_and_complete_settlement() { + let sink = Arc::new(ScriptedSink::new([ + DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + DeliveryOutcome::unavailable(), + ]), + DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + DeliveryOutcome::accepted(), + ]), + ])); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let ((engine, _), clock) = setup_engine_with_sink(signer, sink.clone()); + let push = request_with_policy( + 73, + &["wss://one.example", "wss://two.example"], + SatisfactionClass::Accepted, + TargetPolicy::all(), + ); + block_on(engine.sign_and_enqueue(push.clone())).expect("prepare authored delivery"); + + let first = block_on(engine.deliver_push(push.operation_id())).expect("first attempt"); + assert!(!first.is_replay()); + assert_eq!(first.plan().state(), AuthoredDeliveryState::Retryable); + assert_eq!(first.plan().attempt_count(), 1); + let retry_at = first + .plan() + .retry() + .expect("durable retry") + .not_before_unix_ms(); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::DeliveryDeferred) + ); + let pending = block_on(engine.push_status(push.operation_id())) + .expect("pending status") + .expect("pending operation"); + assert_eq!(pending.settlement().artifacts(), 1); + assert_eq!(pending.settlement().signed(), 1); + assert_eq!(pending.settlement().admitted(), 1); + assert_eq!(pending.settlement().delivery_plans(), 1); + assert_eq!(pending.settlement().delivery_retryable(), 1); + assert!(!pending.settlement().is_settled()); + + clock.0.store(retry_at, Ordering::Relaxed); + let second = block_on(engine.deliver_push(push.operation_id())).expect("retry attempt"); + assert_eq!(second.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(second.plan().attempt_count(), 2); + assert_eq!(second.plan().attempts().len(), 2); + let replay = block_on(engine.deliver_push(push.operation_id())).expect("terminal replay"); + assert!(replay.is_replay()); + assert_eq!(replay.plan(), second.plan()); + let settled = block_on(engine.push_status(push.operation_id())) + .expect("settled status") + .expect("settled operation"); + assert_eq!(settled.settlement().delivery_satisfied(), 1); + assert!(settled.settlement().is_settled()); + assert!(!settled.settlement().has_failures()); + assert!(settled.settlement().is_successful()); + let requests = sink.requests.lock().expect("request log"); + assert_eq!(requests.len(), 2); + assert_eq!(requests[0], requests[1]); +} + +#[test] +fn invalid_authored_delivery_receipt_terminalizes_without_hot_loop() { + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::MismatchedRequest])); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let ((engine, _), _) = setup_engine_with_sink(signer, sink.clone()); + let push = request(74, "wss://one.example"); + block_on(engine.sign_and_enqueue(push.clone())).expect("prepare authored delivery"); + + let failed = block_on(engine.deliver_push(push.operation_id())).expect("durable failure"); + assert_eq!(failed.plan().state(), AuthoredDeliveryState::FailedTerminal); + assert_eq!(failed.plan().attempt_count(), 1); + assert_eq!( + failed.plan().last_failure().expect("failure").code(), + "invalid_transport_contract" + ); + let replay = block_on(engine.deliver_push(push.operation_id())).expect("failure replay"); + assert!(replay.is_replay()); + assert_eq!(replay.plan(), failed.plan()); + assert_eq!(sink.requests.lock().expect("request log").len(), 1); + let status = block_on(engine.push_status(push.operation_id())) + .expect("failed status") + .expect("failed operation"); + assert_eq!(status.settlement().delivery_failed_terminal(), 1); + assert!(status.settlement().is_settled()); + assert!(status.settlement().has_failures()); + assert!(!status.settlement().is_successful()); +} + +#[test] +fn authored_delivery_claims_fence_concurrent_and_stale_workers() { + let pending_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Pending])); + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let ((engine, storage), clock) = setup_engine_with_sink(signer, pending_sink.clone()); + let push = request(75, "wss://one.example"); + block_on(engine.sign_and_enqueue(push.clone())).expect("prepare authored delivery"); + + let mut pending = Box::pin(engine.deliver_push(push.operation_id())).fuse(); + let mut context = std::task::Context::from_waker(noop_waker_ref()); + assert!(pending.poll_unpin(&mut context).is_pending()); + let claimed = block_on(engine.push_status(push.operation_id())) + .expect("claimed status") + .expect("claimed operation"); + let expires_at = claimed + .delivery_plan() + .claim_evidence() + .expect("delivery claim") + .expires_at_unix_ms(); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::WorkClaimConflict) + ); + drop(pending); + + clock.0.store(expires_at + 1, Ordering::Relaxed); + let recovery_sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + ])])); + let capability: Arc<dyn SyncStorage> = storage; + let recovery = Engine::builder( + capability, + clock, + Arc::new(TestIds(AtomicU64::new(220))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(recovery_sink.clone()) + .build() + .expect("recovery engine"); + let recovered = block_on(recovery.deliver_push(push.operation_id())).expect("stale recovery"); + assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!( + recovered.plan().claim_evidence(), + None, + "settlement clears the claim" + ); + assert_eq!(recovery_sink.requests.lock().expect("request log").len(), 1); +} + +#[tokio::test] +async fn sqlite_authored_delivery_retry_survives_reopen() { + let directory = tempfile::tempdir().expect("database directory"); + let paths = Paths::from_directory(directory.path()).expect("paths"); + let push = request_with_policy( + 76, + &["wss://one.example", "wss://two.example"], + SatisfactionClass::Accepted, + TargetPolicy::all(), + ); + let retry_at = { + let store = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([76; 32]).expect("generation"), 1) + .expect("source generation"), + ) + .await + .expect("open SQLite"), + ); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + DeliveryOutcome::unavailable(), + ])])); + let engine = Engine::builder( + store, + Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))), + Arc::new(TestIds(AtomicU64::new(230))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(sink) + .signer(Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + }))) + .build() + .expect("engine"); + engine + .sign_and_enqueue(push.clone()) + .await + .expect("prepare authored delivery"); + let first = engine + .deliver_push(push.operation_id()) + .await + .expect("first attempt"); + assert_eq!(first.plan().attempt_count(), 1); + first + .plan() + .retry() + .expect("retry schedule") + .not_before_unix_ms() + }; + + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + .await + .expect("reopen SQLite"), + ); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + DeliveryOutcome::accepted(), + ])])); + let clock = Arc::new(TestClock(AtomicU64::new(retry_at - 1))); + let engine = Engine::builder( + store, + clock.clone(), + Arc::new(TestIds(AtomicU64::new(240))), + DeadlinePolicy::new(10_000, 10_000, 10_000).expect("deadlines"), + ) + .sink(sink.clone()) + .build() + .expect("reopen engine"); + let reopened = engine + .push_status(push.operation_id()) + .await + .expect("reopened status") + .expect("reopened operation"); + assert_eq!(reopened.delivery_plan().attempt_count(), 1); + assert_eq!( + reopened.delivery_plan().state(), + AuthoredDeliveryState::Retryable + ); + assert_eq!( + engine.deliver_push(push.operation_id()).await, + Err(Error::DeliveryDeferred) + ); + clock.0.store(retry_at, Ordering::Relaxed); + let completed = engine + .deliver_push(push.operation_id()) + .await + .expect("retry after reopen"); + assert_eq!(completed.plan().attempt_count(), 2); + assert_eq!(completed.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(sink.requests.lock().expect("request log").len(), 1); +} + +#[test] fn signer_rejection_challenge_and_timeout_fail_before_enqueue() { for kind in [ SigningErrorKind::SignerRejected,