lib

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

commit f70c151433a82cda12e6060df7c981991ebd9630
parent d783b6ccc602e1a382d4bd2df6b5ac0f2e2f3e75
Author: triesap <tyson@radroots.org>
Date:   Sun,  2 Aug 2026 00:08:33 +0000

storage: complete the conformance capability matrix

- cover every atomic workflow and reliability operation
- implement memory backup restore integrity and close
- distinguish volatile and sqlite backend status truthfully
- prove shared state replay conflicts and rollback behavior

Diffstat:
Mcrates/storage/src/memory.rs | 162++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Mcrates/storage/src/status.rs | 33++++++++++++++++++++++++++++-----
Mcrates/storage/tests/backup.rs | 6++++--
Mcrates/storage/tests/conformance.rs | 16+++++++++++++++-
Mcrates/storage/tests/conformance/suite.rs | 293+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
5 files changed, 496 insertions(+), 14 deletions(-)

diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -7,10 +7,15 @@ use std::sync::{Mutex, MutexGuard}; use crate::{ AtomicStorage, Error, EventStore, Journal, Outbox, PrivateArtifactStore, ProjectionStore, + StorageReliability, atomic::{ AtomicCommit, AtomicCommitDisposition, AtomicCommitId, AtomicCommitOutcome, AtomicCommitReceipt, AtomicWorkflow, }, + backup::{ + BackupId, BackupOperation, BackupPlan, BackupTransition, ReliabilityRevision, + RestoreOperation, RestorePlan, RestoreTransition, + }, event::{ AdmissionDisposition, AdmissionReceipt, AdmissionStage, EventAdmission, EventPage, EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration, @@ -35,7 +40,10 @@ use crate::{ ProjectionHealth, ProjectionId, ProjectionInvalidation, ProjectionStatus, RebuildStage, RebuildTicket, RebuildTransition, }, - status::{EventStoreHealth, EventStoreMode, EventStoreStatus}, + status::{ + EventStoreHealth, EventStoreMode, EventStoreStatus, IntegrityHealth, IntegrityStatus, + ShutdownState, StorageBackend, StorageOpenMode, StorageStatus, WriterPolicy, + }, }; #[derive(Clone)] @@ -55,7 +63,10 @@ struct State { event_index_manifests: Vec<EventIndexManifest>, event_index_checkpoints: Vec<EventIndexCheckpoint>, private_artifacts: Vec<PrivateArtifactMetadata>, + backups: Vec<BackupOperation>, + restores: Vec<RestoreOperation>, atomic_receipts: Vec<AtomicCommitReceipt>, + closed: bool, } /// Bounded deterministic reference backend with no hidden tasks or globals. @@ -77,7 +88,10 @@ impl MemoryStorage { event_index_manifests: Vec::new(), event_index_checkpoints: Vec::new(), private_artifacts: Vec::new(), + backups: Vec::new(), + restores: Vec::new(), atomic_receipts: Vec::new(), + closed: false, }), } } @@ -87,6 +101,14 @@ impl MemoryStorage { } fn state(&self) -> Result<MutexGuard<'_, State>, Error> { + let state = self.state_any()?; + if state.closed { + return Err(Error::BackendUnavailable); + } + Ok(state) + } + + fn state_any(&self) -> Result<MutexGuard<'_, State>, Error> { self.state.lock().map_err(|_| Error::BackendUnavailable) } @@ -279,6 +301,39 @@ impl MemoryStorage { state.projections.push(status.clone()); Ok(status) } + + fn integrity_locked(state: &State) -> Result<IntegrityStatus, Error> { + let members = state + .events + .len() + .checked_add(state.journal.len()) + .and_then(|count| count.checked_add(state.outbox.len())) + .and_then(|count| count.checked_add(state.projections.len())) + .and_then(|count| count.checked_add(state.private_artifacts.len())) + .ok_or(Error::InvalidIntegrityStatus)?; + IntegrityStatus::new( + IntegrityHealth::Healthy, + None, + u32::try_from(members).map_err(|_| Error::InvalidIntegrityStatus)?, + 0, + ) + } + + fn status_locked(state: &State) -> Result<StorageStatus, Error> { + StorageStatus::new( + StorageBackend::Memory, + StorageOpenMode::Create, + WriterPolicy::NoWriter, + if state.closed { + ShutdownState::Closed + } else { + ShutdownState::Open + }, + Self::integrity_locked(state)?, + false, + 0, + ) + } } impl Default for MemoryStorage { @@ -932,6 +987,111 @@ impl PrivateArtifactStore for MemoryStorage { } } +impl StorageReliability for MemoryStorage { + fn begin_backup(&self, plan: BackupPlan) -> BoxFuture<'_, Result<BackupOperation, Error>> { + Box::pin(async move { + let mut state = self.state()?; + if let Some(existing) = state + .backups + .iter() + .find(|operation| operation.plan().backup_id() == plan.backup_id()) + { + return if existing.plan() == &plan { + Ok(existing.clone()) + } else { + Err(Error::ReliabilityRevisionConflict) + }; + } + let operation = BackupOperation::planned(plan); + state.backups.push(operation.clone()); + Ok(operation) + }) + } + + fn transition_backup( + &self, + backup_id: BackupId, + expected_revision: ReliabilityRevision, + transition: BackupTransition, + at_unix_ms: u64, + ) -> BoxFuture<'_, Result<BackupOperation, Error>> { + Box::pin(async move { + let mut state = self.state()?; + let operation = state + .backups + .iter_mut() + .find(|operation| operation.plan().backup_id() == backup_id) + .ok_or(Error::CorruptReliabilityOperation)?; + let next = operation.transition(expected_revision, transition, at_unix_ms)?; + *operation = next.clone(); + Ok(next) + }) + } + + fn begin_restore(&self, plan: RestorePlan) -> BoxFuture<'_, Result<RestoreOperation, Error>> { + Box::pin(async move { + let mut state = self.state()?; + let backup_id = plan.manifest().backup_id(); + if let Some(existing) = state + .restores + .iter() + .find(|operation| operation.plan().manifest().backup_id() == backup_id) + { + return if existing.plan() == &plan { + Ok(existing.clone()) + } else { + Err(Error::ReliabilityRevisionConflict) + }; + } + let operation = RestoreOperation::staging(plan); + state.restores.push(operation.clone()); + Ok(operation) + }) + } + + fn transition_restore( + &self, + backup_id: BackupId, + expected_revision: ReliabilityRevision, + transition: RestoreTransition, + at_unix_ms: u64, + ) -> BoxFuture<'_, Result<RestoreOperation, Error>> { + Box::pin(async move { + let mut state = self.state()?; + let operation = state + .restores + .iter_mut() + .find(|operation| operation.plan().manifest().backup_id() == backup_id) + .ok_or(Error::CorruptReliabilityOperation)?; + let next = operation.transition(expected_revision, transition, at_unix_ms)?; + *operation = next.clone(); + Ok(next) + }) + } + + fn integrity(&self) -> BoxFuture<'_, Result<IntegrityStatus, Error>> { + Box::pin(async move { + let state = self.state_any()?; + Self::integrity_locked(&state) + }) + } + + fn status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { + Box::pin(async move { + let state = self.state_any()?; + Self::status_locked(&state) + }) + } + + fn close(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { + Box::pin(async move { + let mut state = self.state_any()?; + state.closed = true; + Self::status_locked(&state) + }) + } +} + impl AtomicStorage for MemoryStorage { fn commit(&self, request: AtomicCommit) -> BoxFuture<'_, Result<AtomicCommitReceipt, Error>> { Box::pin(async move { diff --git a/crates/storage/src/status.rs b/crates/storage/src/status.rs @@ -2,6 +2,15 @@ use crate::{Error, event::SourceGeneration}; +/// Storage-engine family needed to interpret durability-specific status. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum StorageBackend { + Memory, + Sqlite, +} + #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -85,6 +94,7 @@ impl IntegrityStatus { #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[derive(Clone, Copy, Debug, Eq, PartialEq)] pub struct StorageStatus { + backend: StorageBackend, open_mode: StorageOpenMode, writer_policy: WriterPolicy, shutdown: ShutdownState, @@ -95,6 +105,7 @@ pub struct StorageStatus { impl StorageStatus { pub fn new( + backend: StorageBackend, open_mode: StorageOpenMode, writer_policy: WriterPolicy, shutdown: ShutdownState, @@ -102,14 +113,23 @@ impl StorageStatus { wal_enabled: bool, busy_timeout_ms: u32, ) -> Result<Self, Error> { - if (open_mode == StorageOpenMode::ReadOnly && writer_policy != WriterPolicy::NoWriter) - || (open_mode != StorageOpenMode::ReadOnly - && writer_policy != WriterPolicy::AdvisoryProcessLock) - || (open_mode != StorageOpenMode::ReadOnly && (!wal_enabled || busy_timeout_ms == 0)) - { + let valid_engine_status = match backend { + StorageBackend::Memory => { + writer_policy == WriterPolicy::NoWriter && !wal_enabled && busy_timeout_ms == 0 + } + StorageBackend::Sqlite => { + (open_mode == StorageOpenMode::ReadOnly && writer_policy == WriterPolicy::NoWriter) + || (open_mode != StorageOpenMode::ReadOnly + && writer_policy == WriterPolicy::AdvisoryProcessLock + && wal_enabled + && busy_timeout_ms != 0) + } + }; + if !valid_engine_status { return Err(Error::InvalidStorageStatus); } Ok(Self { + backend, open_mode, writer_policy, shutdown, @@ -118,6 +138,9 @@ impl StorageStatus { busy_timeout_ms, }) } + pub const fn backend(self) -> StorageBackend { + self.backend + } pub const fn open_mode(self) -> StorageOpenMode { self.open_mode } diff --git a/crates/storage/tests/backup.rs b/crates/storage/tests/backup.rs @@ -7,8 +7,8 @@ use radroots_storage::{ RestoreOperation, RestorePlan, RestoreStage, RestoreTransition, }, status::{ - IntegrityHealth, IntegrityStatus, ShutdownState, StorageOpenMode, StorageStatus, - WriterPolicy, + IntegrityHealth, IntegrityStatus, ShutdownState, StorageBackend, StorageOpenMode, + StorageStatus, WriterPolicy, }, }; @@ -173,6 +173,7 @@ fn integrity_and_storage_status_reject_inconsistent_runtime_claims() { let integrity = IntegrityStatus::new(IntegrityHealth::Healthy, Some(100), 2, 0).expect("integrity status"); let status = StorageStatus::new( + StorageBackend::Sqlite, StorageOpenMode::ReadWriteExisting, WriterPolicy::AdvisoryProcessLock, ShutdownState::Open, @@ -189,6 +190,7 @@ fn integrity_and_storage_status_reject_inconsistent_runtime_claims() { ); assert_eq!( StorageStatus::new( + StorageBackend::Sqlite, StorageOpenMode::ReadWriteExisting, WriterPolicy::NoWriter, ShutdownState::Open, diff --git a/crates/storage/tests/conformance.rs b/crates/storage/tests/conformance.rs @@ -3,7 +3,7 @@ mod suite; use radroots_storage::{ AtomicStorage, EventStore, Journal, Outbox, PrivateArtifactStore, ProjectionStore, - memory::MemoryStorage, + StorageReliability, memory::MemoryStorage, }; use suite::StorageConformanceHarness; @@ -37,6 +37,10 @@ impl StorageConformanceHarness for MemoryHarness { fn atomic_storage(&self) -> &dyn AtomicStorage { &self.storage } + + fn reliability(&self) -> &dyn StorageReliability { + &self.storage + } } #[test] @@ -53,3 +57,13 @@ fn memory_backend_atomic_failure_is_all_or_nothing() { fn memory_backend_rejects_identity_and_digest_conflicts() { suite::assert_conflict_conformance(&MemoryHarness::default()); } + +#[test] +fn memory_backend_supports_every_atomic_workflow() { + suite::assert_atomic_workflow_conformance(&MemoryHarness::default()); +} + +#[test] +fn memory_backend_supports_reliability_and_explicit_close() { + suite::assert_reliability_and_close_conformance(&MemoryHarness::default()); +} diff --git a/crates/storage/tests/conformance/suite.rs b/crates/storage/tests/conformance/suite.rs @@ -3,20 +3,37 @@ use radroots_event::{SignedEvent, wire::Nip01EventWire}; use radroots_protocol::runtime::v1::OperationId; use radroots_storage::{ AtomicStorage, EventStore, Journal, Outbox, PrivateArtifactStore, ProjectionStore, - atomic::{AtomicCommit, AtomicCommitDigest, AtomicCommitId, AtomicWorkflow, CommitIngested}, + StorageReliability, + atomic::{ + AtomicCommit, AtomicCommitDigest, AtomicCommitDisposition, AtomicCommitId, AtomicWorkflow, + CommitEnqueued, CommitIngested, CommitSigned, + }, + backup::{ + BackupFormatVersion, BackupId, BackupManifest, BackupMember, BackupMemberKind, BackupPlan, + BackupSecretPolicy, BackupStage, BackupTransition, MemberDigest, MemberVerification, + RestoreMemberStatus, RestorePlan, RestoreStage, RestoreTransition, + }, event::{EventAdmission, EventQuery, EventQueryBounds}, - journal::{IdempotencyDigest, IdempotencyKey, OperationInstanceId, PrepareOperation}, - outbox::{DeliveryPlanDigest, EnqueueDisposition, EnqueueOutboxItem, OutboxItemId}, + journal::{ + IdempotencyDigest, IdempotencyKey, JournalRevision, JournalStage, OperationInstanceId, + PrepareOperation, + }, + outbox::{ + ClaimOutboxItems, DeliveryAttempt, DeliveryAttemptEvidence, DeliveryPlanDigest, + EnqueueDisposition, EnqueueOutboxItem, LeaseId, LeaseOwner, OutboxItemId, OutboxStage, + }, private_artifact::{ ArtifactCommitment, ArtifactKind, ArtifactSchemaId, DurableSecretReference, PrivateArtifactId, PrivateArtifactMetadata, RetentionPolicy, }, projection::{ProjectionCheckpoint, ProjectionGeneration, ProjectionId}, + status::{ShutdownState, StorageBackend}, }; use radroots_transport::{ - DeliveryRequest, Target, TargetSet, TransportId, + DeliveryReceipt, DeliveryRequest, Target, TargetSet, TransportId, + outcome::DeliveryOutcome, policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, - sink::DeliveryPayload, + sink::{DeliveryPayload, DeliveryTargetReceipt}, source::{EventProvenance, ObservedEvent}, }; @@ -28,6 +45,7 @@ pub(crate) trait StorageConformanceHarness { fn projection_store(&self) -> &dyn ProjectionStore; fn private_artifact_store(&self) -> &dyn PrivateArtifactStore; fn atomic_storage(&self) -> &dyn AtomicStorage; + fn reliability(&self) -> &dyn StorageReliability; } pub(crate) fn assert_shared_state_conformance(harness: &impl StorageConformanceHarness) { @@ -169,6 +187,242 @@ pub(crate) fn assert_conflict_conformance(harness: &impl StorageConformanceHarne ); } +pub(crate) fn assert_atomic_workflow_conformance(harness: &impl StorageConformanceHarness) { + let instance = OperationInstanceId::new([10; 16]).expect("operation instance"); + let prepared_request = atomic_commit( + 10, + 10, + 100, + AtomicWorkflow::Prepared(prepare(instance, "atomic", 10)), + ); + let prepared = block_on(harness.atomic_storage().commit(prepared_request.clone())) + .expect("atomic prepare"); + assert_eq!(prepared.disposition(), AtomicCommitDisposition::Committed); + assert_eq!( + block_on(harness.atomic_storage().commit(prepared_request)) + .expect("prepare replay") + .disposition(), + AtomicCommitDisposition::Replay + ); + + let event = signed_event("conformance-atomic"); + let signed = block_on(harness.atomic_storage().commit(atomic_commit( + 11, + 11, + 110, + AtomicWorkflow::Signed(Box::new(CommitSigned::new( + instance, + JournalRevision::INITIAL, + event.clone(), + ))), + ))) + .expect("atomic sign"); + assert_eq!( + signed.outcome().kind(), + radroots_storage::atomic::AtomicWorkflowKind::Signed + ); + + let enqueue = enqueue([11; 16], instance, event.clone()); + let enqueued = block_on( + harness.atomic_storage().commit(atomic_commit( + 12, + 12, + 120, + AtomicWorkflow::Enqueued(Box::new( + CommitEnqueued::new( + instance, + JournalRevision::new(2).expect("journal revision"), + admission(event, 120), + enqueue, + 120, + ) + .expect("enqueue workflow"), + )), + )), + ) + .expect("atomic enqueue"); + assert_eq!( + block_on(harness.journal().operation(instance)) + .expect("journal lookup") + .expect("journal record") + .state() + .stage(), + JournalStage::Committed + ); + assert_eq!( + enqueued.outcome().kind(), + radroots_storage::atomic::AtomicWorkflowKind::Enqueued + ); + + let claimed = block_on( + harness.outbox().claim( + ClaimOutboxItems::new( + LeaseOwner::parse("conformance-worker").expect("lease owner"), + LeaseId::new([12; 16]).expect("lease id"), + 130, + 150, + 1, + ) + .expect("claim request"), + ), + ) + .expect("claim outbox") + .pop() + .expect("claimed item"); + let receipt = DeliveryReceipt::for_request( + claimed.record().request(), + claimed + .record() + .request() + .target_set() + .targets() + .iter() + .cloned() + .map(|target| DeliveryTargetReceipt::attempted(target, DeliveryOutcome::delivered())) + .collect(), + ) + .expect("delivery receipt"); + let evidence = DeliveryAttemptEvidence::new( + claimed.record().item_id(), + claimed.lease().id(), + claimed.record().revision(), + DeliveryAttempt::FIRST, + receipt, + 140, + ) + .expect("delivery evidence"); + let delivered = block_on(harness.atomic_storage().commit(atomic_commit( + 13, + 13, + 140, + AtomicWorkflow::Delivered(Box::new(evidence)), + ))) + .expect("atomic deliver"); + assert_eq!( + delivered.outcome().kind(), + radroots_storage::atomic::AtomicWorkflowKind::Delivered + ); + assert_eq!( + block_on(harness.outbox().item(claimed.record().item_id())) + .expect("outbox lookup") + .expect("outbox item") + .stage(), + OutboxStage::Satisfied + ); + + let ingest_event = signed_event("conformance-ingest"); + let ingested = block_on( + harness.atomic_storage().commit(atomic_commit( + 14, + 14, + 150, + AtomicWorkflow::Ingested(Box::new(CommitIngested::new( + admission(ingest_event, 150), + Some( + ProjectionCheckpoint::new( + ProjectionId::parse("conformance.ingest").expect("projection id"), + ProjectionGeneration::new([14; 32]).expect("projection generation"), + None, + 1, + 150, + ) + .expect("projection checkpoint"), + ), + ))), + )), + ) + .expect("atomic ingest"); + assert_eq!( + ingested.outcome().kind(), + radroots_storage::atomic::AtomicWorkflowKind::Ingested + ); +} + +pub(crate) fn assert_reliability_and_close_conformance(harness: &impl StorageConformanceHarness) { + let backup_id = BackupId::new([15; 16]).expect("backup id"); + let plan = BackupPlan::new( + backup_id, + BackupFormatVersion::V1, + BackupSecretPolicy::ExcludeProtectedStorage, + 100, + ) + .expect("backup plan"); + let planned = block_on(harness.reliability().begin_backup(plan)).expect("begin backup"); + assert_eq!(planned.stage(), BackupStage::Planned); + let manifest = backup_manifest(backup_id); + let captured = block_on(harness.reliability().transition_backup( + backup_id, + planned.revision(), + BackupTransition::Captured(manifest.clone()), + 110, + )) + .expect("capture backup"); + let verified = block_on(harness.reliability().transition_backup( + backup_id, + captured.revision(), + BackupTransition::Verified, + 120, + )) + .expect("verify backup"); + let finalized = block_on(harness.reliability().transition_backup( + backup_id, + verified.revision(), + BackupTransition::Finalize, + 130, + )) + .expect("finalize backup"); + assert_eq!(finalized.stage(), BackupStage::Finalized); + + let restore = block_on( + harness.reliability().begin_restore( + RestorePlan::new(manifest, BackupSecretPolicy::ExcludeProtectedStorage, 140) + .expect("restore plan"), + ), + ) + .expect("begin restore"); + let verifying = block_on(harness.reliability().transition_restore( + backup_id, + restore.revision(), + RestoreTransition::Staged, + 150, + )) + .expect("stage restore"); + let finalizing = block_on(harness.reliability().transition_restore( + backup_id, + verifying.revision(), + RestoreTransition::Verified(vec![ + RestoreMemberStatus::new("memory/state", MemberVerification::Verified) + .expect("member status"), + ]), + 160, + )) + .expect("verify restore"); + let restored = block_on(harness.reliability().transition_restore( + backup_id, + finalizing.revision(), + RestoreTransition::Finalize, + 170, + )) + .expect("finalize restore"); + assert_eq!(restored.stage(), RestoreStage::Finalized); + assert_eq!( + block_on(harness.reliability().status()) + .expect("storage status") + .backend(), + StorageBackend::Memory + ); + assert_eq!( + block_on(harness.reliability().close()) + .expect("close storage") + .shutdown(), + ShutdownState::Closed + ); + assert_eq!( + block_on(harness.event_store().status()), + Err(radroots_storage::Error::BackendUnavailable) + ); +} + fn signed_event(content: &str) -> SignedEvent { let mut wire = Nip01EventWire { id: "0".repeat(64), @@ -250,3 +504,32 @@ fn private_metadata(id: [u8; 16]) -> PrivateArtifactMetadata { ) .expect("private metadata") } + +fn atomic_commit(id: u8, digest: u8, at: u64, workflow: AtomicWorkflow) -> AtomicCommit { + AtomicCommit::new( + AtomicCommitId::new([id; 16]).expect("commit id"), + AtomicCommitDigest::new([digest; 32]), + at, + workflow, + ) + .expect("atomic commit") +} + +fn backup_manifest(backup_id: BackupId) -> BackupManifest { + BackupManifest::new( + BackupFormatVersion::V1, + backup_id, + 105, + BackupSecretPolicy::ExcludeProtectedStorage, + vec![ + BackupMember::new( + "memory/state", + BackupMemberKind::Runtime, + 1, + MemberDigest::new([15; 32]), + ) + .expect("backup member"), + ], + ) + .expect("backup manifest") +}