suite.rs (14530B)
1 use futures_executor::block_on; 2 use radroots_event::{SignedEvent, wire::Nip01EventWire}; 3 use radroots_protocol::runtime::v1::OperationId; 4 use radroots_storage::{ 5 EventStore, Journal, Outbox, ProjectionStore, 6 atomic::{ 7 AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, AtomicStorage, 8 AtomicWorkflow, CommitIngested, 9 }, 10 backup::{ 11 BackupFormatVersion, BackupId, BackupManifest, BackupMember, BackupMemberKind, BackupPlan, 12 BackupSecretPolicy, BackupStage, BackupTransition, MemberDigest, MemberVerification, 13 RestoreMemberStatus, RestorePlan, RestoreStage, RestoreTransition, StorageReliability, 14 }, 15 event::{EventAdmission, EventQuery, EventQueryBounds}, 16 journal::{IdempotencyDigest, IdempotencyKey, OperationInstanceId, PrepareOperation}, 17 outbox::{DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, OutboxItemId}, 18 private_artifact::{ 19 ArtifactCommitment, ArtifactKind, ArtifactSchemaId, DurableSecretReference, 20 PrivateArtifactId, PrivateArtifactMetadata, PrivateArtifactStore, RetentionPolicy, 21 }, 22 projection::{ProjectionCheckpoint, ProjectionGeneration, ProjectionId}, 23 status::{ShutdownState, StorageBackend}, 24 }; 25 use radroots_transport::{ 26 DeliveryRequest, Target, TargetSet, TransportId, 27 policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, 28 sink::DeliveryPayload, 29 source::{EventProvenance, ObservedEvent}, 30 }; 31 32 /// Backend adapter consumed by the shared storage contract assertions. 33 pub(crate) trait StorageConformanceHarness { 34 fn event_store(&self) -> &dyn EventStore; 35 fn journal(&self) -> &dyn Journal; 36 fn outbox(&self) -> &dyn Outbox; 37 fn projection_store(&self) -> &dyn ProjectionStore; 38 fn private_artifact_store(&self) -> &dyn PrivateArtifactStore; 39 fn atomic_storage(&self) -> &dyn AtomicStorage; 40 fn reliability(&self) -> &dyn StorageReliability; 41 } 42 43 pub(crate) fn assert_shared_state_conformance(harness: &impl StorageConformanceHarness) { 44 let event = signed_event("conformance-shared"); 45 let event_id = *event.id(); 46 block_on(harness.event_store().admit(admission(event.clone(), 100))).expect("admit event"); 47 let page = block_on(harness.event_store().query_raw(EventQuery::all( 48 EventQueryBounds::first(10).expect("query bounds"), 49 ))) 50 .expect("query events"); 51 assert_eq!(page.items().len(), 1); 52 assert_eq!(page.items()[0].event().id(), &event_id); 53 54 let instance = OperationInstanceId::new([1; 16]).expect("operation instance"); 55 let prepared = block_on(harness.journal().prepare(prepare(instance, "shared", 1))) 56 .expect("prepare operation"); 57 assert_eq!( 58 block_on(harness.journal().operation(instance)) 59 .expect("journal lookup") 60 .expect("journal record"), 61 *prepared.record() 62 ); 63 64 let outbox = enqueue([2; 16], instance, event); 65 assert_eq!( 66 block_on(harness.outbox().enqueue(outbox.clone())) 67 .expect("enqueue") 68 .disposition(), 69 EnqueueDisposition::Created 70 ); 71 assert_eq!( 72 block_on(harness.outbox().enqueue(outbox)) 73 .expect("enqueue replay") 74 .disposition(), 75 EnqueueDisposition::Replay 76 ); 77 78 let checkpoint = ProjectionCheckpoint::new( 79 ProjectionId::parse("conformance.shared").expect("projection id"), 80 ProjectionGeneration::new([3; 32]).expect("projection generation"), 81 None, 82 1, 83 200, 84 ) 85 .expect("projection checkpoint"); 86 let projection = 87 block_on(harness.projection_store().checkpoint(checkpoint)).expect("store checkpoint"); 88 assert_eq!( 89 projection 90 .checkpoint() 91 .expect("checkpoint") 92 .projected_rows(), 93 1 94 ); 95 96 let metadata = private_metadata([4; 16]); 97 assert_eq!( 98 block_on( 99 harness 100 .private_artifact_store() 101 .put_metadata(metadata.clone()), 102 ) 103 .expect("put private metadata"), 104 metadata 105 ); 106 } 107 108 pub(crate) fn assert_atomic_failure_isolation(harness: &impl StorageConformanceHarness) { 109 let projection_id = ProjectionId::parse("conformance.atomic").expect("projection id"); 110 let generation = ProjectionGeneration::new([5; 32]).expect("projection generation"); 111 let current = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 2, 200) 112 .expect("current checkpoint"); 113 block_on(harness.projection_store().checkpoint(current)).expect("store checkpoint"); 114 let regression = ProjectionCheckpoint::new(projection_id.clone(), generation, None, 1, 201) 115 .expect("regressing checkpoint"); 116 let request = AtomicCommit::new( 117 AtomicCommitId::new([6; 16]).expect("commit id"), 118 AtomicCommitDigest::new([6; 32]), 119 201, 120 AtomicWorkflow::Ingested(Box::new(CommitIngested::new( 121 admission(signed_event("conformance-rollback"), 201), 122 Some(regression), 123 ))), 124 ) 125 .expect("atomic request"); 126 assert_eq!( 127 block_on(harness.atomic_storage().commit(request)), 128 Err(radroots_storage::Error::ProjectionCheckpointRegression) 129 ); 130 assert_eq!( 131 block_on(harness.event_store().status()) 132 .expect("event status") 133 .raw_events(), 134 0 135 ); 136 assert_eq!( 137 block_on(harness.projection_store().status(projection_id)) 138 .expect("projection status") 139 .expect("projection") 140 .checkpoint() 141 .expect("checkpoint") 142 .projected_rows(), 143 2 144 ); 145 assert!( 146 block_on( 147 harness 148 .atomic_storage() 149 .receipt(AtomicCommitId::new([6; 16]).expect("commit id")), 150 ) 151 .expect("receipt lookup") 152 .is_none() 153 ); 154 } 155 156 pub(crate) fn assert_conflict_conformance(harness: &impl StorageConformanceHarness) { 157 let instance = OperationInstanceId::new([7; 16]).expect("operation instance"); 158 block_on(harness.journal().prepare(prepare(instance, "conflict", 7))) 159 .expect("prepare operation"); 160 assert_eq!( 161 block_on(harness.journal().prepare(prepare(instance, "conflict", 8))), 162 Err(radroots_storage::Error::IdempotencyConflict) 163 ); 164 165 let event = signed_event("conformance-conflict"); 166 let first = enqueue([8; 16], instance, event.clone()); 167 block_on(harness.outbox().enqueue(first)).expect("enqueue plan"); 168 let conflicting = EnqueueOutboxItem::new( 169 OutboxItemId::new([8; 16]).expect("item id"), 170 instance, 171 DeliveryPlanDigest::new([9; 32]), 172 delivery_request(event), 173 100, 174 ) 175 .expect("conflicting plan"); 176 assert_eq!( 177 block_on(harness.outbox().enqueue(conflicting)), 178 Err(radroots_storage::Error::OutboxPlanConflict) 179 ); 180 } 181 182 pub(crate) fn assert_atomic_workflow_conformance(harness: &impl StorageConformanceHarness) { 183 let event = signed_event("conformance-ingest"); 184 let request = atomic_commit( 185 14, 186 14, 187 150, 188 AtomicWorkflow::Ingested(Box::new(CommitIngested::new( 189 admission(event, 150), 190 Some( 191 ProjectionCheckpoint::new( 192 ProjectionId::parse("conformance.ingest").expect("projection id"), 193 ProjectionGeneration::new([14; 32]).expect("projection generation"), 194 None, 195 1, 196 150, 197 ) 198 .expect("projection checkpoint"), 199 ), 200 ))), 201 ); 202 let ingested = 203 block_on(harness.atomic_storage().commit(request.clone())).expect("atomic ingest"); 204 assert_eq!(ingested.disposition(), AtomicCommitDisposition::Committed); 205 assert_eq!( 206 ingested.outcome().kind(), 207 radroots_storage::atomic::AtomicWorkflowKind::Ingested 208 ); 209 assert_eq!( 210 block_on(harness.atomic_storage().commit(request)) 211 .expect("atomic replay") 212 .disposition(), 213 AtomicCommitDisposition::Replay 214 ); 215 } 216 217 pub(crate) fn assert_reliability_and_close_conformance(harness: &impl StorageConformanceHarness) { 218 let backup_id = BackupId::new([15; 16]).expect("backup id"); 219 let plan = BackupPlan::new( 220 backup_id, 221 BackupFormatVersion::V1, 222 BackupSecretPolicy::ExcludeProtectedStorage, 223 100, 224 ) 225 .expect("backup plan"); 226 let planned = block_on(harness.reliability().begin_backup(plan)).expect("begin backup"); 227 assert_eq!(planned.stage(), BackupStage::Planned); 228 let manifest = backup_manifest(backup_id); 229 let captured = block_on(harness.reliability().transition_backup( 230 backup_id, 231 planned.revision(), 232 BackupTransition::Captured(manifest.clone()), 233 110, 234 )) 235 .expect("capture backup"); 236 let verified = block_on(harness.reliability().transition_backup( 237 backup_id, 238 captured.revision(), 239 BackupTransition::Verified, 240 120, 241 )) 242 .expect("verify backup"); 243 let finalized = block_on(harness.reliability().transition_backup( 244 backup_id, 245 verified.revision(), 246 BackupTransition::Finalize, 247 130, 248 )) 249 .expect("finalize backup"); 250 assert_eq!(finalized.stage(), BackupStage::Finalized); 251 252 let restore = block_on( 253 harness.reliability().begin_restore( 254 RestorePlan::new(manifest, BackupSecretPolicy::ExcludeProtectedStorage, 140) 255 .expect("restore plan"), 256 ), 257 ) 258 .expect("begin restore"); 259 let verifying = block_on(harness.reliability().transition_restore( 260 backup_id, 261 restore.revision(), 262 RestoreTransition::Staged, 263 150, 264 )) 265 .expect("stage restore"); 266 let finalizing = block_on(harness.reliability().transition_restore( 267 backup_id, 268 verifying.revision(), 269 RestoreTransition::Verified(vec![ 270 RestoreMemberStatus::new("memory/state", MemberVerification::Verified) 271 .expect("member status"), 272 ]), 273 160, 274 )) 275 .expect("verify restore"); 276 let restored = block_on(harness.reliability().transition_restore( 277 backup_id, 278 finalizing.revision(), 279 RestoreTransition::Finalize, 280 170, 281 )) 282 .expect("finalize restore"); 283 assert_eq!(restored.stage(), RestoreStage::Finalized); 284 assert_eq!( 285 block_on(harness.reliability().status()) 286 .expect("storage status") 287 .backend(), 288 StorageBackend::Memory 289 ); 290 assert_eq!( 291 block_on(harness.reliability().close()) 292 .expect("close storage") 293 .shutdown(), 294 ShutdownState::Closed 295 ); 296 assert_eq!( 297 block_on(harness.event_store().status()), 298 Err(radroots_storage::Error::BackendUnavailable) 299 ); 300 } 301 302 fn signed_event(content: &str) -> SignedEvent { 303 let mut wire = Nip01EventWire { 304 id: "0".repeat(64), 305 pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), 306 created_at: 1_800_000_100, 307 kind: 1, 308 tags: vec![], 309 content: content.to_owned(), 310 sig: "42".repeat(64), 311 extra: Default::default(), 312 }; 313 wire.id = wire.computed_event_id().expect("event id").to_hex(); 314 let raw = serde_json::json!({ 315 "id": &wire.id, 316 "pubkey": &wire.pubkey, 317 "created_at": wire.created_at, 318 "kind": wire.kind, 319 "tags": &wire.tags, 320 "content": &wire.content, 321 "sig": &wire.sig, 322 }) 323 .to_string(); 324 SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") 325 } 326 327 fn admission(event: SignedEvent, observed_at_unix_ms: u64) -> EventAdmission { 328 let target = Target::nostr_relay("wss://conformance.example").expect("target"); 329 let provenance = EventProvenance::new( 330 TransportId::NOSTR, 331 target.fingerprint().clone(), 332 observed_at_unix_ms, 333 ) 334 .expect("provenance"); 335 EventAdmission::raw(ObservedEvent::new(event, provenance)) 336 } 337 338 fn prepare(instance: OperationInstanceId, key: &str, digest: u8) -> PrepareOperation { 339 PrepareOperation::new( 340 instance, 341 OperationId::SyncPush, 342 IdempotencyKey::parse(key).expect("idempotency key"), 343 IdempotencyDigest::new([digest; 32]), 344 100, 345 ) 346 .expect("prepare operation") 347 } 348 349 fn delivery_request(event: SignedEvent) -> DeliveryRequest { 350 DeliveryRequest::new( 351 "storage-conformance", 352 DeliveryPayload::new(event), 353 TargetSet::new(vec![ 354 Target::nostr_relay("wss://conformance.example").expect("target"), 355 ]) 356 .expect("target set"), 357 SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), 358 1_000, 359 ) 360 .expect("delivery request") 361 } 362 363 fn enqueue( 364 item_id: [u8; 16], 365 instance: OperationInstanceId, 366 event: SignedEvent, 367 ) -> EnqueueOutboxItem { 368 EnqueueOutboxItem::new( 369 OutboxItemId::new(item_id).expect("item id"), 370 instance, 371 DeliveryPlanDigest::new([2; 32]), 372 delivery_request(event), 373 100, 374 ) 375 .expect("enqueue") 376 } 377 378 fn private_metadata(id: [u8; 16]) -> PrivateArtifactMetadata { 379 PrivateArtifactMetadata::new( 380 PrivateArtifactId::new(id).expect("artifact id"), 381 ArtifactKind::parse("conformance.private").expect("artifact kind"), 382 ArtifactSchemaId::parse("conformance.private.v1").expect("schema id"), 383 ArtifactCommitment::new([4; 32]), 384 32, 385 DurableSecretReference::new("conformance", "opaque-reference", 1) 386 .expect("secret reference"), 387 RetentionPolicy::indefinite(), 388 100, 389 ) 390 .expect("private metadata") 391 } 392 393 fn atomic_commit(id: u8, digest: u8, at: u64, workflow: AtomicWorkflow) -> AtomicCommit { 394 AtomicCommit::new( 395 AtomicCommitId::new([id; 16]).expect("commit id"), 396 AtomicCommitDigest::new([digest; 32]), 397 at, 398 workflow, 399 ) 400 .expect("atomic commit") 401 } 402 403 fn backup_manifest(backup_id: BackupId) -> BackupManifest { 404 BackupManifest::new( 405 BackupFormatVersion::V1, 406 backup_id, 407 105, 408 BackupSecretPolicy::ExcludeProtectedStorage, 409 vec![ 410 BackupMember::new( 411 "memory/state", 412 BackupMemberKind::Runtime, 413 1, 414 MemberDigest::new([15; 32]), 415 ) 416 .expect("backup member"), 417 ], 418 ) 419 .expect("backup manifest") 420 }