lib

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

atomic.rs (5696B)


      1 use futures_executor::block_on;
      2 use radroots_event::{SignedEvent, wire::Nip01EventWire};
      3 use radroots_storage::{
      4     Error,
      5     atomic::{
      6         AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId,
      7         AtomicCommitOutcome, AtomicCommitReceipt, AtomicStorage, AtomicWorkflow,
      8         AtomicWorkflowKind, CommitIngested,
      9     },
     10     event::{
     11         AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPosition,
     12         EventSequence, SourceGeneration,
     13     },
     14     memory::MemoryStorage,
     15 };
     16 use radroots_transport::{
     17     Target, TransportId,
     18     source::{EventProvenance, ObservedEvent},
     19 };
     20 
     21 fn signed_event() -> SignedEvent {
     22     let raw = r#"{"id":"56bfc78223bb2221bad82b539efdec1ade0f56d0eb0e1f592fd387df4b2ceee0","pubkey":"585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df","created_at":1700000001,"kind":0,"tags":[],"content":"{}","sig":"dddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddddd"}"#;
     23     SignedEvent::from_wire_verified_id(Nip01EventWire::parse_json(raw).expect("wire event"), raw)
     24         .expect("signed event")
     25 }
     26 
     27 fn admission(event: SignedEvent) -> EventAdmission {
     28     let target = Target::nostr_relay("wss://relay.example").expect("target");
     29     let provenance = EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), 100)
     30         .expect("provenance");
     31     EventAdmission::raw(ObservedEvent::new(event, provenance))
     32 }
     33 
     34 fn commit_request(id: u8, digest: u8) -> AtomicCommit {
     35     AtomicCommit::new(
     36         AtomicCommitId::new([id; 16]).expect("commit ID"),
     37         AtomicCommitDigest::new([digest; 32]),
     38         100,
     39         AtomicWorkflow::Ingested(Box::new(CommitIngested::new(
     40             admission(signed_event()),
     41             None,
     42         ))),
     43     )
     44     .expect("commit")
     45 }
     46 
     47 fn outcome() -> AtomicCommitOutcome {
     48     let event = signed_event();
     49     AtomicCommitOutcome::Ingested {
     50         admission: AdmissionReceipt::new(
     51             *event.id(),
     52             EventPosition::new(
     53                 SourceGeneration::new([1; 32]).expect("generation"),
     54                 EventSequence::new(1).expect("sequence"),
     55             ),
     56             AdmissionStage::Raw,
     57             AdmissionDisposition::Inserted,
     58         ),
     59         projection: None,
     60     }
     61 }
     62 
     63 #[test]
     64 fn atomic_identity_workflow_and_receipt_models_are_exact() {
     65     assert_eq!(
     66         AtomicCommitId::new([0; 16]),
     67         Err(Error::InvalidAtomicCommitId)
     68     );
     69     let request = commit_request(1, 2);
     70     assert_eq!(request.commit_id().as_bytes(), &[1; 16]);
     71     assert_eq!(request.digest().as_bytes(), &[2; 32]);
     72     assert_eq!(request.requested_at_unix_ms(), 100);
     73     assert_eq!(request.workflow().kind(), AtomicWorkflowKind::Ingested);
     74     let AtomicWorkflow::Ingested(ingested) = request.workflow();
     75     assert_eq!(ingested.admission().event().id(), signed_event().id());
     76     assert_eq!(ingested.projection(), None);
     77 
     78     assert_eq!(
     79         AtomicCommit::new(
     80             request.commit_id(),
     81             request.digest(),
     82             0,
     83             request.workflow().clone(),
     84         ),
     85         Err(Error::InvalidAtomicCommitTimestamp)
     86     );
     87 
     88     let receipt =
     89         AtomicCommitReceipt::new(&request, AtomicCommitDisposition::Committed, 100, outcome())
     90             .expect("receipt");
     91     assert_eq!(receipt.commit_id(), request.commit_id());
     92     assert_eq!(receipt.digest(), request.digest());
     93     assert_eq!(receipt.disposition(), AtomicCommitDisposition::Committed);
     94     assert_eq!(receipt.committed_at_unix_ms(), 100);
     95     assert_eq!(receipt.outcome().kind(), AtomicWorkflowKind::Ingested);
     96     assert_eq!(
     97         AtomicCommitReceipt::new(&request, AtomicCommitDisposition::Committed, 99, outcome(),),
     98         Err(Error::AtomicWorkflowMismatch)
     99     );
    100 }
    101 
    102 #[test]
    103 fn durable_receipt_reconstruction_rejects_timestamp_incoherence() {
    104     let request = commit_request(1, 2);
    105     for (requested_at, committed_at) in [(0, 100), (101, 100)] {
    106         assert_eq!(
    107             AtomicCommitReceipt::from_durable_parts(
    108                 request.commit_id(),
    109                 request.digest(),
    110                 AtomicCommitDisposition::Replay,
    111                 requested_at,
    112                 committed_at,
    113                 AtomicWorkflowKind::Ingested,
    114                 outcome(),
    115             ),
    116             Err(Error::AtomicWorkflowMismatch)
    117         );
    118     }
    119 
    120     let receipt = AtomicCommitReceipt::from_durable_parts(
    121         request.commit_id(),
    122         request.digest(),
    123         AtomicCommitDisposition::Replay,
    124         100,
    125         101,
    126         AtomicWorkflowKind::Ingested,
    127         outcome(),
    128     )
    129     .expect("durable receipt");
    130     assert_eq!(receipt.disposition(), AtomicCommitDisposition::Replay);
    131     assert_eq!(receipt.committed_at_unix_ms(), 101);
    132 }
    133 
    134 #[test]
    135 fn memory_atomic_boundary_replays_exactly_and_conflicts_on_digest_reuse() {
    136     fn accepts_dyn(_: &dyn AtomicStorage) {}
    137 
    138     let store = MemoryStorage::default();
    139     accepts_dyn(&store);
    140     let request = commit_request(1, 2);
    141     assert_eq!(
    142         block_on(store.commit(request.clone()))
    143             .expect("commit")
    144             .disposition(),
    145         AtomicCommitDisposition::Committed
    146     );
    147     assert_eq!(
    148         block_on(store.commit(request.clone()))
    149             .expect("replay")
    150             .disposition(),
    151         AtomicCommitDisposition::Replay
    152     );
    153     assert_eq!(
    154         block_on(store.commit(commit_request(1, 3))),
    155         Err(Error::AtomicCommitConflict)
    156     );
    157     assert_eq!(
    158         block_on(store.receipt(request.commit_id()))
    159             .expect("receipt")
    160             .expect("durable receipt")
    161             .digest(),
    162         request.digest()
    163     );
    164 }