lib

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

delivery_evidence.rs (48672B)


      1 use super::*;
      2 use futures::channel::oneshot;
      3 use radroots_storage::{
      4     authored_atomic::{AuthoredAtomicCommand, RecordDeliveryFact},
      5     authored_delivery::DeliveryAttemptOutcome,
      6 };
      7 use radroots_transport::policy::SatisfactionState;
      8 
      9 type Pending = (
     10     DeliveryRequest,
     11     oneshot::Sender<Result<DeliveryReceipt, SinkFailure>>,
     12 );
     13 
     14 #[derive(Default)]
     15 struct HeldSink {
     16     pending: Mutex<VecDeque<Pending>>,
     17     calls: AtomicUsize,
     18 }
     19 
     20 impl HeldSink {
     21     fn take(&self) -> Pending {
     22         self.pending.lock().unwrap().pop_front().unwrap()
     23     }
     24 }
     25 impl EventSink for HeldSink {
     26     fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
     27         Box::pin(async { panic!("delivery must not probe transport status") })
     28     }
     29     fn deliver(
     30         &self,
     31         request: DeliveryRequest,
     32     ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
     33         Box::pin(async move {
     34             self.calls.fetch_add(1, Ordering::Relaxed);
     35             let (sender, receiver) = oneshot::channel();
     36             self.pending.lock().unwrap().push_back((request, sender));
     37             receiver.await.expect("host retains admitted future")
     38         })
     39     }
     40 }
     41 fn setup(
     42     byte: u8,
     43 ) -> (
     44     Engine,
     45     Arc<MemoryStorage>,
     46     Arc<TestClock>,
     47     Arc<HeldSink>,
     48     PushRequest,
     49 ) {
     50     let sink = Arc::new(HeldSink::default());
     51     let ((engine, storage), clock) = setup_engine_with_sink(
     52         Arc::new(MockSigner::new(SignBehavior::Success {
     53             completed_at_unix_ms: 1_800_000_200_500,
     54         })),
     55         sink.clone(),
     56     );
     57     let push = request(byte, "wss://relay.example");
     58     execute_to_admitted(&engine, &push);
     59     (engine, storage, clock, sink, push)
     60 }
     61 pub(super) fn source_only(storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>) -> Engine {
     62     Engine::builder(
     63         storage,
     64         clock,
     65         Arc::new(TestIds(AtomicU64::new(230))),
     66         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
     67     )
     68     .source(Arc::new(MockSource))
     69     .build()
     70     .unwrap()
     71 }
     72 fn accept((request, sender): Pending) -> DeliveryReceipt {
     73     let result = receipt(
     74         &request,
     75         vec![DeliveryOutcome::accepted(); request.target_set().len()],
     76     )
     77     .unwrap();
     78     sender.send(Ok(result.clone())).unwrap();
     79     result
     80 }
     81 
     82 #[test]
     83 fn capacity_delivery_receipt_failure_retains_claim_and_known_or_unknown_effects() {
     84     for after in [false, true] {
     85         let storage = Arc::new(FaultStorage::new(185));
     86         let sink = Arc::new(HeldSink::default());
     87         let signer = Arc::new(MockSigner::new(SignBehavior::Success {
     88             completed_at_unix_ms: 1_800_000_200_500,
     89         }));
     90         let engine = fault_engine(storage.clone(), signer.clone(), sink.clone());
     91         let push = request(185, "wss://capacity.example");
     92         execute_to_admitted(&engine, &push);
     93         let before = block_on(engine.push_status(push.operation_id()))
     94             .unwrap()
     95             .unwrap();
     96         let mut future = Box::pin(engine.deliver_push(push.operation_id()));
     97         assert!(
     98             future
     99                 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    100                 .is_pending()
    101         );
    102         let pending = sink.take();
    103         assert_eq!(&pending.0, before.delivery_plan().request().unwrap());
    104         *storage.capacity.lock().unwrap() = Some(capacity::Fault {
    105             phase: capacity::Phase::DeliveryFact,
    106             after,
    107         });
    108         let expected = accept(pending);
    109         assert_eq!(block_on(future), Err(Error::StorageSpaceInsufficient));
    110         let failed = block_on(engine.push_status(push.operation_id()))
    111             .unwrap()
    112             .unwrap();
    113         assert_eq!(failed.artifact().signed(), before.artifact().signed());
    114         assert_eq!(
    115             failed.delivery_plan().request(),
    116             before.delivery_plan().request()
    117         );
    118         assert!(!failed.delivery_history().proves_no_issued_attempt());
    119         assert_eq!(failed.delivery_history().has_unresolved_claims(), !after);
    120         assert_eq!(
    121             failed.delivery_plan().delivery_facts().len(),
    122             usize::from(after)
    123         );
    124         if after {
    125             assert_eq!(
    126                 failed.delivery_plan().delivery_facts()[0].outcome(),
    127                 &DeliveryAttemptOutcome::Receipt(expected)
    128             );
    129             assert_eq!(
    130                 block_on(engine.deliver_push(push.operation_id())),
    131                 Err(Error::WorkClaimConflict)
    132             );
    133             let expires_at = failed
    134                 .delivery_plan()
    135                 .claim_evidence()
    136                 .unwrap()
    137                 .expires_at_unix_ms();
    138             let recovery = source_only(
    139                 storage.clone(),
    140                 Arc::new(TestClock(AtomicU64::new(expires_at))),
    141             );
    142             let reconciled = block_on(recovery.deliver_push(push.operation_id())).unwrap();
    143             assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied);
    144         }
    145         assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    146         assert_eq!(signer.calls.load(Ordering::Relaxed), 1);
    147     }
    148 }
    149 
    150 #[test]
    151 fn late_accepted_result_survives_stop_and_terminal_replay_without_sink() {
    152     let (engine, storage, clock, sink, push) = setup(121);
    153     let status = block_on(engine.push_status(push.operation_id()))
    154         .unwrap()
    155         .unwrap();
    156     assert!(status.delivery_history().proves_no_issued_attempt());
    157     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    158     assert!(
    159         future
    160             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    161             .is_pending()
    162     );
    163     let pending = sink.take();
    164     let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap();
    165     assert!(stopped.changed());
    166     assert!(stopped.status().delivery_history().has_unresolved_claims());
    167     assert!(
    168         !stopped
    169             .status()
    170             .delivery_history()
    171             .proves_no_issued_attempt()
    172     );
    173     let stop_at = stopped.status().delivery_plan().stop_requested_at_unix_ms();
    174     let expected = accept(pending);
    175     let delivered = block_on(future).unwrap();
    176     assert_eq!(delivered.plan().state(), AuthoredDeliveryState::Cancelled);
    177     assert_eq!(delivered.plan().attempt_count(), 0);
    178     assert_eq!(
    179         delivered.plan().delivery_satisfaction().unwrap(),
    180         SatisfactionState::Satisfied
    181     );
    182     assert_eq!(
    183         delivered.plan().delivery_facts()[0].outcome(),
    184         &DeliveryAttemptOutcome::Receipt(expected)
    185     );
    186     let status = block_on(engine.push_status(push.operation_id()))
    187         .unwrap()
    188         .unwrap();
    189     assert!(!status.delivery_history().has_unresolved_claims());
    190     assert_eq!(status.delivery_plan().stop_requested_at_unix_ms(), stop_at);
    191     assert!(
    192         !block_on(engine.cancel_push(push.operation_id()))
    193             .unwrap()
    194             .changed()
    195     );
    196     let replay = block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap();
    197     assert!(replay.is_replay());
    198     assert_eq!(replay.plan(), status.delivery_plan());
    199     assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    200 }
    201 
    202 #[test]
    203 fn expired_callback_retains_fact_but_only_fresh_invocation_reconciles() {
    204     let (engine, storage, clock, sink, push) = setup(122);
    205     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    206     assert!(
    207         future
    208             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    209             .is_pending()
    210     );
    211     let before = block_on(engine.push_status(push.operation_id()))
    212         .unwrap()
    213         .unwrap();
    214     clock.0.store(
    215         before
    216             .delivery_plan()
    217             .claim_evidence()
    218             .unwrap()
    219             .expires_at_unix_ms(),
    220         Ordering::Relaxed,
    221     );
    222     accept(sink.take());
    223     let late = block_on(future).unwrap();
    224     assert_eq!(late.plan().revision(), before.delivery_plan().revision());
    225     assert_eq!(late.plan().attempt_count(), 0);
    226     let recovery = source_only(storage, clock);
    227     let recovered = block_on(recovery.deliver_push(push.operation_id())).unwrap();
    228     assert!(recovered.is_replay());
    229     assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied);
    230     assert_eq!(recovered.plan().attempt_count(), 1);
    231     assert_eq!(
    232         recovered.plan().attempts()[0].claim_evidence(),
    233         before.delivery_plan().claim_evidence()
    234     );
    235     assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    236 }
    237 
    238 #[test]
    239 fn superseded_callback_cannot_clear_newer_claim_and_acceptance_never_regresses() {
    240     let (engine, _, clock, sink, push) = setup(123);
    241     let mut first = Box::pin(engine.deliver_push(push.operation_id()));
    242     assert!(
    243         first
    244             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    245             .is_pending()
    246     );
    247     let first_pending = sink.take();
    248     let before = block_on(engine.push_status(push.operation_id()))
    249         .unwrap()
    250         .unwrap();
    251     clock.0.store(
    252         before
    253             .delivery_plan()
    254             .claim_evidence()
    255             .unwrap()
    256             .expires_at_unix_ms(),
    257         Ordering::Relaxed,
    258     );
    259     let mut second = Box::pin(engine.deliver_push(push.operation_id()));
    260     assert!(
    261         second
    262             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    263             .is_pending()
    264     );
    265     let (second_request, second_sender) = sink.take();
    266     assert_eq!(first_pending.0, second_request);
    267     let newer = block_on(engine.push_status(push.operation_id()))
    268         .unwrap()
    269         .unwrap();
    270     accept(first_pending);
    271     let late = block_on(first).unwrap();
    272     assert_eq!(
    273         late.plan().claim_evidence(),
    274         newer.delivery_plan().claim_evidence()
    275     );
    276     assert_eq!(late.plan().revision(), newer.delivery_plan().revision());
    277     assert_eq!(
    278         block_on(engine.deliver_push(push.operation_id())),
    279         Err(Error::WorkClaimConflict)
    280     );
    281     second_sender
    282         .send(Err(SinkFailure::for_request(
    283             &second_request,
    284             "provider_failed",
    285             Retryability::Terminal,
    286             None,
    287             None,
    288             Vec::new(),
    289         )
    290         .unwrap()))
    291         .unwrap();
    292     let complete = block_on(second).unwrap();
    293     assert_eq!(complete.plan().state(), AuthoredDeliveryState::Satisfied);
    294     assert_eq!(complete.plan().attempt_count(), 2);
    295     assert_eq!(complete.plan().delivery_facts().len(), 2);
    296     assert_eq!(sink.calls.load(Ordering::Relaxed), 2);
    297 }
    298 
    299 #[test]
    300 fn clock_loss_retains_raw_result_before_reporting_unavailable() {
    301     for accepted in [true, false] {
    302         let (engine, storage, clock, sink, push) = setup(124 + u8::from(accepted));
    303         let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    304         assert!(
    305             future
    306                 .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    307                 .is_pending()
    308         );
    309         let (request, sender) = sink.take();
    310         let before = block_on(engine.push_status(push.operation_id()))
    311             .unwrap()
    312             .unwrap();
    313         let outcome = if accepted {
    314             DeliveryAttemptOutcome::Receipt(
    315                 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(),
    316             )
    317         } else {
    318             DeliveryAttemptOutcome::SinkFailure(
    319                 SinkFailure::for_request(
    320                     &request,
    321                     "raw_failure",
    322                     Retryability::Retryable,
    323                     None,
    324                     Some("raw diagnostic".into()),
    325                     Vec::new(),
    326                 )
    327                 .unwrap(),
    328             )
    329         };
    330         sender
    331             .send(match &outcome {
    332                 DeliveryAttemptOutcome::Receipt(value) => Ok(value.clone()),
    333                 DeliveryAttemptOutcome::SinkFailure(value) => Err(value.clone()),
    334             })
    335             .unwrap();
    336         clock.0.store(0, Ordering::Relaxed);
    337         assert_eq!(block_on(future), Err(Error::ClockUnavailable));
    338         let status = block_on(engine.push_status(push.operation_id()))
    339             .unwrap()
    340             .unwrap();
    341         assert_eq!(
    342             status.delivery_plan().delivery_facts()[0].outcome(),
    343             &outcome
    344         );
    345         assert_eq!(
    346             status.delivery_plan().revision(),
    347             before.delivery_plan().revision()
    348         );
    349         assert_eq!(status.delivery_plan().attempt_count(), 0);
    350         let expires = status
    351             .delivery_plan()
    352             .claim_evidence()
    353             .unwrap()
    354             .expires_at_unix_ms();
    355         clock.0.store(expires, Ordering::Relaxed);
    356         let recovered =
    357             block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap();
    358         assert_eq!(recovered.plan().attempt_count(), 1);
    359         assert_eq!(recovered.plan().attempts()[0].outcome(), &outcome);
    360         if !accepted {
    361             assert!(recovered.plan().retry().unwrap().not_before_unix_ms() >= expires + 1_000);
    362         }
    363         assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    364     }
    365 }
    366 
    367 #[test]
    368 fn stop_before_delivery_proves_no_issued_attempt_and_stop_after_acceptance_keeps_success() {
    369     let (engine, _, _, sink, push) = setup(126);
    370     let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap();
    371     assert!(
    372         stopped
    373             .status()
    374             .delivery_history()
    375             .proves_no_issued_attempt()
    376     );
    377     assert!(
    378         block_on(engine.deliver_push(push.operation_id()))
    379             .unwrap()
    380             .is_replay()
    381     );
    382     assert_eq!(sink.calls.load(Ordering::Relaxed), 0);
    383     let (engine, _, _, sink, push) = setup(127);
    384     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    385     assert!(
    386         future
    387             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    388             .is_pending()
    389     );
    390     accept(sink.take());
    391     assert_eq!(
    392         block_on(future).unwrap().plan().state(),
    393         AuthoredDeliveryState::Satisfied
    394     );
    395     let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap();
    396     assert!(stopped.changed());
    397     assert!(
    398         stopped
    399             .status()
    400             .delivery_plan()
    401             .stop_requested_at_unix_ms()
    402             .is_some()
    403     );
    404     assert_eq!(
    405         stopped.status().delivery_plan().state(),
    406         AuthoredDeliveryState::Satisfied
    407     );
    408     assert_eq!(
    409         stopped
    410             .status()
    411             .delivery_plan()
    412             .delivery_satisfaction()
    413             .unwrap(),
    414         SatisfactionState::Satisfied
    415     );
    416     assert!(
    417         !block_on(engine.cancel_push(push.operation_id()))
    418             .unwrap()
    419             .changed()
    420     );
    421 }
    422 
    423 #[test]
    424 fn equal_fact_replay_preserves_first_observation_and_conflict_cannot_erase_acceptance() {
    425     let (engine, storage, _, sink, push) = setup(128);
    426     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    427     assert!(
    428         future
    429             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    430             .is_pending()
    431     );
    432     accept(sink.take());
    433     let delivered = block_on(future).unwrap();
    434     let plan = delivered.plan();
    435     let fact = &plan.delivery_facts()[0];
    436     let repeated = AuthoredAtomicCommand::RecordDelivery(
    437         RecordDeliveryFact::new(
    438             plan.plan_id(),
    439             plan.artifact_id(),
    440             fact.claim().clone(),
    441             fact.outcome().clone(),
    442             fact.observed_at_unix_ms() + 10_000,
    443         )
    444         .unwrap(),
    445     );
    446     block_on(storage.execute_authored(repeated)).unwrap();
    447     let changed = DeliveryAttemptOutcome::Receipt(
    448         receipt(
    449             plan.request().unwrap(),
    450             vec![DeliveryOutcome::unavailable()],
    451         )
    452         .unwrap(),
    453     );
    454     let conflicting = AuthoredAtomicCommand::RecordDelivery(
    455         RecordDeliveryFact::new(
    456             plan.plan_id(),
    457             plan.artifact_id(),
    458             fact.claim().clone(),
    459             changed,
    460             fact.observed_at_unix_ms() + 10_000,
    461         )
    462         .unwrap(),
    463     );
    464     assert!(block_on(storage.execute_authored(conflicting)).is_err());
    465     let status = block_on(engine.push_status(push.operation_id()))
    466         .unwrap()
    467         .unwrap();
    468     assert_eq!(status.delivery_plan(), plan);
    469 }
    470 #[test]
    471 fn every_delivery_storage_receipt_is_checked_and_durable_results_remain_recoverable() {
    472     for nth in 1..=3 {
    473         let storage = Arc::new(FaultStorage::new(150 + nth as u8));
    474         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
    475             DeliveryOutcome::accepted(),
    476         ])]));
    477         let engine = fault_engine(
    478             storage.clone(),
    479             Arc::new(MockSigner::new(SignBehavior::Success {
    480                 completed_at_unix_ms: 1_800_000_200_500,
    481             })),
    482             sink.clone(),
    483         );
    484         let push = request(150 + nth as u8, "wss://relay.example");
    485         execute_to_admitted(&engine, &push);
    486         storage.fault_nth_plan(nth);
    487         assert_eq!(
    488             block_on(engine.deliver_push(push.operation_id())),
    489             Err(Error::StorageFailed)
    490         );
    491         let status = block_on(engine.push_status(push.operation_id()))
    492             .unwrap()
    493             .unwrap();
    494         assert_eq!(sink.requests.lock().unwrap().len(), usize::from(nth > 1));
    495         assert_eq!(
    496             status.delivery_plan().delivery_facts().len(),
    497             usize::from(nth > 1)
    498         );
    499         if nth > 1 {
    500             let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_250_000)));
    501             let recovered =
    502                 block_on(source_only(storage, clock).deliver_push(push.operation_id())).unwrap();
    503             assert_eq!(recovered.plan().state(), AuthoredDeliveryState::Satisfied);
    504             assert_eq!(recovered.plan().attempt_count(), 1);
    505         }
    506     }
    507 }
    508 
    509 #[test]
    510 fn legacy_applied_result_gets_exact_marker_without_counting_a_second_attempt() {
    511     legacy_reconciliation(false);
    512 }
    513 
    514 #[test]
    515 fn conflicting_legacy_result_retains_both_observations_without_reconciliation() {
    516     legacy_reconciliation(true);
    517 }
    518 
    519 fn legacy_reconciliation(conflict: bool) {
    520     use core::num::NonZeroU32;
    521     use radroots_storage::{
    522         authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase},
    523         authored_atomic::{ApplyDeliveryAttempt, WorkFence},
    524     };
    525     let (engine, storage, clock, sink, push) = setup(154);
    526     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
    527     assert!(
    528         future
    529             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    530             .is_pending()
    531     );
    532     let (request, _sender) = sink.take();
    533     drop(future);
    534     let status = block_on(engine.push_status(push.operation_id()))
    535         .unwrap()
    536         .unwrap();
    537     let plan = status.delivery_plan();
    538     let claim = plan.claim_evidence().unwrap();
    539     let at = claim.acquired_at_unix_ms() + 100;
    540     let outcome = DeliveryAttemptOutcome::Receipt(
    541         receipt(&request, vec![DeliveryOutcome::unavailable()]).unwrap(),
    542     );
    543     let retry = RetrySchedule::new(
    544         NonZeroU32::MIN,
    545         at + 1000,
    546         WorkFailure::new(
    547             "delivery_pending",
    548             WorkPhase::Delivery,
    549             FailureClass::Retryable,
    550             Some(at + 1000),
    551             None,
    552         )
    553         .unwrap(),
    554     )
    555     .unwrap();
    556     block_on(
    557         storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery(
    558             ApplyDeliveryAttempt::new(
    559                 plan.plan_id(),
    560                 WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap(),
    561                 outcome.clone(),
    562                 Some(retry),
    563                 at,
    564             )
    565             .unwrap(),
    566         )),
    567     )
    568     .unwrap();
    569     block_on(
    570         storage.execute_authored(AuthoredAtomicCommand::RecordDelivery(
    571             RecordDeliveryFact::new(
    572                 plan.plan_id(),
    573                 plan.artifact_id(),
    574                 claim.clone(),
    575                 if conflict {
    576                     DeliveryAttemptOutcome::Receipt(
    577                         receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(),
    578                     )
    579                 } else {
    580                     outcome.clone()
    581                 },
    582                 at + 1,
    583             )
    584             .unwrap(),
    585         )),
    586     )
    587     .unwrap();
    588     clock.0.store(at + 2, Ordering::Relaxed);
    589     let result = block_on(source_only(storage, clock).deliver_push(push.operation_id()));
    590     if conflict {
    591         assert_eq!(result, Err(Error::StorageConflict));
    592         let current = block_on(engine.push_status(push.operation_id()))
    593             .unwrap()
    594             .unwrap();
    595         assert_eq!(current.delivery_plan().attempt_count(), 1);
    596         assert_eq!(current.delivery_plan().attempts()[0].outcome(), &outcome);
    597         assert_eq!(current.delivery_plan().delivery_facts().len(), 1);
    598         assert_eq!(
    599             current.delivery_plan().delivery_satisfaction().unwrap(),
    600             SatisfactionState::Satisfied
    601         );
    602         assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    603         return;
    604     }
    605     let reconciled = result.unwrap();
    606     assert_eq!(reconciled.plan().attempt_count(), 1);
    607     assert_eq!(reconciled.plan().attempts()[0].recorded_at_unix_ms(), at);
    608     assert_eq!(reconciled.plan().attempts()[0].outcome(), &outcome);
    609     assert_eq!(
    610         reconciled.plan().attempts()[0].claim_evidence(),
    611         Some(claim)
    612     );
    613     assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    614 }
    615 
    616 #[tokio::test]
    617 async fn sqlite_late_facts_reopen_after_stop_or_expiry_and_reconcile_once() {
    618     use radroots_storage::authored_atomic::{CancelAuthoredTarget, CancelAuthoredWork};
    619     struct BoundarySink {
    620         storage: Arc<SqliteStorage>,
    621         clock: Arc<TestClock>,
    622         id: radroots_storage::authored_delivery::AuthoredDeliveryPlanId,
    623         stop: bool,
    624     }
    625     impl EventSink for BoundarySink {
    626         fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> {
    627             Box::pin(async { panic!("no probe") })
    628         }
    629         fn deliver(
    630             &self,
    631             request: DeliveryRequest,
    632         ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, SinkFailure>> {
    633             Box::pin(async move {
    634                 let history = self
    635                     .storage
    636                     .authored_delivery_history(self.id)
    637                     .await
    638                     .unwrap()
    639                     .unwrap();
    640                 let plan = history.plan();
    641                 let at = plan.claim_evidence().unwrap().expires_at_unix_ms();
    642                 self.clock.0.store(at, Ordering::Relaxed);
    643                 if self.stop {
    644                     self.storage
    645                         .execute_authored(AuthoredAtomicCommand::Cancel(
    646                             CancelAuthoredWork::new(
    647                                 CancelAuthoredTarget::DeliveryPlan(self.id),
    648                                 plan.revision(),
    649                                 at,
    650                             )
    651                             .unwrap(),
    652                         ))
    653                         .await
    654                         .unwrap();
    655                 }
    656                 Ok(receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap())
    657             })
    658         }
    659     }
    660     for stop in [true, false] {
    661         let directory = tempfile::tempdir().unwrap();
    662         let paths = Paths::from_directory(directory.path()).unwrap();
    663         let storage = Arc::new(
    664             SqliteStorage::open(
    665                 OpenOptions::new(paths.clone(), OpenMode::Create)
    666                     .with_source_generation(SourceGeneration::new([155; 32]).unwrap(), 1)
    667                     .unwrap(),
    668             )
    669             .await
    670             .unwrap(),
    671         );
    672         let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
    673         let push = request(155, "wss://relay.example");
    674         let id = push
    675             .authored_preparation(1_800_000_200_000)
    676             .unwrap()
    677             .delivery_plans()[0]
    678             .plan_id();
    679         let engine = Engine::builder(
    680             storage.clone(),
    681             clock.clone(),
    682             Arc::new(TestIds(AtomicU64::new(10))),
    683             DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
    684         )
    685         .sink(Arc::new(BoundarySink {
    686             storage: storage.clone(),
    687             clock: clock.clone(),
    688             id,
    689             stop,
    690         }))
    691         .signer(Arc::new(MockSigner::new(SignBehavior::Success {
    692             completed_at_unix_ms: 1_800_000_200_500,
    693         })))
    694         .build()
    695         .unwrap();
    696         engine.sign_prepared(push.clone()).await.unwrap();
    697         engine.admit_signed(push.operation_id()).await.unwrap();
    698         let late = engine.deliver_push(push.operation_id()).await.unwrap();
    699         assert_eq!(late.plan().attempt_count(), 0);
    700         assert_eq!(
    701             late.plan().delivery_satisfaction().unwrap(),
    702             SatisfactionState::Satisfied
    703         );
    704         drop(engine);
    705         storage.close().await.unwrap();
    706         let storage = Arc::new(
    707             SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting))
    708                 .await
    709                 .unwrap(),
    710         );
    711         let recovery = source_only(storage.clone(), clock.clone());
    712         let recovered = recovery.deliver_push(push.operation_id()).await.unwrap();
    713         assert_eq!(recovered.plan().attempt_count(), u32::from(!stop));
    714         assert_eq!(
    715             recovered.plan().delivery_facts(),
    716             late.plan().delivery_facts()
    717         );
    718         let status = recovery
    719             .push_status(push.operation_id())
    720             .await
    721             .unwrap()
    722             .unwrap();
    723         assert!(!status.delivery_history().has_unresolved_claims());
    724         drop(recovery);
    725         storage.close().await.unwrap();
    726         let storage = Arc::new(
    727             SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly))
    728                 .await
    729                 .unwrap(),
    730         );
    731         let readonly = source_only(storage.clone(), clock);
    732         let replay = readonly.deliver_push(push.operation_id()).await.unwrap();
    733         assert_eq!(replay.plan(), recovered.plan());
    734         drop(readonly);
    735         storage.close().await.unwrap();
    736     }
    737 }
    738 pub(super) enum Injection {
    739     FactBeforeClaim(Box<AuthoredAtomicCommand>),
    740     StopAfterClaim,
    741     ReplaceAfterClaim,
    742     StopAfterFact,
    743     StopAfterReconcile,
    744 }
    745 
    746 pub(super) async fn inject(storage: &FaultStorage, command: &AuthoredAtomicCommand, before: bool) {
    747     use radroots_storage::{
    748         authored::WorkClaim,
    749         authored_atomic::{
    750             CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget, ClaimAuthoredWork,
    751         },
    752     };
    753     let claim = matches!(command, AuthoredAtomicCommand::Claim(value)
    754         if matches!(value.target(), ClaimAuthoredTarget::DeliveryPlan(_)));
    755     let injection = {
    756         let mut armed = storage.delivery_injection.lock().unwrap();
    757         let ready = match armed.as_ref() {
    758             Some(Injection::FactBeforeClaim(_)) => before && claim,
    759             Some(Injection::StopAfterClaim | Injection::ReplaceAfterClaim) => !before && claim,
    760             Some(Injection::StopAfterFact) => {
    761                 !before && matches!(command, AuthoredAtomicCommand::RecordDelivery(_))
    762             }
    763             Some(Injection::StopAfterReconcile) => {
    764                 !before && matches!(command, AuthoredAtomicCommand::ReconcileDelivery(_))
    765             }
    766             None => false,
    767         };
    768         if ready { armed.take() } else { None }
    769     };
    770     let Some(injection) = injection else {
    771         return;
    772     };
    773     if let Injection::FactBeforeClaim(fact) = injection {
    774         storage.inner.execute_authored(*fact).await.unwrap();
    775         return;
    776     }
    777     let id = match command {
    778         AuthoredAtomicCommand::Claim(value) => match value.target() {
    779             ClaimAuthoredTarget::DeliveryPlan(id) => *id,
    780             _ => panic!("delivery-only injection"),
    781         },
    782         AuthoredAtomicCommand::RecordDelivery(value) => value.plan_id(),
    783         AuthoredAtomicCommand::ReconcileDelivery(value) => value.plan_id(),
    784         _ => panic!("delivery-only injection"),
    785     };
    786     let history = storage
    787         .inner
    788         .authored_delivery_history(id)
    789         .await
    790         .unwrap()
    791         .unwrap();
    792     let plan = history.plan();
    793     let mutation = if matches!(injection, Injection::ReplaceAfterClaim) {
    794         let original = plan.claim_evidence().unwrap();
    795         let at = original.expires_at_unix_ms();
    796         AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    797             ClaimAuthoredTarget::DeliveryPlan(id),
    798             WorkClaim::new(
    799                 [222; 16],
    800                 "competing-worker",
    801                 core::num::NonZeroU64::new(2).unwrap(),
    802                 at,
    803                 at + 10_000,
    804                 plan.revision(),
    805             )
    806             .unwrap(),
    807         ))
    808     } else {
    809         AuthoredAtomicCommand::Cancel(
    810             CancelAuthoredWork::new(
    811                 CancelAuthoredTarget::DeliveryPlan(id),
    812                 plan.revision(),
    813                 plan.updated_at_unix_ms(),
    814             )
    815             .unwrap(),
    816         )
    817     };
    818     storage.inner.execute_authored(mutation).await.unwrap();
    819 }
    820 
    821 #[test]
    822 fn observed_stop_or_replacement_after_claim_prevents_sink_admission() {
    823     for stop in [true, false] {
    824         let storage = Arc::new(FaultStorage::new(160));
    825         let sink = Arc::new(HeldSink::default());
    826         let engine = fault_engine(
    827             storage.clone(),
    828             Arc::new(MockSigner::new(SignBehavior::Success {
    829                 completed_at_unix_ms: 1_800_000_200_500,
    830             })),
    831             sink.clone(),
    832         );
    833         let push = request(160, "wss://relay.example");
    834         execute_to_admitted(&engine, &push);
    835         *storage.delivery_injection.lock().unwrap() = Some(if stop {
    836             Injection::StopAfterClaim
    837         } else {
    838             Injection::ReplaceAfterClaim
    839         });
    840         let result = block_on(engine.deliver_push(push.operation_id()));
    841         if stop {
    842             assert_eq!(
    843                 result.unwrap().plan().state(),
    844                 AuthoredDeliveryState::Cancelled
    845             );
    846         } else {
    847             assert_eq!(result, Err(Error::WorkClaimConflict));
    848         }
    849         assert_eq!(sink.calls.load(Ordering::Relaxed), 0);
    850     }
    851 }
    852 
    853 #[test]
    854 fn fresh_claim_reconciles_fact_arriving_after_initial_status_without_another_effect() {
    855     let storage = Arc::new(FaultStorage::new(161));
    856     let sink = Arc::new(HeldSink::default());
    857     let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
    858     let engine = Engine::builder(
    859         storage.clone(),
    860         clock.clone(),
    861         Arc::new(TestIds(AtomicU64::new(10))),
    862         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
    863     )
    864     .sink(sink.clone())
    865     .signer(Arc::new(MockSigner::new(SignBehavior::Success {
    866         completed_at_unix_ms: 1_800_000_200_500,
    867     })))
    868     .build()
    869     .unwrap();
    870     let push = request(161, "wss://relay.example");
    871     execute_to_admitted(&engine, &push);
    872     let mut first = Box::pin(engine.deliver_push(push.operation_id()));
    873     assert!(
    874         first
    875             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
    876             .is_pending()
    877     );
    878     let (request, _sender) = sink.take();
    879     drop(first);
    880     let status = block_on(engine.push_status(push.operation_id()))
    881         .unwrap()
    882         .unwrap();
    883     let plan = status.delivery_plan();
    884     let claim = plan.claim_evidence().unwrap();
    885     let at = claim.expires_at_unix_ms();
    886     let fact = AuthoredAtomicCommand::RecordDelivery(
    887         RecordDeliveryFact::new(
    888             plan.plan_id(),
    889             plan.artifact_id(),
    890             claim.clone(),
    891             DeliveryAttemptOutcome::Receipt(
    892                 receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(),
    893             ),
    894             at,
    895         )
    896         .unwrap(),
    897     );
    898     *storage.delivery_injection.lock().unwrap() = Some(Injection::FactBeforeClaim(Box::new(fact)));
    899     clock.0.store(at, Ordering::Relaxed);
    900     let reconciled = block_on(engine.deliver_push(push.operation_id())).unwrap();
    901     assert!(reconciled.is_replay());
    902     assert_eq!(reconciled.plan().state(), AuthoredDeliveryState::Satisfied);
    903     assert_eq!(reconciled.plan().attempt_count(), 1);
    904     assert_eq!(sink.calls.load(Ordering::Relaxed), 1);
    905 }
    906 
    907 #[test]
    908 fn receipt_return_reloads_stop_after_fact_or_reconciliation_commit() {
    909     for before_reconcile in [true, false] {
    910         let storage = Arc::new(FaultStorage::new(162));
    911         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
    912             DeliveryOutcome::accepted(),
    913         ])]));
    914         let engine = fault_engine(
    915             storage.clone(),
    916             Arc::new(MockSigner::new(SignBehavior::Success {
    917                 completed_at_unix_ms: 1_800_000_200_500,
    918             })),
    919             sink.clone(),
    920         );
    921         let push = request(162, "wss://relay.example");
    922         execute_to_admitted(&engine, &push);
    923         *storage.delivery_injection.lock().unwrap() = Some(if before_reconcile {
    924             Injection::StopAfterFact
    925         } else {
    926             Injection::StopAfterReconcile
    927         });
    928         let delivered = block_on(engine.deliver_push(push.operation_id())).unwrap();
    929         assert!(delivered.plan().stop_requested_at_unix_ms().is_some());
    930         assert_eq!(
    931             delivered.plan().attempt_count(),
    932             u32::from(!before_reconcile)
    933         );
    934         assert_eq!(
    935             delivered.plan().delivery_satisfaction().unwrap(),
    936             SatisfactionState::Satisfied
    937         );
    938         assert_eq!(
    939             engine
    940                 .retry_decision(delivered.plan(), 1_800_000_250_000)
    941                 .unwrap(),
    942             SyncRetryDecision::Satisfied
    943         );
    944         assert_eq!(sink.requests.lock().unwrap().len(), 1);
    945     }
    946 }
    947 
    948 #[test]
    949 fn forged_stop_receipt_is_not_reported_as_success() {
    950     let storage = Arc::new(FaultStorage::new(163));
    951     let engine = fault_engine(
    952         storage.clone(),
    953         Arc::new(MockSigner::new(SignBehavior::Success {
    954             completed_at_unix_ms: 1_800_000_200_500,
    955         })),
    956         Arc::new(MockSink),
    957     );
    958     let push = request(163, "wss://relay.example");
    959     execute_to_admitted(&engine, &push);
    960     storage.fault_nth_plan(1);
    961     assert_eq!(
    962         block_on(engine.cancel_push(push.operation_id())),
    963         Err(Error::StorageFailed)
    964     );
    965     let status = block_on(engine.push_status(push.operation_id()))
    966         .unwrap()
    967         .unwrap();
    968     assert!(status.delivery_plan().stop_requested_at_unix_ms().is_some());
    969     assert!(status.delivery_history().proves_no_issued_attempt());
    970 }
    971 use radroots_storage::authored_atomic::{AuthoredAtomicOutcome, AuthoredAtomicReceipt};
    972 use radroots_storage::authored_delivery::{AuthoredDeliveryHistory, AuthoredDeliveryPlan};
    973 
    974 #[derive(Clone, Copy, Debug)]
    975 pub(super) enum ReceiptMutation {
    976     Identity,
    977     Plan,
    978     Artifact,
    979     Request,
    980     Created,
    981     Claim,
    982     CommitTime,
    983     Attempts,
    984     State,
    985     Retry,
    986     StopTime,
    987     Revision,
    988 }
    989 pub(super) fn mutate_receipt(
    990     receipt: AuthoredAtomicReceipt,
    991     mutation: ReceiptMutation,
    992 ) -> AuthoredAtomicReceipt {
    993     let AuthoredAtomicOutcome::DeliveryPlan(plan) = receipt.outcome() else {
    994         panic!("delivery receipt");
    995     };
    996     let mut wire = serde_json::to_value(plan).unwrap();
    997     match mutation {
    998         ReceiptMutation::Identity | ReceiptMutation::CommitTime => {}
    999         ReceiptMutation::Plan => wire["plan_id"] = serde_json::json!(vec![241; 16]),
   1000         ReceiptMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![242; 16]),
   1001         ReceiptMutation::Request => {
   1002             let original = plan.request().unwrap();
   1003             let changed = DeliveryRequest::new(
   1004                 "different-request",
   1005                 original.payload().clone(),
   1006                 original.target_set().clone(),
   1007                 original.satisfaction().clone(),
   1008                 original.deadline_unix_ms(),
   1009             )
   1010             .unwrap();
   1011             let other = AuthoredDeliveryPlan::new_bound(
   1012                 plan.plan_id(),
   1013                 plan.artifact_id(),
   1014                 changed,
   1015                 plan.created_at_unix_ms(),
   1016             )
   1017             .unwrap();
   1018             let other = serde_json::to_value(other).unwrap();
   1019             for key in ["request", "intent", "request_digest"] {
   1020                 wire[key] = other[key].clone();
   1021             }
   1022         }
   1023         ReceiptMutation::Created => {
   1024             wire["created_at_unix_ms"] = serde_json::json!(plan.created_at_unix_ms() - 1)
   1025         }
   1026         ReceiptMutation::Claim => wire["claim"]["owner"] = serde_json::json!("another-worker"),
   1027         ReceiptMutation::Attempts => {
   1028             let outcome = DeliveryAttemptOutcome::Receipt(receipt_for_pending(plan));
   1029             wire["attempts"] = serde_json::json!([{"attempt":1,"recorded_at_unix_ms":plan.created_at_unix_ms(),"outcome":outcome,"satisfaction":"pending"}]);
   1030             wire["attempt_count"] = serde_json::json!(1);
   1031         }
   1032         ReceiptMutation::State => {
   1033             wire["state"] = serde_json::json!("pending");
   1034             wire["retry"] = serde_json::Value::Null;
   1035             wire["last_failure"] = serde_json::Value::Null;
   1036         }
   1037         ReceiptMutation::Retry => {
   1038             let at = plan.retry().unwrap().not_before_unix_ms() + 1000;
   1039             wire["retry"]["not_before_unix_ms"] = serde_json::json!(at);
   1040             wire["retry"]["failure"]["retry_after_unix_ms"] = serde_json::json!(at);
   1041             wire["last_failure"] = wire["retry"]["failure"].clone();
   1042         }
   1043         ReceiptMutation::StopTime => {
   1044             wire["stop_requested_at_unix_ms"] =
   1045                 serde_json::json!(plan.stop_requested_at_unix_ms().unwrap() - 1)
   1046         }
   1047         ReceiptMutation::Revision => {
   1048             wire["revision"] = serde_json::json!(plan.revision().get() + 1)
   1049         }
   1050     }
   1051     let changed: AuthoredDeliveryPlan =
   1052         serde_json::from_value(wire).expect("valid but incorrectly bound projection");
   1053     AuthoredAtomicReceipt::from_durable_parts(
   1054         if matches!(mutation, ReceiptMutation::Identity) {
   1055             radroots_storage::atomic::AtomicCommitId::new([243; 16]).unwrap()
   1056         } else {
   1057             receipt.commit_id()
   1058         },
   1059         receipt.digest(),
   1060         receipt.disposition(),
   1061         receipt.committed_at_unix_ms() + u64::from(matches!(mutation, ReceiptMutation::CommitTime)),
   1062         AuthoredAtomicOutcome::DeliveryPlan(changed),
   1063     )
   1064     .unwrap()
   1065 }
   1066 fn receipt_for_pending(plan: &AuthoredDeliveryPlan) -> DeliveryReceipt {
   1067     receipt(
   1068         plan.request().unwrap(),
   1069         vec![DeliveryOutcome::unavailable()],
   1070     )
   1071     .unwrap()
   1072 }
   1073 
   1074 #[test]
   1075 fn claim_projection_must_match_every_original_binding_before_transport() {
   1076     for mutation in [
   1077         ReceiptMutation::Identity,
   1078         ReceiptMutation::Plan,
   1079         ReceiptMutation::Artifact,
   1080         ReceiptMutation::Request,
   1081         ReceiptMutation::Created,
   1082         ReceiptMutation::Claim,
   1083         ReceiptMutation::CommitTime,
   1084         ReceiptMutation::Attempts,
   1085     ] {
   1086         let storage = Arc::new(FaultStorage::new(171));
   1087         let sink = Arc::new(HeldSink::default());
   1088         let engine = fault_engine(
   1089             storage.clone(),
   1090             Arc::new(MockSigner::new(SignBehavior::Success {
   1091                 completed_at_unix_ms: 1_800_000_200_500,
   1092             })),
   1093             sink.clone(),
   1094         );
   1095         let push = request(171, "wss://relay.example");
   1096         execute_to_admitted(&engine, &push);
   1097         *storage.receipt_mutation.lock().unwrap() = Some(mutation);
   1098         assert_eq!(
   1099             block_on(engine.deliver_push(push.operation_id())),
   1100             Err(Error::StorageFailed),
   1101             "{mutation:?}"
   1102         );
   1103         assert_eq!(sink.calls.load(Ordering::Relaxed), 0);
   1104     }
   1105 }
   1106 
   1107 #[test]
   1108 fn stop_projection_must_match_identity_request_time_and_revision() {
   1109     for mutation in [
   1110         ReceiptMutation::Identity,
   1111         ReceiptMutation::Plan,
   1112         ReceiptMutation::Artifact,
   1113         ReceiptMutation::Request,
   1114         ReceiptMutation::StopTime,
   1115         ReceiptMutation::Revision,
   1116     ] {
   1117         let storage = Arc::new(FaultStorage::new(172));
   1118         let engine = fault_engine(
   1119             storage.clone(),
   1120             Arc::new(MockSigner::new(SignBehavior::Success {
   1121                 completed_at_unix_ms: 1_800_000_200_500,
   1122             })),
   1123             Arc::new(MockSink),
   1124         );
   1125         let push = request(172, "wss://relay.example");
   1126         execute_to_admitted(&engine, &push);
   1127         *storage.receipt_mutation.lock().unwrap() = Some(mutation);
   1128         assert_eq!(
   1129             block_on(engine.cancel_push(push.operation_id())),
   1130             Err(Error::StorageFailed),
   1131             "{mutation:?}"
   1132         );
   1133         let current = block_on(engine.push_status(push.operation_id()))
   1134             .unwrap()
   1135             .unwrap();
   1136         assert!(current.delivery_history().proves_no_issued_attempt());
   1137         assert!(
   1138             current
   1139                 .delivery_plan()
   1140                 .stop_requested_at_unix_ms()
   1141                 .is_some()
   1142         );
   1143     }
   1144 }
   1145 
   1146 #[test]
   1147 fn claim_projection_cannot_replace_retained_retry_state_or_backoff() {
   1148     for mutation in [ReceiptMutation::State, ReceiptMutation::Retry] {
   1149         let storage = Arc::new(FaultStorage::new(173));
   1150         let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000)));
   1151         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Failure(
   1152             Retryability::Retryable,
   1153         )]));
   1154         let engine = Engine::builder(
   1155             storage.clone(),
   1156             clock.clone(),
   1157             Arc::new(TestIds(AtomicU64::new(10))),
   1158             DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
   1159         )
   1160         .sink(sink.clone())
   1161         .signer(Arc::new(MockSigner::new(SignBehavior::Success {
   1162             completed_at_unix_ms: 1_800_000_200_500,
   1163         })))
   1164         .build()
   1165         .unwrap();
   1166         let push = request(173, "wss://relay.example");
   1167         execute_to_admitted(&engine, &push);
   1168         let first = block_on(engine.deliver_push(push.operation_id())).unwrap();
   1169         clock.0.store(
   1170             first.plan().retry().unwrap().not_before_unix_ms(),
   1171             Ordering::Relaxed,
   1172         );
   1173         *storage.receipt_mutation.lock().unwrap() = Some(mutation);
   1174         assert_eq!(
   1175             block_on(engine.deliver_push(push.operation_id())),
   1176             Err(Error::StorageFailed),
   1177             "{mutation:?}"
   1178         );
   1179         assert_eq!(sink.requests.lock().unwrap().len(), 1);
   1180     }
   1181 }
   1182 
   1183 #[derive(Clone, Copy)]
   1184 pub(super) enum HistoryMutation {
   1185     Plan,
   1186     Artifact,
   1187     MissingFact,
   1188     MissingStop,
   1189 }
   1190 pub(super) fn mutate_history(
   1191     history: AuthoredDeliveryHistory,
   1192     mutation: HistoryMutation,
   1193 ) -> AuthoredDeliveryHistory {
   1194     let mut wire = serde_json::to_value(history.plan()).unwrap();
   1195     match mutation {
   1196         HistoryMutation::Plan => wire["plan_id"] = serde_json::json!(vec![244; 16]),
   1197         HistoryMutation::Artifact => wire["artifact_id"] = serde_json::json!(vec![245; 16]),
   1198         HistoryMutation::MissingFact => wire["delivery_facts"] = serde_json::json!([]),
   1199         HistoryMutation::MissingStop => {
   1200             wire["stop_requested_at_unix_ms"] = serde_json::Value::Null;
   1201             wire["state"] = serde_json::json!("pending");
   1202         }
   1203     }
   1204     AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), None).unwrap()
   1205 }
   1206 
   1207 #[test]
   1208 fn wrong_history_and_lost_committed_fact_or_stop_never_become_success() {
   1209     for (nth, mutation, stop) in [
   1210         (1, HistoryMutation::Plan, false),
   1211         (1, HistoryMutation::Artifact, false),
   1212         (2, HistoryMutation::Plan, false),
   1213         (3, HistoryMutation::MissingFact, false),
   1214         (2, HistoryMutation::MissingStop, true),
   1215     ] {
   1216         let storage = Arc::new(FaultStorage::new(174));
   1217         let sink = Arc::new(ScriptedSink::new([DeliveryBehavior::Outcomes(vec![
   1218             DeliveryOutcome::accepted(),
   1219         ])]));
   1220         let engine = fault_engine(
   1221             storage.clone(),
   1222             Arc::new(MockSigner::new(SignBehavior::Success {
   1223                 completed_at_unix_ms: 1_800_000_200_500,
   1224             })),
   1225             sink.clone(),
   1226         );
   1227         let push = request(174, "wss://relay.example");
   1228         execute_to_admitted(&engine, &push);
   1229         *storage.history_mutation.lock().unwrap() = Some((nth, mutation));
   1230         if stop {
   1231             assert_eq!(
   1232                 block_on(engine.cancel_push(push.operation_id())),
   1233                 Err(Error::StorageFailed)
   1234             );
   1235         } else {
   1236             assert_eq!(
   1237                 block_on(engine.deliver_push(push.operation_id())),
   1238                 Err(Error::StorageFailed)
   1239             );
   1240         }
   1241         let actual = block_on(engine.push_status(push.operation_id()))
   1242             .unwrap()
   1243             .unwrap();
   1244         assert_eq!(
   1245             actual.delivery_plan().delivery_facts().len(),
   1246             usize::from(nth == 3)
   1247         );
   1248     }
   1249 }
   1250 
   1251 #[test]
   1252 fn clock_expiry_before_sink_and_future_fact_before_reconciliation_fail_closed() {
   1253     struct ExpiringClock(AtomicUsize);
   1254     impl Clock for ExpiringClock {
   1255         fn now_unix_ms(&self) -> Result<u64, Error> {
   1256             Ok(1_800_000_220_000 + 10_000 * self.0.fetch_add(1, Ordering::Relaxed) as u64)
   1257         }
   1258     }
   1259     let (engine, storage, _, sink, push) = setup(175);
   1260     let boundary = Engine::builder(
   1261         storage,
   1262         Arc::new(ExpiringClock(AtomicUsize::new(0))),
   1263         Arc::new(TestIds(AtomicU64::new(230))),
   1264         DeadlinePolicy::new(10_000, 10_000, 10_000).unwrap(),
   1265     )
   1266     .sink(sink.clone())
   1267     .build()
   1268     .unwrap();
   1269     assert_eq!(
   1270         block_on(boundary.deliver_push(push.operation_id())),
   1271         Err(Error::WorkClaimConflict)
   1272     );
   1273     assert_eq!(sink.calls.load(Ordering::Relaxed), 0);
   1274     assert_eq!(
   1275         block_on(engine.deliver_push(SyncId::new([231; 16]).unwrap())),
   1276         Err(Error::StorageFailed)
   1277     );
   1278     let (engine, storage, clock, sink, push) = setup(176);
   1279     let mut future = Box::pin(engine.deliver_push(push.operation_id()));
   1280     assert!(
   1281         future
   1282             .poll_unpin(&mut std::task::Context::from_waker(noop_waker_ref()))
   1283             .is_pending()
   1284     );
   1285     let (request, _sender) = sink.take();
   1286     drop(future);
   1287     let status = block_on(engine.push_status(push.operation_id()))
   1288         .unwrap()
   1289         .unwrap();
   1290     let plan = status.delivery_plan();
   1291     let claim = plan.claim_evidence().unwrap();
   1292     let at = claim.expires_at_unix_ms();
   1293     block_on(
   1294         storage.execute_authored(AuthoredAtomicCommand::RecordDelivery(
   1295             RecordDeliveryFact::new(
   1296                 plan.plan_id(),
   1297                 plan.artifact_id(),
   1298                 claim.clone(),
   1299                 DeliveryAttemptOutcome::Receipt(
   1300                     receipt(&request, vec![DeliveryOutcome::accepted()]).unwrap(),
   1301                 ),
   1302                 at + 1000,
   1303             )
   1304             .unwrap(),
   1305         )),
   1306     )
   1307     .unwrap();
   1308     clock.0.store(at, Ordering::Relaxed);
   1309     assert_eq!(
   1310         block_on(engine.deliver_push(push.operation_id())),
   1311         Err(Error::ClockUnavailable)
   1312     );
   1313     let current = block_on(engine.push_status(push.operation_id()))
   1314         .unwrap()
   1315         .unwrap();
   1316     assert_eq!(current.delivery_plan().attempt_count(), 0);
   1317     assert_eq!(current.delivery_plan().delivery_facts().len(), 1);
   1318 }
   1319 
   1320 #[test]
   1321 fn unbound_preparation_and_stop_keep_honest_retry_and_no_issued_proof() {
   1322     let signer = Arc::new(MockSigner::new(SignBehavior::Success {
   1323         completed_at_unix_ms: 1_800_000_200_500,
   1324     }));
   1325     let (engine, _) = setup_engine(signer);
   1326     let push = request(177, "wss://relay.example");
   1327     block_on(engine.prepare_push(push.clone())).unwrap();
   1328     let status = block_on(engine.push_status(push.operation_id()))
   1329         .unwrap()
   1330         .unwrap();
   1331     assert!(status.delivery_plan().request().is_none());
   1332     assert!(status.delivery_history().proves_no_issued_attempt());
   1333     assert_eq!(
   1334         engine
   1335             .retry_decision(status.delivery_plan(), 1_800_000_200_001)
   1336             .unwrap(),
   1337         SyncRetryDecision::Ready
   1338     );
   1339     let stopped = block_on(engine.cancel_push(push.operation_id())).unwrap();
   1340     assert!(
   1341         stopped
   1342             .status()
   1343             .delivery_history()
   1344             .proves_no_issued_attempt()
   1345     );
   1346     assert_eq!(
   1347         engine
   1348             .retry_decision(stopped.status().delivery_plan(), 1_800_000_200_010)
   1349             .unwrap(),
   1350         SyncRetryDecision::Exhausted
   1351     );
   1352 }