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 }