lib

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

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 }