lib

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

commit 144e8bdb06d67d19f915cd55ac1fc0c5e89460a9
parent 01b77455004e7a9003df8005ab966bc288a6d81c
Author: triesap <tyson@radroots.org>
Date:   Sat,  1 Aug 2026 19:21:53 +0000

storage: define the canonical event storage SPI

- define monotonic raw verified and visible admissions
- add generation-aware bounded event and provenance queries
- expose dyn-safe backend-neutral event store operations
- prove transitions and queries with a memory conformance harness

Diffstat:
MCargo.lock | 2++
Mcrates/storage/Cargo.toml | 12+++++++++++-
Mcrates/storage/src/event.rs | 557+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage/src/lib.rs | 2++
Mcrates/storage/src/status.rs | 81+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage/tests/event_store.rs | 425+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
6 files changed, 1078 insertions(+), 1 deletion(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -5123,11 +5123,13 @@ dependencies = [ name = "radroots_storage" version = "0.1.0-alpha" dependencies = [ + "futures-executor", "radroots_event", "radroots_protocol", "radroots_trade", "radroots_transport", "serde", + "serde_json", ] [[package]] diff --git a/crates/storage/Cargo.toml b/crates/storage/Cargo.toml @@ -17,7 +17,13 @@ name = "radroots_storage" [features] default = ["memory", "serde"] memory = [] -serde = ["dep:serde"] +serde = [ + "dep:serde", + "radroots_event/serde", + "radroots_protocol/serde", + "radroots_trade/serde", + "radroots_transport/serde", +] [dependencies] radroots_event = { workspace = true, default-features = false } @@ -29,5 +35,9 @@ serde = { workspace = true, default-features = false, features = [ "derive", ], optional = true } +[dev-dependencies] +futures-executor = { workspace = true } +serde_json = { workspace = true, features = ["std"] } + [lints] workspace = true diff --git a/crates/storage/src/event.rs b/crates/storage/src/event.rs @@ -1 +1,558 @@ //! Canonical event persistence contracts. + +use core::fmt; +use radroots_event::{EventId, SignedEvent, VerifiedEvent, admission::VisibleEvent}; +use radroots_transport::{ + BoxFuture, + source::{EventProvenance, ObservedEvent}, +}; +use std::collections::BTreeSet; + +use crate::status::EventStoreStatus; + +/// Maximum events returned by one storage query. +pub const EVENT_QUERY_LIMIT_MAX: u16 = 1_000; +/// Maximum explicit event identifiers in one storage query. +pub const EVENT_QUERY_ID_MAX: usize = 256; +/// Opaque identity of one append-only canonical event source. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct SourceGeneration([u8; 32]); + +impl SourceGeneration { + /// Creates a generation from host-provided entropy. + pub const fn new(bytes: [u8; 32]) -> Result<Self, Error> { + if is_all_zero(&bytes) { + return Err(Error::InvalidSourceGeneration); + } + Ok(Self(bytes)) + } + + /// Returns the opaque generation bytes. + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +const fn is_all_zero(bytes: &[u8; 32]) -> bool { + let mut index = 0; + while index < bytes.len() { + if bytes[index] != 0 { + return false; + } + index += 1; + } + true +} + +/// Non-zero sequence within one source generation. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] +pub struct EventSequence(u64); + +impl EventSequence { + /// Creates a non-zero source-local sequence. + pub const fn new(value: u64) -> Result<Self, Error> { + if value == 0 { + return Err(Error::InvalidEventSequence); + } + Ok(Self(value)) + } + + /// Returns the source-local sequence. + pub const fn get(self) -> u64 { + self.0 + } +} + +/// Stable location of an event within one source generation. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct EventPosition { + generation: SourceGeneration, + sequence: EventSequence, +} + +impl EventPosition { + /// Creates a source position. + pub const fn new(generation: SourceGeneration, sequence: EventSequence) -> Self { + Self { + generation, + sequence, + } + } + + /// Returns the source generation. + pub const fn generation(self) -> SourceGeneration { + self.generation + } + + /// Returns the generation-local sequence. + pub const fn sequence(self) -> EventSequence { + self.sequence + } +} + +/// Cursor after which a query resumes. +pub type EventCursor = EventPosition; + +/// Validated bounds for a canonical event query. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub struct EventQueryBounds { + limit: u16, + after: Option<EventCursor>, +} + +impl EventQueryBounds { + /// Creates bounds for a first-page query. + pub const fn first(limit: u16) -> Result<Self, Error> { + if limit == 0 || limit > EVENT_QUERY_LIMIT_MAX { + return Err(Error::InvalidEventQueryLimit); + } + Ok(Self { limit, after: None }) + } + + /// Resumes strictly after a prior cursor. + #[must_use] + pub const fn after(mut self, cursor: EventCursor) -> Self { + self.after = Some(cursor); + self + } + + /// Returns the maximum number of records. + pub const fn limit(self) -> u16 { + self.limit + } + + /// Returns the optional exclusive cursor. + pub const fn cursor(self) -> Option<EventCursor> { + self.after + } +} + +/// Bounded event selection; an empty identifier set selects every event. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EventQuery { + bounds: EventQueryBounds, + event_ids: Vec<EventId>, +} + +impl EventQuery { + /// Selects all events under the supplied bounds. + pub const fn all(bounds: EventQueryBounds) -> Self { + Self { + bounds, + event_ids: Vec::new(), + } + } + + /// Selects a bounded, duplicate-free set of event identifiers. + pub fn for_ids(bounds: EventQueryBounds, event_ids: Vec<EventId>) -> Result<Self, Error> { + if event_ids.is_empty() { + return Err(Error::EmptyEventQueryIds); + } + if event_ids.len() > EVENT_QUERY_ID_MAX { + return Err(Error::TooManyEventQueryIds); + } + let unique = event_ids.iter().collect::<BTreeSet<_>>(); + if unique.len() != event_ids.len() { + return Err(Error::DuplicateEventQueryId); + } + Ok(Self { bounds, event_ids }) + } + + /// Returns the query bounds. + pub const fn bounds(&self) -> EventQueryBounds { + self.bounds + } + + /// Returns the selected identifiers; empty means all. + pub fn event_ids(&self) -> &[EventId] { + self.event_ids.as_slice() + } + + /// Reports whether an identifier is selected. + pub fn selects(&self, event_id: &EventId) -> bool { + self.event_ids.is_empty() || self.event_ids.contains(event_id) + } +} + +/// Durable event admission stage. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[cfg_attr(feature = "serde", serde(rename_all = "snake_case"))] +#[derive(Clone, Copy, Debug, Eq, Ord, PartialEq, PartialOrd)] +pub enum AdmissionStage { + /// Structurally valid and ID-checked, but not signature verified. + Raw, + /// Canonical identifier and signature verified. + Verified, + /// Contract-admitted and visibility-authorized. + Visible, +} + +/// One canonical event admission with exact transport provenance. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EventAdmission { + observed: ObservedEvent, + state: AdmissionState, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +enum AdmissionState { + Raw, + Verified(VerifiedEvent), + Visible(VisibleEvent), +} + +impl EventAdmission { + /// Retains an observed signed event without claiming signature verification. + pub const fn raw(observed: ObservedEvent) -> Self { + Self { + observed, + state: AdmissionState::Raw, + } + } + + /// Retains a verified event after proving it matches the observed payload. + pub fn verified(observed: ObservedEvent, verified: VerifiedEvent) -> Result<Self, Error> { + if observed.event().envelope() != verified.event() { + return Err(Error::AdmissionEventMismatch); + } + Ok(Self { + observed, + state: AdmissionState::Verified(verified), + }) + } + + /// Retains a visible event after proving it matches the observed payload. + pub fn visible(observed: ObservedEvent, visible: VisibleEvent) -> Result<Self, Error> { + if observed.event().envelope() != visible.event() { + return Err(Error::AdmissionEventMismatch); + } + Ok(Self { + observed, + state: AdmissionState::Visible(visible), + }) + } + + /// Returns the durable stage represented by this admission. + pub const fn stage(&self) -> AdmissionStage { + match self.state { + AdmissionState::Raw => AdmissionStage::Raw, + AdmissionState::Verified(_) => AdmissionStage::Verified, + AdmissionState::Visible(_) => AdmissionStage::Visible, + } + } + + /// Returns the exact observed signed event. + pub const fn event(&self) -> &SignedEvent { + self.observed.event() + } + + /// Returns the event identifier. + pub fn event_id(&self) -> &EventId { + self.event().id() + } + + /// Returns the transport observation attached to this admission. + pub const fn provenance(&self) -> &EventProvenance { + self.observed.provenance() + } + + /// Returns the verified event when this admission reached verification. + pub const fn verified_event(&self) -> Option<&VerifiedEvent> { + match &self.state { + AdmissionState::Raw => None, + AdmissionState::Verified(event) => Some(event), + AdmissionState::Visible(event) => { + Some(event.admitted_event().validated_event().verified_event()) + } + } + } + + /// Returns the visible event when visibility was authorized. + pub const fn visible_event(&self) -> Option<&VisibleEvent> { + match &self.state { + AdmissionState::Visible(event) => Some(event), + AdmissionState::Raw | AdmissionState::Verified(_) => None, + } + } +} + +/// Persistence result for one idempotent admission. +#[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 AdmissionDisposition { + Inserted, + Advanced, + Duplicate, +} + +/// Request-bound durable admission receipt. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct AdmissionReceipt { + event_id: EventId, + position: EventPosition, + stage: AdmissionStage, + disposition: AdmissionDisposition, +} + +impl AdmissionReceipt { + /// Creates a backend receipt from validated durable state. + pub const fn new( + event_id: EventId, + position: EventPosition, + stage: AdmissionStage, + disposition: AdmissionDisposition, + ) -> Self { + Self { + event_id, + position, + stage, + disposition, + } + } + + pub const fn event_id(&self) -> &EventId { + &self.event_id + } + + pub const fn position(&self) -> EventPosition { + self.position + } + + pub const fn stage(&self) -> AdmissionStage { + self.stage + } + + pub const fn disposition(&self) -> AdmissionDisposition { + self.disposition + } +} + +/// Raw event returned from canonical storage. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct StoredRawEvent { + position: EventPosition, + event: SignedEvent, + stage: AdmissionStage, +} + +impl StoredRawEvent { + pub const fn new(position: EventPosition, event: SignedEvent, stage: AdmissionStage) -> Self { + Self { + position, + event, + stage, + } + } + + pub const fn position(&self) -> EventPosition { + self.position + } + + pub const fn event(&self) -> &SignedEvent { + &self.event + } + + pub const fn stage(&self) -> AdmissionStage { + self.stage + } +} + +/// Signature-verified event returned from canonical storage. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct StoredVerifiedEvent { + position: EventPosition, + event: VerifiedEvent, +} + +impl StoredVerifiedEvent { + pub const fn new(position: EventPosition, event: VerifiedEvent) -> Self { + Self { position, event } + } + + pub const fn position(&self) -> EventPosition { + self.position + } + + pub const fn event(&self) -> &VerifiedEvent { + &self.event + } +} + +/// Visibility-authorized event returned from canonical storage. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct StoredVisibleEvent { + position: EventPosition, + event: VisibleEvent, +} + +impl StoredVisibleEvent { + pub const fn new(position: EventPosition, event: VisibleEvent) -> Self { + Self { position, event } + } + + pub const fn position(&self) -> EventPosition { + self.position + } + + pub const fn event(&self) -> &VisibleEvent { + &self.event + } +} + +/// One bounded, generation-consistent page. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EventPage<T> { + generation: SourceGeneration, + items: Vec<T>, + next: Option<EventCursor>, +} + +impl<T> EventPage<T> { + /// Creates a page and enforces its caller-supplied item bound. + pub fn new( + generation: SourceGeneration, + items: Vec<T>, + next: Option<EventCursor>, + bounds: EventQueryBounds, + ) -> Result<Self, Error> { + if items.len() > usize::from(bounds.limit()) { + return Err(Error::EventPageLimitExceeded); + } + if let Some(cursor) = next + && cursor.generation() != generation + { + return Err(Error::CursorGenerationMismatch); + } + Ok(Self { + generation, + items, + next, + }) + } + + pub const fn generation(&self) -> SourceGeneration { + self.generation + } + + pub fn items(&self) -> &[T] { + self.items.as_slice() + } + + pub const fn next_cursor(&self) -> Option<EventCursor> { + self.next + } +} + +/// Provenance retained for one event observation. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct StoredEventProvenance { + position: EventPosition, + provenance: EventProvenance, +} + +impl StoredEventProvenance { + pub const fn new(position: EventPosition, provenance: EventProvenance) -> Self { + Self { + position, + provenance, + } + } + + pub const fn position(&self) -> EventPosition { + self.position + } + + pub const fn provenance(&self) -> &EventProvenance { + &self.provenance + } +} + +/// Backend-neutral canonical event storage SPI. +/// +/// Implementations are dyn-compatible and return `Send` futures. They may not +/// expose backend transactions, handles, SQL, or filesystem paths. Admission is +/// idempotent for an identical event and provenance observation; a stage may +/// advance but never regress. Query cursors are bound to one source generation. +pub trait EventStore: Send + Sync { + /// Returns passive event-store status without initiating maintenance. + fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>>; + + /// Durably admits or advances one canonical event observation. + fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>>; + + /// Queries retained raw events. + fn query_raw( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>>; + + /// Queries signature-verified events. + fn query_verified( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>>; + + /// Queries visibility-authorized events. + fn query_visible( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>>; + + /// Queries bounded provenance for one event. + fn query_provenance( + &self, + event_id: EventId, + bounds: EventQueryBounds, + ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>>; +} + +/// Stable, secret-safe event-storage failure. +#[non_exhaustive] +#[derive(Clone, Debug, Eq, PartialEq)] +pub enum Error { + InvalidSourceGeneration, + InvalidEventSequence, + InvalidEventQueryLimit, + EmptyEventQueryIds, + TooManyEventQueryIds, + DuplicateEventQueryId, + AdmissionEventMismatch, + AdmissionRegression, + EventConflict, + EventPageLimitExceeded, + CursorGenerationMismatch, + SourceGenerationChanged, + EventNotFound, + CorruptStoredEvent, + BackendUnavailable, +} + +impl fmt::Display for Error { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::InvalidSourceGeneration => "storage source generation is invalid", + Self::InvalidEventSequence => "storage event sequence is invalid", + Self::InvalidEventQueryLimit => "storage event query limit is invalid", + Self::EmptyEventQueryIds => "storage event id query is empty", + Self::TooManyEventQueryIds => "storage event id query exceeds its limit", + Self::DuplicateEventQueryId => "storage event id query contains a duplicate", + Self::AdmissionEventMismatch => "storage admission event identities do not match", + Self::AdmissionRegression => "storage admission cannot regress event state", + Self::EventConflict => "storage contains conflicting data for the event id", + Self::EventPageLimitExceeded => "storage event page exceeds its requested limit", + Self::CursorGenerationMismatch => "storage cursor belongs to another source generation", + Self::SourceGenerationChanged => "storage source generation changed", + Self::EventNotFound => "storage event was not found", + Self::CorruptStoredEvent => "storage event data is corrupt", + Self::BackendUnavailable => "storage backend is unavailable", + }) + } +} + +impl std::error::Error for Error {} diff --git a/crates/storage/src/lib.rs b/crates/storage/src/lib.rs @@ -12,3 +12,5 @@ pub mod outbox; pub mod private_artifact; pub mod projection; pub mod status; + +pub use event::{Error, EventStore}; diff --git a/crates/storage/src/status.rs b/crates/storage/src/status.rs @@ -1 +1,82 @@ //! Storage capability, health, and integrity status contracts. + +use crate::event::{Error, SourceGeneration}; + +/// Current event-store operating mode. +#[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 EventStoreMode { + ReadOnly, + ReadWrite, +} + +/// Current health of the canonical event source. +#[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 EventStoreHealth { + Available, + Degraded, + Unavailable, +} + +/// Passive event-store capability and cardinality report. +#[cfg_attr(feature = "serde", derive(serde::Serialize, serde::Deserialize))] +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct EventStoreStatus { + generation: SourceGeneration, + mode: EventStoreMode, + health: EventStoreHealth, + raw_events: u64, + verified_events: u64, + visible_events: u64, +} + +impl EventStoreStatus { + /// Creates a consistent event-store status report. + pub const fn new( + generation: SourceGeneration, + mode: EventStoreMode, + health: EventStoreHealth, + raw_events: u64, + verified_events: u64, + visible_events: u64, + ) -> Result<Self, Error> { + if verified_events > raw_events || visible_events > verified_events { + return Err(Error::CorruptStoredEvent); + } + Ok(Self { + generation, + mode, + health, + raw_events, + verified_events, + visible_events, + }) + } + + pub const fn generation(&self) -> SourceGeneration { + self.generation + } + + pub const fn mode(&self) -> EventStoreMode { + self.mode + } + + pub const fn health(&self) -> EventStoreHealth { + self.health + } + + pub const fn raw_events(&self) -> u64 { + self.raw_events + } + + pub const fn verified_events(&self) -> u64 { + self.verified_events + } + + pub const fn visible_events(&self) -> u64 { + self.visible_events + } +} diff --git a/crates/storage/tests/event_store.rs b/crates/storage/tests/event_store.rs @@ -0,0 +1,425 @@ +use futures_executor::block_on; +use radroots_event::{ + EventId, SignedEvent, VerifiedEvent, + admission::{AdmissionPolicy, RawEvent, VisibilityPolicy, VisibleEvent}, + wire::Nip01EventWire, +}; +use radroots_storage::{ + EventStore, + event::{ + AdmissionDisposition, AdmissionReceipt, AdmissionStage, Error, EventAdmission, EventPage, + EventPosition, EventQuery, EventQueryBounds, EventSequence, SourceGeneration, + StoredEventProvenance, StoredRawEvent, StoredVerifiedEvent, StoredVisibleEvent, + }, + status::{EventStoreHealth, EventStoreMode, EventStoreStatus}, +}; +use radroots_transport::{ + BoxFuture, Target, TransportId, + source::{EventProvenance, ObservedEvent}, +}; +use std::sync::Mutex; + +#[derive(Clone)] +struct Entry { + position: EventPosition, + admission: EventAdmission, + provenance: Vec<EventProvenance>, +} + +struct MemoryEventStore { + generation: SourceGeneration, + entries: Mutex<Vec<Entry>>, +} + +impl MemoryEventStore { + fn new() -> Self { + Self { + generation: SourceGeneration::new([7; 32]).expect("non-zero generation"), + entries: Mutex::new(Vec::new()), + } + } + + fn selected(&self, query: &EventQuery) -> Result<Vec<Entry>, Error> { + if let Some(cursor) = query.bounds().cursor() + && cursor.generation() != self.generation + { + return Err(Error::SourceGenerationChanged); + } + let after = query + .bounds() + .cursor() + .map_or(0, |cursor| cursor.sequence().get()); + Ok(self + .entries + .lock() + .expect("test store lock") + .iter() + .filter(|entry| { + entry.position.sequence().get() > after && query.selects(entry.admission.event_id()) + }) + .take(usize::from(query.bounds().limit())) + .cloned() + .collect()) + } +} + +impl EventStore for MemoryEventStore { + fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> { + Box::pin(async move { + let entries = self.entries.lock().expect("test store lock"); + let raw = entries.len() as u64; + let verified = entries + .iter() + .filter(|entry| entry.admission.stage() >= AdmissionStage::Verified) + .count() as u64; + let visible = entries + .iter() + .filter(|entry| entry.admission.stage() == AdmissionStage::Visible) + .count() as u64; + EventStoreStatus::new( + self.generation, + EventStoreMode::ReadWrite, + EventStoreHealth::Available, + raw, + verified, + visible, + ) + }) + } + + fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> { + Box::pin(async move { + let mut entries = self.entries.lock().expect("test store lock"); + if let Some(entry) = entries + .iter_mut() + .find(|entry| entry.admission.event_id() == admission.event_id()) + { + if entry.admission.event() != admission.event() { + return Err(Error::EventConflict); + } + if admission.stage() < entry.admission.stage() { + return Err(Error::AdmissionRegression); + } + let disposition = if admission.stage() == entry.admission.stage() { + AdmissionDisposition::Duplicate + } else { + AdmissionDisposition::Advanced + }; + if !entry.provenance.contains(admission.provenance()) { + entry.provenance.push(admission.provenance().clone()); + } + entry.admission = admission; + return Ok(AdmissionReceipt::new( + *entry.admission.event_id(), + entry.position, + entry.admission.stage(), + disposition, + )); + } + + let sequence = EventSequence::new(entries.len() as u64 + 1)?; + let position = EventPosition::new(self.generation, sequence); + let receipt = AdmissionReceipt::new( + *admission.event_id(), + position, + admission.stage(), + AdmissionDisposition::Inserted, + ); + let provenance = vec![admission.provenance().clone()]; + entries.push(Entry { + position, + admission, + provenance, + }); + Ok(receipt) + }) + } + + fn query_raw( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> { + Box::pin(async move { + let items = self + .selected(&query)? + .into_iter() + .map(|entry| { + StoredRawEvent::new( + entry.position, + entry.admission.event().clone(), + entry.admission.stage(), + ) + }) + .collect(); + EventPage::new(self.generation, items, None, query.bounds()) + }) + } + + fn query_verified( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> { + Box::pin(async move { + let items = self + .selected(&query)? + .into_iter() + .filter_map(|entry| { + entry + .admission + .verified_event() + .cloned() + .map(|event| StoredVerifiedEvent::new(entry.position, event)) + }) + .collect(); + EventPage::new(self.generation, items, None, query.bounds()) + }) + } + + fn query_visible( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> { + Box::pin(async move { + let items = self + .selected(&query)? + .into_iter() + .filter_map(|entry| { + entry + .admission + .visible_event() + .cloned() + .map(|event| StoredVisibleEvent::new(entry.position, event)) + }) + .collect(); + EventPage::new(self.generation, items, None, query.bounds()) + }) + } + + fn query_provenance( + &self, + event_id: EventId, + bounds: EventQueryBounds, + ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> { + Box::pin(async move { + let entries = self.entries.lock().expect("test store lock"); + let entry = entries + .iter() + .find(|entry| entry.admission.event_id() == &event_id) + .ok_or(Error::EventNotFound)?; + let items = entry + .provenance + .iter() + .take(usize::from(bounds.limit())) + .cloned() + .map(|provenance| StoredEventProvenance::new(entry.position, provenance)) + .collect(); + EventPage::new(self.generation, items, None, bounds) + }) + } +} + +struct Allow; + +impl radroots_event::admission::SignatureVerifier for Allow { + fn verify_signature( + &self, + _event: &radroots_event::Event, + ) -> Result<(), radroots_event::Error> { + Ok(()) + } +} + +impl AdmissionPolicy for Allow { + type Error = core::convert::Infallible; + + fn policy_id(&self) -> &'static str { + "test.storage.admission.v1" + } + + fn admit( + &self, + _event: &radroots_event::admission::ContractValidatedEvent, + ) -> Result<(), Self::Error> { + Ok(()) + } +} + +impl VisibilityPolicy for Allow { + type Error = core::convert::Infallible; + + fn policy_id(&self) -> &'static str { + "test.storage.visibility.v1" + } + + fn make_visible( + &self, + _event: &radroots_event::admission::AdmittedEvent, + ) -> Result<(), Self::Error> { + Ok(()) + } +} + +fn signed_event_with_signature(signature_byte: &str) -> SignedEvent { + let mut wire = Nip01EventWire { + id: "0".repeat(64), + pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), + created_at: 1_800_000_100, + kind: 0, + tags: vec![], + content: "{\"display_name\":\"Moss Street Farm\",\"bot\":false}".to_owned(), + sig: signature_byte.repeat(64), + extra: Default::default(), + }; + wire.id = wire + .computed_event_id() + .expect("canonical event id") + .to_hex(); + let raw_json = serde_json::to_string(&wire).expect("event JSON"); + SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") +} + +fn signed_event() -> SignedEvent { + signed_event_with_signature("42") +} + +fn observed(event: SignedEvent, observed_at: u64) -> ObservedEvent { + let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("relay target"); + let provenance = EventProvenance::new( + TransportId::NOSTR, + target.fingerprint().clone(), + observed_at, + ) + .expect("provenance"); + ObservedEvent::new(event, provenance) +} + +fn verified(event: &SignedEvent) -> VerifiedEvent { + RawEvent::new(event.envelope().clone()) + .verify_id() + .expect("event id") + .verify_signature(&Allow) + .expect("signature") +} + +fn visible(event: &SignedEvent) -> VisibleEvent { + verified(event) + .validate_contract() + .expect("contract") + .admit_with(&Allow) + .expect("admission") + .make_visible_with(&Allow) + .expect("visibility") +} + +#[test] +fn event_store_is_dyn_compatible_and_enforces_monotonic_admission() { + let store = MemoryEventStore::new(); + let dynamic: &dyn EventStore = &store; + let event = signed_event(); + + let inserted = block_on(dynamic.admit(EventAdmission::raw(observed(event.clone(), 1)))) + .expect("raw insertion"); + assert_eq!(inserted.disposition(), AdmissionDisposition::Inserted); + assert_eq!(inserted.stage(), AdmissionStage::Raw); + + let advanced = block_on( + dynamic.admit( + EventAdmission::verified(observed(event.clone(), 2), verified(&event)) + .expect("verified admission"), + ), + ) + .expect("verified advancement"); + assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced); + assert_eq!(advanced.position(), inserted.position()); + + let visible_event = visible(&event); + let advanced = block_on( + dynamic.admit( + EventAdmission::visible(observed(event.clone(), 3), visible_event) + .expect("visible admission"), + ), + ) + .expect("visible advancement"); + assert_eq!(advanced.stage(), AdmissionStage::Visible); + + let duplicate = block_on( + dynamic.admit( + EventAdmission::visible(observed(event.clone(), 3), visible(&event)) + .expect("duplicate visible admission"), + ), + ) + .expect("idempotent duplicate"); + assert_eq!(duplicate.disposition(), AdmissionDisposition::Duplicate); + + assert_eq!( + block_on(dynamic.admit(EventAdmission::raw(observed(event, 4)))), + Err(Error::AdmissionRegression) + ); + + let conflicting = signed_event_with_signature("43"); + assert_eq!( + EventAdmission::verified(observed(conflicting.clone(), 5), verified(&signed_event())), + Err(Error::AdmissionEventMismatch) + ); + assert_eq!( + block_on(dynamic.admit(EventAdmission::raw(observed(conflicting, 6)))), + Err(Error::EventConflict) + ); +} + +#[test] +fn event_queries_preserve_stage_generation_bounds_and_provenance() { + let store = MemoryEventStore::new(); + let event = signed_event(); + let event_id = *event.id(); + block_on( + store.admit( + EventAdmission::visible(observed(event.clone(), 5), visible(&event)) + .expect("visible admission"), + ), + ) + .expect("insert visible event"); + + let bounds = EventQueryBounds::first(1).expect("query bounds"); + let query = EventQuery::for_ids(bounds, vec![event_id]).expect("id query"); + let raw = block_on(store.query_raw(query.clone())).expect("raw page"); + let verified = block_on(store.query_verified(query.clone())).expect("verified page"); + let visible = block_on(store.query_visible(query)).expect("visible page"); + let provenance = block_on(store.query_provenance(event_id, bounds)).expect("provenance page"); + + assert_eq!(raw.items().len(), 1); + assert_eq!(raw.items()[0].stage(), AdmissionStage::Visible); + assert_eq!(verified.items().len(), 1); + assert_eq!(visible.items().len(), 1); + assert_eq!(provenance.items().len(), 1); + assert_eq!(raw.generation(), store.generation); + + let status = block_on(store.status()).expect("status"); + assert_eq!(status.raw_events(), 1); + assert_eq!(status.verified_events(), 1); + assert_eq!(status.visible_events(), 1); +} + +#[test] +fn bounds_generations_and_status_reject_invalid_state() { + assert_eq!( + SourceGeneration::new([0; 32]), + Err(Error::InvalidSourceGeneration) + ); + assert_eq!(EventSequence::new(0), Err(Error::InvalidEventSequence)); + assert_eq!( + EventQueryBounds::first(0), + Err(Error::InvalidEventQueryLimit) + ); + assert_eq!( + EventStoreStatus::new( + SourceGeneration::new([1; 32]).expect("generation"), + EventStoreMode::ReadOnly, + EventStoreHealth::Degraded, + 1, + 2, + 0, + ), + Err(Error::CorruptStoredEvent) + ); +}