lib

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

authored_delivery_facts.rs (27097B)


      1 use futures_executor::block_on;
      2 use radroots_storage::{
      3     Error,
      4     atomic::{AtomicCommitDigest, AtomicCommitDisposition},
      5     authored::WorkClaim,
      6     authored_atomic::{
      7         ApplyDeliveryAttempt, AuthoredAtomicCommand, AuthoredAtomicOutcome, AuthoredAtomicReceipt,
      8         AuthoredAtomicStorage, CancelAuthoredTarget, CancelAuthoredWork, ClaimAuthoredTarget,
      9         ClaimAuthoredWork, RecordDeliveryFact, WorkFence,
     10     },
     11     authored_delivery::{
     12         AuthoredDeliveryPlan, AuthoredDeliveryState, DELIVERY_PLAN_ATTEMPTS_MAX,
     13         DeliveryAttemptOutcome,
     14     },
     15     event::SourceGeneration,
     16     memory::MemoryStorage,
     17 };
     18 use radroots_transport::{
     19     DeliveryReceipt, SinkFailure,
     20     outcome::{DeliveryOutcome, Retryability},
     21     policy::SatisfactionState,
     22     sink::DeliveryTargetReceipt,
     23 };
     24 use std::num::NonZeroU64;
     25 
     26 #[path = "authored_signing/fixture.rs"]
     27 mod fixture;
     28 use fixture::*;
     29 
     30 #[path = "authored_delivery/reconciliation_tests.rs"]
     31 mod reconciliation;
     32 
     33 fn plan(storage: &MemoryStorage) -> AuthoredDeliveryPlan {
     34     block_on(storage.authored_delivery_plan(ids().2))
     35         .unwrap()
     36         .unwrap()
     37 }
     38 
     39 fn prepared() -> (MemoryStorage, WorkClaim, AuthoredAtomicReceipt) {
     40     let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap());
     41     let (command, event) = prepare();
     42     block_on(storage.execute_authored(command)).unwrap();
     43     let signing = claim(NonZeroU64::MIN, 1, 11);
     44     block_on(
     45         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
     46             ClaimAuthoredTarget::ArtifactSigning(ids().1),
     47             signing.clone(),
     48         ))),
     49     )
     50     .unwrap();
     51     block_on(storage.execute_authored(record(event, signing, 12))).unwrap();
     52     let delivery = claim(plan(&storage).revision(), 2, 13);
     53     let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::Claim(
     54         ClaimAuthoredWork::new(ClaimAuthoredTarget::DeliveryPlan(ids().2), delivery.clone()),
     55     )))
     56     .unwrap();
     57     (storage, delivery, receipt)
     58 }
     59 
     60 fn outcome(plan: &AuthoredDeliveryPlan, accepted: bool) -> DeliveryAttemptOutcome {
     61     let request = plan.request().unwrap();
     62     DeliveryAttemptOutcome::Receipt(
     63         DeliveryReceipt::for_request(
     64             request,
     65             request
     66                 .target_set()
     67                 .targets()
     68                 .iter()
     69                 .cloned()
     70                 .map(|target| {
     71                     DeliveryTargetReceipt::attempted(
     72                         target,
     73                         if accepted {
     74                             DeliveryOutcome::accepted()
     75                         } else {
     76                             DeliveryOutcome::unavailable()
     77                         },
     78                     )
     79                 })
     80                 .collect(),
     81         )
     82         .unwrap(),
     83     )
     84 }
     85 
     86 fn fact(
     87     plan: &AuthoredDeliveryPlan,
     88     claim: WorkClaim,
     89     accepted: bool,
     90     at: u64,
     91 ) -> AuthoredAtomicCommand {
     92     AuthoredAtomicCommand::RecordDelivery(
     93         RecordDeliveryFact::new(
     94             plan.plan_id(),
     95             plan.artifact_id(),
     96             claim,
     97             outcome(plan, accepted),
     98             at,
     99         )
    100         .unwrap(),
    101     )
    102 }
    103 
    104 fn stop(storage: &MemoryStorage, at: u64) {
    105     block_on(
    106         storage.execute_authored(AuthoredAtomicCommand::Cancel(
    107             CancelAuthoredWork::new(
    108                 CancelAuthoredTarget::DeliveryPlan(ids().2),
    109                 plan(storage).revision(),
    110                 at,
    111             )
    112             .unwrap(),
    113         )),
    114     )
    115     .unwrap();
    116 }
    117 
    118 #[test]
    119 fn expired_cancelled_delivery_retains_exact_evidence_and_idempotent_first_time() {
    120     let (storage, active, original) = prepared();
    121     stop(&storage, 20);
    122     let before = plan(&storage);
    123     let command = fact(&before, active.clone(), true, 50);
    124     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
    125     assert!(receipt.matches_command(&command));
    126     let after = plan(&storage);
    127     assert_eq!(after.state(), AuthoredDeliveryState::Cancelled);
    128     assert_eq!(after.stop_requested_at_unix_ms(), Some(20));
    129     assert_eq!(after.request().unwrap().payload().event().raw_json(), RAW);
    130     assert_eq!(
    131         after.delivery_satisfaction().unwrap(),
    132         SatisfactionState::Satisfied
    133     );
    134     assert_eq!(after.revision(), before.revision());
    135     assert_eq!(after.updated_at_unix_ms(), before.updated_at_unix_ms());
    136     assert_eq!(after.attempt_count(), 0);
    137     assert_eq!(after.delivery_facts()[0].claim(), &active);
    138     assert_eq!(after.delivery_facts()[0].observed_at_unix_ms(), 50);
    139     let replay_command = fact(&after, active.clone(), true, 90);
    140     assert_eq!(replay_command.commit_id(), command.commit_id());
    141     let replay = block_on(storage.execute_authored(replay_command)).unwrap();
    142     assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
    143     assert_eq!(replay.committed_at_unix_ms(), 50);
    144     assert_eq!(plan(&storage), after);
    145     assert_eq!(
    146         block_on(storage.authored_receipt(original.commit_id()))
    147             .unwrap()
    148             .unwrap(),
    149         original
    150     );
    151     stop(&storage, 99);
    152     assert_eq!(plan(&storage), after);
    153     assert!(block_on(storage.execute_authored(fact(&after, active.clone(), false, 95))).is_err());
    154     let fenced = AuthoredAtomicCommand::ApplyDelivery(
    155         ApplyDeliveryAttempt::new(
    156             ids().2,
    157             WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(),
    158             outcome(&after, true),
    159             None,
    160             50,
    161         )
    162         .unwrap(),
    163     );
    164     assert_eq!(
    165         block_on(storage.execute_authored(fenced)),
    166         Err(Error::DeliveryPlanClaimConflict)
    167     );
    168     assert!(
    169         block_on(
    170             storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    171                 ClaimAuthoredTarget::DeliveryPlan(ids().2),
    172                 claim(after.revision(), 9, 100)
    173             )))
    174         )
    175         .is_err()
    176     );
    177 }
    178 
    179 #[test]
    180 fn superseded_result_does_not_steal_new_claim_or_regress_acceptance() {
    181     let (storage, old, _) = prepared();
    182     let newer = claim(plan(&storage).revision(), 3, 40);
    183     block_on(
    184         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    185             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    186             newer.clone(),
    187         ))),
    188     )
    189     .unwrap();
    190     let before = plan(&storage);
    191     block_on(storage.execute_authored(fact(&before, old, true, 30))).unwrap();
    192     let after = plan(&storage);
    193     assert_eq!(after.claim_evidence(), Some(&newer));
    194     assert_eq!(after.revision(), before.revision());
    195     assert_eq!(after.updated_at_unix_ms(), 40);
    196     let partial = SinkFailure::for_request(
    197         after.request().unwrap(),
    198         "sink_lost",
    199         Retryability::Retryable,
    200         Some(65),
    201         Some("connection lost".into()),
    202         vec![],
    203     )
    204     .unwrap();
    205     let command = AuthoredAtomicCommand::RecordDelivery(
    206         RecordDeliveryFact::new(
    207             ids().2,
    208             ids().1,
    209             newer.clone(),
    210             DeliveryAttemptOutcome::SinkFailure(partial.clone()),
    211             45,
    212         )
    213         .unwrap(),
    214     );
    215     block_on(storage.execute_authored(command)).unwrap();
    216     assert_eq!(
    217         plan(&storage).delivery_satisfaction().unwrap(),
    218         SatisfactionState::Satisfied
    219     );
    220     assert_eq!(plan(&storage).claim_evidence(), Some(&newer));
    221     let retry = radroots_storage::authored::RetrySchedule::new(
    222         std::num::NonZeroU32::MIN,
    223         65,
    224         radroots_storage::authored::WorkFailure::new(
    225             "sink_lost",
    226             radroots_storage::authored::WorkPhase::Delivery,
    227             radroots_storage::authored::FailureClass::Retryable,
    228             Some(65),
    229             Some("connection lost".into()),
    230         )
    231         .unwrap(),
    232     )
    233     .unwrap();
    234     let fenced = AuthoredAtomicCommand::ApplyDelivery(
    235         ApplyDeliveryAttempt::new(
    236             ids().2,
    237             WorkFence::new(*newer.token(), newer.generation(), newer.row_revision()).unwrap(),
    238             DeliveryAttemptOutcome::SinkFailure(partial),
    239             Some(retry.clone()),
    240             46,
    241         )
    242         .unwrap(),
    243     );
    244     block_on(storage.execute_authored(fenced)).unwrap();
    245     let transitioned = plan(&storage);
    246     assert_eq!(transitioned.state(), AuthoredDeliveryState::Retryable);
    247     assert_eq!(transitioned.retry(), Some(&retry));
    248     assert_eq!(transitioned.attempt_count(), 1);
    249     assert_eq!(
    250         transitioned.delivery_satisfaction().unwrap(),
    251         SatisfactionState::Satisfied
    252     );
    253     stop(&storage, 47);
    254     assert_eq!(
    255         plan(&storage).delivery_satisfaction().unwrap(),
    256         SatisfactionState::Satisfied
    257     );
    258 }
    259 
    260 #[test]
    261 fn backend_provenance_rejects_every_forged_complete_claim_and_wrong_artifact() {
    262     let (storage, original, _) = prepared();
    263     let before = plan(&storage);
    264     for changed in 0..6 {
    265         let forged = WorkClaim::new(
    266             if changed == 0 {
    267                 [9; 16]
    268             } else {
    269                 *original.token()
    270             },
    271             if changed == 1 {
    272                 "forged"
    273             } else {
    274                 original.owner()
    275             },
    276             if changed == 2 {
    277                 NonZeroU64::new(99).unwrap()
    278             } else {
    279                 original.generation()
    280             },
    281             if changed == 3 {
    282                 14
    283             } else {
    284                 original.acquired_at_unix_ms()
    285             },
    286             if changed == 4 {
    287                 90
    288             } else {
    289                 original.expires_at_unix_ms()
    290             },
    291             if changed == 5 {
    292                 NonZeroU64::new(99).unwrap()
    293             } else {
    294                 original.row_revision()
    295             },
    296         )
    297         .unwrap();
    298         let command = fact(&before, forged, true, 100);
    299         assert_eq!(
    300             block_on(storage.execute_authored(command.clone())),
    301             Err(Error::AtomicWorkflowMismatch)
    302         );
    303         assert!(
    304             block_on(storage.authored_receipt(command.commit_id()))
    305                 .unwrap()
    306                 .is_none()
    307         );
    308         assert_eq!(plan(&storage), before);
    309     }
    310     let wrong = AuthoredAtomicCommand::RecordDelivery(
    311         RecordDeliveryFact::new(
    312             ids().2,
    313             radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(),
    314             original.clone(),
    315             outcome(&before, true),
    316             50,
    317         )
    318         .unwrap(),
    319     );
    320     assert_eq!(
    321         block_on(storage.execute_authored(wrong)),
    322         Err(Error::AtomicWorkflowMismatch)
    323     );
    324     assert!(
    325         RecordDeliveryFact::new(ids().2, ids().1, original, outcome(&before, true), 12).is_err()
    326     );
    327 }
    328 
    329 #[test]
    330 fn late_fact_rejects_rebound_raw_request_and_caller_receipt_with_wrong_time() {
    331     let (storage, active, original) = prepared();
    332     let before = plan(&storage);
    333     let value =
    334         RecordDeliveryFact::new(ids().2, ids().1, active.clone(), outcome(&before, true), 50)
    335             .unwrap();
    336     let request = before
    337         .intent()
    338         .materialize(radroots_transport::sink::DeliveryPayload::new(event(
    339             OTHER_RAW,
    340         )))
    341         .unwrap();
    342     let mut wire = serde_json::to_value(&before).unwrap();
    343     wire["request"] = serde_json::to_value(request).unwrap();
    344     let mut rebound: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    345     let saved = rebound.clone();
    346     assert_eq!(
    347         value.apply_to(&mut rebound, &original),
    348         Err(Error::AtomicWorkflowMismatch)
    349     );
    350     assert_eq!(rebound, saved);
    351     let wrong = AuthoredAtomicReceipt::from_durable_parts(
    352         original.commit_id(),
    353         original.digest(),
    354         AtomicCommitDisposition::Committed,
    355         14,
    356         original.outcome().clone(),
    357     )
    358     .unwrap();
    359     let mut unchanged = before.clone();
    360     assert_eq!(
    361         value.apply_to(&mut unchanged, &wrong),
    362         Err(Error::AtomicWorkflowMismatch)
    363     );
    364     assert_eq!(unchanged, before);
    365 }
    366 
    367 #[test]
    368 fn invalid_result_binding_rolls_back_and_changed_failure_details_have_distinct_ids() {
    369     let (storage, active, _) = prepared();
    370     let before = plan(&storage);
    371     let request = before.request().unwrap();
    372     let wrong = radroots_transport::DeliveryRequest::new(
    373         "wrong-request",
    374         request.payload().clone(),
    375         request.target_set().clone(),
    376         request.satisfaction().clone(),
    377         request.deadline_unix_ms(),
    378     )
    379     .unwrap();
    380     let wrong = DeliveryReceipt::for_request(
    381         &wrong,
    382         wrong
    383             .target_set()
    384             .targets()
    385             .iter()
    386             .cloned()
    387             .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()))
    388             .collect(),
    389     )
    390     .unwrap();
    391     let command = AuthoredAtomicCommand::RecordDelivery(
    392         RecordDeliveryFact::new(
    393             ids().2,
    394             ids().1,
    395             active.clone(),
    396             DeliveryAttemptOutcome::Receipt(wrong),
    397             50,
    398         )
    399         .unwrap(),
    400     );
    401     assert_eq!(
    402         block_on(storage.execute_authored(command.clone())),
    403         Err(Error::InvalidAuthoredDeliveryPlan)
    404     );
    405     assert_eq!(plan(&storage), before);
    406     assert!(
    407         block_on(storage.authored_receipt(command.commit_id()))
    408             .unwrap()
    409             .is_none()
    410     );
    411     let mut identities = std::collections::BTreeSet::new();
    412     for retry in [None, Some(65), Some(66)] {
    413         for message in [
    414             None,
    415             Some("connection lost".to_owned()),
    416             Some("connection closed".to_owned()),
    417         ] {
    418             let failure = SinkFailure::for_request(
    419                 request,
    420                 "sink_lost",
    421                 Retryability::Retryable,
    422                 retry,
    423                 message,
    424                 vec![],
    425             )
    426             .unwrap();
    427             let command = AuthoredAtomicCommand::RecordDelivery(
    428                 RecordDeliveryFact::new(
    429                     ids().2,
    430                     ids().1,
    431                     active.clone(),
    432                     DeliveryAttemptOutcome::SinkFailure(failure),
    433                     50,
    434                 )
    435                 .unwrap(),
    436             );
    437             assert!(identities.insert(*command.commit_id().as_bytes()));
    438         }
    439     }
    440     assert_eq!(identities.len(), 9);
    441 }
    442 
    443 #[test]
    444 fn stop_after_satisfaction_and_legacy_cancelled_snapshot_preserve_intent() {
    445     let (storage, active, _) = prepared();
    446     let before = plan(&storage);
    447     let command = AuthoredAtomicCommand::ApplyDelivery(
    448         ApplyDeliveryAttempt::new(
    449             ids().2,
    450             WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap(),
    451             outcome(&before, true),
    452             None,
    453             14,
    454         )
    455         .unwrap(),
    456     );
    457     block_on(storage.execute_authored(command)).unwrap();
    458     stop(&storage, 15);
    459     let after = plan(&storage);
    460     assert_eq!(after.state(), AuthoredDeliveryState::Satisfied);
    461     assert_eq!(after.stop_requested_at_unix_ms(), Some(15));
    462     assert_eq!(
    463         after.delivery_satisfaction().unwrap(),
    464         SatisfactionState::Satisfied
    465     );
    466     let (storage, _, _) = prepared();
    467     stop(&storage, 20);
    468     let expected = plan(&storage);
    469     let mut legacy = serde_json::to_value(&expected).unwrap();
    470     legacy.as_object_mut().unwrap().remove("delivery_facts");
    471     legacy
    472         .as_object_mut()
    473         .unwrap()
    474         .remove("stop_requested_at_unix_ms");
    475     assert_eq!(
    476         serde_json::from_value::<AuthoredDeliveryPlan>(legacy).unwrap(),
    477         expected
    478     );
    479     let mut pending = before.clone();
    480     assert!(pending.request_stop(9).is_err());
    481     assert_eq!(pending, before);
    482     assert!(pending.request_stop(12).is_err());
    483     assert_eq!(pending, before);
    484 }
    485 
    486 #[test]
    487 fn receipt_binding_rejects_wrong_outcome_digest_and_observation_time() {
    488     let (storage, active, _) = prepared();
    489     let before = plan(&storage);
    490     let command = fact(&before, active, true, 50);
    491     assert!(
    492         AuthoredAtomicReceipt::new(
    493             &command,
    494             AtomicCommitDisposition::Committed,
    495             50,
    496             AuthoredAtomicOutcome::DeliveryPlan(before.clone())
    497         )
    498         .is_err()
    499     );
    500     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
    501     assert!(
    502         AuthoredAtomicReceipt::new(
    503             &command,
    504             AtomicCommitDisposition::Committed,
    505             49,
    506             receipt.outcome().clone()
    507         )
    508         .is_err()
    509     );
    510     let wrong = AuthoredAtomicReceipt::from_durable_parts(
    511         receipt.commit_id(),
    512         AtomicCommitDigest::new([99; 32]),
    513         AtomicCommitDisposition::Committed,
    514         50,
    515         receipt.outcome().clone(),
    516     )
    517     .unwrap();
    518     assert!(!wrong.matches_command(&command));
    519     let wrong = AuthoredAtomicReceipt::from_durable_parts(
    520         receipt.commit_id(),
    521         receipt.digest(),
    522         AtomicCommitDisposition::Committed,
    523         50,
    524         AuthoredAtomicOutcome::Artifact(
    525             block_on(storage.authored_artifact(ids().1))
    526                 .unwrap()
    527                 .unwrap(),
    528         ),
    529     )
    530     .unwrap();
    531     assert!(!wrong.matches_command(&command));
    532 }
    533 
    534 #[test]
    535 fn fact_capacity_and_structural_corruption_fail_without_evicting_evidence() {
    536     let (storage, active, original) = prepared();
    537     let before = plan(&storage);
    538     let command = fact(&before, active.clone(), true, 50);
    539     block_on(storage.execute_authored(command)).unwrap();
    540     let after = plan(&storage);
    541     let base = serde_json::to_value(&after).unwrap();
    542     for (path, value) in [
    543         ("observed_at_unix_ms", serde_json::json!(1)),
    544         ("claim", serde_json::json!(null)),
    545         ("outcome", serde_json::json!(null)),
    546     ] {
    547         let mut corrupt = base.clone();
    548         corrupt["delivery_facts"][0][path] = value;
    549         assert!(serde_json::from_value::<AuthoredDeliveryPlan>(corrupt).is_err());
    550     }
    551     let mut duplicate = base.clone();
    552     duplicate["delivery_facts"]
    553         .as_array_mut()
    554         .unwrap()
    555         .push(base["delivery_facts"][0].clone());
    556     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err());
    557     let mut full = base.clone();
    558     full["revision"] = serde_json::json!(2048);
    559     full["claim"] = serde_json::Value::Null;
    560     let entries = (1..=DELIVERY_PLAN_ATTEMPTS_MAX)
    561         .map(|index| {
    562             let mut entry = base["delivery_facts"][0].clone();
    563             let claim = WorkClaim::new(
    564                 [7; 16],
    565                 "capacity",
    566                 NonZeroU64::new(u64::from(index)).unwrap(),
    567                 13,
    568                 33,
    569                 NonZeroU64::new(u64::from(index)).unwrap(),
    570             )
    571             .unwrap();
    572             entry["claim"] = serde_json::to_value(claim).unwrap();
    573             entry
    574         })
    575         .collect::<Vec<_>>();
    576     full["delivery_facts"] = serde_json::to_value(entries).unwrap();
    577     let mut full_plan: AuthoredDeliveryPlan = serde_json::from_value(full.clone()).unwrap();
    578     let checkpoint = full_plan.clone();
    579     let value =
    580         RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap();
    581     assert_eq!(
    582         value.apply_to(&mut full_plan, &original),
    583         Err(Error::DeliveryAttemptOverflow)
    584     );
    585     assert_eq!(full_plan, checkpoint);
    586     assert!(
    587         full_plan
    588             .claim(claim(full_plan.revision(), 9, 60), 60)
    589             .is_err()
    590     );
    591     let mut replay = full.clone();
    592     replay["delivery_facts"][0] = base["delivery_facts"][0].clone();
    593     let mut replay_plan: AuthoredDeliveryPlan = serde_json::from_value(replay).unwrap();
    594     let checkpoint = replay_plan.clone();
    595     value.apply_to(&mut replay_plan, &original).unwrap();
    596     assert_eq!(replay_plan, checkpoint);
    597     full["delivery_facts"]
    598         .as_array_mut()
    599         .unwrap()
    600         .push(base["delivery_facts"][0].clone());
    601     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(full).is_err());
    602 }
    603 
    604 #[test]
    605 fn corrupt_original_claim_receipts_and_rebound_current_rows_are_rejected() {
    606     let (storage, active, original) = prepared();
    607     let before = plan(&storage);
    608     let value =
    609         RecordDeliveryFact::new(ids().2, ids().1, active, outcome(&before, true), 50).unwrap();
    610     let AuthoredAtomicOutcome::DeliveryPlan(original_plan) = original.outcome() else {
    611         unreachable!()
    612     };
    613     for updates in [
    614         vec![("plan_id", serde_json::json!(vec![9_u8; 16]))],
    615         vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))],
    616         vec![("created_at_unix_ms", serde_json::json!(9))],
    617         vec![("request", serde_json::Value::Null)],
    618         vec![("claim", serde_json::Value::Null)],
    619     ] {
    620         let mut wire = serde_json::to_value(original_plan).unwrap();
    621         for (key, replacement) in updates {
    622             wire[key] = replacement;
    623         }
    624         let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    625         let forged = AuthoredAtomicReceipt::from_durable_parts(
    626             original.commit_id(),
    627             original.digest(),
    628             AtomicCommitDisposition::Committed,
    629             original.committed_at_unix_ms(),
    630             AuthoredAtomicOutcome::DeliveryPlan(altered),
    631         )
    632         .unwrap();
    633         let mut current = before.clone();
    634         assert_eq!(
    635             value.apply_to(&mut current, &forged),
    636             Err(Error::AtomicWorkflowMismatch)
    637         );
    638         assert_eq!(current, before);
    639     }
    640     for updates in [
    641         vec![("plan_id", serde_json::json!(vec![9_u8; 16]))],
    642         vec![("artifact_id", serde_json::json!(vec![9_u8; 16]))],
    643         vec![("created_at_unix_ms", serde_json::json!(9))],
    644         vec![
    645             ("claim", serde_json::Value::Null),
    646             ("revision", serde_json::json!(2)),
    647         ],
    648         vec![
    649             ("claim", serde_json::Value::Null),
    650             ("updated_at_unix_ms", serde_json::json!(12)),
    651         ],
    652     ] {
    653         let mut wire = serde_json::to_value(&before).unwrap();
    654         for (key, replacement) in updates {
    655             wire[key] = replacement;
    656         }
    657         let mut current: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    658         let unchanged = current.clone();
    659         assert_eq!(
    660             value.apply_to(&mut current, &original),
    661             Err(Error::AtomicWorkflowMismatch)
    662         );
    663         assert_eq!(current, unchanged);
    664     }
    665     let wrong_kind = AuthoredAtomicReceipt::from_durable_parts(
    666         original.commit_id(),
    667         original.digest(),
    668         AtomicCommitDisposition::Committed,
    669         original.committed_at_unix_ms(),
    670         AuthoredAtomicOutcome::Artifact(
    671             block_on(storage.authored_artifact(ids().1))
    672                 .unwrap()
    673                 .unwrap(),
    674         ),
    675     )
    676     .unwrap();
    677     let mut current = before.clone();
    678     assert_eq!(
    679         value.apply_to(&mut current, &wrong_kind),
    680         Err(Error::AtomicWorkflowMismatch)
    681     );
    682     let wrong_digest = AuthoredAtomicReceipt::from_durable_parts(
    683         original.commit_id(),
    684         AtomicCommitDigest::new([99; 32]),
    685         AtomicCommitDisposition::Committed,
    686         original.committed_at_unix_ms(),
    687         original.outcome().clone(),
    688     )
    689     .unwrap();
    690     assert_eq!(
    691         value.apply_to(&mut current, &wrong_digest),
    692         Err(Error::AtomicWorkflowMismatch)
    693     );
    694     assert_eq!(current, before);
    695 }
    696 
    697 #[test]
    698 fn forged_fact_snapshots_cannot_bypass_stop_time_claim_or_request_invariants() {
    699     let (storage, active, _) = prepared();
    700     let before = plan(&storage);
    701     block_on(storage.execute_authored(fact(&before, active, true, 50))).unwrap();
    702     let current = plan(&storage);
    703     let base = serde_json::to_value(&current).unwrap();
    704     for at in [9, 12, 99] {
    705         let mut wire = base.clone();
    706         wire["stop_requested_at_unix_ms"] = serde_json::json!(at);
    707         assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
    708     }
    709     for (at, revision) in [(9, 2), (15, 2), (13, 3)] {
    710         let forged = WorkClaim::new(
    711             [8; 16],
    712             "invalid-fact",
    713             NonZeroU64::MIN,
    714             at,
    715             at + 20,
    716             NonZeroU64::new(revision).unwrap(),
    717         )
    718         .unwrap();
    719         let mut wire = base.clone();
    720         wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap();
    721         assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
    722     }
    723     let mut wire = base;
    724     wire["request"] = serde_json::Value::Null;
    725     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(wire).is_err());
    726 }
    727 
    728 #[test]
    729 fn receipt_matching_checks_each_plan_fact_and_monotonic_row_time() {
    730     let (storage, active, _) = prepared();
    731     let before = plan(&storage);
    732     let command = fact(&before, active, true, 50);
    733     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
    734     let current = plan(&storage);
    735     let base = serde_json::to_value(&current).unwrap();
    736     for key in ["plan_id", "artifact_id"] {
    737         let mut wire = base.clone();
    738         wire[key] = serde_json::json!(vec![9_u8; 16]);
    739         let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    740         let forged = AuthoredAtomicReceipt::from_durable_parts(
    741             receipt.commit_id(),
    742             receipt.digest(),
    743             AtomicCommitDisposition::Committed,
    744             50,
    745             AuthoredAtomicOutcome::DeliveryPlan(altered.clone()),
    746         )
    747         .unwrap();
    748         assert!(!forged.matches_command(&command));
    749         assert!(
    750             AuthoredAtomicReceipt::new(
    751                 &command,
    752                 AtomicCommitDisposition::Committed,
    753                 50,
    754                 AuthoredAtomicOutcome::DeliveryPlan(altered)
    755             )
    756             .is_err()
    757         );
    758     }
    759     for change_claim in [false, true] {
    760         let mut wire = base.clone();
    761         if change_claim {
    762             let forged = WorkClaim::new(
    763                 [8; 16],
    764                 "different",
    765                 NonZeroU64::MIN,
    766                 13,
    767                 33,
    768                 NonZeroU64::new(2).unwrap(),
    769             )
    770             .unwrap();
    771             wire["delivery_facts"][0]["claim"] = serde_json::to_value(forged).unwrap();
    772         } else {
    773             wire["delivery_facts"][0]["outcome"] =
    774                 serde_json::to_value(outcome(&current, false)).unwrap();
    775         }
    776         let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    777         let forged = AuthoredAtomicReceipt::from_durable_parts(
    778             receipt.commit_id(),
    779             receipt.digest(),
    780             AtomicCommitDisposition::Committed,
    781             50,
    782             AuthoredAtomicOutcome::DeliveryPlan(altered),
    783         )
    784         .unwrap();
    785         assert!(!forged.matches_command(&command));
    786     }
    787     let mut stopped = current;
    788     stopped.request_stop(100).unwrap();
    789     assert!(
    790         AuthoredAtomicReceipt::new(
    791             &command,
    792             AtomicCommitDisposition::Committed,
    793             60,
    794             AuthoredAtomicOutcome::DeliveryPlan(stopped)
    795         )
    796         .is_err()
    797     );
    798 }