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 }