lib

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

event_store.rs (20552B)


      1 use futures_executor::block_on;
      2 use radroots_event::{
      3     EventId, SignedEvent, VerifiedEvent,
      4     admission::{AdmissionPolicy, RawEvent, VisibilityPolicy, VisibleEvent},
      5     wire::Nip01EventWire,
      6 };
      7 use radroots_storage::{
      8     Error, EventStore,
      9     event::{
     10         AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage,
     11         EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration,
     12         StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent,
     13         VisibilityInput, VisibilitySnapshot, evaluate_visibility,
     14     },
     15     status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
     16 };
     17 use radroots_transport::{
     18     BoxFuture, Target, TransportId,
     19     source::{EventProvenance, ObservedEvent},
     20 };
     21 use std::sync::Mutex;
     22 
     23 #[derive(Clone)]
     24 struct Entry {
     25     position: EventPosition,
     26     admission: EventAdmission,
     27     provenance: Vec<EventProvenance>,
     28 }
     29 
     30 struct MemoryEventStore {
     31     generation: SourceGeneration,
     32     entries: Mutex<Vec<Entry>>,
     33 }
     34 
     35 impl MemoryEventStore {
     36     fn new() -> Self {
     37         Self {
     38             generation: SourceGeneration::new([7; 32]).expect("non-zero generation"),
     39             entries: Mutex::new(Vec::new()),
     40         }
     41     }
     42 
     43     fn selected(&self, query: &EventQuery) -> Result<Vec<Entry>, Error> {
     44         if let Some(cursor) = query.bounds().cursor()
     45             && cursor.generation() != self.generation
     46         {
     47             return Err(Error::SourceGenerationChanged);
     48         }
     49         let after = query
     50             .bounds()
     51             .cursor()
     52             .map_or(0, |cursor| cursor.sequence().get());
     53         Ok(self
     54             .entries
     55             .lock()
     56             .expect("test store lock")
     57             .iter()
     58             .filter(|entry| {
     59                 entry.position.sequence().get() > after && query.selects(entry.admission.event_id())
     60             })
     61             .take(usize::from(query.bounds().limit()))
     62             .cloned()
     63             .collect())
     64     }
     65 }
     66 
     67 impl EventStore for MemoryEventStore {
     68     fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> {
     69         Box::pin(async move {
     70             let entries = self.entries.lock().expect("test store lock");
     71             let raw = entries.len() as u64;
     72             let verified = entries
     73                 .iter()
     74                 .filter(|entry| entry.admission.stage() >= AdmissionStage::Verified)
     75                 .count() as u64;
     76             let visible = entries
     77                 .iter()
     78                 .filter(|entry| entry.admission.stage() == AdmissionStage::Visible)
     79                 .count() as u64;
     80             EventStoreStatus::new(
     81                 self.generation,
     82                 EventStoreMode::ReadWrite,
     83                 EventStoreHealth::Available,
     84                 raw,
     85                 verified,
     86                 visible,
     87             )
     88         })
     89     }
     90 
     91     fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> {
     92         Box::pin(async move {
     93             let mut entries = self.entries.lock().expect("test store lock");
     94             if let Some(entry) = entries
     95                 .iter_mut()
     96                 .find(|entry| entry.admission.event_id() == admission.event_id())
     97             {
     98                 if entry.admission.event() != admission.event() {
     99                     return Err(Error::EventConflict);
    100                 }
    101                 if admission.stage() < entry.admission.stage() {
    102                     return Err(Error::AdmissionRegression);
    103                 }
    104                 let disposition = if admission.stage() == entry.admission.stage() {
    105                     AdmissionDisposition::Duplicate
    106                 } else {
    107                     AdmissionDisposition::Advanced
    108                 };
    109                 if !entry.provenance.contains(admission.provenance()) {
    110                     entry.provenance.push(admission.provenance().clone());
    111                 }
    112                 entry.admission = admission;
    113                 return Ok(AdmissionReceipt::new(
    114                     *entry.admission.event_id(),
    115                     entry.position,
    116                     entry.admission.stage(),
    117                     disposition,
    118                 ));
    119             }
    120 
    121             let sequence = EventSequence::new(entries.len() as u64 + 1)?;
    122             let position = EventPosition::new(self.generation, sequence);
    123             let receipt = AdmissionReceipt::new(
    124                 *admission.event_id(),
    125                 position,
    126                 admission.stage(),
    127                 AdmissionDisposition::Inserted,
    128             );
    129             let provenance = vec![admission.provenance().clone()];
    130             entries.push(Entry {
    131                 position,
    132                 admission,
    133                 provenance,
    134             });
    135             Ok(receipt)
    136         })
    137     }
    138 
    139     fn query_raw(
    140         &self,
    141         query: EventQuery,
    142     ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> {
    143         Box::pin(async move {
    144             let items = self
    145                 .selected(&query)?
    146                 .into_iter()
    147                 .map(|entry| {
    148                     StoredRawEvent::new(
    149                         entry.position,
    150                         entry.admission.event().clone(),
    151                         entry.admission.stage(),
    152                     )
    153                 })
    154                 .collect();
    155             EventPage::new(self.generation, items, None, query.bounds())
    156         })
    157     }
    158 
    159     fn query_verified(
    160         &self,
    161         query: EventQuery,
    162     ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> {
    163         Box::pin(async move {
    164             let items = self
    165                 .selected(&query)?
    166                 .into_iter()
    167                 .filter_map(|entry| {
    168                     (entry.admission.stage() >= AdmissionStage::Verified).then(|| {
    169                         StoredVerifiedEvent::new(entry.position, entry.admission.event().clone())
    170                     })
    171                 })
    172                 .collect();
    173             EventPage::new(self.generation, items, None, query.bounds())
    174         })
    175     }
    176 
    177     fn query_visible(
    178         &self,
    179         query: EventQuery,
    180     ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> {
    181         Box::pin(async move {
    182             let items = self
    183                 .selected(&query)?
    184                 .into_iter()
    185                 .filter_map(|entry| {
    186                     (entry.admission.stage() == AdmissionStage::Visible).then(|| {
    187                         StoredVisibleEvent::new(entry.position, entry.admission.event().clone())
    188                     })
    189                 })
    190                 .collect();
    191             EventPage::new(self.generation, items, None, query.bounds())
    192         })
    193     }
    194 
    195     fn rebuild_visibility(&self) -> BoxFuture<'_, Result<VisibilitySnapshot, Error>> {
    196         Box::pin(async move {
    197             let entries = self.entries.lock().expect("test store lock");
    198             evaluate_visibility(
    199                 self.generation,
    200                 entries.iter().map(|entry| {
    201                     VisibilityInput::new(
    202                         entry.position,
    203                         entry.admission.event(),
    204                         entry.admission.stage(),
    205                     )
    206                 }),
    207             )
    208             .map(|evaluation| evaluation.into_snapshot())
    209         })
    210     }
    211 
    212     fn query_provenance(
    213         &self,
    214         event_id: EventId,
    215         bounds: EventQueryBounds,
    216     ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> {
    217         Box::pin(async move {
    218             let entries = self.entries.lock().expect("test store lock");
    219             let entry = entries
    220                 .iter()
    221                 .find(|entry| entry.admission.event_id() == &event_id)
    222                 .ok_or(Error::EventNotFound)?;
    223             let items = entry
    224                 .provenance
    225                 .iter()
    226                 .take(usize::from(bounds.limit()))
    227                 .cloned()
    228                 .map(|provenance| StoredEventProvenance::new(entry.position, provenance))
    229                 .collect();
    230             EventPage::new(self.generation, items, None, bounds)
    231         })
    232     }
    233 }
    234 
    235 struct Allow;
    236 
    237 impl radroots_event::admission::SignatureVerifier for Allow {
    238     fn verify_signature(
    239         &self,
    240         _event: &radroots_event::Event,
    241     ) -> Result<(), radroots_event::Error> {
    242         Ok(())
    243     }
    244 }
    245 
    246 impl AdmissionPolicy for Allow {
    247     type Error = core::convert::Infallible;
    248 
    249     fn policy_id(&self) -> &'static str {
    250         "test.storage.admission.v1"
    251     }
    252 
    253     fn admit(
    254         &self,
    255         _event: &radroots_event::admission::ContractValidatedEvent,
    256     ) -> Result<(), Self::Error> {
    257         Ok(())
    258     }
    259 }
    260 
    261 impl VisibilityPolicy for Allow {
    262     type Error = core::convert::Infallible;
    263 
    264     fn policy_id(&self) -> &'static str {
    265         "test.storage.visibility.v1"
    266     }
    267 
    268     fn make_visible(
    269         &self,
    270         _event: &radroots_event::admission::AdmittedEvent,
    271     ) -> Result<(), Self::Error> {
    272         Ok(())
    273     }
    274 }
    275 
    276 fn signed_event_with_signature(signature_byte: &str) -> SignedEvent {
    277     let mut wire = Nip01EventWire {
    278         id: "0".repeat(64),
    279         pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(),
    280         created_at: 1_800_000_100,
    281         kind: 0,
    282         tags: vec![],
    283         content: "{\"display_name\":\"Moss Street Farm\",\"bot\":false}".to_owned(),
    284         sig: signature_byte.repeat(64),
    285         extra: Default::default(),
    286     };
    287     wire.id = wire
    288         .computed_event_id()
    289         .expect("canonical event id")
    290         .to_hex();
    291     let raw_json = serde_json::json!({
    292         "id": &wire.id,
    293         "pubkey": &wire.pubkey,
    294         "created_at": wire.created_at,
    295         "kind": wire.kind,
    296         "tags": &wire.tags,
    297         "content": &wire.content,
    298         "sig": &wire.sig,
    299     })
    300     .to_string();
    301     SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event")
    302 }
    303 
    304 fn signed_event() -> SignedEvent {
    305     signed_event_with_signature("42")
    306 }
    307 
    308 fn observed(event: SignedEvent, observed_at: u64) -> ObservedEvent {
    309     let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("relay target");
    310     let provenance = EventProvenance::new(
    311         TransportId::NOSTR,
    312         target.fingerprint().clone(),
    313         observed_at,
    314     )
    315     .expect("provenance");
    316     ObservedEvent::new(event, provenance)
    317 }
    318 
    319 fn verified(event: &SignedEvent) -> VerifiedEvent {
    320     RawEvent::new(event.envelope().clone())
    321         .verify_id()
    322         .expect("event id")
    323         .verify_signature(&Allow)
    324         .expect("signature")
    325 }
    326 
    327 fn visible(event: &SignedEvent) -> VisibleEvent {
    328     verified(event)
    329         .validate_contract()
    330         .expect("contract")
    331         .admit_with(&Allow)
    332         .expect("admission")
    333         .make_visible_with(&Allow)
    334         .expect("visibility")
    335 }
    336 
    337 #[test]
    338 fn event_store_is_dyn_compatible_and_enforces_monotonic_admission() {
    339     let store = MemoryEventStore::new();
    340     let dynamic: &dyn EventStore = &store;
    341     let event = signed_event();
    342 
    343     let inserted = block_on(dynamic.admit(EventAdmission::raw(observed(event.clone(), 1))))
    344         .expect("raw insertion");
    345     assert_eq!(inserted.disposition(), AdmissionDisposition::Inserted);
    346     assert_eq!(inserted.stage(), AdmissionStage::Raw);
    347 
    348     let advanced = block_on(
    349         dynamic.admit(
    350             EventAdmission::verified(observed(event.clone(), 2), verified(&event))
    351                 .expect("verified admission"),
    352         ),
    353     )
    354     .expect("verified advancement");
    355     assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced);
    356     assert_eq!(advanced.position(), inserted.position());
    357 
    358     let visible_event = visible(&event);
    359     let advanced = block_on(
    360         dynamic.admit(
    361             EventAdmission::visible(observed(event.clone(), 3), visible_event)
    362                 .expect("visible admission"),
    363         ),
    364     )
    365     .expect("visible advancement");
    366     assert_eq!(advanced.stage(), AdmissionStage::Visible);
    367 
    368     let duplicate = block_on(
    369         dynamic.admit(
    370             EventAdmission::visible(observed(event.clone(), 3), visible(&event))
    371                 .expect("duplicate visible admission"),
    372         ),
    373     )
    374     .expect("idempotent duplicate");
    375     assert_eq!(duplicate.disposition(), AdmissionDisposition::Duplicate);
    376 
    377     assert_eq!(
    378         block_on(dynamic.admit(EventAdmission::raw(observed(event, 4)))),
    379         Err(Error::AdmissionRegression)
    380     );
    381 
    382     let conflicting = signed_event_with_signature("43");
    383     assert_eq!(
    384         EventAdmission::verified(observed(conflicting.clone(), 5), verified(&signed_event())),
    385         Err(Error::AdmissionEventMismatch)
    386     );
    387     assert_eq!(
    388         block_on(dynamic.admit(EventAdmission::raw(observed(conflicting, 6)))),
    389         Err(Error::EventConflict)
    390     );
    391 }
    392 
    393 #[test]
    394 fn event_queries_preserve_stage_generation_bounds_and_provenance() {
    395     let store = MemoryEventStore::new();
    396     let event = signed_event();
    397     let event_id = *event.id();
    398     block_on(
    399         store.admit(
    400             EventAdmission::visible(observed(event.clone(), 5), visible(&event))
    401                 .expect("visible admission"),
    402         ),
    403     )
    404     .expect("insert visible event");
    405 
    406     let bounds = EventQueryBounds::first(1).expect("query bounds");
    407     let query = EventQuery::for_ids(bounds, vec![event_id]).expect("id query");
    408     let raw = block_on(store.query_raw(query.clone())).expect("raw page");
    409     let verified = block_on(store.query_verified(query.clone())).expect("verified page");
    410     let visible = block_on(store.query_visible(query)).expect("visible page");
    411     let provenance = block_on(store.query_provenance(event_id, bounds)).expect("provenance page");
    412 
    413     assert_eq!(raw.items().len(), 1);
    414     assert_eq!(raw.items()[0].stage(), AdmissionStage::Visible);
    415     assert_eq!(verified.items().len(), 1);
    416     assert_eq!(visible.items().len(), 1);
    417     assert_eq!(provenance.items().len(), 1);
    418     assert_eq!(raw.generation(), store.generation);
    419 
    420     let status = block_on(store.status()).expect("status");
    421     assert_eq!(status.raw_events(), 1);
    422     assert_eq!(status.verified_events(), 1);
    423     assert_eq!(status.visible_events(), 1);
    424 }
    425 
    426 #[test]
    427 fn bounds_generations_and_status_reject_invalid_state() {
    428     assert_eq!(
    429         SourceGeneration::new([0; 32]),
    430         Err(Error::InvalidSourceGeneration)
    431     );
    432     assert_eq!(EventSequence::new(0), Err(Error::InvalidEventSequence));
    433     assert_eq!(
    434         EventQueryBounds::first(0),
    435         Err(Error::InvalidEventQueryLimit)
    436     );
    437     assert_eq!(
    438         EventStoreStatus::new(
    439             SourceGeneration::new([1; 32]).expect("generation"),
    440             EventStoreMode::ReadOnly,
    441             EventStoreHealth::Degraded,
    442             1,
    443             2,
    444             0,
    445         ),
    446         Err(Error::CorruptStoredEvent)
    447     );
    448 }
    449 
    450 #[test]
    451 fn event_value_models_cover_bounds_accessors_and_durable_reconstruction() {
    452     let generation = SourceGeneration::new([1; 32]).expect("generation");
    453     let other_generation = SourceGeneration::new([2; 32]).expect("other generation");
    454     let sequence = EventSequence::new(1).expect("sequence");
    455     let position = EventPosition::new(generation, sequence);
    456     assert_eq!(generation.as_bytes(), &[1; 32]);
    457     assert_eq!(sequence.get(), 1);
    458     assert_eq!(position.generation(), generation);
    459     assert_eq!(position.sequence(), sequence);
    460 
    461     assert_eq!(
    462         EventQueryBounds::first(radroots_storage::event::EVENT_QUERY_LIMIT_MAX + 1),
    463         Err(Error::InvalidEventQueryLimit)
    464     );
    465     let bounds = EventQueryBounds::first(1).expect("bounds").after(position);
    466     assert_eq!(bounds.limit(), 1);
    467     assert_eq!(bounds.cursor(), Some(position));
    468     let event = signed_event();
    469     let event_id = *event.id();
    470     assert_eq!(
    471         EventQuery::for_ids(bounds, Vec::new()),
    472         Err(Error::EmptyEventQueryIds)
    473     );
    474     assert_eq!(
    475         EventQuery::for_ids(bounds, vec![event_id, event_id]),
    476         Err(Error::DuplicateEventQueryId)
    477     );
    478     assert_eq!(
    479         EventQuery::for_ids(
    480             bounds,
    481             vec![event_id; radroots_storage::event::EVENT_QUERY_ID_MAX + 1]
    482         ),
    483         Err(Error::TooManyEventQueryIds)
    484     );
    485     let all = EventQuery::all(bounds);
    486     assert!(all.event_ids().is_empty());
    487     assert!(all.selects(&event_id));
    488     let selected = EventQuery::for_ids(bounds, vec![event_id]).expect("selected query");
    489     assert_eq!(selected.bounds(), bounds);
    490     assert_eq!(selected.event_ids(), &[event_id]);
    491     assert!(selected.selects(&event_id));
    492     let other_event_id = EventId::parse("f".repeat(64)).expect("other event id");
    493     assert!(!selected.selects(&other_event_id));
    494 
    495     let raw_admission = EventAdmission::raw(observed(event.clone(), 1));
    496     assert_eq!(raw_admission.stage(), AdmissionStage::Raw);
    497     assert_eq!(raw_admission.event(), &event);
    498     assert_eq!(raw_admission.event_id(), &event_id);
    499     assert_eq!(raw_admission.provenance().observed_at_unix_ms(), 1);
    500     assert!(raw_admission.verified_event().is_none());
    501     assert!(raw_admission.visible_event().is_none());
    502     let verified_admission =
    503         EventAdmission::verified(observed(event.clone(), 2), verified(&event)).expect("verified");
    504     assert!(verified_admission.verified_event().is_some());
    505     assert!(verified_admission.visible_event().is_none());
    506     let visible_admission =
    507         EventAdmission::visible(observed(event.clone(), 3), visible(&event)).expect("visible");
    508     assert!(visible_admission.verified_event().is_some());
    509     assert!(visible_admission.visible_event().is_some());
    510     assert_eq!(
    511         EventAdmission::visible(
    512             observed(signed_event_with_signature("43"), 4),
    513             visible(&event)
    514         ),
    515         Err(Error::AdmissionEventMismatch)
    516     );
    517 
    518     let receipt = AdmissionReceipt::new(
    519         event_id,
    520         position,
    521         AdmissionStage::Raw,
    522         AdmissionDisposition::Inserted,
    523     );
    524     assert_eq!(receipt.event_id(), &event_id);
    525     assert_eq!(receipt.position(), position);
    526     assert_eq!(receipt.stage(), AdmissionStage::Raw);
    527     assert_eq!(receipt.disposition(), AdmissionDisposition::Inserted);
    528     let stored_raw = StoredRawEvent::new(position, event.clone(), AdmissionStage::Raw);
    529     assert_eq!(stored_raw.position(), position);
    530     assert_eq!(stored_raw.event(), &event);
    531     assert_eq!(stored_raw.stage(), AdmissionStage::Raw);
    532     let stored_verified = StoredVerifiedEvent::new(position, event.clone());
    533     assert_eq!(stored_verified.position(), position);
    534     assert_eq!(stored_verified.event(), &event);
    535     let stored_visible = StoredVisibleEvent::new(position, event);
    536     assert_eq!(stored_visible.position(), position);
    537     assert_eq!(stored_visible.event().id(), &event_id);
    538 
    539     assert_eq!(
    540         EventPage::new(
    541             generation,
    542             vec![1, 2],
    543             None,
    544             EventQueryBounds::first(1).unwrap()
    545         ),
    546         Err(Error::EventPageLimitExceeded)
    547     );
    548     assert_eq!(
    549         EventPage::<u8>::new(
    550             generation,
    551             vec![],
    552             Some(EventPosition::new(other_generation, sequence)),
    553             EventQueryBounds::first(1).unwrap(),
    554         ),
    555         Err(Error::CursorGenerationMismatch)
    556     );
    557     let page = EventPage::new(generation, vec![1], Some(position), bounds).expect("page");
    558     assert_eq!(page.generation(), generation);
    559     assert_eq!(page.items(), &[1]);
    560     assert_eq!(page.next_cursor(), Some(position));
    561 
    562     let provenance = observed(signed_event(), 5).provenance().clone();
    563     let stored = StoredEventProvenance::new(position, provenance.clone());
    564     assert_eq!(stored.position(), position);
    565     assert_eq!(stored.provenance(), &provenance);
    566     let reconstructed = StoredEventProvenance::from_stored_parts(
    567         position,
    568         "nostr",
    569         provenance.target().as_str(),
    570         5,
    571         Some("cursor"),
    572     )
    573     .expect("stored provenance");
    574     assert_eq!(
    575         reconstructed.provenance().cursor().unwrap().as_str(),
    576         "cursor"
    577     );
    578     for (transport, target, observed_at, cursor) in [
    579         ("BAD ID", provenance.target().as_str(), 5, None),
    580         ("nostr", "bad", 5, None),
    581         ("nostr", provenance.target().as_str(), 0, None),
    582         ("nostr", provenance.target().as_str(), 5, Some(" bad")),
    583     ] {
    584         assert_eq!(
    585             StoredEventProvenance::from_stored_parts(
    586                 position,
    587                 transport,
    588                 target,
    589                 observed_at,
    590                 cursor,
    591             ),
    592             Err(Error::CorruptStoredEvent)
    593         );
    594     }
    595 }