lib

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

memory.rs (35103B)


      1 #![cfg(feature = "memory")]
      2 
      3 use futures_executor::block_on;
      4 use radroots_event::{
      5     SignedEvent,
      6     admission::{AdmissionPolicy, RawEvent, SignatureVerifier, VisibilityPolicy},
      7     envelope::EventEnvelope,
      8     wire::Nip01EventWire,
      9 };
     10 use radroots_protocol::runtime::v1::OperationId;
     11 use radroots_storage::{
     12     Error, EventStore, Journal, Outbox, ProjectionStore,
     13     atomic::{
     14         AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, AtomicStorage,
     15         AtomicWorkflow, CommitIngested,
     16     },
     17     event::{
     18         EventAdmission, EventPosition, EventQuery, EventQueryBounds, EventSequence,
     19         SourceGeneration,
     20     },
     21     journal::{
     22         IdempotencyDigest, IdempotencyKey, JournalStage, OperationInstanceId, PrepareOperation,
     23         RECOVERABLE_QUERY_LIMIT_MAX,
     24     },
     25     memory::MemoryStorage,
     26     outbox::{
     27         ClaimOutboxItems, DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, LeaseId,
     28         LeaseOwner, OutboxItemId, OutboxStage,
     29     },
     30     private_artifact::{
     31         ArtifactCommitment, ArtifactKind, ArtifactSchemaId, DurableSecretReference,
     32         EXPIRED_ARTIFACT_QUERY_LIMIT_MAX, PrivateArtifactId, PrivateArtifactMetadata,
     33         PrivateArtifactRevision, PrivateArtifactStore, RetentionPolicy,
     34     },
     35     projection::{
     36         InvalidationReason, ProjectionCheckpoint, ProjectionDocument, ProjectionGeneration,
     37         ProjectionHealth, ProjectionId, ProjectionInvalidation, ProjectionRevision,
     38         ProjectionSnapshot, RawSourceDigest, RebuildStage, RebuildTicket, RebuildTicketId,
     39         RebuildTransition,
     40     },
     41 };
     42 use radroots_transport::{
     43     DeliveryRequest, Target, TargetSet, TransportId,
     44     policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy},
     45     sink::DeliveryPayload,
     46     source::{EventProvenance, ObservedEvent},
     47 };
     48 
     49 fn signed_event() -> SignedEvent {
     50     signed_event_with(1_800_000_100, 0, vec![], "memory-backend")
     51 }
     52 
     53 fn signed_event_with(
     54     created_at: u64,
     55     kind: u32,
     56     tags: Vec<Vec<String>>,
     57     content: &str,
     58 ) -> SignedEvent {
     59     let mut wire = Nip01EventWire {
     60         id: "0".repeat(64),
     61         pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
     62         created_at,
     63         kind,
     64         tags,
     65         content: content.to_owned(),
     66         sig: "42".repeat(64),
     67         extra: Default::default(),
     68     };
     69     wire.id = wire.computed_event_id().expect("event id").to_hex();
     70     let raw = serde_json::json!({
     71         "id": &wire.id,
     72         "pubkey": &wire.pubkey,
     73         "created_at": wire.created_at,
     74         "kind": wire.kind,
     75         "tags": &wire.tags,
     76         "content": &wire.content,
     77         "sig": &wire.sig,
     78     })
     79     .to_string();
     80     SignedEvent::from_wire_verified_id(wire, raw).expect("signed event")
     81 }
     82 
     83 struct Allow;
     84 
     85 impl SignatureVerifier for Allow {
     86     fn verify_signature(&self, _event: &EventEnvelope) -> Result<(), radroots_event::Error> {
     87         Ok(())
     88     }
     89 }
     90 
     91 impl AdmissionPolicy for Allow {
     92     type Error = core::convert::Infallible;
     93 
     94     fn policy_id(&self) -> &'static str {
     95         "test.storage-memory.admission.v1"
     96     }
     97 
     98     fn admit(
     99         &self,
    100         _event: &radroots_event::admission::ContractValidatedEvent,
    101     ) -> Result<(), Self::Error> {
    102         Ok(())
    103     }
    104 }
    105 
    106 impl VisibilityPolicy for Allow {
    107     type Error = core::convert::Infallible;
    108 
    109     fn policy_id(&self) -> &'static str {
    110         "test.storage-memory.visibility.v1"
    111     }
    112 
    113     fn make_visible(
    114         &self,
    115         _event: &radroots_event::admission::AdmittedEvent,
    116     ) -> Result<(), Self::Error> {
    117         Ok(())
    118     }
    119 }
    120 
    121 fn admission(event: SignedEvent, at: u64) -> EventAdmission {
    122     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    123     let provenance = EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at)
    124         .expect("provenance");
    125     EventAdmission::raw(ObservedEvent::new(event, provenance))
    126 }
    127 
    128 fn visible_admission(event: SignedEvent, at: u64) -> EventAdmission {
    129     let verified = RawEvent::new(event.envelope().clone())
    130         .verify_id()
    131         .expect("event id")
    132         .verify_signature(&Allow)
    133         .expect("signature");
    134     let validated = if event.envelope().kind_u32() == 5 {
    135         verified
    136             .validate_contract_for_admission("radroots.social.deletion_request.v1")
    137             .expect("admission-selected contract")
    138     } else {
    139         verified.validate_contract().expect("contract")
    140     };
    141     let visible = validated
    142         .admit_with(&Allow)
    143         .expect("admission")
    144         .make_visible_with(&Allow)
    145         .expect("visibility");
    146     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target");
    147     let provenance = EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at)
    148         .expect("provenance");
    149     EventAdmission::visible(ObservedEvent::new(event, provenance), visible)
    150         .expect("visible admission")
    151 }
    152 
    153 #[test]
    154 fn verified_replacement_changes_visibility_without_increasing_raw_count() {
    155     let store = MemoryStorage::new(SourceGeneration::new([19; 32]).unwrap());
    156     let old = signed_event_with(10, 0, vec![], r#"{"display_name":"Old Farm","bot":false}"#);
    157     let newer = signed_event_with(20, 0, vec![], "malformed profile");
    158     block_on(store.admit(visible_admission(old.clone(), 10))).unwrap();
    159     let raw = admission(newer.clone(), 20);
    160     block_on(store.admit(raw.clone())).unwrap();
    161     let before = block_on(store.rebuild_visibility()).unwrap();
    162     assert_eq!(before.visible_event_ids(), &[*old.id()]);
    163     let verified = RawEvent::new(newer.envelope().clone())
    164         .verify_id()
    165         .unwrap()
    166         .verify_signature(&Allow)
    167         .unwrap();
    168     let advanced = block_on(
    169         store.admit(
    170             EventAdmission::verified(
    171                 ObservedEvent::new(newer.clone(), raw.provenance().clone()),
    172                 verified,
    173             )
    174             .unwrap(),
    175         ),
    176     )
    177     .unwrap();
    178     assert_eq!(
    179         advanced.disposition(),
    180         radroots_storage::event::AdmissionDisposition::Advanced
    181     );
    182     let after = block_on(store.rebuild_visibility()).unwrap();
    183     assert_ne!(before.digest(), after.digest());
    184     assert_eq!(after.current_heads()[0].event_id, *newer.id());
    185     assert!(after.visible_event_ids().is_empty());
    186     assert_eq!(
    187         block_on(EventStore::status(&store)).unwrap().raw_events(),
    188         2
    189     );
    190     assert!(
    191         block_on(store.query_visible(EventQuery::all(EventQueryBounds::first(10).unwrap())))
    192             .unwrap()
    193             .items()
    194             .is_empty()
    195     );
    196 }
    197 
    198 #[test]
    199 fn memory_visibility_rebuild_is_current_delete_aware_and_atomic_parity_safe() {
    200     let generation = SourceGeneration::new([17; 32]).expect("generation");
    201     let direct = MemoryStorage::new(generation);
    202     let atomic_store = MemoryStorage::new(generation);
    203     let old = signed_event_with(
    204         1_800_000_100,
    205         0,
    206         vec![],
    207         r#"{"display_name":"Old Farm","bot":false}"#,
    208     );
    209     let current = signed_event_with(
    210         1_800_000_200,
    211         0,
    212         vec![],
    213         r#"{"display_name":"Current Farm","bot":false}"#,
    214     );
    215     let deletion = signed_event_with(
    216         1_800_000_300,
    217         5,
    218         vec![vec!["e".to_owned(), current.id().to_hex()]],
    219         "retired profile",
    220     );
    221     let admissions = [
    222         visible_admission(old.clone(), 100),
    223         visible_admission(current.clone(), 200),
    224         visible_admission(deletion.clone(), 300),
    225     ];
    226 
    227     for admission in admissions.clone() {
    228         block_on(direct.admit(admission)).expect("direct admission");
    229     }
    230     for (index, admission) in admissions.into_iter().enumerate() {
    231         let identity = u8::try_from(index).expect("commit identity") + 20;
    232         block_on(atomic_store.commit(atomic(
    233             identity,
    234             identity,
    235             AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission, None))),
    236         )))
    237         .expect("atomic admission");
    238     }
    239 
    240     let snapshot = block_on(direct.rebuild_visibility()).expect("visibility rebuild");
    241     let atomic_snapshot =
    242         block_on(atomic_store.rebuild_visibility()).expect("atomic visibility rebuild");
    243     assert_eq!(snapshot, atomic_snapshot);
    244     assert_eq!(snapshot.current_heads()[0].event_id, *current.id());
    245     assert_eq!(snapshot.visible_event_ids(), &[*deletion.id()]);
    246     assert_eq!(snapshot.suppressed_event_ids(), &[*current.id()]);
    247     assert_eq!(snapshot.superseded_event_ids(), &[*old.id()]);
    248     let page = block_on(direct.query_visible(EventQuery::all(
    249         EventQueryBounds::first(10).expect("bounds"),
    250     )))
    251     .expect("visible page");
    252     assert_eq!(page.items().len(), 1);
    253     assert_eq!(page.items()[0].event().id(), deletion.id());
    254     assert_eq!(
    255         block_on(EventStore::status(&direct))
    256             .expect("status")
    257             .visible_events(),
    258         1
    259     );
    260 }
    261 
    262 fn prepare(instance: OperationInstanceId) -> PrepareOperation {
    263     PrepareOperation::new(
    264         instance,
    265         OperationId::SyncPush,
    266         IdempotencyKey::parse("memory-operation").expect("key"),
    267         IdempotencyDigest::new([3; 32]),
    268         100,
    269     )
    270     .expect("prepare")
    271 }
    272 
    273 fn atomic(id: u8, digest: u8, workflow: AtomicWorkflow) -> AtomicCommit {
    274     AtomicCommit::new(
    275         AtomicCommitId::new([id; 16]).expect("commit id"),
    276         AtomicCommitDigest::new([digest; 32]),
    277         200,
    278         workflow,
    279     )
    280     .expect("atomic commit")
    281 }
    282 
    283 fn delivery_request(event: SignedEvent) -> DeliveryRequest {
    284     DeliveryRequest::new(
    285         "memory-delivery",
    286         DeliveryPayload::new(event),
    287         TargetSet::new(vec![
    288             Target::nostr_relay("wss://relay.example").expect("target"),
    289         ])
    290         .expect("target set"),
    291         SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    292         1_000,
    293     )
    294     .expect("delivery request")
    295 }
    296 
    297 fn private_metadata() -> PrivateArtifactMetadata {
    298     PrivateArtifactMetadata::new(
    299         PrivateArtifactId::new([9; 16]).expect("artifact id"),
    300         ArtifactKind::parse("memory.private").expect("kind"),
    301         ArtifactSchemaId::parse("memory.private.v1").expect("schema"),
    302         ArtifactCommitment::new([8; 32]),
    303         64,
    304         DurableSecretReference::new("memory", "caller-owned-key", 1).expect("secret reference"),
    305         RetentionPolicy::new(None, Some(300)).expect("retention"),
    306         100,
    307     )
    308     .expect("metadata")
    309 }
    310 
    311 #[test]
    312 fn memory_event_and_journal_implement_the_canonical_spis() {
    313     let store = MemoryStorage::new(SourceGeneration::new([7; 32]).expect("generation"));
    314     let event = signed_event();
    315     let event_id = *event.id();
    316     block_on(store.admit(admission(event, 10))).expect("admit");
    317     let raw = block_on(store.query_raw(EventQuery::all(
    318         EventQueryBounds::first(10).expect("bounds"),
    319     )))
    320     .expect("query");
    321     assert_eq!(raw.items().len(), 1);
    322     assert_eq!(raw.items()[0].event().id(), &event_id);
    323     assert_eq!(
    324         block_on(EventStore::status(&store))
    325             .expect("status")
    326             .raw_events(),
    327         1
    328     );
    329 
    330     let instance = OperationInstanceId::new([1; 16]).expect("instance");
    331     let prepared = block_on(store.prepare(prepare(instance))).expect("prepare");
    332     let signed = block_on(
    333         store.transition(radroots_storage::journal::JournalTransition::signed(
    334             instance,
    335             prepared.record().revision(),
    336             event_id,
    337         )),
    338     )
    339     .expect("signed");
    340     assert_eq!(signed.state().stage(), JournalStage::Signed);
    341     assert_eq!(
    342         block_on(store.by_idempotency_key(
    343             OperationId::SyncPush,
    344             IdempotencyKey::parse("memory-operation").expect("key"),
    345         ))
    346         .expect("lookup")
    347         .expect("record")
    348         .instance_id(),
    349         instance
    350     );
    351 }
    352 
    353 #[test]
    354 fn memory_atomic_ingest_commits_event_state_and_replays() {
    355     let store = MemoryStorage::default();
    356     let event = signed_event();
    357     let ingest = atomic(
    358         1,
    359         1,
    360         AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(event, 20), None))),
    361     );
    362     let committed = block_on(store.commit(ingest.clone())).expect("atomic ingest");
    363     assert_eq!(committed.disposition(), AtomicCommitDisposition::Committed);
    364     assert_eq!(
    365         block_on(store.commit(ingest))
    366             .expect("atomic replay")
    367             .disposition(),
    368         AtomicCommitDisposition::Replay
    369     );
    370     assert_eq!(
    371         block_on(EventStore::status(&store))
    372             .expect("status")
    373             .raw_events(),
    374         1
    375     );
    376 }
    377 
    378 #[test]
    379 fn memory_outbox_uses_caller_supplied_time_and_lease_identity() {
    380     let store = MemoryStorage::default();
    381     let item = EnqueueOutboxItem::new(
    382         OutboxItemId::new([5; 16]).expect("item id"),
    383         OperationInstanceId::new([5; 16]).expect("instance"),
    384         DeliveryPlanDigest::new([5; 32]),
    385         delivery_request(signed_event()),
    386         100,
    387     )
    388     .expect("enqueue");
    389     assert_eq!(
    390         block_on(store.enqueue(item.clone()))
    391             .expect("enqueue")
    392             .disposition(),
    393         EnqueueDisposition::Created
    394     );
    395     assert_eq!(
    396         block_on(store.enqueue(item)).expect("replay").disposition(),
    397         EnqueueDisposition::Replay
    398     );
    399     let claimed = block_on(
    400         store.claim(
    401             ClaimOutboxItems::new(
    402                 LeaseOwner::parse("memory-worker").expect("owner"),
    403                 LeaseId::new([6; 16]).expect("lease seed"),
    404                 200,
    405                 250,
    406                 1,
    407             )
    408             .expect("claim"),
    409         ),
    410     )
    411     .expect("claim");
    412     assert_eq!(claimed.len(), 1);
    413     assert_eq!(claimed[0].record().stage(), OutboxStage::Leased);
    414     assert_eq!(block_on(Outbox::status(&store)).expect("status").leased, 1);
    415 }
    416 
    417 #[test]
    418 fn memory_projection_rebuild_and_private_metadata_share_deterministic_state() {
    419     let store = MemoryStorage::default();
    420     let projection_id = ProjectionId::parse("memory.test").expect("projection id");
    421     let initial_generation = ProjectionGeneration::new([4; 32]).expect("generation");
    422     let replacement_generation = ProjectionGeneration::new([5; 32]).expect("generation");
    423     let checkpoint =
    424         ProjectionCheckpoint::new(projection_id.clone(), initial_generation, None, 0, 200)
    425             .expect("checkpoint");
    426     assert_eq!(
    427         block_on(store.checkpoint(checkpoint))
    428             .expect("checkpoint")
    429             .health(),
    430         ProjectionHealth::Ready
    431     );
    432     let invalidation = ProjectionInvalidation::new(
    433         projection_id,
    434         initial_generation,
    435         replacement_generation,
    436         InvalidationReason::ProjectionGenerationChanged,
    437         210,
    438     )
    439     .expect("invalidation");
    440     assert_eq!(
    441         block_on(store.invalidate(invalidation.clone()))
    442             .expect("invalidate")
    443             .health(),
    444         ProjectionHealth::Invalidated
    445     );
    446     let ticket_id = RebuildTicketId::new([7; 16]).expect("ticket id");
    447     block_on(
    448         store.request_rebuild(
    449             RebuildTicket::requested(
    450                 ticket_id,
    451                 invalidation,
    452                 store.generation(),
    453                 None,
    454                 RawSourceDigest::new([8; 32]),
    455             )
    456             .expect("ticket"),
    457         ),
    458     )
    459     .expect("request rebuild");
    460     let running = block_on(store.transition_rebuild(RebuildTransition::start(
    461         ticket_id,
    462         ProjectionRevision::INITIAL,
    463         220,
    464     )))
    465     .expect("start rebuild");
    466     assert_eq!(running.stage(), RebuildStage::Running);
    467 
    468     let metadata = private_metadata();
    469     block_on(store.put_metadata(metadata.clone())).expect("put metadata");
    470     assert_eq!(
    471         block_on(store.expired(299, 1)).expect("not expired").len(),
    472         0
    473     );
    474     assert_eq!(block_on(store.expired(300, 1)).expect("expired").len(), 1);
    475     block_on(store.mark_expired(
    476         metadata.artifact_id(),
    477         PrivateArtifactRevision::INITIAL,
    478         300,
    479     ))
    480     .expect("mark expired");
    481     assert_eq!(
    482         block_on(PrivateArtifactStore::status(&store))
    483             .expect("status")
    484             .expired,
    485         1
    486     );
    487 }
    488 
    489 #[test]
    490 fn atomic_projection_failure_leaves_event_and_checkpoint_unchanged() {
    491     let store = MemoryStorage::default();
    492     let projection_id = ProjectionId::parse("memory.atomic").expect("projection id");
    493     let generation = ProjectionGeneration::new([4; 32]).expect("projection generation");
    494     let first = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 2, 200)
    495         .expect("checkpoint");
    496     block_on(store.checkpoint(first)).expect("initial checkpoint");
    497     let regressing = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 1, 201)
    498         .expect("checkpoint");
    499     let request = atomic(
    500         4,
    501         4,
    502         AtomicWorkflow::Ingested(Box::new(CommitIngested::new(
    503             admission(signed_event(), 20),
    504             Some(regressing),
    505         ))),
    506     );
    507     assert_eq!(
    508         block_on(store.commit(request)),
    509         Err(radroots_storage::Error::ProjectionCheckpointRegression)
    510     );
    511     assert_eq!(
    512         block_on(EventStore::status(&store))
    513             .expect("status")
    514             .raw_events(),
    515         0
    516     );
    517     assert_eq!(
    518         block_on(ProjectionStore::status(&store, projection_id))
    519             .expect("projection status")
    520             .expect("projection")
    521             .checkpoint()
    522             .expect("checkpoint")
    523             .projected_rows(),
    524         2
    525     );
    526 }
    527 
    528 #[test]
    529 fn memory_event_journal_and_closed_state_fail_closed() {
    530     let generation = SourceGeneration::new([7; 32]).unwrap();
    531     let store = MemoryStorage::new(generation);
    532     assert_eq!(store.generation(), generation);
    533     let event = signed_event();
    534     let inserted = block_on(store.admit(admission(event.clone(), 10))).unwrap();
    535     assert_eq!(
    536         block_on(store.admit(admission(event.clone(), 10)))
    537             .unwrap()
    538             .disposition(),
    539         radroots_storage::event::AdmissionDisposition::Duplicate
    540     );
    541     let wrong_cursor = EventPosition::new(
    542         SourceGeneration::new([8; 32]).unwrap(),
    543         EventSequence::new(1).unwrap(),
    544     );
    545     let query = EventQuery::all(EventQueryBounds::first(1).unwrap().after(wrong_cursor));
    546     assert_eq!(
    547         block_on(store.query_raw(query)),
    548         Err(Error::SourceGenerationChanged)
    549     );
    550     assert_eq!(
    551         block_on(store.query_provenance(
    552             *event.id(),
    553             EventQueryBounds::first(1).unwrap().after(wrong_cursor),
    554         )),
    555         Err(Error::SourceGenerationChanged)
    556     );
    557     let after = EventQueryBounds::first(1)
    558         .unwrap()
    559         .after(inserted.position());
    560     assert!(
    561         block_on(store.query_provenance(*event.id(), after))
    562             .unwrap()
    563             .items()
    564             .is_empty()
    565     );
    566     assert_eq!(
    567         block_on(store.query_provenance(
    568             radroots_event::EventId::parse("f".repeat(64)).unwrap(),
    569             EventQueryBounds::first(1).unwrap(),
    570         )),
    571         Err(Error::EventNotFound)
    572     );
    573 
    574     let instance = OperationInstanceId::new([1; 16]).unwrap();
    575     let operation = prepare(instance);
    576     block_on(store.prepare(operation.clone())).unwrap();
    577     assert_eq!(
    578         block_on(
    579             store.prepare(
    580                 PrepareOperation::new(
    581                     instance,
    582                     OperationId::SyncPush,
    583                     IdempotencyKey::parse("memory-operation").unwrap(),
    584                     IdempotencyDigest::new([4; 32]),
    585                     100,
    586                 )
    587                 .unwrap()
    588             )
    589         ),
    590         Err(Error::IdempotencyConflict)
    591     );
    592     assert_eq!(
    593         block_on(
    594             store.prepare(
    595                 PrepareOperation::new(
    596                     instance,
    597                     OperationId::SyncPush,
    598                     IdempotencyKey::parse("different-key").unwrap(),
    599                     IdempotencyDigest::new([3; 32]),
    600                     100,
    601                 )
    602                 .unwrap()
    603             )
    604         ),
    605         Err(Error::OperationIdentityMismatch)
    606     );
    607     assert!(
    608         block_on(store.operation(OperationInstanceId::new([2; 16]).unwrap()))
    609             .unwrap()
    610             .is_none()
    611     );
    612     assert!(
    613         block_on(store.by_idempotency_key(
    614             OperationId::SyncPull,
    615             IdempotencyKey::parse("memory-operation").unwrap()
    616         ))
    617         .unwrap()
    618         .is_none()
    619     );
    620     assert_eq!(
    621         block_on(store.recoverable(0)),
    622         Err(Error::InvalidJournalQueryLimit)
    623     );
    624 
    625     block_on(radroots_storage::backup::StorageReliability::close(&store)).unwrap();
    626     assert_eq!(
    627         block_on(EventStore::status(&store)),
    628         Err(Error::BackendUnavailable)
    629     );
    630     assert_eq!(
    631         block_on(store.admit(admission(event, 11))),
    632         Err(Error::BackendUnavailable)
    633     );
    634 }
    635 
    636 #[test]
    637 fn memory_outbox_conflict_and_claim_matrix_is_complete() {
    638     let store = MemoryStorage::default();
    639     let make_item = |item: u8, instance: u8, digest: u8, created: u64| {
    640         EnqueueOutboxItem::new(
    641             OutboxItemId::new([item; 16]).unwrap(),
    642             OperationInstanceId::new([instance; 16]).unwrap(),
    643             DeliveryPlanDigest::new([digest; 32]),
    644             delivery_request(signed_event()),
    645             created,
    646         )
    647         .unwrap()
    648     };
    649     block_on(store.enqueue(make_item(1, 1, 1, 100))).unwrap();
    650     assert_eq!(
    651         block_on(store.enqueue(make_item(1, 1, 2, 100))),
    652         Err(Error::OutboxPlanConflict)
    653     );
    654     assert_eq!(
    655         block_on(store.enqueue(make_item(2, 1, 1, 100))),
    656         Err(Error::OutboxPlanConflict)
    657     );
    658     assert!(
    659         block_on(store.item(OutboxItemId::new([9; 16]).unwrap()))
    660             .unwrap()
    661             .is_none()
    662     );
    663     let first_claim = ClaimOutboxItems::new(
    664         LeaseOwner::parse("worker").unwrap(),
    665         LeaseId::new([3; 16]).unwrap(),
    666         200,
    667         300,
    668         1,
    669     )
    670     .unwrap();
    671     assert_eq!(block_on(store.claim(first_claim)).unwrap().len(), 1);
    672     let concurrent = ClaimOutboxItems::new(
    673         LeaseOwner::parse("worker-two").unwrap(),
    674         LeaseId::new([4; 16]).unwrap(),
    675         250,
    676         350,
    677         1,
    678     )
    679     .unwrap();
    680     assert!(block_on(store.claim(concurrent)).unwrap().is_empty());
    681     assert_eq!(
    682         block_on(store.release(
    683             OutboxItemId::new([9; 16]).unwrap(),
    684             LeaseId::new([4; 16]).unwrap(),
    685             radroots_storage::outbox::OutboxRevision::INITIAL,
    686             260,
    687             None,
    688         )),
    689         Err(Error::OutboxItemNotFound)
    690     );
    691 }
    692 
    693 #[test]
    694 fn memory_projection_and_private_artifact_conflict_matrix_is_complete() {
    695     let store = MemoryStorage::default();
    696     let projection_id = ProjectionId::parse("memory.matrix").unwrap();
    697     let initial = ProjectionGeneration::new([4; 32]).unwrap();
    698     let replacement = ProjectionGeneration::new([5; 32]).unwrap();
    699     let initial_checkpoint =
    700         ProjectionCheckpoint::new(projection_id.clone(), initial, None, 1, 100).unwrap();
    701     block_on(store.checkpoint(initial_checkpoint.clone())).unwrap();
    702     assert_eq!(
    703         block_on(store.checkpoint(
    704             ProjectionCheckpoint::new(projection_id.clone(), replacement, None, 2, 101).unwrap()
    705         )),
    706         Err(Error::ProjectionCheckpointMismatch)
    707     );
    708     assert_eq!(
    709         block_on(store.checkpoint(
    710             ProjectionCheckpoint::new(projection_id.clone(), initial, None, 0, 101).unwrap()
    711         )),
    712         Err(Error::ProjectionCheckpointRegression)
    713     );
    714     let missing = ProjectionInvalidation::new(
    715         ProjectionId::parse("missing").unwrap(),
    716         initial,
    717         replacement,
    718         InvalidationReason::OperatorRequested,
    719         110,
    720     )
    721     .unwrap();
    722     assert_eq!(
    723         block_on(store.invalidate(missing)),
    724         Err(Error::ProjectionCheckpointMismatch)
    725     );
    726     let wrong = ProjectionInvalidation::new(
    727         projection_id.clone(),
    728         replacement,
    729         ProjectionGeneration::new([6; 32]).unwrap(),
    730         InvalidationReason::OperatorRequested,
    731         110,
    732     )
    733     .unwrap();
    734     assert_eq!(
    735         block_on(store.invalidate(wrong)),
    736         Err(Error::ProjectionCheckpointMismatch)
    737     );
    738     let invalidation = ProjectionInvalidation::new(
    739         projection_id.clone(),
    740         initial,
    741         replacement,
    742         InvalidationReason::OperatorRequested,
    743         110,
    744     )
    745     .unwrap();
    746     block_on(store.invalidate(invalidation.clone())).unwrap();
    747     assert_eq!(
    748         block_on(store.invalidate(invalidation.clone()))
    749             .unwrap()
    750             .health(),
    751         ProjectionHealth::Invalidated
    752     );
    753     let conflicting_invalidation = ProjectionInvalidation::new(
    754         projection_id.clone(),
    755         initial,
    756         ProjectionGeneration::new([6; 32]).unwrap(),
    757         InvalidationReason::ProjectionGenerationChanged,
    758         111,
    759     )
    760     .unwrap();
    761     assert_eq!(
    762         block_on(store.invalidate(conflicting_invalidation)),
    763         Err(Error::ProjectionRevisionConflict)
    764     );
    765     assert!(
    766         block_on(store.invalidation(projection_id.clone(), replacement))
    767             .unwrap()
    768             .is_some()
    769     );
    770     assert!(
    771         block_on(store.invalidation(
    772             projection_id.clone(),
    773             ProjectionGeneration::new([9; 32]).unwrap(),
    774         ))
    775         .unwrap()
    776         .is_none()
    777     );
    778     let ticket = RebuildTicket::requested(
    779         RebuildTicketId::new([7; 16]).unwrap(),
    780         invalidation.clone(),
    781         store.generation(),
    782         None,
    783         RawSourceDigest::new([8; 32]),
    784     )
    785     .unwrap();
    786     assert_eq!(
    787         block_on(store.request_rebuild(ticket.clone())).unwrap(),
    788         ticket
    789     );
    790     let conflicting_ticket = RebuildTicket::requested(
    791         ticket.ticket_id(),
    792         invalidation,
    793         store.generation(),
    794         None,
    795         RawSourceDigest::new([9; 32]),
    796     )
    797     .unwrap();
    798     assert_eq!(
    799         block_on(store.request_rebuild(conflicting_ticket)),
    800         Err(Error::ProjectionRevisionConflict)
    801     );
    802     assert_eq!(
    803         block_on(store.request_rebuild(ticket.clone())).unwrap(),
    804         ticket
    805     );
    806     assert!(
    807         block_on(store.rebuild(ticket.ticket_id()))
    808             .unwrap()
    809             .is_some()
    810     );
    811 
    812     let metadata = private_metadata();
    813     assert_eq!(
    814         block_on(store.put_metadata(metadata.clone())).unwrap(),
    815         metadata
    816     );
    817     assert_eq!(
    818         block_on(store.put_metadata(metadata.clone())).unwrap(),
    819         metadata
    820     );
    821     let conflict = PrivateArtifactMetadata::new(
    822         metadata.artifact_id(),
    823         ArtifactKind::parse("memory.private").unwrap(),
    824         ArtifactSchemaId::parse("memory.private.v1").unwrap(),
    825         ArtifactCommitment::new([9; 32]),
    826         64,
    827         DurableSecretReference::new("memory", "caller-owned-key", 1).unwrap(),
    828         RetentionPolicy::new(None, Some(300)).unwrap(),
    829         100,
    830     )
    831     .unwrap();
    832     assert_eq!(
    833         block_on(store.put_metadata(conflict)),
    834         Err(Error::PrivateArtifactConflict)
    835     );
    836     assert!(
    837         block_on(store.metadata(PrivateArtifactId::new([8; 16]).unwrap()))
    838             .unwrap()
    839             .is_none()
    840     );
    841     assert_eq!(
    842         block_on(store.expired(0, 1)),
    843         Err(Error::InvalidExpiredArtifactQueryLimit)
    844     );
    845     assert_eq!(
    846         block_on(store.expired(1, 0)),
    847         Err(Error::InvalidExpiredArtifactQueryLimit)
    848     );
    849 }
    850 
    851 #[test]
    852 fn memory_branch_guards_distinguish_each_identity_and_bound() {
    853     let store = MemoryStorage::new(SourceGeneration::new([7; 32]).unwrap());
    854     let event = signed_event();
    855     block_on(store.admit(admission(event.clone(), 10))).unwrap();
    856     let other_id = radroots_event::EventId::parse("f".repeat(64)).unwrap();
    857     let selected = block_on(store.query_raw(
    858         EventQuery::for_ids(EventQueryBounds::first(10).unwrap(), vec![other_id]).unwrap(),
    859     ))
    860     .unwrap();
    861     assert!(selected.items().is_empty());
    862 
    863     let instance = OperationInstanceId::new([1; 16]).unwrap();
    864     block_on(store.prepare(prepare(instance))).unwrap();
    865     let different_operation = PrepareOperation::new(
    866         OperationInstanceId::new([2; 16]).unwrap(),
    867         OperationId::SyncPull,
    868         IdempotencyKey::parse("memory-operation").unwrap(),
    869         IdempotencyDigest::new([3; 32]),
    870         100,
    871     )
    872     .unwrap();
    873     assert_eq!(
    874         block_on(store.prepare(different_operation)),
    875         Err(Error::IdempotencyConflict)
    876     );
    877     let different_instance = PrepareOperation::new(
    878         OperationInstanceId::new([2; 16]).unwrap(),
    879         OperationId::SyncPush,
    880         IdempotencyKey::parse("memory-operation").unwrap(),
    881         IdempotencyDigest::new([3; 32]),
    882         100,
    883     )
    884     .unwrap();
    885     assert_eq!(
    886         block_on(store.prepare(different_instance)),
    887         Err(Error::IdempotencyConflict)
    888     );
    889     assert_eq!(
    890         block_on(store.recoverable(RECOVERABLE_QUERY_LIMIT_MAX + 1)),
    891         Err(Error::InvalidJournalQueryLimit)
    892     );
    893 
    894     let private = MemoryStorage::default();
    895     block_on(private.put_metadata(private_metadata())).unwrap();
    896     assert_eq!(
    897         block_on(private.expired(1, EXPIRED_ARTIFACT_QUERY_LIMIT_MAX + 1)),
    898         Err(Error::InvalidExpiredArtifactQueryLimit)
    899     );
    900     block_on(private.mark_expired(
    901         PrivateArtifactId::new([9; 16]).unwrap(),
    902         PrivateArtifactRevision::INITIAL,
    903         300,
    904     ))
    905     .unwrap();
    906     assert!(block_on(private.expired(400, 1)).unwrap().is_empty());
    907 }
    908 
    909 #[test]
    910 fn memory_outbox_replay_checks_each_durable_plan_field_and_retry_window() {
    911     let make_request = |id: &str| {
    912         DeliveryRequest::new(
    913             id,
    914             DeliveryPayload::new(signed_event()),
    915             TargetSet::new(vec![Target::nostr_relay("wss://relay.example").unwrap()]).unwrap(),
    916             SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()),
    917             1_000,
    918         )
    919         .unwrap()
    920     };
    921     let make_item = |instance: u8, digest: u8, request: DeliveryRequest, created: u64| {
    922         EnqueueOutboxItem::new(
    923             OutboxItemId::new([1; 16]).unwrap(),
    924             OperationInstanceId::new([instance; 16]).unwrap(),
    925             DeliveryPlanDigest::new([digest; 32]),
    926             request,
    927             created,
    928         )
    929         .unwrap()
    930     };
    931 
    932     for conflict in [
    933         make_item(2, 1, make_request("memory-delivery"), 100),
    934         make_item(1, 2, make_request("memory-delivery"), 100),
    935         make_item(1, 1, make_request("different-delivery"), 100),
    936         make_item(1, 1, make_request("memory-delivery"), 101),
    937     ] {
    938         let store = MemoryStorage::default();
    939         block_on(store.enqueue(make_item(1, 1, make_request("memory-delivery"), 100))).unwrap();
    940         assert_eq!(
    941             block_on(store.enqueue(conflict)),
    942             Err(Error::OutboxPlanConflict)
    943         );
    944     }
    945 
    946     let store = MemoryStorage::default();
    947     block_on(store.enqueue(make_item(1, 1, make_request("memory-delivery"), 100))).unwrap();
    948     let lease_seed = LeaseId::new([3; 16]).unwrap();
    949     let claimed = block_on(
    950         store.claim(
    951             ClaimOutboxItems::new(
    952                 LeaseOwner::parse("worker").unwrap(),
    953                 lease_seed,
    954                 200,
    955                 250,
    956                 1,
    957             )
    958             .unwrap(),
    959         ),
    960     )
    961     .unwrap();
    962     let record = &claimed[0];
    963     block_on(store.release(
    964         record.record().item_id(),
    965         record.lease().id(),
    966         record.record().revision(),
    967         210,
    968         Some(300),
    969     ))
    970     .unwrap();
    971     let early = ClaimOutboxItems::new(
    972         LeaseOwner::parse("worker").unwrap(),
    973         LeaseId::new([4; 16]).unwrap(),
    974         299,
    975         350,
    976         1,
    977     )
    978     .unwrap();
    979     assert!(block_on(store.claim(early)).unwrap().is_empty());
    980 }
    981 
    982 #[test]
    983 fn memory_atomic_replay_rejects_a_changed_digest_without_mutation() {
    984     let store = MemoryStorage::default();
    985     let event = signed_event();
    986     let committed = atomic(
    987         1,
    988         1,
    989         AtomicWorkflow::Ingested(Box::new(CommitIngested::new(
    990             admission(event.clone(), 20),
    991             None,
    992         ))),
    993     );
    994     block_on(store.commit(committed)).unwrap();
    995     let conflict = atomic(
    996         1,
    997         2,
    998         AtomicWorkflow::Ingested(Box::new(CommitIngested::new(admission(event, 20), None))),
    999     );
   1000     assert_eq!(
   1001         block_on(store.commit(conflict)),
   1002         Err(Error::AtomicCommitConflict)
   1003     );
   1004     assert_eq!(
   1005         block_on(EventStore::status(&store)).unwrap().raw_events(),
   1006         1
   1007     );
   1008 }
   1009 
   1010 #[test]
   1011 fn memory_materialized_documents_replace_and_snapshots_remain_immutable() {
   1012     let store = MemoryStorage::default();
   1013     let projection_id = ProjectionId::parse("memory.today").unwrap();
   1014     let generation = ProjectionGeneration::new([21; 32]).unwrap();
   1015     block_on(store.put_projection_document(
   1016         projection_id.clone(),
   1017         generation,
   1018         ProjectionDocument::new("context.one".into(), vec![1]).unwrap(),
   1019     ))
   1020     .unwrap();
   1021     block_on(store.put_projection_document(
   1022         projection_id.clone(),
   1023         generation,
   1024         ProjectionDocument::new("context.one".into(), vec![2]).unwrap(),
   1025     ))
   1026     .unwrap();
   1027     assert_eq!(
   1028         block_on(store.projection_document(
   1029             projection_id.clone(),
   1030             generation,
   1031             "context.one".into(),
   1032         ))
   1033         .unwrap()
   1034         .unwrap()
   1035         .value(),
   1036         [2]
   1037     );
   1038     for (candidate_id, candidate_generation, candidate_key) in [
   1039         (
   1040             ProjectionId::parse("memory.other").unwrap(),
   1041             generation,
   1042             "context.one".to_owned(),
   1043         ),
   1044         (
   1045             projection_id.clone(),
   1046             ProjectionGeneration::new([23; 32]).unwrap(),
   1047             "context.one".to_owned(),
   1048         ),
   1049         (
   1050             projection_id.clone(),
   1051             generation,
   1052             "context.other".to_owned(),
   1053         ),
   1054     ] {
   1055         assert!(
   1056             block_on(store.projection_document(candidate_id, candidate_generation, candidate_key,))
   1057                 .unwrap()
   1058                 .is_none()
   1059         );
   1060     }
   1061     let snapshot =
   1062         ProjectionSnapshot::new(projection_id.clone(), [22; 32], generation, 100, vec![3]).unwrap();
   1063     block_on(store.put_projection_snapshot(snapshot.clone())).unwrap();
   1064     block_on(store.put_projection_snapshot(snapshot.clone())).unwrap();
   1065     assert_eq!(
   1066         block_on(store.projection_snapshot(projection_id.clone(), [22; 32])).unwrap(),
   1067         Some(snapshot)
   1068     );
   1069     assert_eq!(
   1070         block_on(store.put_projection_snapshot(
   1071             ProjectionSnapshot::new(projection_id, [22; 32], generation, 100, vec![4]).unwrap()
   1072         )),
   1073         Err(Error::CorruptProjectionDocument)
   1074     );
   1075     assert!(
   1076         block_on(
   1077             store.projection_snapshot(ProjectionId::parse("memory.other").unwrap(), [22; 32],)
   1078         )
   1079         .unwrap()
   1080         .is_none()
   1081     );
   1082     assert!(
   1083         block_on(
   1084             store.projection_snapshot(ProjectionId::parse("memory.today").unwrap(), [23; 32],)
   1085         )
   1086         .unwrap()
   1087         .is_none()
   1088     );
   1089 }