lib

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

commit c9d75ca260a9191dd91a44581c1daa0460e366ab
parent eff47745c38195d2e48674d6b98236b714dcb880
Author: triesap <tyson@radroots.org>
Date:   Mon, 14 Sep 2026 17:13:12 +0000

sync: retain late delivery facts without extending worker authority

- Persist raw delivery outcomes before clock-dependent retry reconciliation
- Validate exact claims and receipts and expose bounded durable history
- Preserve stop intent and reconcile legacy attempts without duplicate effects
- Verify concurrency, reopen, API and unchanged owner coverage gates

Diffstat:
Mcontracts/api_baselines/radroots_sync.txt | 5++++-
Acontracts/architecture/decisions/sync_delivery_evidence.v1.json | 13+++++++++++++
Mcrates/sync/README.md | 16++++++++++++++++
Mcrates/sync/src/push.rs | 276+++++++++++++++++++++++++++-----------------------------------------------------
Acrates/sync/src/push/delivery.rs | 334+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/sync/src/status.rs | 8++++++++
Mcrates/sync/tests/push_enqueue.rs | 52+++++++++++++++++++++++++++++++++++++++++++++++++++-
Acrates/sync/tests/push_enqueue/delivery_evidence.rs | 1284+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
8 files changed, 1802 insertions(+), 186 deletions(-)

diff --git a/contracts/api_baselines/radroots_sync.txt b/contracts/api_baselines/radroots_sync.txt @@ -215,6 +215,7 @@ pub fn radroots_sync::push::PushRequest::fmt(&self, &mut core::fmt::Formatter<'_ pub struct radroots_sync::push::PushStatus impl radroots_sync::push::PushStatus pub const fn radroots_sync::push::PushStatus::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact +pub const fn radroots_sync::push::PushStatus::delivery_history(&self) -> &radroots_storage::authored_delivery::history::AuthoredDeliveryHistory pub const fn radroots_sync::push::PushStatus::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan pub const fn radroots_sync::push::PushStatus::operation(&self) -> &radroots_storage::authored::AuthoredOperation pub const fn radroots_sync::push::PushStatus::settlement(&self) -> radroots_storage::authored::OperationSettlement @@ -289,7 +290,6 @@ pub struct radroots_sync::Engine impl radroots_sync::Engine pub async fn radroots_sync::Engine::admit_signed(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::AdmissionRunReceipt, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::cancel_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::PushCancellationReceipt, radroots_sync::policy::Error> -pub async fn radroots_sync::Engine::deliver_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::prepare_push(&self, radroots_sync::push::PushRequest) -> core::result::Result<radroots_sync::push::PushPreparation, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::push_status(&self, radroots_sync::policy::SyncId) -> core::result::Result<core::option::Option<radroots_sync::push::PushStatus>, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::sign_prepared(&self, radroots_sync::push::PushRequest) -> core::result::Result<radroots_sync::push::SigningRunReceipt, radroots_sync::policy::Error> @@ -303,6 +303,8 @@ pub fn radroots_sync::Engine::sink(&self) -> core::option::Option<&dyn radroots_ pub fn radroots_sync::Engine::source(&self) -> core::option::Option<&dyn radroots_transport::source::EventSource> pub fn radroots_sync::Engine::storage(&self) -> &dyn radroots_sync::policy::SyncStorage impl radroots_sync::Engine +pub async fn radroots_sync::Engine::deliver_push(&self, radroots_sync::policy::SyncId) -> core::result::Result<radroots_sync::push::DeliveryExecutionReceipt, radroots_sync::policy::Error> +impl radroots_sync::Engine pub async fn radroots_sync::Engine::ingest(&self, radroots_transport::source::ObservedEvent, &dyn radroots_sync::ingest::AdmissionPolicy) -> core::result::Result<radroots_sync::ingest::IngestReceipt, radroots_sync::policy::Error> pub async fn radroots_sync::Engine::ingest_batch(&self, alloc::vec::Vec<radroots_transport::source::ObservedEvent>, &dyn radroots_sync::ingest::AdmissionPolicy) -> radroots_sync::ingest::IngestBatchReceipt impl radroots_sync::Engine @@ -366,6 +368,7 @@ pub fn radroots_sync::push::PushRequest::fmt(&self, &mut core::fmt::Formatter<'_ pub struct radroots_sync::PushStatus impl radroots_sync::push::PushStatus pub const fn radroots_sync::push::PushStatus::artifact(&self) -> &radroots_storage::authored::AuthoredArtifact +pub const fn radroots_sync::push::PushStatus::delivery_history(&self) -> &radroots_storage::authored_delivery::history::AuthoredDeliveryHistory pub const fn radroots_sync::push::PushStatus::delivery_plan(&self) -> &radroots_storage::authored_delivery::AuthoredDeliveryPlan pub const fn radroots_sync::push::PushStatus::operation(&self) -> &radroots_storage::authored::AuthoredOperation pub const fn radroots_sync::push::PushStatus::settlement(&self) -> radroots_storage::authored::OperationSettlement diff --git a/contracts/architecture/decisions/sync_delivery_evidence.v1.json b/contracts/architecture/decisions/sync_delivery_evidence.v1.json @@ -0,0 +1,13 @@ +{ + "schema": "radroots.sync-delivery-evidence.v1", + "status": "approved", + "owner": "radroots_sync", + "producer_contracts": ["authored_delivery_reconciliation.v1.json"], + "binding": "Every transport invocation uses the exact persisted signed request under a backend-owned original claim. Validate raw final sink results and retain them with RecordDeliveryFact before retry normalization. Invalid adapter contracts produce the existing bound invalid_transport_contract failure, never fabricated acceptance.", + "authority": "Validate each claim, fact and reconciliation receipt against its exact command and reload current backend history. Late callbacks may reconcile only under their own still-current unexpired claim. Stop, expiry or supersession never erases a valid fact and never grants that callback fresh scheduling authority. A fresh call reconciles pending facts as a complete bounded action before any subsequent transport invocation.", + "history": "Expose the consistent bounded AuthoredDeliveryHistory through PushStatus. Proven no-issued-attempt, unresolved or incomplete history, raw acceptance, first stop and scheduling settlement remain distinct. Legacy settlement fields retain scheduling semantics. Missing backend history fails closed; no client infers no effect from an empty attempt list.", + "retry": "Preserve raw outcomes, provider retry bounds, the existing delay and attempt cap. Count new facts once; recognize legacy applications using both the original claim attempt ordinal and valid lease interval. Reconciliation rejects conflicting legacy observations without deleting either fact. Retry and deadline exhaustion preserve immutable intent for explicit host recovery.", + "clock": "A validated non-expiring delivery result is retained even when the post-effect host clock fails, using the known valid pre-effect timestamp only as a causal lower bound. Return ClockUnavailable after retaining it and do not reconcile scheduling in that callback. Do not claim a measured response timestamp or alter event time. Fresh reconciliation cannot advance host time to a future fact. Strict signer evidence and expiring authorization clock requirements are unchanged.", + "lifetime": "Sync creates no executor, worker, timer or hidden retry. The host must retain and poll an admitted future to deliver a late result. Dropping it leaves durable unresolved provenance; cancellation is not remote rollback.", + "scope": "Standalone shared Sync only. Tera native adoption, UI status and whole-operation stop are separate checkpoints. No credential mutation, release publication or deployment." +} diff --git a/crates/sync/README.md b/crates/sync/README.md @@ -45,3 +45,19 @@ unresolved; recovery follows the declared replay capability and never invents a new preimage, event timestamp, or key. Sync creates no worker or timer to retain a discarded future. Delivery-wide stop reconciliation is a separate contract from authored signing evidence. + +Authored delivery retains validated raw results against their original durable +claim even after stop, expiry or replacement of that claim. A late callback can +change scheduling only under its original still-current lease. Fresh calls +reconcile pending facts before any further delivery; reconciliation invokes no +transport. Every retry uses the same persisted signed request and target policy. +Stop prevents further local work and preserves unknown and accepted effects. + +`PushStatus::delivery_history` exposes consistent bounded claim provenance. +Use its explicit no-issued proof and unresolved-history classification together +with the plan's cumulative delivery satisfaction and first stop. The existing +settlement counters describe scheduling state; they do not prove remote absence. +After a post-delivery clock failure, Sync retains the raw non-expiring result +with the known pre-effect time as a causal lower bound and returns +`ClockUnavailable` without retry scheduling. This is not a measured response +time and does not change strict signing or expiring authorization requirements. diff --git a/crates/sync/src/push.rs b/crates/sync/src/push.rs @@ -19,14 +19,14 @@ use radroots_storage::{ OperationSettlement, RetrySchedule, SigningState, WorkClaim, WorkFailure, WorkPhase, }, authored_atomic::{ - ApplyAdmissionResult, ApplyDeliveryAttempt, ApplyWorkFailure, AuthoredAtomicCommand, - AuthoredAtomicOutcome, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, - ClaimAuthoredTarget, ClaimAuthoredWork, PrepareAuthoredOperation, RecordSignedArtifact, - WorkFence, + ApplyAdmissionResult, ApplyWorkFailure, AuthoredAtomicCommand, AuthoredAtomicOutcome, + AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, + ClaimAuthoredWork, PrepareAuthoredOperation, RecordSignedArtifact, WorkFence, }, authored_delivery::{ - AuthoredDeliveryIntent, AuthoredDeliveryPlan, AuthoredDeliveryPlanId, - AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, DeliveryAttemptOutcome, + AuthoredDeliveryHistory, AuthoredDeliveryIntent, AuthoredDeliveryPlan, + AuthoredDeliveryPlanId, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX, + DeliveryAttemptOutcome, }, event::{AdmissionDisposition, EventAdmission, EventStore}, journal::{IdempotencyKey, OperationInstanceId}, @@ -46,6 +46,8 @@ use crate::{ policy::{Error, OperationKind, SyncId}, }; +mod delivery; + 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"; @@ -206,7 +208,7 @@ impl PushPreparation { pub struct PushStatus { operation: AuthoredOperation, artifact: AuthoredArtifact, - delivery_plan: AuthoredDeliveryPlan, + delivery_history: AuthoredDeliveryHistory, settlement: OperationSettlement, } @@ -218,7 +220,11 @@ impl PushStatus { &self.artifact } pub const fn delivery_plan(&self) -> &AuthoredDeliveryPlan { - &self.delivery_plan + self.delivery_history.plan() + } + /// Consistent original-claim history; missing provenance remains unknown. + pub const fn delivery_history(&self) -> &AuthoredDeliveryHistory { + &self.delivery_history } pub const fn settlement(&self) -> OperationSettlement { self.settlement @@ -352,13 +358,17 @@ impl Engine { .await .map_err(map_storage_error)? .ok_or(Error::StorageFailed)?; - let delivery_plan = self + let delivery_history = self .storage - .authored_delivery_plan(expected_plan_id) + .authored_delivery_history(expected_plan_id) .await .map_err(map_storage_error)? .ok_or(Error::StorageFailed)?; - if artifact.operation_id() != operation.operation_id() + delivery_history.validate().map_err(map_storage_error)?; + let delivery_plan = delivery_history.plan(); + if artifact.artifact_id() != expected_artifact_id + || delivery_plan.plan_id() != expected_plan_id + || artifact.operation_id() != operation.operation_id() || delivery_plan.artifact_id() != artifact.artifact_id() { return Err(Error::StorageFailed); @@ -366,13 +376,13 @@ impl Engine { let settlement = OperationSettlement::evaluate_complete( &operation, core::slice::from_ref(&artifact), - core::slice::from_ref(&delivery_plan), + core::slice::from_ref(delivery_plan), ) .map_err(map_storage_error)?; Ok(Some(PushStatus { operation, artifact, - delivery_plan, + delivery_history, settlement, })) } @@ -393,7 +403,7 @@ impl Engine { .clock .now_unix_ms()? .max(status.artifact.updated_at_unix_ms()) - .max(status.delivery_plan.updated_at_unix_ms()); + .max(status.delivery_plan().updated_at_unix_ms()); let artifact_target = match status.artifact.signing_state() { SigningState::Planned | SigningState::Retryable => Some( @@ -429,26 +439,44 @@ impl Engine { .ok_or(Error::StorageFailed)?; } - if matches!( - status.delivery_plan.state(), - AuthoredDeliveryState::Pending | AuthoredDeliveryState::Retryable - ) { - self.storage - .execute_authored(AuthoredAtomicCommand::Cancel( - CancelAuthoredWork::new( - CancelAuthoredTarget::DeliveryPlan(status.delivery_plan.plan_id()), - status.delivery_plan.revision(), - now, - ) - .map_err(map_storage_error)?, - )) + if status.delivery_plan().stop_requested_at_unix_ms().is_none() { + let command = AuthoredAtomicCommand::Cancel( + CancelAuthoredWork::new( + CancelAuthoredTarget::DeliveryPlan(status.delivery_plan().plan_id()), + status.delivery_plan().revision(), + now, + ) + .map_err(map_storage_error)?, + ); + let receipt = self + .storage + .execute_authored(command.clone()) .await .map_err(map_storage_error)?; + if !receipt.matches_command(&command) { + return Err(Error::StorageFailed); + } + let AuthoredAtomicOutcome::DeliveryPlan(stopped) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + if stopped.validate().is_err() + || stopped.plan_id() != status.delivery_plan().plan_id() + || stopped.artifact_id() != status.artifact.artifact_id() + || stopped.request() != status.delivery_plan().request() + || stopped.stop_requested_at_unix_ms() != Some(now) + || status.delivery_plan().revision().get().checked_add(1) + != Some(stopped.revision().get()) + { + return Err(Error::StorageFailed); + } changed = true; status = self .push_status(operation_id) .await? .ok_or(Error::StorageFailed)?; + if status.delivery_plan().stop_requested_at_unix_ms().is_none() { + return Err(Error::StorageFailed); + } } Ok(PushCancellationReceipt { status, changed }) @@ -461,7 +489,7 @@ impl Engine { .push_status(request.operation_id) .await? .ok_or(Error::StorageFailed)?; - if status.delivery_plan.state() == AuthoredDeliveryState::Cancelled { + if status.delivery_plan().state() == AuthoredDeliveryState::Cancelled { self.cancel_push(request.operation_id).await?; return Err(Error::SigningCancelled); } @@ -587,7 +615,7 @@ impl Engine { return Err(Error::StorageFailed); } if expected_request.cancellation_signal().is_cancelled() - || status.delivery_plan.state() == AuthoredDeliveryState::Cancelled + || status.delivery_plan().state() == AuthoredDeliveryState::Cancelled { let stopped = self.cancel_push(request.operation_id).await?; if stopped @@ -715,6 +743,7 @@ impl Engine { .push_status(operation_id) .await? .ok_or(Error::StorageFailed)?; + let delivery_stopped = status.delivery_plan().stop_requested_at_unix_ms().is_some(); let artifact = status.artifact; if artifact.signing_state() != SigningState::Signed { return Err(Error::InvalidSignerOutput); @@ -725,7 +754,7 @@ impl Engine { replay: true, }); } - if status.delivery_plan.state() == AuthoredDeliveryState::Cancelled { + if delivery_stopped { self.cancel_push(operation_id).await?; return Err(Error::AdmissionFailed); } @@ -849,141 +878,6 @@ 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, @@ -1131,16 +1025,11 @@ fn normalize_sink_failure( } fn delivery_retry_schedule( - plan: &AuthoredDeliveryPlan, + attempt: NonZeroU32, 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); } @@ -1667,26 +1556,45 @@ mod tests { .unwrap(); let receipt_outcome = DeliveryAttemptOutcome::Receipt(receipt); assert!( - delivery_retry_schedule(&plan, &receipt_outcome, SatisfactionState::Satisfied, 100) - .unwrap() - .is_none() + delivery_retry_schedule( + NonZeroU32::new(plan.attempt_count() + 1).unwrap(), + &receipt_outcome, + SatisfactionState::Satisfied, + 100 + ) + .unwrap() + .is_none() ); - let schedule = - delivery_retry_schedule(&plan, &receipt_outcome, SatisfactionState::Pending, 100) - .unwrap() - .unwrap(); + let schedule = delivery_retry_schedule( + NonZeroU32::new(plan.attempt_count() + 1).unwrap(), + &receipt_outcome, + SatisfactionState::Pending, + 100, + ) + .unwrap() + .unwrap(); assert_eq!(schedule.not_before_unix_ms(), 1_100); let retry_outcome = DeliveryAttemptOutcome::SinkFailure(normalized); assert!( - delivery_retry_schedule(&plan, &retry_outcome, SatisfactionState::Pending, 100) - .unwrap() - .is_some() + delivery_retry_schedule( + NonZeroU32::new(plan.attempt_count() + 1).unwrap(), + &retry_outcome, + SatisfactionState::Pending, + 100 + ) + .unwrap() + .is_some() ); let terminal_outcome = DeliveryAttemptOutcome::SinkFailure(terminal); assert!( - delivery_retry_schedule(&plan, &terminal_outcome, SatisfactionState::Pending, 100) - .unwrap() - .is_none() + delivery_retry_schedule( + NonZeroU32::new(plan.attempt_count() + 1).unwrap(), + &terminal_outcome, + SatisfactionState::Pending, + 100 + ) + .unwrap() + .is_none() ); for policy in [ diff --git a/crates/sync/src/push/delivery.rs b/crates/sync/src/push/delivery.rs @@ -0,0 +1,334 @@ +//! Delivery observations outlive the lease that admitted their external effect. + +use super::*; +use radroots_storage::authored_atomic::{ReconcileDeliveryFacts, RecordDeliveryFact}; + +impl Engine { + /// Reconciles retained evidence or executes one bounded delivery attempt. + /// + /// A late result retains its original raw binding without authorizing a + /// stale worker to change scheduling. Hosts must keep polling this future + /// to deliver late results; dropping it leaves durable unresolved work. + pub async fn deliver_push( + &self, + operation_id: SyncId, + ) -> Result<DeliveryExecutionReceipt, Error> { + let status = self.push_status(operation_id).await?.ok_or_else(|| { + if self.sink.is_none() { + Error::MissingSink + } else { + Error::StorageFailed + } + })?; + let plan = status.delivery_plan(); + if plan.state().is_terminal() { + return Ok(DeliveryExecutionReceipt { + plan: plan.clone(), + replay: true, + }); + } + if status.artifact.signing_state() != SigningState::Signed { + return Err(Error::InvalidSignerOutput); + } + if !status.artifact.admission_state().is_admitted() { + return Err(Error::AdmissionFailed); + } + let now = self.delivery_now()?.max(plan.updated_at_unix_ms()); + // This invocation is fresh scheduling authority, not a late callback. + // Reconciliation is a complete bounded action; never also send a retry. + if plan.pending_delivery_facts().next().is_some() { + let plan = self + .reconcile_delivery(status.delivery_history(), None, now) + .await?; + return Ok(DeliveryExecutionReceipt { plan, replay: true }); + } + 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 sink = self.sink.as_deref().ok_or(Error::MissingSink)?; + let request = plan.request().cloned().ok_or(Error::InvalidSignerOutput)?; + let claimed = self.claim_delivery_plan(plan, now).await?; + let claim = claimed + .claim_evidence() + .cloned() + .ok_or(Error::StorageFailed)?; + // Re-read after claim admission: an observed stop or newer worker must + // prevent this callback from initiating an effect. + let current = self.delivery_history(claimed.plan_id()).await?; + if current.plan().state().is_terminal() { + return Ok(DeliveryExecutionReceipt { + plan: current.plan().clone(), + replay: true, + }); + } + let execution_started_at = self + .delivery_now()? + .max(current.plan().updated_at_unix_ms()); + if current.plan().claim_evidence() != Some(&claim) + || execution_started_at >= claim.expires_at_unix_ms() + { + return Err(Error::WorkClaimConflict); + } + if current.plan().pending_delivery_facts().next().is_some() { + let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) + .map_err(map_storage_error)?; + let plan = self + .reconcile_delivery(&current, Some(fence), execution_started_at) + .await?; + return Ok(DeliveryExecutionReceipt { plan, replay: true }); + } + let outcome = 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) + } + Err(failure) if failure.validate_for_request(&request).is_ok() => { + DeliveryAttemptOutcome::SinkFailure(failure) + } + Ok(_) | Err(_) => { + DeliveryAttemptOutcome::SinkFailure(SinkFailure::invalid_contract(&request)) + } + } + }; + let observed = self.delivery_now().map(|at| at.max(execution_started_at)); + // A known pre-effect lower bound suffices for the non-expiring delivery + // fact. It does not claim a measured response time or grant retry rights. + // Strict signing/expiring authorization observation rules are separate. + let command = AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + claimed.plan_id(), + claimed.artifact_id(), + claim.clone(), + outcome.clone(), + observed.unwrap_or(execution_started_at), + ) + .map_err(map_storage_error)?, + ); + let receipt = self + .storage + .execute_authored(command.clone()) + .await + .map_err(map_storage_error)?; + if !receipt.matches_command(&command) { + return Err(Error::StorageFailed); + } + let current = self.delivery_history(claimed.plan_id()).await?; + if !current + .plan() + .delivery_facts() + .iter() + .any(|fact| fact.claim() == &claim && fact.outcome() == &outcome) + { + return Err(Error::StorageFailed); + } + let at = observed?; + let plan = if !current.plan().state().is_terminal() + && current.plan().claim_evidence() == Some(&claim) + && at < claim.expires_at_unix_ms() + { + let fence = WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()) + .map_err(map_storage_error)?; + self.reconcile_delivery(&current, Some(fence), at).await? + } else { + // Acceptance remains visible even when scheduling was stopped, + // expired or superseded. Only a fresh caller may reconcile later. + current.plan().clone() + }; + Ok(DeliveryExecutionReceipt { + plan, + replay: false, + }) + } + + fn delivery_now(&self) -> Result<u64, Error> { + match self.clock.now_unix_ms()? { + 0 => Err(Error::ClockUnavailable), + at => Ok(at), + } + } + + async fn delivery_history( + &self, + id: AuthoredDeliveryPlanId, + ) -> Result<AuthoredDeliveryHistory, Error> { + let history = self + .storage + .authored_delivery_history(id) + .await + .map_err(map_storage_error)? + .ok_or(Error::StorageFailed)?; + history.validate().map_err(map_storage_error)?; + if history.plan().plan_id() != id { + return Err(Error::StorageFailed); + } + Ok(history) + } + + async fn reconcile_delivery( + &self, + history: &AuthoredDeliveryHistory, + fence: Option<WorkFence>, + at: u64, + ) -> Result<AuthoredDeliveryPlan, Error> { + let plan = history.plan(); + let at = at.max(plan.updated_at_unix_ms()); + // Never advance host time to a future fact to bypass a competing lease. + if plan + .delivery_facts() + .iter() + .any(|fact| fact.observed_at_unix_ms() > at) + { + return Err(Error::ClockUnavailable); + } + let retry = reconciliation_retry(history, at)?; + let reconcile = + ReconcileDeliveryFacts::new(plan, fence, retry, at).map_err(map_storage_error)?; + reconcile.apply_to(history).map_err(map_claim_error)?; + let command = AuthoredAtomicCommand::ReconcileDelivery(reconcile); + let receipt = self + .storage + .execute_authored(command.clone()) + .await + .map_err(map_claim_error)?; + if !receipt.matches_command(&command) { + return Err(Error::StorageFailed); + } + // A receipt replay is historical; expose current stop/facts/authority. + Ok(self.delivery_history(plan.plan_id()).await?.plan().clone()) + } + + async fn claim_delivery_plan( + &self, + plan: &AuthoredDeliveryPlan, + acquired_at: u64, + ) -> Result<AuthoredDeliveryPlan, Error> { + let generation = plan + .claim_evidence() + .map_or(1, |claim| claim.generation().get().saturating_add(1)); + 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, + NonZeroU64::new(generation).ok_or(Error::StorageFailed)?, + acquired_at, + expires_at, + plan.revision(), + ) + .map_err(map_storage_error)?; + let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(plan.plan_id()), + claim.clone(), + )); + let receipt = self + .storage + .execute_authored(command.clone()) + .await + .map_err(map_claim_error)?; + if !receipt.matches_command(&command) { + return Err(Error::StorageFailed); + } + let AuthoredAtomicOutcome::DeliveryPlan(claimed) = receipt.outcome() else { + return Err(Error::StorageFailed); + }; + if claimed.validate().is_err() + || claimed.plan_id() != plan.plan_id() + || claimed.artifact_id() != plan.artifact_id() + || claimed.request() != plan.request() + || claimed.created_at_unix_ms() != plan.created_at_unix_ms() + || claimed.claim_evidence() != Some(&claim) + || plan.revision().get().checked_add(1) != Some(claimed.revision().get()) + || claimed.updated_at_unix_ms() != acquired_at + || receipt.committed_at_unix_ms() != acquired_at + || claimed.attempts() != plan.attempts() + || claimed.state() != plan.state() + || claimed.retry() != plan.retry() + { + return Err(Error::StorageFailed); + } + Ok(claimed.clone()) + } +} + +fn reconciliation_retry( + history: &AuthoredDeliveryHistory, + at: u64, +) -> Result<Option<RetrySchedule>, Error> { + history + .require_pending_fact_provenance() + .map_err(map_storage_error)?; + let plan = history.plan(); + let mut count = plan.attempt_count(); + let mut matched_legacy = Vec::new(); + let mut last = plan.attempts().last().map(|attempt| attempt.outcome()); + for fact in plan.pending_delivery_facts() { + let original = history + .claims() + .iter() + .find(|entry| entry.claim() == fact.claim()) + .ok_or(Error::StorageFailed)?; + let legacy = plan.attempts().iter().find(|attempt| { + !matched_legacy.contains(&attempt.attempt()) + && attempt.claim_evidence().is_none() + && original.prior_attempt_count().checked_add(1) == Some(attempt.attempt().get()) + && attempt.recorded_at_unix_ms() >= fact.claim().acquired_at_unix_ms() + && attempt.recorded_at_unix_ms() < fact.claim().expires_at_unix_ms() + }); + if let Some(legacy) = legacy { + if legacy.outcome() != fact.outcome() { + return Err(Error::StorageConflict); + } + matched_legacy.push(legacy.attempt()); + } else { + count = count.checked_add(1).ok_or(Error::StorageFailed)?; + last = Some(fact.outcome()); + } + } + if count > DELIVERY_PLAN_ATTEMPTS_MAX { + return Err(Error::StorageFailed); + } + let satisfaction = plan.delivery_satisfaction().map_err(map_storage_error)?; + if satisfaction != SatisfactionState::Pending || count == DELIVERY_PLAN_ATTEMPTS_MAX { + return Ok(None); + } + let last = last.ok_or(Error::StorageFailed)?; + let normalized; + let last = if let DeliveryAttemptOutcome::SinkFailure(failure) = last { + normalized = DeliveryAttemptOutcome::SinkFailure(normalize_sink_failure( + plan.request().ok_or(Error::InvalidDeliveryRequest)?, + failure.clone(), + at, + )?); + &normalized + } else { + last + }; + delivery_retry_schedule( + NonZeroU32::new(count).ok_or(Error::StorageFailed)?, + last, + satisfaction, + at, + ) +} diff --git a/crates/sync/src/status.rs b/crates/sync/src/status.rs @@ -222,6 +222,14 @@ impl Engine { if now_unix_ms == 0 { return Err(Error::ClockUnavailable); } + if plan.request().is_some() + && plan + .delivery_satisfaction() + .map_err(|_| Error::StorageFailed)? + == radroots_transport::policy::SatisfactionState::Satisfied + { + return Ok(SyncRetryDecision::Satisfied); + } match plan.state() { AuthoredDeliveryState::Satisfied => return Ok(SyncRetryDecision::Satisfied), AuthoredDeliveryState::Exhausted diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -57,10 +57,16 @@ const CONTENT: &str = "frozen-content"; #[path = "push_enqueue/signing_evidence.rs"] mod signing_evidence; +#[path = "push_enqueue/delivery_evidence.rs"] +mod delivery_evidence; + struct MockSink; struct FaultStorage { inner: Arc<MemoryStorage>, + delivery_injection: Mutex<Option<delivery_evidence::Injection>>, + receipt_mutation: Mutex<Option<delivery_evidence::ReceiptMutation>>, + history_mutation: Mutex<Option<(usize, delivery_evidence::HistoryMutation)>>, fault_kind: AtomicUsize, remaining: AtomicUsize, admit_error: AtomicUsize, @@ -70,6 +76,9 @@ struct FaultStorage { impl FaultStorage { fn new(generation: u8) -> Self { Self { + delivery_injection: Mutex::new(None), + receipt_mutation: Mutex::new(None), + history_mutation: Mutex::new(None), inner: Arc::new(MemoryStorage::new( SourceGeneration::new([generation; 32]).expect("generation"), )), @@ -476,8 +485,18 @@ impl AuthoredAtomicStorage for FaultStorage { Result<radroots_storage::authored_atomic::AuthoredAtomicReceipt, radroots_storage::Error>, > { Box::pin(async move { + delivery_evidence::inject(self, &command, true).await; let receipt = - AuthoredAtomicStorage::execute_authored(self.inner.as_ref(), command).await?; + AuthoredAtomicStorage::execute_authored(self.inner.as_ref(), command.clone()) + .await?; + delivery_evidence::inject(self, &command, false).await; + if matches!( + receipt.outcome(), + radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(_) + ) && let Some(mutation) = self.receipt_mutation.lock().unwrap().take() + { + return Ok(delivery_evidence::mutate_receipt(receipt, mutation)); + } if let radroots_storage::authored_atomic::AuthoredAtomicOutcome::Prepared { .. } = receipt.outcome() { @@ -586,6 +605,37 @@ impl AuthoredAtomicStorage for FaultStorage { > { AuthoredAtomicStorage::authored_delivery_plan(self.inner.as_ref(), id) } + fn authored_delivery_history( + &self, + id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId, + ) -> radroots_transport::BoxFuture< + '_, + Result< + Option<radroots_storage::authored_delivery::AuthoredDeliveryHistory>, + radroots_storage::Error, + >, + > { + Box::pin(async move { + let history = self.inner.authored_delivery_history(id).await?; + let mutation = { + let mut armed = self.history_mutation.lock().unwrap(); + if let Some((remaining, _)) = armed.as_mut() { + *remaining -= 1; + } + if armed.as_ref().is_some_and(|(remaining, _)| *remaining == 0) { + armed.take().map(|(_, mutation)| mutation) + } else { + None + } + }; + Ok(match (history, mutation) { + (Some(history), Some(mutation)) => { + Some(delivery_evidence::mutate_history(history, mutation)) + } + (history, _) => history, + }) + }) + } } struct MockSource; diff --git a/crates/sync/tests/push_enqueue/delivery_evidence.rs b/crates/sync/tests/push_enqueue/delivery_evidence.rs @@ -0,0 +1,1284 @@ +use super::*; +use futures::channel::oneshot; +use radroots_storage::{ + authored_atomic::{AuthoredAtomicCommand, RecordDeliveryFact}, + authored_delivery::DeliveryAttemptOutcome, +}; +use radroots_transport::policy::SatisfactionState; + +type Pending = ( + DeliveryRequest, + oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>, +); + +#[derive(Default)] +struct HeldSink { + pending: Mutex<VecDeque<Pending>>, + calls: AtomicUsize, +} + +impl HeldSink { + fn take(&self) -> Pending { + self.pending.lock().unwrap().pop_front().unwrap() + } +} +impl EventSink for HeldSink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { panic!("delivery must not probe transport status") }) + } + fn deliver( + &self, + request: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + self.calls.fetch_add(1, Ordering::Relaxed); + let (sender, receiver) = oneshot::channel(); + self.pending.lock().unwrap().push_back((request, sender)); + receiver.await.expect("host retains admitted future") + }) + } +} +fn setup( + byte: u8, +) -> ( + Engine, + Arc<MemoryStorage>, + Arc<TestClock>, + Arc<HeldSink>, + PushRequest, +) { + let sink = Arc::new(HeldSink::default()); + let ((engine, storage), clock) = setup_engine_with_sink( + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(byte, "wss://relay.example"); + execute_to_admitted(&engine, &push); + (engine, storage, clock, sink, push) +} +fn source_only(storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>) -> Engine { + Engine::builder( + storage, + clock, + Arc::new(TestIds(AtomicU64::new(230))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .source(Arc::new(MockSource)) + .build() + .unwrap() +} +fn accept((request, sender): Pending) -> DeliveryReceipt { + let result = receipt( + &request, + vec![DeliveryOutcome::accepted(); request.target_set().len()], + ) + .unwrap(); + sender.send(Ok(result.clone())).unwrap(); + result +} + +#[test] +fn late_accepted_result_survives_stop_and_terminal_replay_without_sink() { + let (engine, storage, clock, sink, push) = setup(121); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(status.delivery_history().proves_no_issued_attempt()); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let pending = sink.take(); + let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); + assert!(stopped.changed()); + assert!(stopped.status().delivery_history().has_unresolved_claims()); + assert!( + !stopped + .status() + .delivery_history() + .proves_no_issued_attempt() + ); + let stop_at = stopped.status().delivery_plan().stop_requested_at_unix_ms(); + let expected = accept(pending); + let delivered = block_on(future).unwrap(); + assert_eq!(delivered.plan().state(), AuthoredDeliveryState::Cancelled); + assert_eq!(delivered.plan().attempt_count(), 0); + assert_eq!( + delivered.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + assert_eq!( + delivered.plan().delivery_facts()[0].outcome(), + &DeliveryAttemptOutcome::Receipt(expected) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(!status.delivery_history().has_unresolved_claims()); + assert_eq!(status.delivery_plan().stop_requested_at_unix_ms(), stop_at); + assert!( + !block_on(engine.cancel_push(push.operation_id())) + .unwrap() + .changed() + ); + let replay = block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); + assert!(replay.is_replay()); + assert_eq!(replay.plan(), status.delivery_plan()); + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); +} + +#[test] +fn expired_callback_retains_fact_but_only_fresh_invocation_reconciles() { + let (engine, storage, clock, sink, push) = setup(122); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + clock.0.store( + before + .delivery_plan() + .claim_evidence() + .unwrap() + .expires_at_unix_ms(), + Ordering::Relaxed, + ); + accept(sink.take()); + let late = block_on(future).unwrap(); + assert_eq!(late.plan().revision(), before.delivery_plan().revision()); + assert_eq!(late.plan().attempt_count(), 0); + let recovery = source_only(storage, clock); + let recovered = block_on(recovery.deliver_push(push.operation_id())).unwrap(); + assert!(recovered.is_replay()); + assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(recovered.plan().attempt_count(), 1); + assert_eq!( + recovered.plan().attempts()[0].claim_evidence(), + before.delivery_plan().claim_evidence() + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); +} + +#[test] +fn superseded_callback_cannot_clear_newer_claim_and_acceptance_never_regresses() { + let (engine, _, clock, sink, push) = setup(123); + let mut first = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + first + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let first_pending = sink.take(); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + clock.0.store( + before + .delivery_plan() + .claim_evidence() + .unwrap() + .expires_at_unix_ms(), + Ordering::Relaxed, + ); + let mut second = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + second + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (second_request, second_sender) = sink.take(); + assert_eq!(first_pending.0, second_request); + let newer = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + accept(first_pending); + let late = block_on(first).unwrap(); + assert_eq!( + late.plan().claim_evidence(), + newer.delivery_plan().claim_evidence() + ); + assert_eq!(late.plan().revision(), newer.delivery_plan().revision()); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::WorkClaimConflict) + ); + second_sender + .send(Err(SinkFailure::for_request( + &second_request, + "provider_failed", + Retryability::Terminal, + None, + None, + Vec::new(), + ) + .unwrap())) + .unwrap(); + let complete = block_on(second).unwrap(); + assert_eq!(complete.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(complete.plan().attempt_count(), 2); + assert_eq!(complete.plan().delivery_facts().len(), 2); + assert_eq!(sink.calls.load(Ordering::Relaxed), 2); +} + +#[test] +fn clock_loss_retains_raw_result_before_reporting_unavailable() { + for accepted in [true, false] { + let (engine, storage, clock, sink, push) = setup(124 + u8::from(accepted)); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (request, sender) = sink.take(); + let before = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let outcome = if accepted { + DeliveryAttemptOutcome::Receipt( + receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), + ) + } else { + DeliveryAttemptOutcome::SinkFailure( + SinkFailure::for_request( + &request, + "raw_failure", + Retryability::Retryable, + None, + Some("raw diagnostic".into()), + Vec::new(), + ) + .unwrap(), + ) + }; + sender + .send(match &outcome { + DeliveryAttemptOutcome::Receipt(value) => Ok(value.clone()), + DeliveryAttemptOutcome::SinkFailure(value) => Err(value.clone()), + }) + .unwrap(); + clock.0.store(0, Ordering::Relaxed); + assert_eq!(block_on(future), Err(Error::ClockUnavailable)); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!( + status.delivery_plan().delivery_facts()[0].outcome(), + &outcome + ); + assert_eq!( + status.delivery_plan().revision(), + before.delivery_plan().revision() + ); + assert_eq!(status.delivery_plan().attempt_count(), 0); + let expires = status + .delivery_plan() + .claim_evidence() + .unwrap() + .expires_at_unix_ms(); + clock.0.store(expires, Ordering::Relaxed); + let recovered = + block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); + assert_eq!(recovered.plan().attempt_count(), 1); + assert_eq!(recovered.plan().attempts()[0].outcome(), &outcome); + if !accepted { + assert!(recovered.plan().retry().unwrap().not_before_unix_ms() >= expires + 1_000); + } + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); + } +} + +#[test] +fn stop_before_delivery_proves_no_issued_attempt_and_stop_after_acceptance_keeps_success() { + let (engine, _, _, sink, push) = setup(126); + let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); + assert!( + stopped + .status() + .delivery_history() + .proves_no_issued_attempt() + ); + assert!( + block_on(engine.deliver_push(push.operation_id())) + .unwrap() + .is_replay() + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 0); + let (engine, _, _, sink, push) = setup(127); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + accept(sink.take()); + assert_eq!( + block_on(future).unwrap().plan().state(), + AuthoredDeliveryState::Satisfied + ); + let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); + assert!(stopped.changed()); + assert!( + stopped + .status() + .delivery_plan() + .stop_requested_at_unix_ms() + .is_some() + ); + assert_eq!( + stopped.status().delivery_plan().state(), + AuthoredDeliveryState::Satisfied + ); + assert_eq!( + stopped + .status() + .delivery_plan() + .delivery_satisfaction() + .unwrap(), + SatisfactionState::Satisfied + ); + assert!( + !block_on(engine.cancel_push(push.operation_id())) + .unwrap() + .changed() + ); +} + +#[test] +fn equal_fact_replay_preserves_first_observation_and_conflict_cannot_erase_acceptance() { + let (engine, storage, _, sink, push) = setup(128); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + accept(sink.take()); + let delivered = block_on(future).unwrap(); + let plan = delivered.plan(); + let fact = &plan.delivery_facts()[0]; + let repeated = AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + fact.claim().clone(), + fact.outcome().clone(), + fact.observed_at_unix_ms() + 10_000, + ) + .unwrap(), + ); + block_on(storage.execute_authored(repeated)).unwrap(); + let changed = DeliveryAttemptOutcome::Receipt( + receipt( + plan.request().unwrap(), + vec![DeliveryOutcome::unavailable()], + ) + .unwrap(), + ); + let conflicting = AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + fact.claim().clone(), + changed, + fact.observed_at_unix_ms() + 10_000, + ) + .unwrap(), + ); + assert!(block_on(storage.execute_authored(conflicting)).is_err()); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(status.delivery_plan(), plan); +} +#[test] +fn every_delivery_storage_receipt_is_checked_and_durable_results_remain_recoverable() { + for nth in 1..=3 { + let storage = Arc::new(FaultStorage::new(150 + nth as u8)); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + ])])); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(150 + nth as u8, "wss://relay.example"); + execute_to_admitted(&engine, &push); + storage.fault_nth_plan(nth); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::StorageFailed) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(sink.requests.lock().unwrap().len(), usize::from(nth > 1)); + assert_eq!( + status.delivery_plan().delivery_facts().len(), + usize::from(nth > 1) + ); + if nth > 1 { + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_250_000))); + let recovered = + block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap(); + assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(recovered.plan().attempt_count(), 1); + } + } +} + +#[test] +fn legacy_applied_result_gets_exact_marker_without_counting_a_second_attempt() { + legacy_reconciliation(false); +} + +#[test] +fn conflicting_legacy_result_retains_both_observations_without_reconciliation() { + legacy_reconciliation(true); +} + +fn legacy_reconciliation(conflict: bool) { + use core::num::NonZeroU32; + use radroots_storage::{ + authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase}, + authored_atomic::{ApplyDeliveryAttempt, WorkFence}, + }; + let (engine, storage, clock, sink, push) = setup(154); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (request, _sender) = sink.take(); + drop(future); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let plan = status.delivery_plan(); + let claim = plan.claim_evidence().unwrap(); + let at = claim.acquired_at_unix_ms() + 100; + let outcome = DeliveryAttemptOutcome::Receipt( + receipt(&request, vec![DeliveryOutcome::unavailable()]).unwrap(), + ); + let retry = RetrySchedule::new( + NonZeroU32::MIN, + at + 1000, + WorkFailure::new( + "delivery_pending", + WorkPhase::Delivery, + FailureClass::Retryable, + Some(at + 1000), + None, + ) + .unwrap(), + ) + .unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery( + ApplyDeliveryAttempt::new( + plan.plan_id(), + WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap(), + outcome.clone(), + Some(retry), + at, + ) + .unwrap(), + )), + ) + .unwrap(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + claim.clone(), + if conflict { + DeliveryAttemptOutcome::Receipt( + receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), + ) + } else { + outcome.clone() + }, + at + 1, + ) + .unwrap(), + )), + ) + .unwrap(); + clock.0.store(at + 2, Ordering::Relaxed); + let result = block_on(source_only(storage, clock).deliver_push(push.operation_id())); + if conflict { + assert_eq!(result, Err(Error::StorageConflict)); + let current = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(current.delivery_plan().attempt_count(), 1); + assert_eq!(current.delivery_plan().attempts()[0].outcome(), &outcome); + assert_eq!(current.delivery_plan().delivery_facts().len(), 1); + assert_eq!( + current.delivery_plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); + return; + } + let reconciled = result.unwrap(); + assert_eq!(reconciled.plan().attempt_count(), 1); + assert_eq!(reconciled.plan().attempts()[0].recorded_at_unix_ms(), at); + assert_eq!(reconciled.plan().attempts()[0].outcome(), &outcome); + assert_eq!( + reconciled.plan().attempts()[0].claim_evidence(), + Some(claim) + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); +} + +#[tokio::test] +async fn sqlite_late_facts_reopen_after_stop_or_expiry_and_reconcile_once() { + use radroots_storage::authored_atomic::{CancelAuthoredTarget, CancelAuthoredWork}; + struct BoundarySink { + storage: Arc<SqliteStorage>, + clock: Arc<TestClock>, + id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId, + stop: bool, + } + impl EventSink for BoundarySink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { panic!("no probe") }) + } + fn deliver( + &self, + request: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> { + Box::pin(async move { + let history = self + .storage + .authored_delivery_history(self.id) + .await + .unwrap() + .unwrap(); + let plan = history.plan(); + let at = plan.claim_evidence().unwrap().expires_at_unix_ms(); + self.clock.0.store(at, Ordering::Relaxed); + if self.stop { + self.storage + .execute_authored(AuthoredAtomicCommand::Cancel( + CancelAuthoredWork::new( + CancelAuthoredTarget::DeliveryPlan(self.id), + plan.revision(), + at, + ) + .unwrap(), + )) + .await + .unwrap(); + } + Ok(receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap()) + }) + } + } + for stop in [true, false] { + let directory = tempfile::tempdir().unwrap(); + let paths = Paths::from_directory(directory.path()).unwrap(); + let storage = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([155; 32]).unwrap(), 1) + .unwrap(), + ) + .await + .unwrap(), + ); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let push = request(155, "wss://relay.example"); + let id = push + .authored_preparation(1_800_000_200_000) + .unwrap() + .delivery_plans()[0] + .plan_id(); + let engine = Engine::builder( + storage.clone(), + clock.clone(), + Arc::new(TestIds(AtomicU64::new(10))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(Arc::new(BoundarySink { + storage: storage.clone(), + clock: clock.clone(), + id, + stop, + })) + .signer(Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + }))) + .build() + .unwrap(); + engine.sign_prepared(push.clone()).await.unwrap(); + engine.admit_signed(push.operation_id()).await.unwrap(); + let late = engine.deliver_push(push.operation_id()).await.unwrap(); + assert_eq!(late.plan().attempt_count(), 0); + assert_eq!( + late.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + drop(engine); + storage.close().await.unwrap(); + let storage = Arc::new( + SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting)) + .await + .unwrap(), + ); + let recovery = source_only(storage.clone(), clock.clone()); + let recovered = recovery.deliver_push(push.operation_id()).await.unwrap(); + assert_eq!(recovered.plan().attempt_count(), u32::from(!stop)); + assert_eq!( + recovered.plan().delivery_facts(), + late.plan().delivery_facts() + ); + let status = recovery + .push_status(push.operation_id()) + .await + .unwrap() + .unwrap(); + assert!(!status.delivery_history().has_unresolved_claims()); + drop(recovery); + storage.close().await.unwrap(); + let storage = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly)) + .await + .unwrap(), + ); + let readonly = source_only(storage.clone(), clock); + let replay = readonly.deliver_push(push.operation_id()).await.unwrap(); + assert_eq!(replay.plan(), recovered.plan()); + drop(readonly); + storage.close().await.unwrap(); + } +} +pub(super) enum Injection { + FactBeforeClaim(Box<AuthoredAtomicCommand>), + StopAfterClaim, + ReplaceAfterClaim, + StopAfterFact, + StopAfterReconcile, +} + +pub(super) async fn inject(storage: &FaultStorage, command: &AuthoredAtomicCommand, before: bool) { + use radroots_storage::{ + authored::WorkClaim, + authored_atomic::{ + CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, ClaimAuthoredWork, + }, + }; + let claim = matches!(command, AuthoredAtomicCommand::Claim(value) + if matches!(value.target(), ClaimAuthoredTarget::DeliveryPlan(_))); + let injection = { + let mut armed = storage.delivery_injection.lock().unwrap(); + let ready = match armed.as_ref() { + Some(Injection::FactBeforeClaim(_)) => before && claim, + Some(Injection::StopAfterClaim | Injection::ReplaceAfterClaim) => !before && claim, + Some(Injection::StopAfterFact) => { + !before && matches!(command, AuthoredAtomicCommand::RecordDelivery(_)) + } + Some(Injection::StopAfterReconcile) => { + !before && matches!(command, AuthoredAtomicCommand::ReconcileDelivery(_)) + } + None => false, + }; + if ready { armed.take() } else { None } + }; + let Some(injection) = injection else { + return; + }; + if let Injection::FactBeforeClaim(fact) = injection { + storage.inner.execute_authored(*fact).await.unwrap(); + return; + } + let id = match command { + AuthoredAtomicCommand::Claim(value) => match value.target() { + ClaimAuthoredTarget::DeliveryPlan(id) => *id, + _ => panic!("delivery-only injection"), + }, + AuthoredAtomicCommand::RecordDelivery(value) => value.plan_id(), + AuthoredAtomicCommand::ReconcileDelivery(value) => value.plan_id(), + _ => panic!("delivery-only injection"), + }; + let history = storage + .inner + .authored_delivery_history(id) + .await + .unwrap() + .unwrap(); + let plan = history.plan(); + let mutation = if matches!(injection, Injection::ReplaceAfterClaim) { + let original = plan.claim_evidence().unwrap(); + let at = original.expires_at_unix_ms(); + AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new( + ClaimAuthoredTarget::DeliveryPlan(id), + WorkClaim::new( + [222; 16], + "competing-worker", + core::num::NonZeroU64::new(2).unwrap(), + at, + at + 10_000, + plan.revision(), + ) + .unwrap(), + )) + } else { + AuthoredAtomicCommand::Cancel( + CancelAuthoredWork::new( + CancelAuthoredTarget::DeliveryPlan(id), + plan.revision(), + plan.updated_at_unix_ms(), + ) + .unwrap(), + ) + }; + storage.inner.execute_authored(mutation).await.unwrap(); +} + +#[test] +fn observed_stop_or_replacement_after_claim_prevents_sink_admission() { + for stop in [true, false] { + let storage = Arc::new(FaultStorage::new(160)); + let sink = Arc::new(HeldSink::default()); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(160, "wss://relay.example"); + execute_to_admitted(&engine, &push); + *storage.delivery_injection.lock().unwrap() = Some(if stop { + Injection::StopAfterClaim + } else { + Injection::ReplaceAfterClaim + }); + let result = block_on(engine.deliver_push(push.operation_id())); + if stop { + assert_eq!( + result.unwrap().plan().state(), + AuthoredDeliveryState::Cancelled + ); + } else { + assert_eq!(result, Err(Error::WorkClaimConflict)); + } + assert_eq!(sink.calls.load(Ordering::Relaxed), 0); + } +} + +#[test] +fn fresh_claim_reconciles_fact_arriving_after_initial_status_without_another_effect() { + let storage = Arc::new(FaultStorage::new(161)); + let sink = Arc::new(HeldSink::default()); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let engine = Engine::builder( + storage.clone(), + clock.clone(), + Arc::new(TestIds(AtomicU64::new(10))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(sink.clone()) + .signer(Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + }))) + .build() + .unwrap(); + let push = request(161, "wss://relay.example"); + execute_to_admitted(&engine, &push); + let mut first = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + first + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (request, _sender) = sink.take(); + drop(first); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let plan = status.delivery_plan(); + let claim = plan.claim_evidence().unwrap(); + let at = claim.expires_at_unix_ms(); + let fact = AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + claim.clone(), + DeliveryAttemptOutcome::Receipt( + receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), + ), + at, + ) + .unwrap(), + ); + *storage.delivery_injection.lock().unwrap() = Some(Injection::FactBeforeClaim(Box::new(fact))); + clock.0.store(at, Ordering::Relaxed); + let reconciled = block_on(engine.deliver_push(push.operation_id())).unwrap(); + assert!(reconciled.is_replay()); + assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied); + assert_eq!(reconciled.plan().attempt_count(), 1); + assert_eq!(sink.calls.load(Ordering::Relaxed), 1); +} + +#[test] +fn receipt_return_reloads_stop_after_fact_or_reconciliation_commit() { + for before_reconcile in [true, false] { + let storage = Arc::new(FaultStorage::new(162)); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + ])])); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(162, "wss://relay.example"); + execute_to_admitted(&engine, &push); + *storage.delivery_injection.lock().unwrap() = Some(if before_reconcile { + Injection::StopAfterFact + } else { + Injection::StopAfterReconcile + }); + let delivered = block_on(engine.deliver_push(push.operation_id())).unwrap(); + assert!(delivered.plan().stop_requested_at_unix_ms().is_some()); + assert_eq!( + delivered.plan().attempt_count(), + u32::from(!before_reconcile) + ); + assert_eq!( + delivered.plan().delivery_satisfaction().unwrap(), + SatisfactionState::Satisfied + ); + assert_eq!( + engine + .retry_decision(delivered.plan(), 1_800_000_250_000) + .unwrap(), + SyncRetryDecision::Satisfied + ); + assert_eq!(sink.requests.lock().unwrap().len(), 1); + } +} + +#[test] +fn forged_stop_receipt_is_not_reported_as_success() { + let storage = Arc::new(FaultStorage::new(163)); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + Arc::new(MockSink), + ); + let push = request(163, "wss://relay.example"); + execute_to_admitted(&engine, &push); + storage.fault_nth_plan(1); + assert_eq!( + block_on(engine.cancel_push(push.operation_id())), + Err(Error::StorageFailed) + ); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(status.delivery_plan().stop_requested_at_unix_ms().is_some()); + assert!(status.delivery_history().proves_no_issued_attempt()); +} +use radroots_storage::authored_atomic::{AuthoredAtomicOutcome, AuthoredAtomicReceipt}; +use radroots_storage::authored_delivery::{AuthoredDeliveryHistory, AuthoredDeliveryPlan}; + +#[derive(Clone, Copy, Debug)] +pub(super) enum ReceiptMutation { + Identity, + Plan, + Artifact, + Request, + Created, + Claim, + CommitTime, + Attempts, + State, + Retry, + StopTime, + Revision, +} +pub(super) fn mutate_receipt( + receipt: AuthoredAtomicReceipt, + mutation: ReceiptMutation, +) -> AuthoredAtomicReceipt { + let AuthoredAtomicOutcome::DeliveryPlan(plan) = receipt.outcome() else { + panic!("delivery receipt"); + }; + let mut wire = serde_json::to_value(plan).unwrap(); + match mutation { + ReceiptMutation::Identity | ReceiptMutation::CommitTime => {} + ReceiptMutation::Plan => wire["plan_id"] = serde_json::json!(vec![241; 16]), + ReceiptMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![242; 16]), + ReceiptMutation::Request => { + let original = plan.request().unwrap(); + let changed = DeliveryRequest::new( + "different-request", + original.payload().clone(), + original.target_set().clone(), + original.satisfaction().clone(), + original.deadline_unix_ms(), + ) + .unwrap(); + let other = AuthoredDeliveryPlan::new_bound( + plan.plan_id(), + plan.artifact_id(), + changed, + plan.created_at_unix_ms(), + ) + .unwrap(); + let other = serde_json::to_value(other).unwrap(); + for key in ["request", "intent", "request_digest"] { + wire[key] = other[key].clone(); + } + } + ReceiptMutation::Created => { + wire["created_at_unix_ms"] = serde_json::json!(plan.created_at_unix_ms() - 1) + } + ReceiptMutation::Claim => wire["claim"]["owner"] = serde_json::json!("another-worker"), + ReceiptMutation::Attempts => { + let outcome = DeliveryAttemptOutcome::Receipt(receipt_for_pending(plan)); + wire["attempts"] = serde_json::json!([{"attempt":1,"recorded_at_unix_ms":plan.created_at_unix_ms(),"outcome":outcome,"satisfaction":"pending"}]); + wire["attempt_count"] = serde_json::json!(1); + } + ReceiptMutation::State => { + wire["state"] = serde_json::json!("pending"); + wire["retry"] = serde_json::Value::Null; + wire["last_failure"] = serde_json::Value::Null; + } + ReceiptMutation::Retry => { + let at = plan.retry().unwrap().not_before_unix_ms() + 1000; + wire["retry"]["not_before_unix_ms"] = serde_json::json!(at); + wire["retry"]["failure"]["retry_after_unix_ms"] = serde_json::json!(at); + wire["last_failure"] = wire["retry"]["failure"].clone(); + } + ReceiptMutation::StopTime => { + wire["stop_requested_at_unix_ms"] = + serde_json::json!(plan.stop_requested_at_unix_ms().unwrap() - 1) + } + ReceiptMutation::Revision => { + wire["revision"] = serde_json::json!(plan.revision().get() + 1) + } + } + let changed: AuthoredDeliveryPlan = + serde_json::from_value(wire).expect("valid but incorrectly bound projection"); + AuthoredAtomicReceipt::from_durable_parts( + if matches!(mutation, ReceiptMutation::Identity) { + radroots_storage::atomic::AtomicCommitId::new([243; 16]).unwrap() + } else { + receipt.commit_id() + }, + receipt.digest(), + receipt.disposition(), + receipt.committed_at_unix_ms() + u64::from(matches!(mutation, ReceiptMutation::CommitTime)), + AuthoredAtomicOutcome::DeliveryPlan(changed), + ) + .unwrap() +} +fn receipt_for_pending(plan: &AuthoredDeliveryPlan) -> DeliveryReceipt { + receipt( + plan.request().unwrap(), + vec![DeliveryOutcome::unavailable()], + ) + .unwrap() +} + +#[test] +fn claim_projection_must_match_every_original_binding_before_transport() { + for mutation in [ + ReceiptMutation::Identity, + ReceiptMutation::Plan, + ReceiptMutation::Artifact, + ReceiptMutation::Request, + ReceiptMutation::Created, + ReceiptMutation::Claim, + ReceiptMutation::CommitTime, + ReceiptMutation::Attempts, + ] { + let storage = Arc::new(FaultStorage::new(171)); + let sink = Arc::new(HeldSink::default()); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(171, "wss://relay.example"); + execute_to_admitted(&engine, &push); + *storage.receipt_mutation.lock().unwrap() = Some(mutation); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::StorageFailed), + "{mutation:?}" + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 0); + } +} + +#[test] +fn stop_projection_must_match_identity_request_time_and_revision() { + for mutation in [ + ReceiptMutation::Identity, + ReceiptMutation::Plan, + ReceiptMutation::Artifact, + ReceiptMutation::Request, + ReceiptMutation::StopTime, + ReceiptMutation::Revision, + ] { + let storage = Arc::new(FaultStorage::new(172)); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + Arc::new(MockSink), + ); + let push = request(172, "wss://relay.example"); + execute_to_admitted(&engine, &push); + *storage.receipt_mutation.lock().unwrap() = Some(mutation); + assert_eq!( + block_on(engine.cancel_push(push.operation_id())), + Err(Error::StorageFailed), + "{mutation:?}" + ); + let current = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(current.delivery_history().proves_no_issued_attempt()); + assert!( + current + .delivery_plan() + .stop_requested_at_unix_ms() + .is_some() + ); + } +} + +#[test] +fn claim_projection_cannot_replace_retained_retry_state_or_backoff() { + for mutation in [ReceiptMutation::State, ReceiptMutation::Retry] { + let storage = Arc::new(FaultStorage::new(173)); + let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Failure( + Retryability::Retryable, + )])); + let engine = Engine::builder( + storage.clone(), + clock.clone(), + Arc::new(TestIds(AtomicU64::new(10))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(sink.clone()) + .signer(Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + }))) + .build() + .unwrap(); + let push = request(173, "wss://relay.example"); + execute_to_admitted(&engine, &push); + let first = block_on(engine.deliver_push(push.operation_id())).unwrap(); + clock.0.store( + first.plan().retry().unwrap().not_before_unix_ms(), + Ordering::Relaxed, + ); + *storage.receipt_mutation.lock().unwrap() = Some(mutation); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::StorageFailed), + "{mutation:?}" + ); + assert_eq!(sink.requests.lock().unwrap().len(), 1); + } +} + +#[derive(Clone, Copy)] +pub(super) enum HistoryMutation { + Plan, + Artifact, + MissingFact, + MissingStop, +} +pub(super) fn mutate_history( + history: AuthoredDeliveryHistory, + mutation: HistoryMutation, +) -> AuthoredDeliveryHistory { + let mut wire = serde_json::to_value(history.plan()).unwrap(); + match mutation { + HistoryMutation::Plan => wire["plan_id"] = serde_json::json!(vec![244; 16]), + HistoryMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![245; 16]), + HistoryMutation::MissingFact => wire["delivery_facts"] = serde_json::json!([]), + HistoryMutation::MissingStop => { + wire["stop_requested_at_unix_ms"] = serde_json::Value::Null; + wire["state"] = serde_json::json!("pending"); + } + } + AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), None).unwrap() +} + +#[test] +fn wrong_history_and_lost_committed_fact_or_stop_never_become_success() { + for (nth, mutation, stop) in [ + (1, HistoryMutation::Plan, false), + (1, HistoryMutation::Artifact, false), + (2, HistoryMutation::Plan, false), + (3, HistoryMutation::MissingFact, false), + (2, HistoryMutation::MissingStop, true), + ] { + let storage = Arc::new(FaultStorage::new(174)); + let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![ + DeliveryOutcome::accepted(), + ])])); + let engine = fault_engine( + storage.clone(), + Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })), + sink.clone(), + ); + let push = request(174, "wss://relay.example"); + execute_to_admitted(&engine, &push); + *storage.history_mutation.lock().unwrap() = Some((nth, mutation)); + if stop { + assert_eq!( + block_on(engine.cancel_push(push.operation_id())), + Err(Error::StorageFailed) + ); + } else { + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::StorageFailed) + ); + } + let actual = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!( + actual.delivery_plan().delivery_facts().len(), + usize::from(nth == 3) + ); + } +} + +#[test] +fn clock_expiry_before_sink_and_future_fact_before_reconciliation_fail_closed() { + struct ExpiringClock(AtomicUsize); + impl Clock for ExpiringClock { + fn now_unix_ms(&self) -> Result<u64, Error> { + Ok(1_800_000_220_000 + 10_000 * self.0.fetch_add(1, Ordering::Relaxed) as u64) + } + } + let (engine, storage, _, sink, push) = setup(175); + let boundary = Engine::builder( + storage, + Arc::new(ExpiringClock(AtomicUsize::new(0))), + Arc::new(TestIds(AtomicU64::new(230))), + DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(), + ) + .sink(sink.clone()) + .build() + .unwrap(); + assert_eq!( + block_on(boundary.deliver_push(push.operation_id())), + Err(Error::WorkClaimConflict) + ); + assert_eq!(sink.calls.load(Ordering::Relaxed), 0); + assert_eq!( + block_on(engine.deliver_push(SyncId::new([231; 16]).unwrap())), + Err(Error::StorageFailed) + ); + let (engine, storage, clock, sink, push) = setup(176); + let mut future = Box::pin(engine.deliver_push(push.operation_id())); + assert!( + future + .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref())) + .is_pending() + ); + let (request, _sender) = sink.take(); + drop(future); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + let plan = status.delivery_plan(); + let claim = plan.claim_evidence().unwrap(); + let at = claim.expires_at_unix_ms(); + block_on( + storage.execute_authored(AuthoredAtomicCommand::RecordDelivery( + RecordDeliveryFact::new( + plan.plan_id(), + plan.artifact_id(), + claim.clone(), + DeliveryAttemptOutcome::Receipt( + receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(), + ), + at + 1000, + ) + .unwrap(), + )), + ) + .unwrap(); + clock.0.store(at, Ordering::Relaxed); + assert_eq!( + block_on(engine.deliver_push(push.operation_id())), + Err(Error::ClockUnavailable) + ); + let current = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert_eq!(current.delivery_plan().attempt_count(), 0); + assert_eq!(current.delivery_plan().delivery_facts().len(), 1); +} + +#[test] +fn unbound_preparation_and_stop_keep_honest_retry_and_no_issued_proof() { + let signer = Arc::new(MockSigner::new(SignBehavior::Success { + completed_at_unix_ms: 1_800_000_200_500, + })); + let (engine, _) = setup_engine(signer); + let push = request(177, "wss://relay.example"); + block_on(engine.prepare_push(push.clone())).unwrap(); + let status = block_on(engine.push_status(push.operation_id())) + .unwrap() + .unwrap(); + assert!(status.delivery_plan().request().is_none()); + assert!(status.delivery_history().proves_no_issued_attempt()); + assert_eq!( + engine + .retry_decision(status.delivery_plan(), 1_800_000_200_001) + .unwrap(), + SyncRetryDecision::Ready + ); + let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap(); + assert!( + stopped + .status() + .delivery_history() + .proves_no_issued_attempt() + ); + assert_eq!( + engine + .retry_decision(stopped.status().delivery_plan(), 1_800_000_200_010) + .unwrap(), + SyncRetryDecision::Exhausted + ); +}