lib

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

authored_signed_facts.rs (16451B)


      1 use std::num::NonZeroU64;
      2 
      3 use futures_executor::block_on;
      4 use radroots_event::SignedEvent;
      5 use radroots_storage::{
      6     Error,
      7     atomic::AtomicCommitDisposition,
      8     authored::{
      9         AdmissionState, AuthoredArtifact, FailureClass, SigningState, WorkClaim, WorkFailure,
     10         WorkPhase,
     11     },
     12     authored_atomic::{
     13         ApplySignedArtifact, ApplyWorkFailure, AuthoredAtomicCommand, AuthoredAtomicReceipt,
     14         AuthoredAtomicStorage, AuthoredWorkTarget, CancelAuthoredTarget, CancelAuthoredWork,
     15         ClaimAuthoredTarget, ClaimAuthoredWork, RecordSignedArtifact, WorkFence,
     16     },
     17     authored_delivery::AuthoredDeliveryState,
     18     event::SourceGeneration,
     19     memory::MemoryStorage,
     20 };
     21 
     22 #[path = "authored_signing/fixture.rs"]
     23 mod fixture;
     24 use fixture::*;
     25 
     26 fn prepared() -> (MemoryStorage, SignedEvent, WorkClaim, AuthoredAtomicReceipt) {
     27     let storage = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap());
     28     let (preparation, event) = prepare();
     29     block_on(storage.execute_authored(preparation)).unwrap();
     30     let active = claim(NonZeroU64::new(1).unwrap(), 4, 11);
     31     let receipt = block_on(storage.execute_authored(AuthoredAtomicCommand::Claim(
     32         ClaimAuthoredWork::new(
     33             ClaimAuthoredTarget::ArtifactSigning(ids().1),
     34             active.clone(),
     35         ),
     36     )))
     37     .unwrap();
     38     (storage, event, active, receipt)
     39 }
     40 
     41 fn artifact(storage: &MemoryStorage) -> AuthoredArtifact {
     42     block_on(storage.authored_artifact(ids().1))
     43         .unwrap()
     44         .unwrap()
     45 }
     46 
     47 fn fence(active: &WorkClaim) -> WorkFence {
     48     WorkFence::new(*active.token(), active.generation(), active.row_revision()).unwrap()
     49 }
     50 
     51 #[test]
     52 fn expired_claim_retains_exact_facts_without_relaxing_active_fences() {
     53     let (storage, event, active, original) = prepared();
     54     let stale = AuthoredAtomicCommand::ApplySigned(
     55         ApplySignedArtifact::new(ids().1, fence(&active), event.clone(), 40).unwrap(),
     56     );
     57     assert_eq!(
     58         block_on(storage.execute_authored(stale.clone())),
     59         Err(Error::DeliveryPlanClaimConflict)
     60     );
     61     let fact = record(event.clone(), active.clone(), 40);
     62     assert_ne!(fact.commit_id(), stale.commit_id());
     63     let receipt = block_on(storage.execute_authored(fact.clone())).unwrap();
     64     assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed);
     65     let signed = artifact(&storage);
     66     assert_eq!(signed.signing_state(), SigningState::Signed);
     67     assert_eq!(signed.signed().unwrap().event().raw_json(), RAW);
     68     assert!(signed.signing_claim().is_none());
     69     let delivery = block_on(storage.authored_delivery_plan(ids().2))
     70         .unwrap()
     71         .unwrap();
     72     assert_eq!(delivery.request().unwrap().payload().event(), &event);
     73     assert!(delivery.attempts().is_empty());
     74     assert_eq!(
     75         block_on(storage.authored_receipt(original.commit_id()))
     76             .unwrap()
     77             .unwrap(),
     78         original
     79     );
     80     let later = record(event, active, 90);
     81     assert_eq!(later.commit_id(), fact.commit_id());
     82     let replay = block_on(storage.execute_authored(later)).unwrap();
     83     assert_eq!(replay.disposition(), AtomicCommitDisposition::Replay);
     84     assert_eq!(replay.committed_at_unix_ms(), 40);
     85     assert_eq!(artifact(&storage), signed);
     86 }
     87 
     88 #[test]
     89 fn cancellation_and_terminal_failure_survive_late_signed_bytes() {
     90     for cancelled in [false, true] {
     91         let (storage, event, active, _) = prepared();
     92         let stop = if cancelled {
     93             AuthoredAtomicCommand::Cancel(
     94                 CancelAuthoredWork::new(
     95                     CancelAuthoredTarget::ArtifactSigning(ids().1),
     96                     artifact(&storage).revision(),
     97                     20,
     98                 )
     99                 .unwrap(),
    100             )
    101         } else {
    102             AuthoredAtomicCommand::ApplyFailure(
    103                 ApplyWorkFailure::new(
    104                     AuthoredWorkTarget::Artifact(ids().1),
    105                     fence(&active),
    106                     WorkFailure::new(
    107                         "signing_stopped",
    108                         WorkPhase::Signing,
    109                         FailureClass::Terminal,
    110                         None,
    111                         None,
    112                     )
    113                     .unwrap(),
    114                     None,
    115                     20,
    116                 )
    117                 .unwrap(),
    118             )
    119         };
    120         let stop_receipt = block_on(storage.execute_authored(stop)).unwrap();
    121         let stopped = artifact(&storage);
    122         // An older observation still records the fact without moving row time back.
    123         let fact_receipt = block_on(storage.execute_authored(record(event, active, 15))).unwrap();
    124         assert_eq!(fact_receipt.committed_at_unix_ms(), 20);
    125         let retained = artifact(&storage);
    126         assert_eq!(retained.signing_state(), stopped.signing_state());
    127         assert_eq!(retained.last_failure(), stopped.last_failure());
    128         assert_eq!(retained.updated_at_unix_ms(), 20);
    129         assert_eq!(retained.admission_state(), AdmissionState::Pending);
    130         assert_eq!(retained.signed().unwrap().event().raw_json(), RAW);
    131         let delivery = block_on(storage.authored_delivery_plan(ids().2))
    132             .unwrap()
    133             .unwrap();
    134         assert!(delivery.request().is_none());
    135         for target in [
    136             ClaimAuthoredTarget::ArtifactSigning(ids().1),
    137             ClaimAuthoredTarget::ArtifactAdmission(ids().1),
    138             ClaimAuthoredTarget::DeliveryPlan(ids().2),
    139         ] {
    140             let revision = if matches!(target, ClaimAuthoredTarget::DeliveryPlan(_)) {
    141                 delivery.revision()
    142             } else {
    143                 retained.revision()
    144             };
    145             assert!(
    146                 block_on(storage.execute_authored(AuthoredAtomicCommand::Claim(
    147                     ClaimAuthoredWork::new(target, claim(revision, 8, 50)),
    148                 )))
    149                 .is_err()
    150             );
    151         }
    152         assert_eq!(artifact(&storage), retained);
    153         assert_eq!(
    154             block_on(storage.authored_receipt(stop_receipt.commit_id()))
    155                 .unwrap()
    156                 .unwrap(),
    157             stop_receipt
    158         );
    159         let operation = block_on(storage.authored_operation(ids().0))
    160             .unwrap()
    161             .unwrap();
    162         let settlement = radroots_storage::authored::OperationSettlement::evaluate(
    163             &operation,
    164             std::slice::from_ref(&retained),
    165         )
    166         .unwrap();
    167         assert_eq!(settlement.signed(), 1);
    168         assert_eq!(settlement.cancelled(), u16::from(cancelled));
    169         assert_eq!(settlement.failed_terminal(), u16::from(!cancelled));
    170         #[cfg(feature = "serde")]
    171         assert_eq!(
    172             serde_json::from_str::<AuthoredArtifact>(&serde_json::to_string(&retained).unwrap())
    173                 .unwrap(),
    174             retained
    175         );
    176     }
    177 }
    178 
    179 #[test]
    180 fn superseded_attempts_cannot_overwrite_first_exact_bytes_or_restart_stopped_delivery() {
    181     let (storage, event, first, _) = prepared();
    182     let second = claim(artifact(&storage).revision(), 5, 40);
    183     block_on(
    184         storage.execute_authored(AuthoredAtomicCommand::Claim(ClaimAuthoredWork::new(
    185             ClaimAuthoredTarget::ArtifactSigning(ids().1),
    186             second.clone(),
    187         ))),
    188     )
    189     .unwrap();
    190     let plan = block_on(storage.authored_delivery_plan(ids().2))
    191         .unwrap()
    192         .unwrap();
    193     block_on(
    194         storage.execute_authored(AuthoredAtomicCommand::Cancel(
    195             CancelAuthoredWork::new(
    196                 CancelAuthoredTarget::DeliveryPlan(ids().2),
    197                 plan.revision(),
    198                 41,
    199             )
    200             .unwrap(),
    201         )),
    202     )
    203     .unwrap();
    204     block_on(storage.execute_authored(record(event.clone(), first.clone(), 42))).unwrap();
    205     let retained = artifact(&storage);
    206     assert!(retained.signing_claim().is_none());
    207     block_on(storage.execute_authored(record(event, second.clone(), 43))).unwrap();
    208     assert_eq!(artifact(&storage), retained);
    209     let alternate = fixture::event(&format!(" {RAW} "));
    210     assert_eq!(alternate.id(), retained.signed().unwrap().event().id());
    211     assert_eq!(
    212         block_on(storage.execute_authored(record(alternate, second, 44))),
    213         Err(Error::AtomicCommitConflict)
    214     );
    215     assert_eq!(artifact(&storage), retained);
    216     let plan = block_on(storage.authored_delivery_plan(ids().2))
    217         .unwrap()
    218         .unwrap();
    219     assert_eq!(plan.state(), AuthoredDeliveryState::Cancelled);
    220     assert!(plan.request().is_none());
    221     assert!(plan.attempts().is_empty());
    222 }
    223 
    224 #[test]
    225 fn unrelated_plan_or_attempt_provenance_rolls_back_without_fact_receipts() {
    226     let (storage, event, active, original) = prepared();
    227     let before = artifact(&storage);
    228     let wrong_claims = [
    229         WorkClaim::new(
    230             [9; 16],
    231             active.owner(),
    232             active.generation(),
    233             11,
    234             31,
    235             active.row_revision(),
    236         )
    237         .unwrap(),
    238         WorkClaim::new(
    239             *active.token(),
    240             "different-worker",
    241             active.generation(),
    242             11,
    243             31,
    244             active.row_revision(),
    245         )
    246         .unwrap(),
    247         WorkClaim::new(
    248             *active.token(),
    249             active.owner(),
    250             NonZeroU64::new(9).unwrap(),
    251             11,
    252             31,
    253             active.row_revision(),
    254         )
    255         .unwrap(),
    256         WorkClaim::new(
    257             *active.token(),
    258             active.owner(),
    259             active.generation(),
    260             12,
    261             31,
    262             active.row_revision(),
    263         )
    264         .unwrap(),
    265         WorkClaim::new(
    266             *active.token(),
    267             active.owner(),
    268             active.generation(),
    269             11,
    270             32,
    271             active.row_revision(),
    272         )
    273         .unwrap(),
    274         WorkClaim::new(
    275             *active.token(),
    276             active.owner(),
    277             active.generation(),
    278             11,
    279             31,
    280             NonZeroU64::new(9).unwrap(),
    281         )
    282         .unwrap(),
    283     ];
    284     let expected = record(event.clone(), active.clone(), 50);
    285     for wrong in wrong_claims {
    286         let command = record(event.clone(), wrong, 50);
    287         assert_ne!(command.commit_id(), expected.commit_id());
    288         assert!(block_on(storage.execute_authored(command.clone())).is_err());
    289         assert!(
    290             block_on(storage.authored_receipt(command.commit_id()))
    291                 .unwrap()
    292                 .is_none()
    293         );
    294     }
    295     let wrong_operation = RecordSignedArtifact::new(
    296         radroots_storage::journal::OperationInstanceId::new([9; 16]).unwrap(),
    297         ids().1,
    298         active.clone(),
    299         event.clone(),
    300         50,
    301     )
    302     .unwrap();
    303     assert!(
    304         block_on(storage.execute_authored(AuthoredAtomicCommand::RecordSigned(wrong_operation)))
    305             .is_err()
    306     );
    307     let mismatch = record(fixture::event(OTHER_RAW), active.clone(), 50);
    308     assert!(block_on(storage.execute_authored(mismatch.clone())).is_err());
    309     assert!(
    310         block_on(storage.authored_receipt(mismatch.commit_id()))
    311             .unwrap()
    312             .is_none()
    313     );
    314     assert_eq!(artifact(&storage), before);
    315     let AuthoredAtomicCommand::RecordSigned(value) = expected else {
    316         unreachable!()
    317     };
    318     assert_eq!(value.operation_id(), ids().0);
    319     assert_eq!(value.artifact_id(), ids().1);
    320     assert_eq!(value.claim(), &active);
    321     assert_eq!(value.event(), &event);
    322     assert_eq!(value.observed_at_unix_ms(), 50);
    323     assert_eq!(value.clone(), value);
    324     let (preparation, _) = prepare();
    325     let wrong_outcome = block_on(storage.execute_authored(preparation)).unwrap();
    326     assert!(value.apply_to(&mut before.clone(), &wrong_outcome).is_err());
    327     assert!(value.apply_to(&mut before.clone(), &original).is_ok());
    328     #[cfg(feature = "serde")]
    329     {
    330         let mut regressed = serde_json::to_value(&before).unwrap();
    331         regressed["signing_claim"] = serde_json::Value::Null;
    332         regressed["updated_at_unix_ms"] = serde_json::Value::from(10);
    333         let mut regressed = serde_json::from_value::<AuthoredArtifact>(regressed).unwrap();
    334         assert!(value.apply_to(&mut regressed, &original).is_err());
    335     }
    336 }
    337 
    338 #[test]
    339 fn invalid_signature_and_pre_attempt_observation_never_become_facts() {
    340     let (_, event, active, _) = prepared();
    341     for at in [0, 10] {
    342         assert!(
    343             RecordSignedArtifact::new(ids().0, ids().1, active.clone(), event.clone(), at).is_err()
    344         );
    345     }
    346     let mut value: serde_json::Value = serde_json::from_str(RAW).unwrap();
    347     value["sig"] = serde_json::Value::String("ff".repeat(64));
    348     let hostile = fixture::event(&serde_json::to_string(&value).unwrap());
    349     assert_eq!(
    350         RecordSignedArtifact::new(ids().0, ids().1, active, hostile, 50),
    351         Err(Error::InvalidAuthoredArtifact)
    352     );
    353 }
    354 
    355 #[test]
    356 fn indeterminate_signing_is_resolved_by_original_verified_evidence() {
    357     let (storage, event, active, _) = prepared();
    358     block_on(
    359         storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(
    360             ApplyWorkFailure::new(
    361                 AuthoredWorkTarget::Artifact(ids().1),
    362                 fence(&active),
    363                 WorkFailure::new(
    364                     "signing_unknown",
    365                     WorkPhase::Signing,
    366                     FailureClass::Indeterminate,
    367                     None,
    368                     None,
    369                 )
    370                 .unwrap(),
    371                 None,
    372                 20,
    373             )
    374             .unwrap(),
    375         )),
    376     )
    377     .unwrap();
    378     assert_eq!(
    379         artifact(&storage).signing_state(),
    380         SigningState::Indeterminate
    381     );
    382     block_on(storage.execute_authored(record(event, active, 40))).unwrap();
    383     let retained = artifact(&storage);
    384     assert_eq!(retained.signing_state(), SigningState::Signed);
    385     assert!(retained.last_failure().is_none());
    386     assert_eq!(retained.signed().unwrap().event().raw_json(), RAW);
    387 }
    388 
    389 #[test]
    390 fn signed_fact_receipts_require_exact_outcome_and_monotonic_commit_time() {
    391     let (storage, event, active, _) = prepared();
    392     let command = record(event, active, 40);
    393     let unsigned = artifact(&storage);
    394     assert!(
    395         AuthoredAtomicReceipt::new(
    396             &command,
    397             AtomicCommitDisposition::Committed,
    398             40,
    399             radroots_storage::authored_atomic::AuthoredAtomicOutcome::Artifact(unsigned),
    400         )
    401         .is_err()
    402     );
    403     let receipt = block_on(storage.execute_authored(command.clone())).unwrap();
    404     assert!(receipt.matches_command(&command));
    405     assert!(
    406         AuthoredAtomicReceipt::new(
    407             &command,
    408             AtomicCommitDisposition::Committed,
    409             39,
    410             receipt.outcome().clone(),
    411         )
    412         .is_err()
    413     );
    414     let wrong = radroots_storage::authored_atomic::AuthoredAtomicOutcome::DeliveryPlan(
    415         block_on(storage.authored_delivery_plan(ids().2))
    416             .unwrap()
    417             .unwrap(),
    418     );
    419     assert!(
    420         AuthoredAtomicReceipt::new(
    421             &command,
    422             AtomicCommitDisposition::Committed,
    423             40,
    424             wrong.clone()
    425         )
    426         .is_err()
    427     );
    428     let malformed = AuthoredAtomicReceipt::from_durable_parts(
    429         command.commit_id(),
    430         command.digest(),
    431         AtomicCommitDisposition::Committed,
    432         40,
    433         wrong,
    434     )
    435     .unwrap();
    436     assert!(!malformed.matches_command(&command));
    437 }
    438 
    439 #[cfg(feature = "serde")]
    440 #[test]
    441 fn stopped_signed_snapshots_cannot_erase_the_required_terminal_failure() {
    442     let (storage, event, active, _) = prepared();
    443     block_on(
    444         storage.execute_authored(AuthoredAtomicCommand::ApplyFailure(
    445             ApplyWorkFailure::new(
    446                 AuthoredWorkTarget::Artifact(ids().1),
    447                 fence(&active),
    448                 WorkFailure::new(
    449                     "signing_stopped",
    450                     WorkPhase::Signing,
    451                     FailureClass::Terminal,
    452                     None,
    453                     None,
    454                 )
    455                 .unwrap(),
    456                 None,
    457                 20,
    458             )
    459             .unwrap(),
    460         )),
    461     )
    462     .unwrap();
    463     block_on(storage.execute_authored(record(event, active, 40))).unwrap();
    464     let retained = artifact(&storage);
    465     let mut forged = serde_json::to_value(&retained).unwrap();
    466     forged["last_failure"] = serde_json::Value::Null;
    467     assert!(serde_json::from_value::<AuthoredArtifact>(forged).is_err());
    468 }