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 }