lib

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

commit d096d9e8e2ee058e75d0d6684240887da244c546
parent d293ea7fb79da801ab579582f21cbab59d64dd15
Author: triesap <tyson@radroots.org>
Date:   Mon,  3 Aug 2026 09:15:34 +0000

sync: add end-to-end sync reliability tests

- define the exact runtime storage capability sync consumes
- expose passive storage status independently of backup workflows
- run durable retry scenarios against memory and SQLite
- prove SQLite retry recovery across a complete reopen boundary

Diffstat:
MCargo.lock | 4++++
Mcrates/storage/src/memory.rs | 12+++++++++++-
Mcrates/storage/src/status.rs | 7+++++++
Mcrates/storage_sqlite/src/status.rs | 9++++++++-
Mcrates/sync/Cargo.toml | 4++++
Mcrates/sync/src/engine.rs | 9++++-----
Mcrates/sync/src/policy.rs | 20+++++++++++++++++---
Mcrates/sync/src/status.rs | 9++++++---
Mcrates/sync/tests/engine_composition.rs | 8+++-----
Mcrates/sync/tests/ingest.rs | 6+++---
Mcrates/sync/tests/projection.rs | 6+++---
Mcrates/sync/tests/pull.rs | 6+++---
Mcrates/sync/tests/push_enqueue.rs | 6+++---
Acrates/sync/tests/reliability_scenarios.rs | 273+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
14 files changed, 349 insertions(+), 30 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -5162,11 +5162,15 @@ dependencies = [ "radroots_protocol", "radroots_signing", "radroots_storage", + "radroots_storage_sqlite", "radroots_trade", "radroots_transport", "secp256k1", "serde", + "serde_json", "sha2", + "tempfile", + "tokio", ] [[package]] diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -41,7 +41,8 @@ use crate::{ }, status::{ EventStoreHealth, EventStoreMode, EventStoreStatus, IntegrityHealth, IntegrityStatus, - ShutdownState, StorageBackend, StorageOpenMode, StorageStatus, WriterPolicy, + ShutdownState, StorageBackend, StorageOpenMode, StorageStatus, StorageStatusProvider, + WriterPolicy, }, }; @@ -1132,6 +1133,15 @@ impl StorageReliability for MemoryStorage { } } +impl StorageStatusProvider for MemoryStorage { + fn storage_status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { + Box::pin(async move { + let state = self.state_any()?; + 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 @@ -1,7 +1,14 @@ //! Storage capability, health, and integrity status contracts. +use radroots_transport::BoxFuture; + use crate::{Error, event::SourceGeneration}; +/// Passive backend-level status capability independent of backup workflows. +pub trait StorageStatusProvider: Send + Sync { + fn storage_status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>>; +} + /// 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"))] diff --git a/crates/storage_sqlite/src/status.rs b/crates/storage_sqlite/src/status.rs @@ -10,9 +10,10 @@ use std::{ use radroots_storage::{ Error, + outbox::BoxFuture, status::{ IntegrityStatus, ShutdownState, StorageBackend, StorageOpenMode, StorageStatus, - WriterPolicy, + StorageStatusProvider, WriterPolicy, }, }; @@ -184,6 +185,12 @@ impl SqliteStorage { } } +impl StorageStatusProvider for SqliteStorage { + fn storage_status(&self) -> BoxFuture<'_, Result<StorageStatus, Error>> { + Box::pin(async move { SqliteStorage::storage_status(self).await }) + } +} + const fn storage_open_mode(mode: OpenMode) -> StorageOpenMode { match mode { OpenMode::ReadOnly => StorageOpenMode::ReadOnly, diff --git a/crates/sync/Cargo.toml b/crates/sync/Cargo.toml @@ -42,7 +42,11 @@ sha2 = { workspace = true, default-features = false } futures = { workspace = true } futures-executor = { workspace = true } radroots_storage = { workspace = true, features = ["memory"] } +radroots_storage_sqlite = { workspace = true } secp256k1 = { workspace = true } +serde_json = { workspace = true, features = ["std"] } +tempfile = { workspace = true } +tokio = { workspace = true, features = ["macros", "rt"] } [lints] workspace = true diff --git a/crates/sync/src/engine.rs b/crates/sync/src/engine.rs @@ -1,15 +1,14 @@ use std::sync::Arc; use radroots_signing::Signer; -use radroots_storage::Storage; use radroots_transport::{EventSink, EventSource}; -use crate::policy::{Clock, DeadlinePolicy, EngineBuilder, IdSource}; +use crate::policy::{Clock, DeadlinePolicy, EngineBuilder, IdSource, SyncStorage}; /// Injected, executor-neutral synchronization composition boundary. #[derive(Clone)] pub struct Engine { - pub(crate) storage: Arc<dyn Storage>, + pub(crate) storage: Arc<dyn SyncStorage>, pub(crate) source: Option<Arc<dyn EventSource>>, pub(crate) sink: Option<Arc<dyn EventSink>>, pub(crate) signer: Option<Arc<dyn Signer>>, @@ -21,7 +20,7 @@ pub struct Engine { impl Engine { /// Starts an explicit capability builder around required host policies. pub fn builder( - storage: Arc<dyn Storage>, + storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>, ids: Arc<dyn IdSource>, deadlines: DeadlinePolicy, @@ -30,7 +29,7 @@ impl Engine { } /// Returns the canonical storage capability. - pub fn storage(&self) -> &dyn Storage { + pub fn storage(&self) -> &dyn SyncStorage { self.storage.as_ref() } diff --git a/crates/sync/src/policy.rs b/crates/sync/src/policy.rs @@ -3,13 +3,27 @@ use std::sync::Arc; use radroots_signing::Signer; -use radroots_storage::Storage; +use radroots_storage::{ + EventStore, Journal, Outbox, ProjectionStore, atomic::AtomicStorage, + status::StorageStatusProvider, +}; use radroots_transport::{EventSink, EventSource}; use crate::Engine; const MAX_OPERATION_TIMEOUT_MS: u64 = 86_400_000; +/// Exact backend-neutral storage capability required by sync orchestration. +pub trait SyncStorage: + EventStore + Journal + Outbox + ProjectionStore + AtomicStorage + StorageStatusProvider +{ +} + +impl<T> SyncStorage for T where + T: EventStore + Journal + Outbox + ProjectionStore + AtomicStorage + StorageStatusProvider +{ +} + /// Sync operation class used for identity and deadline policy. #[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] #[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] @@ -115,7 +129,7 @@ const fn valid_timeout(value: u64) -> bool { /// Builder for an [`Engine`] with explicit optional transport capabilities. pub struct EngineBuilder { - storage: Arc<dyn Storage>, + storage: Arc<dyn SyncStorage>, source: Option<Arc<dyn EventSource>>, sink: Option<Arc<dyn EventSink>>, signer: Option<Arc<dyn Signer>>, @@ -126,7 +140,7 @@ pub struct EngineBuilder { impl EngineBuilder { pub(crate) fn new( - storage: Arc<dyn Storage>, + storage: Arc<dyn SyncStorage>, clock: Arc<dyn Clock>, ids: Arc<dyn IdSource>, deadlines: DeadlinePolicy, diff --git a/crates/sync/src/status.rs b/crates/sync/src/status.rs @@ -8,10 +8,13 @@ use radroots_protocol::runtime::v1::{ }; use radroots_signing::{SignerStatus, status::SignerAvailability}; use radroots_storage::{ - BackupSource, EventStore, Outbox, ProjectionStore, + EventStore, Outbox, ProjectionStore, outbox::{OutboxRecord, OutboxStage, OutboxStatus}, projection::{ProjectionHealth, ProjectionId, ProjectionStatus}, - status::{EventStoreHealth, EventStoreStatus, IntegrityHealth, ShutdownState, StorageStatus}, + status::{ + EventStoreHealth, EventStoreStatus, IntegrityHealth, ShutdownState, StorageStatus, + StorageStatusProvider, + }, }; use radroots_transport::{SinkStatus, SourceStatus, capability::Availability}; @@ -168,7 +171,7 @@ impl Engine { { return Err(Error::InvalidStatusRequest); } - let storage = BackupSource::status(self.storage.as_ref()) + let storage = StorageStatusProvider::storage_status(self.storage.as_ref()) .await .map_err(|_| Error::StorageFailed)?; let events = EventStore::status(self.storage.as_ref()) diff --git a/crates/sync/tests/engine_composition.rs b/crates/sync/tests/engine_composition.rs @@ -3,12 +3,10 @@ use std::sync::Arc; use futures_executor::block_on; use radroots_protocol::runtime::v1::{OPERATION_SCHEMA_VERSION, SyncCapabilityState, SyncHealth}; use radroots_signing::{Error as SigningError, SignReceipt, SignRequest, Signer, SignerStatus}; -use radroots_storage::{ - Storage, event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId, -}; +use radroots_storage::{event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId}; use radroots_sync::{ Engine, - policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, }; use radroots_transport::{ DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, EventSource, FetchPage, @@ -24,7 +22,7 @@ struct FixedClock; struct FixedIds; type TestDependencies = ( - Arc<dyn Storage>, + Arc<dyn SyncStorage>, Arc<dyn Clock>, Arc<dyn IdSource>, DeadlinePolicy, diff --git a/crates/sync/tests/ingest.rs b/crates/sync/tests/ingest.rs @@ -6,14 +6,14 @@ use std::sync::{ use futures_executor::block_on; use radroots_event::{SignedEvent, draft::SignedEventParts}; use radroots_storage::{ - EventStore, Storage, + EventStore, event::{AdmissionDisposition, AdmissionStage, EventQuery, EventQueryBounds, SourceGeneration}, memory::MemoryStorage, }; use radroots_sync::{ Engine, ingest::{AdmissionDecision, AdmissionPolicy, RegistryPolicy}, - policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, }; use radroots_transport::{ Error as TransportError, EventSource, FetchPage, FetchRequest, SourceStatus, Target, @@ -91,7 +91,7 @@ fn setup_engine_with_ids(ids: Arc<dyn IdSource>) -> (Engine, Arc<MemoryStorage>) let storage = Arc::new(MemoryStorage::new( SourceGeneration::new([7; 32]).expect("source generation"), )); - let storage_capability: Arc<dyn Storage> = storage.clone(); + let storage_capability: Arc<dyn SyncStorage> = storage.clone(); let engine = Engine::builder( storage_capability, Arc::new(FixedClock), diff --git a/crates/sync/tests/projection.rs b/crates/sync/tests/projection.rs @@ -12,14 +12,14 @@ use radroots_event::{ wire::compute_canonical_nip01_event_id, }; use radroots_storage::{ - EventStore, ProjectionStore, Storage, + EventStore, ProjectionStore, event::{EventAdmission, SourceGeneration, StoredVisibleEvent}, memory::MemoryStorage, projection::{ProjectionGeneration, ProjectionHealth, ProjectionId}, }; use radroots_sync::{ Engine, - policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, projection::{Reducer, ReducerError, RefreshKind, RefreshRequest, RefreshState}, }; use radroots_transport::{ @@ -131,7 +131,7 @@ fn setup() -> (Engine, Arc<MemoryStorage>, ProjectionId) { let storage = Arc::new(MemoryStorage::new( SourceGeneration::new([9; 32]).expect("generation"), )); - let storage_capability: Arc<dyn Storage> = storage.clone(); + let storage_capability: Arc<dyn SyncStorage> = storage.clone(); let engine = Engine::builder( storage_capability, Arc::new(TestClock(AtomicU64::new(1_000))), diff --git a/crates/sync/tests/pull.rs b/crates/sync/tests/pull.rs @@ -5,11 +5,11 @@ use std::{ use futures_executor::block_on; use radroots_event::{SignedEvent, draft::SignedEventParts}; -use radroots_storage::{Storage, event::SourceGeneration, memory::MemoryStorage}; +use radroots_storage::{event::SourceGeneration, memory::MemoryStorage}; use radroots_sync::{ Engine, PullRequest, ingest::RegistryPolicy, - policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, pull::{PULL_MAX_PAGES, PullTermination}, }; use radroots_transport::{ @@ -173,7 +173,7 @@ fn observed(signature: &str, observed_at: u64) -> ObservedEvent { } fn engine(source: Arc<ScriptedSource>, clock: Arc<dyn Clock>, timeout_ms: u64) -> Engine { - let storage: Arc<dyn Storage> = Arc::new(MemoryStorage::new( + let storage: Arc<dyn SyncStorage> = Arc::new(MemoryStorage::new( SourceGeneration::new([8; 32]).expect("generation"), )); Engine::builder( diff --git a/crates/sync/tests/push_enqueue.rs b/crates/sync/tests/push_enqueue.rs @@ -15,7 +15,7 @@ use radroots_signing::{ actor::ActorSource, error::Kind as SigningErrorKind, request::CancellationPolicy, }; use radroots_storage::{ - EventStore, Journal, Outbox, Storage, + EventStore, Journal, Outbox, event::{EventQuery, EventQueryBounds, SourceGeneration}, journal::{IdempotencyKey, JournalStage, OperationInstanceId}, memory::MemoryStorage, @@ -23,7 +23,7 @@ use radroots_storage::{ }; use radroots_sync::{ Engine, PushRequest, - policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId}, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, push::DeliveryRunRequest, }; use radroots_transport::{ @@ -291,7 +291,7 @@ fn setup_engine_with_sink( let storage = Arc::new(MemoryStorage::new( SourceGeneration::new([6; 32]).expect("generation"), )); - let capability: Arc<dyn Storage> = storage.clone(); + let capability: Arc<dyn SyncStorage> = storage.clone(); let clock = Arc::new(TestClock(AtomicU64::new(1_800_000_200_000))); let engine = Engine::builder( capability, diff --git a/crates/sync/tests/reliability_scenarios.rs b/crates/sync/tests/reliability_scenarios.rs @@ -0,0 +1,273 @@ +use std::sync::{ + Arc, + atomic::{AtomicU64, AtomicUsize, Ordering}, +}; + +use futures_executor::block_on; +use radroots_event::{SignedEvent, wire::Nip01EventWire}; +use radroots_protocol::runtime::v1::OperationId; +use radroots_storage::{ + Journal, Outbox, + event::SourceGeneration, + journal::{IdempotencyDigest, IdempotencyKey, OperationInstanceId, PrepareOperation}, + memory::MemoryStorage, + outbox::{DeliveryPlanDigest, EnqueueOutboxItem, LeaseOwner, OutboxItemId, OutboxStage}, +}; +use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage}; +use radroots_sync::{ + Engine, + policy::{Clock, DeadlinePolicy, Error, IdSource, OperationKind, SyncId, SyncStorage}, + push::DeliveryRunRequest, +}; +use radroots_transport::{ + DeliveryReceipt, DeliveryRequest, Error as TransportError, EventSink, SinkStatus, Target, + TargetSet, TransportId, + capability::{Availability, Maturity, SinkCapabilities}, + outcome::DeliveryOutcome, + policy::{SatisfactionClass, SatisfactionPolicy, TargetPolicy}, + sink::{DeliveryPayload, DeliveryTargetReceipt}, +}; + +struct SequenceClock(AtomicU64); + +impl Clock for SequenceClock { + fn now_unix_ms(&self) -> Result<u64, Error> { + Ok(self.0.fetch_add(1, Ordering::Relaxed)) + } +} + +struct SequenceIds(AtomicU64); + +impl IdSource for SequenceIds { + fn next_id(&self, _operation: OperationKind) -> Result<SyncId, Error> { + let byte = u8::try_from(self.0.fetch_add(1, Ordering::Relaxed)) + .map_err(|_| Error::InvalidSyncId)?; + SyncId::new([byte; 16]) + } +} + +struct RecoverySink(AtomicUsize); + +impl EventSink for RecoverySink { + fn status(&self) -> radroots_transport::BoxFuture<'_, Result<SinkStatus, TransportError>> { + Box::pin(async { + Ok(SinkStatus::new( + TransportId::NOSTR, + true, + Maturity::Stable, + Availability::Available, + SinkCapabilities::DELIVER, + "ready", + )) + }) + } + + fn deliver( + &self, + request: DeliveryRequest, + ) -> radroots_transport::BoxFuture<'_, Result<DeliveryReceipt, TransportError>> { + let call = self.0.fetch_add(1, Ordering::Relaxed); + Box::pin(async move { + if call == 0 { + return Err(TransportError::UnsupportedOperation); + } + DeliveryReceipt::for_request( + &request, + request + .target_set() + .targets() + .iter() + .cloned() + .map(|target| { + DeliveryTargetReceipt::attempted(target, DeliveryOutcome::accepted()) + }) + .collect(), + ) + }) + } +} + +fn signed_event() -> SignedEvent { + let mut wire = Nip01EventWire { + id: "0".repeat(64), + pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), + created_at: 1_800_000_100, + kind: 0, + tags: vec![], + content: "reliability-scenario".to_owned(), + sig: "42".repeat(64), + extra: Default::default(), + }; + wire.id = wire.computed_event_id().expect("event id").to_hex(); + let raw = serde_json::json!({ + "id": &wire.id, + "pubkey": &wire.pubkey, + "created_at": wire.created_at, + "kind": wire.kind, + "tags": &wire.tags, + "content": &wire.content, + "sig": &wire.sig, + }) + .to_string(); + SignedEvent::from_wire_verified_id(wire, raw).expect("signed event") +} + +fn enqueue_request() -> EnqueueOutboxItem { + let target = Target::new(TransportId::NOSTR, "wss://reliability.example").expect("target"); + let request = DeliveryRequest::new( + "sync-reliability", + DeliveryPayload::new(signed_event()), + TargetSet::new(vec![target]).expect("targets"), + SatisfactionPolicy::new(SatisfactionClass::Accepted, TargetPolicy::all()), + 10_000, + ) + .expect("request"); + EnqueueOutboxItem::new( + OutboxItemId::new([5; 16]).expect("item id"), + OperationInstanceId::new([6; 16]).expect("operation id"), + DeliveryPlanDigest::new([7; 32]), + request, + 10, + ) + .expect("enqueue") +} + +fn engine( + storage: Arc<dyn SyncStorage>, + sink: Arc<RecoverySink>, + clock_start: u64, + id_start: u64, +) -> Engine { + Engine::builder( + storage, + Arc::new(SequenceClock(AtomicU64::new(clock_start))), + Arc::new(SequenceIds(AtomicU64::new(id_start))), + DeadlinePolicy::new(1_000, 1_000, 1_000).expect("deadlines"), + ) + .sink(sink) + .build() + .expect("engine") +} + +fn run(seed: u8) -> DeliveryRunRequest { + DeliveryRunRequest::new( + LeaseOwner::parse("sync-reliability").expect("owner"), + SyncId::new([seed; 16]).expect("seed"), + 100, + 1, + ) + .expect("run") +} + +async fn retry_scenario(storage: Arc<dyn SyncStorage>) { + Journal::prepare( + storage.as_ref(), + PrepareOperation::new( + OperationInstanceId::new([6; 16]).expect("operation id"), + OperationId::SyncPush, + IdempotencyKey::parse("sync-reliability").expect("idempotency key"), + IdempotencyDigest::new([8; 32]), + 9, + ) + .expect("prepare operation"), + ) + .await + .expect("prepare"); + Outbox::enqueue(storage.as_ref(), enqueue_request()) + .await + .expect("enqueue"); + let sync = engine( + storage, + Arc::new(RecoverySink(AtomicUsize::new(0))), + 100, + 20, + ); + let first = sync.deliver_pending(run(30)).await.expect("first attempt"); + assert_eq!( + first.outcomes()[0].as_ref().expect("durable retry").stage(), + OutboxStage::Retryable + ); + let second = sync.deliver_pending(run(31)).await.expect("retry"); + assert_eq!( + second.outcomes()[0] + .as_ref() + .expect("durable success") + .stage(), + OutboxStage::Satisfied + ); + let status = sync.status(&[]).await.expect("status").to_protocol(); + assert_eq!(status.outbox.satisfied, 1); +} + +#[test] +fn memory_runs_the_durable_retry_scenario() { + let storage: Arc<dyn SyncStorage> = Arc::new(MemoryStorage::new( + SourceGeneration::new([9; 32]).expect("generation"), + )); + block_on(retry_scenario(storage)); +} + +#[tokio::test] +async fn sqlite_recovers_retryable_work_after_reopen() { + let directory = tempfile::tempdir().expect("directory"); + let paths = Paths::from_directory(directory.path()).expect("paths"); + { + let store = Arc::new( + SqliteStorage::open( + OpenOptions::new(paths.clone(), OpenMode::Create) + .with_source_generation(SourceGeneration::new([10; 32]).expect("generation"), 1) + .expect("source generation"), + ) + .await + .expect("open"), + ); + Journal::prepare( + store.as_ref(), + PrepareOperation::new( + OperationInstanceId::new([6; 16]).expect("operation id"), + OperationId::SyncPush, + IdempotencyKey::parse("sync-reliability").expect("idempotency key"), + IdempotencyDigest::new([8; 32]), + 9, + ) + .expect("prepare operation"), + ) + .await + .expect("prepare"); + Outbox::enqueue(store.as_ref(), enqueue_request()) + .await + .expect("enqueue"); + let storage: Arc<dyn SyncStorage> = store; + let sync = engine( + storage, + Arc::new(RecoverySink(AtomicUsize::new(0))), + 100, + 20, + ); + let first = sync.deliver_pending(run(40)).await.expect("first attempt"); + assert_eq!( + first.outcomes()[0].as_ref().expect("retryable").stage(), + OutboxStage::Retryable + ); + } + let store = Arc::new( + SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting)) + .await + .expect("reopen"), + ); + let storage: Arc<dyn SyncStorage> = store; + let sync = engine( + storage, + Arc::new(RecoverySink(AtomicUsize::new(1))), + 200, + 40, + ); + let recovered = sync + .deliver_pending(run(41)) + .await + .expect("recovered retry"); + assert_eq!( + recovered.outcomes()[0].as_ref().expect("satisfied").stage(), + OutboxStage::Satisfied + ); +}