lib

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

reconciliation_tests.rs (37438B)


      1 use super::*;
      2 use radroots_storage::{
      3     authored::{FailureClass, RetrySchedule, WorkFailure, WorkPhase},
      4     authored_atomic::ReconcileDeliveryFacts,
      5     authored_delivery::AuthoredDeliveryHistory,
      6 };
      7 use std::num::NonZeroU32;
      8 
      9 fn history(storage: &MemoryStorage) -> AuthoredDeliveryHistory {
     10     block_on(storage.authored_delivery_history(ids().2))
     11         .unwrap()
     12         .unwrap()
     13 }
     14 
     15 fn fence(claim: &WorkClaim) -> WorkFence {
     16     WorkFence::new(*claim.token(), claim.generation(), claim.row_revision()).unwrap()
     17 }
     18 
     19 fn reconcile(
     20     plan: &AuthoredDeliveryPlan,
     21     authority: Option<WorkFence>,
     22     retry: Option<RetrySchedule>,
     23     at: u64,
     24 ) -> AuthoredAtomicCommand {
     25     AuthoredAtomicCommand::ReconcileDelivery(
     26         ReconcileDeliveryFacts::new(plan, authority, retry, at).unwrap(),
     27     )
     28 }
     29 
     30 fn retry(attempt: u32, at: u64) -> RetrySchedule {
     31     RetrySchedule::new(
     32         NonZeroU32::new(attempt).unwrap(),
     33         at,
     34         WorkFailure::new(
     35             "delivery_pending",
     36             WorkPhase::Delivery,
     37             FailureClass::Retryable,
     38             Some(at),
     39             None,
     40         )
     41         .unwrap(),
     42     )
     43     .unwrap()
     44 }
     45 
     46 #[test]
     47 fn history_distinguishes_no_issued_work_from_stopped_unresolved_work() {
     48     let untouched = MemoryStorage::new(SourceGeneration::new([71; 32]).unwrap());
     49     assert!(
     50         block_on(untouched.authored_delivery_history(ids().2))
     51             .unwrap()
     52             .is_none()
     53     );
     54     block_on(untouched.execute_authored(prepare().0)).unwrap();
     55     let before = history(&untouched);
     56     assert!(before.is_complete());
     57     assert!(before.proves_no_issued_attempt());
     58     assert!(!before.has_unresolved_claims());
     59     stop(&untouched, 20);
     60     assert!(history(&untouched).proves_no_issued_attempt());
     61 
     62     let (storage, active, _) = prepared();
     63     let issued = history(&storage);
     64     assert_eq!(issued.claims().len(), 1);
     65     assert_eq!(issued.claims()[0].claim(), &active);
     66     assert_eq!(issued.claims()[0].prior_attempt_count(), 0);
     67     assert!(issued.has_unresolved_claims());
     68     assert!(!issued.proves_no_issued_attempt());
     69     stop(&storage, 20);
     70     assert!(history(&storage).has_unresolved_claims());
     71     block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap();
     72     let observed = history(&storage);
     73     assert!(!observed.has_unresolved_claims());
     74     assert!(!observed.proves_no_issued_attempt());
     75     assert_eq!(observed.plan().state(), AuthoredDeliveryState::Cancelled);
     76     assert_eq!(
     77         observed.plan().delivery_satisfaction().unwrap(),
     78         SatisfactionState::Satisfied
     79     );
     80 }
     81 
     82 #[test]
     83 fn exact_current_fence_reconciles_once_and_preserves_raw_fact_provenance() {
     84     let (storage, active, original) = prepared();
     85     block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap();
     86     let before = plan(&storage);
     87     let command = reconcile(&before, Some(fence(&active)), None, 15);
     88     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
     89     let after = plan(&storage);
     90     assert!(receipt.matches_command(&command));
     91     assert_eq!(after.state(), AuthoredDeliveryState::Satisfied);
     92     assert_eq!(after.attempt_count(), 1);
     93     assert_eq!(after.attempts()[0].claim_evidence(), Some(&active));
     94     assert_eq!(
     95         after.attempts()[0].outcome(),
     96         before.delivery_facts()[0].outcome()
     97     );
     98     assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 15);
     99     assert_eq!(after.delivery_facts(), before.delivery_facts());
    100     assert_eq!(after.pending_delivery_facts().count(), 0);
    101     assert!(!history(&storage).has_unresolved_claims());
    102     let replay = block_on(storage.execute_authored(command)).unwrap();
    103     assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
    104     assert_eq!(plan(&storage), after);
    105     assert_eq!(
    106         block_on(storage.authored_receipt(original.commit_id()))
    107             .unwrap()
    108             .unwrap(),
    109         original
    110     );
    111     assert_eq!(
    112         ReconcileDeliveryFacts::new(&after, None, None, 16),
    113         Err(Error::AtomicWorkflowMismatch)
    114     );
    115 }
    116 
    117 #[test]
    118 fn distinct_current_reconciliation_survives_legacy_generation_identity_reuse() {
    119     let (storage, first, _) = prepared();
    120     let first_plan = plan(&storage);
    121     block_on(storage.execute_authored(fact(&first_plan, first.clone(), false, 14))).unwrap();
    122     let first_command = reconcile(&plan(&storage), Some(fence(&first)), Some(retry(1, 18)), 15);
    123     block_on(storage.execute_authored(first_command.clone())).unwrap();
    124     let second = claim(plan(&storage).revision(), 2, 20);
    125     block_on(
    126         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    127             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    128             second.clone(),
    129         ))),
    130     )
    131     .unwrap();
    132     block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 21))).unwrap();
    133     let second_command = reconcile(
    134         &plan(&storage),
    135         Some(fence(&second)),
    136         Some(retry(2, 25)),
    137         22,
    138     );
    139     assert_ne!(first_command.commit_id(), second_command.commit_id());
    140     let legacy = |claim: &WorkClaim, at| {
    141         AuthoredAtomicCommand::ApplyDelivery(
    142             ApplyDeliveryAttempt::new(ids().2, fence(claim), outcome(&first_plan, false), None, at)
    143                 .unwrap(),
    144         )
    145     };
    146     assert_eq!(
    147         legacy(&first, 15).commit_id(),
    148         legacy(&second, 22).commit_id()
    149     );
    150     block_on(storage.execute_authored(second_command)).unwrap();
    151     let after = plan(&storage);
    152     assert_eq!(after.state(), AuthoredDeliveryState::Retryable);
    153     assert_eq!(after.attempt_count(), 2);
    154     assert_eq!(after.retry().unwrap().not_before_unix_ms(), 25);
    155     assert_eq!(after.attempts()[0].claim_evidence(), Some(&first));
    156     assert_eq!(after.attempts()[1].claim_evidence(), Some(&second));
    157     assert!(!history(&storage).has_unresolved_claims());
    158 }
    159 
    160 #[test]
    161 fn late_fact_cannot_take_a_newer_lease_and_its_marker_cannot_resolve_that_lease() {
    162     let (storage, old, _) = prepared();
    163     let newer = claim(plan(&storage).revision(), 3, 40);
    164     block_on(
    165         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    166             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    167             newer.clone(),
    168         ))),
    169     )
    170     .unwrap();
    171     block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap();
    172     let before = plan(&storage);
    173     for authority in [None, Some(fence(&old))] {
    174         let command = reconcile(&before, authority, None, 51);
    175         assert_eq!(
    176             block_on(storage.execute_authored(command)),
    177             Err(Error::DeliveryPlanClaimConflict)
    178         );
    179         assert_eq!(plan(&storage), before);
    180     }
    181     block_on(storage.execute_authored(reconcile(&before, Some(fence(&newer)), None, 51))).unwrap();
    182     let after = history(&storage);
    183     assert_eq!(after.plan().attempts()[0].claim_evidence(), Some(&old));
    184     assert!(
    185         after.has_unresolved_claims(),
    186         "the new worker has no final result yet"
    187     );
    188     assert_eq!(
    189         after.plan().delivery_satisfaction().unwrap(),
    190         SatisfactionState::Satisfied
    191     );
    192 }
    193 
    194 #[test]
    195 fn fresh_observer_can_reconcile_after_expiry_but_expired_worker_cannot() {
    196     let (storage, old, _) = prepared();
    197     block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap();
    198     let before = plan(&storage);
    199     assert_eq!(
    200         block_on(storage.execute_authored(reconcile(&before, Some(fence(&old)), None, 51))),
    201         Err(Error::DeliveryPlanClaimConflict)
    202     );
    203     block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap();
    204     assert_eq!(plan(&storage).state(), AuthoredDeliveryState::Satisfied);
    205     assert_eq!(plan(&storage).attempts()[0].claim_evidence(), Some(&old));
    206 }
    207 
    208 #[test]
    209 fn exact_fact_set_and_revision_races_fail_without_partial_scheduling() {
    210     let (storage, old, _) = prepared();
    211     let newer = claim(plan(&storage).revision(), 3, 40);
    212     block_on(
    213         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    214             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    215             newer.clone(),
    216         ))),
    217     )
    218     .unwrap();
    219     block_on(storage.execute_authored(fact(&plan(&storage), old.clone(), true, 50))).unwrap();
    220     let before = plan(&storage);
    221     let stale = reconcile(&before, Some(fence(&newer)), None, 55);
    222     block_on(storage.execute_authored(fact(&before, newer.clone(), false, 51))).unwrap();
    223     let current = plan(&storage);
    224     assert_eq!(current.revision(), before.revision());
    225     assert_eq!(
    226         block_on(storage.execute_authored(stale)),
    227         Err(Error::DeliveryPlanClaimConflict)
    228     );
    229     assert_eq!(plan(&storage), current);
    230     block_on(storage.execute_authored(reconcile(&current, Some(fence(&newer)), None, 55))).unwrap();
    231     let settled = plan(&storage);
    232     assert_eq!(settled.attempt_count(), 2);
    233     assert_eq!(settled.state(), AuthoredDeliveryState::Satisfied);
    234     assert_eq!(settled.attempts()[0].claim_evidence(), Some(&old));
    235     assert_eq!(settled.attempts()[1].claim_evidence(), Some(&newer));
    236     assert!(
    237         settled
    238             .attempts()
    239             .iter()
    240             .all(|attempt| attempt.recorded_at_unix_ms() == 55)
    241     );
    242 
    243     let (stopped, active, _) = prepared();
    244     block_on(stopped.execute_authored(fact(&plan(&stopped), active, true, 50))).unwrap();
    245     let stale = reconcile(&plan(&stopped), None, None, 51);
    246     stop(&stopped, 52);
    247     let before = plan(&stopped);
    248     assert_eq!(
    249         block_on(stopped.execute_authored(stale)),
    250         Err(Error::DeliveryPlanClaimConflict)
    251     );
    252     assert_eq!(
    253         block_on(stopped.execute_authored(reconcile(&before, None, None, 53))),
    254         Err(Error::InvalidAuthoredTransition)
    255     );
    256     assert_eq!(plan(&stopped), before);
    257 }
    258 
    259 #[test]
    260 fn historical_application_is_resolved_without_rewriting_its_snapshot_shape() {
    261     let (storage, active, _) = prepared();
    262     let command = AuthoredAtomicCommand::ApplyDelivery(
    263         ApplyDeliveryAttempt::new(
    264             ids().2,
    265             fence(&active),
    266             outcome(&plan(&storage), true),
    267             None,
    268             14,
    269         )
    270         .unwrap(),
    271     );
    272     block_on(storage.execute_authored(command)).unwrap();
    273     let legacy = history(&storage);
    274     assert!(!legacy.has_unresolved_claims());
    275     let attempt = &legacy.plan().attempts()[0];
    276     assert!(attempt.claim_evidence().is_none());
    277     assert!(
    278         !serde_json::to_value(attempt)
    279             .unwrap()
    280             .as_object()
    281             .unwrap()
    282             .contains_key("claim")
    283     );
    284     let restored: AuthoredDeliveryPlan =
    285         serde_json::from_slice(&serde_json::to_vec(legacy.plan()).unwrap()).unwrap();
    286     assert_eq!(restored, *legacy.plan());
    287 }
    288 
    289 #[test]
    290 fn missing_truncated_or_malformed_history_never_proves_absence() {
    291     let storage = MemoryStorage::new(SourceGeneration::new([72; 32]).unwrap());
    292     let original = block_on(storage.execute_authored(prepare().0)).unwrap();
    293     let initial = plan(&storage);
    294     let unknown = AuthoredDeliveryHistory::new(initial.clone(), None).unwrap();
    295     assert!(!unknown.is_complete());
    296     assert!(!unknown.proves_no_issued_attempt());
    297     assert!(unknown.has_unresolved_claims());
    298     let mut bounded = AuthoredDeliveryHistory::new(initial, Some(&original)).unwrap();
    299     bounded.mark_truncated();
    300     assert!(bounded.is_truncated());
    301     assert!(!bounded.proves_no_issued_attempt());
    302     assert!(bounded.has_unresolved_claims());
    303     assert_eq!(
    304         bounded.require_pending_fact_provenance(),
    305         Err(Error::DeliveryAttemptOverflow)
    306     );
    307     assert_eq!(
    308         bounded.push_claim(&original),
    309         Err(Error::AtomicWorkflowMismatch)
    310     );
    311 
    312     let (storage, _, issued) = prepared();
    313     let mut duplicate = history(&storage);
    314     assert_eq!(
    315         duplicate.push_claim(&issued),
    316         Err(Error::AtomicWorkflowMismatch)
    317     );
    318     assert!(AuthoredDeliveryHistory::new(plan(&storage), Some(&issued)).is_err());
    319     let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap();
    320     assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch));
    321 }
    322 
    323 #[test]
    324 fn forged_or_duplicate_reconciliation_markers_fail_snapshot_validation() {
    325     let (storage, active, _) = prepared();
    326     block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), true, 14))).unwrap();
    327     block_on(storage.execute_authored(reconcile(&plan(&storage), Some(fence(&active)), None, 15)))
    328         .unwrap();
    329     let valid = serde_json::to_value(plan(&storage)).unwrap();
    330     let mut missing_fact = valid.clone();
    331     missing_fact["delivery_facts"] = serde_json::json!([]);
    332     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(missing_fact).is_err());
    333     let mut forged = valid.clone();
    334     forged["attempts"][0]["claim"] = serde_json::to_value(claim(NonZeroU64::MIN, 9, 13)).unwrap();
    335     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(forged).is_err());
    336     let mut duplicate = valid;
    337     let mut second = duplicate["attempts"][0].clone();
    338     second["attempt"] = serde_json::json!(2);
    339     duplicate["attempts"].as_array_mut().unwrap().push(second);
    340     duplicate["attempt_count"] = serde_json::json!(2);
    341     assert!(serde_json::from_value::<AuthoredDeliveryPlan>(duplicate).is_err());
    342 }
    343 
    344 #[test]
    345 fn issued_claim_limit_rejects_new_work_atomically_and_keeps_original_replay() {
    346     let (storage, _, first) = prepared();
    347     for index in 1..DELIVERY_PLAN_ATTEMPTS_MAX {
    348         let at = 40 + u64::from(index) * 21;
    349         let active = WorkClaim::new(
    350             [7; 16],
    351             "bounded-worker",
    352             NonZeroU64::new(u64::from(index) + 2).unwrap(),
    353             at,
    354             at + 20,
    355             plan(&storage).revision(),
    356         )
    357         .unwrap();
    358         block_on(
    359             storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    360                 ClaimAuthoredTarget::DeliveryPlan(ids().2),
    361                 active,
    362             ))),
    363         )
    364         .unwrap();
    365     }
    366     let before = plan(&storage);
    367     let at = 40 + u64::from(DELIVERY_PLAN_ATTEMPTS_MAX) * 21;
    368     let active = WorkClaim::new(
    369         [8; 16],
    370         "overflow-worker",
    371         NonZeroU64::new(2048).unwrap(),
    372         at,
    373         at + 20,
    374         before.revision(),
    375     )
    376     .unwrap();
    377     let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    378         ClaimAuthoredTarget::DeliveryPlan(ids().2),
    379         active,
    380     ));
    381     assert_eq!(
    382         block_on(storage.execute_authored(command.clone())),
    383         Err(Error::DeliveryAttemptOverflow)
    384     );
    385     assert_eq!(plan(&storage), before);
    386     assert!(
    387         block_on(storage.authored_receipt(command.commit_id()))
    388             .unwrap()
    389             .is_none()
    390     );
    391     let history = history(&storage);
    392     assert_eq!(history.claims().len(), DELIVERY_PLAN_ATTEMPTS_MAX as usize);
    393     assert!(history.is_complete());
    394     assert!(history.has_unresolved_claims());
    395     assert_eq!(
    396         history.clone().push_claim(&first),
    397         Err(Error::DeliveryAttemptOverflow)
    398     );
    399     let AuthoredAtomicOutcome::DeliveryPlan(original) = first.outcome() else {
    400         unreachable!()
    401     };
    402     let original = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    403         ClaimAuthoredTarget::DeliveryPlan(ids().2),
    404         original.claim_evidence().unwrap().clone(),
    405     ));
    406     assert_eq!(
    407         block_on(storage.execute_authored(original))
    408             .unwrap()
    409             .disposition(),
    410         AtomicCommitDisposition::Replay
    411     );
    412     assert_eq!(plan(&storage), before);
    413 }
    414 
    415 #[test]
    416 fn reconciliation_rejects_future_facts_zero_time_and_invalid_retry_without_mutation() {
    417     for accepted in [false, true] {
    418         let (storage, active, _) = prepared();
    419         block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 30)))
    420             .unwrap();
    421         let before = plan(&storage);
    422         assert_eq!(
    423             ReconcileDeliveryFacts::new(&before, None, None, 0),
    424             Err(Error::AtomicWorkflowMismatch)
    425         );
    426         assert_eq!(
    427             block_on(storage.execute_authored(reconcile(&before, Some(fence(&active)), None, 29))),
    428             Err(Error::AtomicWorkflowMismatch)
    429         );
    430         let invalid = if accepted { Some(retry(1, 60)) } else { None };
    431         assert_eq!(
    432             block_on(storage.execute_authored(reconcile(&before, None, invalid, 50))),
    433             Err(Error::InvalidRetrySchedule)
    434         );
    435         if !accepted {
    436             assert_eq!(
    437                 block_on(storage.execute_authored(reconcile(
    438                     &before,
    439                     None,
    440                     Some(retry(2, 60)),
    441                     50
    442                 ))),
    443                 Err(Error::InvalidRetrySchedule)
    444             );
    445         }
    446         assert_eq!(plan(&storage), before);
    447     }
    448 }
    449 
    450 #[test]
    451 fn raw_sink_failure_controls_retry_diagnostic_and_provider_backoff() {
    452     let (storage, active, _) = prepared();
    453     let failure = SinkFailure::for_request(
    454         plan(&storage).request().unwrap(),
    455         "sink_lost",
    456         Retryability::Retryable,
    457         Some(65),
    458         Some("connection lost".into()),
    459         vec![],
    460     )
    461     .unwrap();
    462     block_on(
    463         storage.execute_authored(AuthoredAtomicCommand::RecordDelivery(
    464             RecordDeliveryFact::new(
    465                 ids().2,
    466                 ids().1,
    467                 active,
    468                 DeliveryAttemptOutcome::SinkFailure(failure.clone()),
    469                 50,
    470             )
    471             .unwrap(),
    472         )),
    473     )
    474     .unwrap();
    475     let before = plan(&storage);
    476     for (code, diagnostic, at) in [
    477         ("other", Some("connection lost"), 65),
    478         ("sink_lost", None, 65),
    479         ("sink_lost", Some("connection lost"), 64),
    480     ] {
    481         let retry = RetrySchedule::new(
    482             NonZeroU32::MIN,
    483             at,
    484             WorkFailure::new(
    485                 code,
    486                 WorkPhase::Delivery,
    487                 FailureClass::Retryable,
    488                 None,
    489                 diagnostic.map(str::to_owned),
    490             )
    491             .unwrap(),
    492         )
    493         .unwrap();
    494         assert_eq!(
    495             block_on(storage.execute_authored(reconcile(&before, None, Some(retry), 51))),
    496             Err(Error::InvalidRetrySchedule)
    497         );
    498         assert_eq!(plan(&storage), before);
    499     }
    500     let retry = RetrySchedule::new(
    501         NonZeroU32::MIN,
    502         65,
    503         WorkFailure::new(
    504             "sink_lost",
    505             WorkPhase::Delivery,
    506             FailureClass::Retryable,
    507             Some(65),
    508             Some("connection lost".into()),
    509         )
    510         .unwrap(),
    511     )
    512     .unwrap();
    513     block_on(storage.execute_authored(reconcile(&before, None, Some(retry.clone()), 51))).unwrap();
    514     let after = plan(&storage);
    515     assert_eq!(after.state(), AuthoredDeliveryState::Retryable);
    516     assert_eq!(after.retry(), Some(&retry));
    517     assert_eq!(
    518         after.attempts()[0].outcome(),
    519         &DeliveryAttemptOutcome::SinkFailure(failure)
    520     );
    521 }
    522 
    523 #[test]
    524 fn matching_late_fact_marks_legacy_attempt_without_inventing_another_attempt() {
    525     for accepted in [false, true] {
    526         let (storage, active, _) = prepared();
    527         let command = AuthoredAtomicCommand::ApplyDelivery(
    528             ApplyDeliveryAttempt::new(
    529                 ids().2,
    530                 fence(&active),
    531                 outcome(&plan(&storage), false),
    532                 Some(retry(1, 18)),
    533                 14,
    534             )
    535             .unwrap(),
    536         );
    537         let original = block_on(storage.execute_authored(command)).unwrap();
    538         block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), accepted, 50)))
    539             .unwrap();
    540         let before = plan(&storage);
    541         let command = reconcile(&before, None, Some(retry(1, 60)), 51);
    542         if accepted {
    543             assert_eq!(
    544                 block_on(storage.execute_authored(command)),
    545                 Err(Error::AtomicWorkflowMismatch)
    546             );
    547             assert_eq!(plan(&storage), before);
    548             // Conflicting raw facts remain factual evidence, not a replacement
    549             // for the immutable earlier attempt or authority for another effect.
    550             assert_eq!(
    551                 before.delivery_satisfaction().unwrap(),
    552                 SatisfactionState::Satisfied
    553             );
    554         } else {
    555             block_on(storage.execute_authored(command.clone())).unwrap();
    556             let after = plan(&storage);
    557             assert_eq!(after.attempt_count(), 1);
    558             assert_eq!(after.attempts()[0].recorded_at_unix_ms(), 14);
    559             assert_eq!(after.attempts()[0].claim_evidence(), Some(&active));
    560             assert_eq!(after.pending_delivery_facts().count(), 0);
    561             assert_eq!(
    562                 block_on(storage.execute_authored(command))
    563                     .unwrap()
    564                     .disposition(),
    565                 AtomicCommitDisposition::Replay
    566             );
    567         }
    568         assert_eq!(
    569             block_on(storage.authored_receipt(original.commit_id()))
    570                 .unwrap()
    571                 .unwrap(),
    572             original
    573         );
    574     }
    575 }
    576 
    577 #[test]
    578 fn partial_acceptance_and_terminal_evidence_settle_without_retry_authority() {
    579     for (retryability, partial, expected) in [
    580         (
    581             Retryability::Terminal,
    582             None,
    583             AuthoredDeliveryState::FailedTerminal,
    584         ),
    585         (
    586             Retryability::Retryable,
    587             Some(true),
    588             AuthoredDeliveryState::Satisfied,
    589         ),
    590         (
    591             Retryability::Retryable,
    592             Some(false),
    593             AuthoredDeliveryState::Exhausted,
    594         ),
    595         (
    596             Retryability::Terminal,
    597             Some(false),
    598             AuthoredDeliveryState::Exhausted,
    599         ),
    600     ] {
    601         let (storage, active, _) = prepared();
    602         let before = plan(&storage);
    603         let evidence = partial
    604             .map(|accepted| {
    605                 DeliveryTargetReceipt::attempted(
    606                     before.request().unwrap().target_set().targets()[0].clone(),
    607                     if accepted {
    608                         DeliveryOutcome::accepted()
    609                     } else {
    610                         DeliveryOutcome::rejected()
    611                     },
    612                 )
    613             })
    614             .into_iter()
    615             .collect();
    616         let failure = SinkFailure::for_request(
    617             before.request().unwrap(),
    618             "sink_lost",
    619             retryability,
    620             None,
    621             None,
    622             evidence,
    623         )
    624         .unwrap();
    625         block_on(
    626             storage.execute_authored(AuthoredAtomicCommand::RecordDelivery(
    627                 RecordDeliveryFact::new(
    628                     ids().2,
    629                     ids().1,
    630                     active,
    631                     DeliveryAttemptOutcome::SinkFailure(failure),
    632                     50,
    633                 )
    634                 .unwrap(),
    635             )),
    636         )
    637         .unwrap();
    638         let before = plan(&storage);
    639         assert_eq!(
    640             block_on(storage.execute_authored(reconcile(&before, None, Some(retry(1, 60)), 51))),
    641             Err(Error::InvalidRetrySchedule)
    642         );
    643         assert_eq!(plan(&storage), before);
    644         block_on(storage.execute_authored(reconcile(&before, None, None, 51))).unwrap();
    645         assert_eq!(plan(&storage).state(), expected);
    646         assert!(plan(&storage).retry().is_none());
    647     }
    648 }
    649 
    650 #[test]
    651 fn earlier_unresolved_claim_cannot_adopt_a_new_workers_legacy_attempt() {
    652     let (storage, first, _) = prepared();
    653     let second = claim(plan(&storage).revision(), 3, 40);
    654     block_on(
    655         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    656             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    657             second.clone(),
    658         ))),
    659     )
    660     .unwrap();
    661     block_on(storage.execute_authored(fact(&plan(&storage), first, false, 41))).unwrap();
    662     block_on(storage.execute_authored(reconcile(
    663         &plan(&storage),
    664         Some(fence(&second)),
    665         Some(retry(1, 43)),
    666         42,
    667     )))
    668     .unwrap();
    669     let third = claim(plan(&storage).revision(), 4, 44);
    670     block_on(
    671         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    672             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    673             third.clone(),
    674         ))),
    675     )
    676     .unwrap();
    677     let legacy = AuthoredAtomicCommand::ApplyDelivery(
    678         ApplyDeliveryAttempt::new(
    679             ids().2,
    680             fence(&third),
    681             outcome(&plan(&storage), false),
    682             Some(retry(2, 47)),
    683             45,
    684         )
    685         .unwrap(),
    686     );
    687     block_on(storage.execute_authored(legacy)).unwrap();
    688     assert!(history(&storage).has_unresolved_claims());
    689     block_on(storage.execute_authored(fact(&plan(&storage), second.clone(), false, 50))).unwrap();
    690     block_on(storage.execute_authored(reconcile(&plan(&storage), None, Some(retry(3, 65)), 61)))
    691         .unwrap();
    692     let after = plan(&storage);
    693     assert_eq!(after.attempt_count(), 3);
    694     assert!(after.attempts()[1].claim_evidence().is_none());
    695     assert_eq!(after.attempts()[2].claim_evidence(), Some(&second));
    696     assert!(!history(&storage).has_unresolved_claims());
    697 }
    698 
    699 fn restored_receipt(
    700     original: &AuthoredAtomicReceipt,
    701     plan: AuthoredDeliveryPlan,
    702     at: u64,
    703 ) -> AuthoredAtomicReceipt {
    704     AuthoredAtomicReceipt::from_durable_parts(
    705         original.commit_id(),
    706         original.digest(),
    707         AtomicCommitDisposition::Committed,
    708         at,
    709         AuthoredAtomicOutcome::DeliveryPlan(plan),
    710     )
    711     .unwrap()
    712 }
    713 
    714 #[test]
    715 fn original_claim_history_rejects_each_mismatched_binding() {
    716     let (storage, _, original) = prepared();
    717     let current = plan(&storage);
    718     for field in ["plan_id", "artifact_id", "request", "created_at_unix_ms"] {
    719         let mut wire = serde_json::to_value(&current).unwrap();
    720         match field {
    721             "plan_id" => {
    722                 wire[field] = serde_json::to_value(
    723                     radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16])
    724                         .unwrap(),
    725                 )
    726                 .unwrap()
    727             }
    728             "artifact_id" => {
    729                 wire[field] = serde_json::to_value(
    730                     radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(),
    731                 )
    732                 .unwrap()
    733             }
    734             "request" => wire[field] = serde_json::Value::Null,
    735             _ => wire[field] = serde_json::json!(9),
    736         }
    737         let altered = serde_json::from_value(wire).unwrap();
    738         let forged = restored_receipt(&original, altered, 13);
    739         let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap();
    740         assert_eq!(
    741             history.push_claim(&forged),
    742             Err(Error::AtomicWorkflowMismatch),
    743             "{field}"
    744         );
    745         assert!(history.claims().is_empty());
    746     }
    747     let mut rebound = serde_json::to_value(&current).unwrap();
    748     rebound["request"] = serde_json::to_value(
    749         current
    750             .intent()
    751             .materialize(radroots_transport::sink::DeliveryPayload::new(event(
    752                 OTHER_RAW,
    753             )))
    754             .unwrap(),
    755     )
    756     .unwrap();
    757     let mut history =
    758         AuthoredDeliveryHistory::new(serde_json::from_value(rebound).unwrap(), None).unwrap();
    759     assert_eq!(
    760         history.push_claim(&original),
    761         Err(Error::AtomicWorkflowMismatch)
    762     );
    763     let mut history = AuthoredDeliveryHistory::new(current.clone(), None).unwrap();
    764     assert_eq!(
    765         history.push_claim(&restored_receipt(&original, current.clone(), 14)),
    766         Err(Error::AtomicWorkflowMismatch)
    767     );
    768     let corrupt_digest = AuthoredAtomicReceipt::from_durable_parts(
    769         original.commit_id(),
    770         AtomicCommitDigest::new([99; 32]),
    771         AtomicCommitDisposition::Committed,
    772         13,
    773         original.outcome().clone(),
    774     )
    775     .unwrap();
    776     assert_eq!(
    777         history.push_claim(&corrupt_digest),
    778         Err(Error::AtomicWorkflowMismatch)
    779     );
    780     // A structurally valid future original cannot explain the current row.
    781     for (revision, at) in [
    782         (current.revision().get() + 1, 13),
    783         (current.revision().get(), 14),
    784     ] {
    785         let mut wire = serde_json::to_value(&current).unwrap();
    786         let active = claim(NonZeroU64::new(revision - 1).unwrap(), 8, at);
    787         wire["revision"] = serde_json::json!(revision);
    788         wire["updated_at_unix_ms"] = serde_json::json!(at);
    789         wire["claim"] = serde_json::to_value(&active).unwrap();
    790         let future = serde_json::from_value(wire).unwrap();
    791         let command = AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    792             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    793             active,
    794         ));
    795         let receipt = AuthoredAtomicReceipt::new(
    796             &command,
    797             AtomicCommitDisposition::Committed,
    798             at,
    799             AuthoredAtomicOutcome::DeliveryPlan(future),
    800         )
    801         .unwrap();
    802         assert_eq!(
    803             history.push_claim(&receipt),
    804             Err(Error::AtomicWorkflowMismatch)
    805         );
    806     }
    807 }
    808 
    809 #[test]
    810 fn preparation_history_requires_exact_initial_identity_and_monotonic_row() {
    811     let storage = MemoryStorage::new(SourceGeneration::new([73; 32]).unwrap());
    812     let original = block_on(storage.execute_authored(prepare().0)).unwrap();
    813     let initial = plan(&storage);
    814     for field in ["plan_id", "artifact_id", "created_at_unix_ms"] {
    815         let mut wire = serde_json::to_value(&initial).unwrap();
    816         match field {
    817             "plan_id" => {
    818                 wire[field] = serde_json::to_value(
    819                     radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16])
    820                         .unwrap(),
    821                 )
    822                 .unwrap()
    823             }
    824             "artifact_id" => {
    825                 wire[field] = serde_json::to_value(
    826                     radroots_storage::authored::AuthoredArtifactId::new([99; 16]).unwrap(),
    827                 )
    828                 .unwrap()
    829             }
    830             _ => wire[field] = serde_json::json!(9),
    831         }
    832         assert!(
    833             AuthoredDeliveryHistory::new(serde_json::from_value(wire).unwrap(), Some(&original))
    834                 .is_err(),
    835             "{field}"
    836         );
    837     }
    838     let intent = radroots_storage::authored_delivery::AuthoredDeliveryIntent::new(
    839         "different-intent",
    840         initial.intent().target_set().clone(),
    841         initial.intent().satisfaction().clone(),
    842         100,
    843     )
    844     .unwrap();
    845     let changed = AuthoredDeliveryPlan::new(ids().2, ids().1, intent, 10).unwrap();
    846     assert!(AuthoredDeliveryHistory::new(changed, Some(&original)).is_err());
    847     for field in ["revision", "updated_at_unix_ms"] {
    848         let mut wire = serde_json::to_value(&initial).unwrap();
    849         wire[field] = serde_json::json!(12);
    850         let altered = serde_json::from_value(wire).unwrap();
    851         let AuthoredAtomicOutcome::Prepared {
    852             operation,
    853             artifacts,
    854             ..
    855         } = original.outcome()
    856         else {
    857             unreachable!()
    858         };
    859         let receipt = AuthoredAtomicReceipt::from_durable_parts(
    860             original.commit_id(),
    861             original.digest(),
    862             AtomicCommitDisposition::Committed,
    863             12,
    864             AuthoredAtomicOutcome::Prepared {
    865                 operation: operation.clone(),
    866                 artifacts: artifacts.clone(),
    867                 delivery_plans: vec![altered],
    868             },
    869         )
    870         .unwrap();
    871         assert!(
    872             AuthoredDeliveryHistory::new(initial.clone(), Some(&receipt)).is_err(),
    873             "{field}"
    874         );
    875     }
    876 }
    877 
    878 #[test]
    879 fn reconciliation_receipts_cannot_substitute_a_different_plan_or_result() {
    880     let (storage, active, _) = prepared();
    881     block_on(storage.execute_authored(fact(&plan(&storage), active.clone(), false, 14))).unwrap();
    882     let command = reconcile(
    883         &plan(&storage),
    884         Some(fence(&active)),
    885         Some(retry(1, 20)),
    886         15,
    887     );
    888     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
    889     let after = plan(&storage);
    890     for field in [
    891         "plan_id",
    892         "revision",
    893         "updated_at_unix_ms",
    894         "claim",
    895         "stop",
    896         "pending",
    897         "facts",
    898         "retry",
    899     ] {
    900         let mut wire = serde_json::to_value(&after).unwrap();
    901         match field {
    902             "plan_id" => {
    903                 wire[field] = serde_json::to_value(
    904                     radroots_storage::authored_delivery::AuthoredDeliveryPlanId::new([99; 16])
    905                         .unwrap(),
    906                 )
    907                 .unwrap()
    908             }
    909             "revision" => wire[field] = serde_json::json!(after.revision().get() + 1),
    910             "updated_at_unix_ms" => wire[field] = serde_json::json!(16),
    911             "claim" => {
    912                 wire[field] = serde_json::to_value(claim(
    913                     NonZeroU64::new(after.revision().get() - 1).unwrap(),
    914                     9,
    915                     15,
    916                 ))
    917                 .unwrap();
    918             }
    919             "stop" => {
    920                 wire["stop_requested_at_unix_ms"] = serde_json::json!(15);
    921                 wire["state"] = serde_json::json!("cancelled");
    922                 wire["retry"] = serde_json::Value::Null;
    923                 wire["last_failure"] = serde_json::Value::Null;
    924             }
    925             "pending" => wire["attempts"][0]["claim"] = serde_json::Value::Null,
    926             "facts" => wire["delivery_facts"][0]["observed_at_unix_ms"] = serde_json::json!(15),
    927             _ => {
    928                 let other = retry(1, 21);
    929                 wire["retry"] = serde_json::to_value(&other).unwrap();
    930                 wire["last_failure"] = serde_json::to_value(other.failure()).unwrap();
    931             }
    932         }
    933         let altered: AuthoredDeliveryPlan = serde_json::from_value(wire).unwrap();
    934         let wrong = restored_receipt(&receipt, altered.clone(), 15);
    935         assert!(!wrong.matches_command(&command), "{field}");
    936         assert_eq!(
    937             AuthoredAtomicReceipt::new(
    938                 &command,
    939                 AtomicCommitDisposition::Committed,
    940                 15,
    941                 AuthoredAtomicOutcome::DeliveryPlan(altered)
    942             ),
    943             Err(Error::AtomicWorkflowMismatch),
    944             "{field}"
    945         );
    946     }
    947     let artifact = block_on(storage.authored_artifact(ids().1))
    948         .unwrap()
    949         .unwrap();
    950     let wrong = AuthoredAtomicReceipt::from_durable_parts(
    951         receipt.commit_id(),
    952         receipt.digest(),
    953         AtomicCommitDisposition::Committed,
    954         15,
    955         AuthoredAtomicOutcome::Artifact(artifact.clone()),
    956     )
    957     .unwrap();
    958     assert!(!wrong.matches_command(&command));
    959     assert!(
    960         AuthoredAtomicReceipt::new(
    961             &command,
    962             AtomicCommitDisposition::Committed,
    963             15,
    964             AuthoredAtomicOutcome::Artifact(artifact)
    965         )
    966         .is_err()
    967     );
    968 }
    969 
    970 #[test]
    971 fn missing_fact_claims_and_prepopulated_preparation_remain_uncertain() {
    972     let (storage, active, _) = prepared();
    973     let original = block_on(storage.authored_receipt(prepare().0.commit_id()))
    974         .unwrap()
    975         .unwrap();
    976     let AuthoredAtomicOutcome::Prepared {
    977         operation,
    978         artifacts,
    979         ..
    980     } = original.outcome()
    981     else {
    982         unreachable!()
    983     };
    984     let prepopulated = AuthoredAtomicReceipt::from_durable_parts(
    985         original.commit_id(),
    986         original.digest(),
    987         AtomicCommitDisposition::Committed,
    988         13,
    989         AuthoredAtomicOutcome::Prepared {
    990             operation: operation.clone(),
    991             artifacts: artifacts.clone(),
    992             delivery_plans: vec![plan(&storage)],
    993         },
    994     )
    995     .unwrap();
    996     let legacy = AuthoredDeliveryHistory::new(plan(&storage), Some(&prepopulated)).unwrap();
    997     assert!(!legacy.is_complete());
    998     assert!(!legacy.proves_no_issued_attempt());
    999     assert!(legacy.has_unresolved_claims());
   1000     block_on(storage.execute_authored(fact(&plan(&storage), active, true, 50))).unwrap();
   1001     stop(&storage, 51);
   1002     let unknown = AuthoredDeliveryHistory::new(plan(&storage), None).unwrap();
   1003     assert_eq!(
   1004         unknown.require_pending_fact_provenance(),
   1005         Err(Error::AtomicWorkflowMismatch)
   1006     );
   1007     let incomplete = AuthoredDeliveryHistory::new(plan(&storage), Some(&original)).unwrap();
   1008     assert_eq!(incomplete.validate(), Err(Error::AtomicWorkflowMismatch));
   1009 }
   1010 
   1011 #[test]
   1012 fn later_legacy_attempt_outside_old_lease_cannot_resolve_old_claim() {
   1013     let (storage, _, _) = prepared();
   1014     let newer = claim(plan(&storage).revision(), 3, 40);
   1015     block_on(
   1016         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
   1017             ClaimAuthoredTarget::DeliveryPlan(ids().2),
   1018             newer.clone(),
   1019         ))),
   1020     )
   1021     .unwrap();
   1022     block_on(
   1023         storage.execute_authored(AuthoredAtomicCommand::ApplyDelivery(
   1024             ApplyDeliveryAttempt::new(
   1025                 ids().2,
   1026                 fence(&newer),
   1027                 outcome(&plan(&storage), true),
   1028                 None,
   1029                 41,
   1030             )
   1031             .unwrap(),
   1032         )),
   1033     )
   1034     .unwrap();
   1035     let history = history(&storage);
   1036     assert_eq!(history.plan().attempt_count(), 1);
   1037     assert_eq!(history.plan().state(), AuthoredDeliveryState::Satisfied);
   1038     assert!(history.has_unresolved_claims());
   1039 }