lib

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

commit 7c92ca8ec8948c9696474c5cd4d6fec15f094aef
parent aba7ccd97cffb3017cf81b96203134df160da0d7
Author: triesap <tyson@radroots.org>
Date:   Sun,  2 Aug 2026 17:20:55 +0000

storage-sqlite: port canonical event storage

- implement the backend-neutral event SPI over private unified SQLite state
- preserve exact provenance, monotonic admission, generation cursors, and bounded queries
- return signed events with durable stage evidence without forging authorization typestates
- verify corruption guards plus NIP-09 and source-maintenance compatibility vectors

Diffstat:
MCargo.lock | 3+++
Mcontracts/storage/runtime_schema_v1.toml | 42+++++++++++++++++++++++++++++++++++++++---
Mcrates/event_codec/src/codec.rs | 6++++++
Mcrates/event_codec/src/decode.rs | 20++++++++++++++++++++
Mcrates/storage/src/event.rs | 54+++++++++++++++++++++++++++++++++++++++++++-----------
Mcrates/storage/src/memory.rs | 16++++++----------
Mcrates/storage/tests/event_store.rs | 16++++++----------
Mcrates/storage_sqlite/Cargo.toml | 7+++++--
Acrates/storage_sqlite/src/event/mod.rs | 773+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/src/lib.rs | 3+++
Acrates/storage_sqlite/src/migration/runtime/0002_canonical_event_storage.up.sql | 44++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/src/migration/runtime/mod.rs | 103+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------------
Mcrates/storage_sqlite/tests/package_boundary.rs | 3+++
13 files changed, 1028 insertions(+), 62 deletions(-)

diff --git a/Cargo.lock b/Cargo.lock @@ -5136,10 +5136,13 @@ dependencies = [ name = "radroots_storage_sqlite" version = "0.1.0-alpha" dependencies = [ + "radroots_event", "radroots_event_codec", "radroots_secrets", "radroots_storage", + "radroots_transport", "serde", + "serde_json", "sha2", "sqlx", "tempfile", diff --git a/contracts/storage/runtime_schema_v1.toml b/contracts/storage/runtime_schema_v1.toml @@ -1,9 +1,9 @@ schema_version = 1 database = "runtime.sqlite" minimum_version = 1 -current_version = 1 -migration_name = "runtime_authority" -migration_sha256 = "3b869122dd5bd58f4a15e7a71fd1377879640b36cd496d7b3f15278ef1e128c9" +current_version = 2 +migration_name = "canonical_event_storage" +migration_sha256 = "35b036ba84eff7135665c4ae42fa8232d8bacd8115805b03011a3eb76423f4b8" forward_only = true raw_sql_public = false @@ -59,3 +59,39 @@ owned_objects = [ "radroots_runtime_source_generations_delete_guard", "radroots_runtime_source_generations_identity_guard", ] + +[[migrations]] +version = 2 +name = "canonical_event_storage" +sha256 = "35b036ba84eff7135665c4ae42fa8232d8bacd8115805b03011a3eb76423f4b8" + +owned_objects = [ + "radroots_runtime_atomic_commits", + "radroots_runtime_delivery_evidence", + "radroots_runtime_delivery_evidence_item_idx", + "radroots_runtime_event_index_checkpoints", + "radroots_runtime_event_index_manifests", + "radroots_runtime_event_index_shards", + "radroots_runtime_event_provenance", + "radroots_runtime_event_provenance_observed_idx", + "radroots_runtime_events", + "radroots_runtime_events_admission_idx", + "radroots_runtime_events_delete_guard", + "radroots_runtime_events_event_id_idx", + "radroots_runtime_events_raw_update_guard", + "radroots_runtime_journal_idempotency_idx", + "radroots_runtime_journal_operations", + "radroots_runtime_journal_recovery_idx", + "radroots_runtime_outbox_items", + "radroots_runtime_outbox_ready_idx", + "radroots_runtime_outbox_targets", + "radroots_runtime_projection_checkpoints", + "radroots_runtime_projection_invalidations", + "radroots_runtime_projection_rebuilds", + "radroots_runtime_projection_rebuilds_stage_idx", + "radroots_runtime_source_generations", + "radroots_runtime_source_generations_active_idx", + "radroots_runtime_source_generations_delete_guard", + "radroots_runtime_source_generations_identity_guard", + "radroots_runtime_source_generations_sequence_guard", +] diff --git a/crates/event_codec/src/codec.rs b/crates/event_codec/src/codec.rs @@ -22,6 +22,12 @@ impl Codec { decode::event(raw_json) } + /// Decodes canonical NIP-01 JSON and verifies its declared event ID. + #[cfg(feature = "json")] + pub fn decode_signed_event(raw_json: &str) -> Result<radroots_event::SignedEvent, DecodeError> { + decode::signed_event(raw_json) + } + /// Encodes a native event envelope as compact NIP-01 JSON. #[cfg(feature = "json")] pub fn encode_event( diff --git a/crates/event_codec/src/decode.rs b/crates/event_codec/src/decode.rs @@ -123,6 +123,7 @@ pub mod trade { #[cfg(feature = "json")] use radroots_event::{ + SignedEvent, admission::RawEvent, wire::{DEFAULT_RAW_JSON_MAX_BYTES, Nip01EventWire}, }; @@ -213,6 +214,25 @@ pub fn event(raw_json: &str) -> Result<RawEvent, DecodeError> { Ok(RawEvent::new(envelope)) } +/// Decodes compact NIP-01 JSON into an ID-verified signed event. +#[cfg(feature = "json")] +pub fn signed_event(raw_json: &str) -> Result<SignedEvent, DecodeError> { + let max = MAX_EVENT_JSON_BYTES; + let actual = raw_json.len(); + if actual > max { + return Err(DecodeError::InputTooLarge { max, actual }); + } + let wire = Nip01EventWire::parse_json(raw_json).map_err(map_wire_error)?; + SignedEvent::from_wire_verified_id(wire, raw_json).map_err(|error| match error { + radroots_event::draft::SignedEventError::Envelope(error) => { + DecodeError::InvalidEnvelope(error) + } + radroots_event::draft::SignedEventError::Wire(error) + | radroots_event::draft::SignedEventError::RawJson(error) => map_wire_error(error), + radroots_event::draft::SignedEventError::RawJsonMismatch => DecodeError::InvalidJson, + }) +} + #[cfg(feature = "json")] fn map_wire_error(error: EventWireError) -> DecodeError { match error { diff --git a/crates/storage/src/event.rs b/crates/storage/src/event.rs @@ -1,9 +1,12 @@ //! Canonical event persistence contracts. -use radroots_event::{EventId, SignedEvent, VerifiedEvent, admission::VisibleEvent}; +pub use radroots_event::EventId; +use radroots_event::{SignedEvent, VerifiedEvent, admission::VisibleEvent}; +pub use radroots_transport::BoxFuture; use radroots_transport::{ - BoxFuture, - source::{EventProvenance, ObservedEvent}, + TransportId, + source::{EventProvenance, FetchCursor, ObservedEvent}, + target::TargetFingerprint, }; use std::collections::BTreeSet; @@ -362,15 +365,19 @@ impl StoredRawEvent { } } -/// Signature-verified event returned from canonical storage. +/// Canonical event returned with durable signature-verification evidence. +/// +/// The signed event is returned rather than forging an in-memory verification +/// typestate from persisted bytes. [`EventStore`] guarantees that only records +/// durably admitted at [`AdmissionStage::Verified`] or later appear here. #[derive(Clone, Debug, Eq, PartialEq)] pub struct StoredVerifiedEvent { position: EventPosition, - event: VerifiedEvent, + event: SignedEvent, } impl StoredVerifiedEvent { - pub const fn new(position: EventPosition, event: VerifiedEvent) -> Self { + pub const fn new(position: EventPosition, event: SignedEvent) -> Self { Self { position, event } } @@ -378,20 +385,24 @@ impl StoredVerifiedEvent { self.position } - pub const fn event(&self) -> &VerifiedEvent { + pub const fn event(&self) -> &SignedEvent { &self.event } } -/// Visibility-authorized event returned from canonical storage. +/// Canonical event returned with durable visibility evidence. +/// +/// The signed event is returned rather than rerunning a host authorization +/// policy during a storage read. [`EventStore`] guarantees that only records +/// durably admitted at [`AdmissionStage::Visible`] appear here. #[derive(Clone, Debug, Eq, PartialEq)] pub struct StoredVisibleEvent { position: EventPosition, - event: VisibleEvent, + event: SignedEvent, } impl StoredVisibleEvent { - pub const fn new(position: EventPosition, event: VisibleEvent) -> Self { + pub const fn new(position: EventPosition, event: SignedEvent) -> Self { Self { position, event } } @@ -399,7 +410,7 @@ impl StoredVisibleEvent { self.position } - pub const fn event(&self) -> &VisibleEvent { + pub const fn event(&self) -> &SignedEvent { &self.event } } @@ -470,6 +481,27 @@ impl StoredEventProvenance { pub const fn provenance(&self) -> &EventProvenance { &self.provenance } + + /// Reconstructs validated backend-neutral provenance from durable fields. + pub fn from_stored_parts( + position: EventPosition, + transport_id: &str, + target_fingerprint: &str, + observed_at_unix_ms: u64, + cursor: Option<&str>, + ) -> Result<Self, Error> { + let transport_id = + TransportId::parse(transport_id).map_err(|_| Error::CorruptStoredEvent)?; + let target = + TargetFingerprint::parse(target_fingerprint).map_err(|_| Error::CorruptStoredEvent)?; + let mut provenance = EventProvenance::new(transport_id, target, observed_at_unix_ms) + .map_err(|_| Error::CorruptStoredEvent)?; + if let Some(cursor) = cursor { + provenance = provenance + .with_cursor(FetchCursor::parse(cursor).map_err(|_| Error::CorruptStoredEvent)?); + } + Ok(Self::new(position, provenance)) + } } /// Backend-neutral canonical event storage SPI. diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -411,11 +411,9 @@ impl EventStore for MemoryStorage { .selected(&state, &query)? .into_iter() .filter_map(|entry| { - entry - .admission - .verified_event() - .cloned() - .map(|event| StoredVerifiedEvent::new(entry.position, event)) + (entry.admission.stage() >= AdmissionStage::Verified).then(|| { + StoredVerifiedEvent::new(entry.position, entry.admission.event().clone()) + }) }) .collect(); EventPage::new(self.generation, items, None, query.bounds()) @@ -432,11 +430,9 @@ impl EventStore for MemoryStorage { .selected(&state, &query)? .into_iter() .filter_map(|entry| { - entry - .admission - .visible_event() - .cloned() - .map(|event| StoredVisibleEvent::new(entry.position, event)) + (entry.admission.stage() == AdmissionStage::Visible).then(|| { + StoredVisibleEvent::new(entry.position, entry.admission.event().clone()) + }) }) .collect(); EventPage::new(self.generation, items, None, query.bounds()) diff --git a/crates/storage/tests/event_store.rs b/crates/storage/tests/event_store.rs @@ -164,11 +164,9 @@ impl EventStore for MemoryEventStore { .selected(&query)? .into_iter() .filter_map(|entry| { - entry - .admission - .verified_event() - .cloned() - .map(|event| StoredVerifiedEvent::new(entry.position, event)) + (entry.admission.stage() >= AdmissionStage::Verified).then(|| { + StoredVerifiedEvent::new(entry.position, entry.admission.event().clone()) + }) }) .collect(); EventPage::new(self.generation, items, None, query.bounds()) @@ -184,11 +182,9 @@ impl EventStore for MemoryEventStore { .selected(&query)? .into_iter() .filter_map(|entry| { - entry - .admission - .visible_event() - .cloned() - .map(|event| StoredVisibleEvent::new(entry.position, event)) + (entry.admission.stage() == AdmissionStage::Visible).then(|| { + StoredVisibleEvent::new(entry.position, entry.admission.event().clone()) + }) }) .collect(); EventPage::new(self.generation, items, None, query.bounds()) diff --git a/crates/storage_sqlite/Cargo.toml b/crates/storage_sqlite/Cargo.toml @@ -15,14 +15,17 @@ publish = false name = "radroots_storage_sqlite" [dependencies] -radroots_event_codec = { workspace = true, default-features = false } +radroots_event_codec = { workspace = true, default-features = false, features = ["json", "std"] } radroots_secrets = { workspace = true, default-features = false } radroots_storage = { workspace = true, default-features = false } +sqlx = { workspace = true, features = ["runtime-tokio", "sqlite-bundled"] } [dev-dependencies] +radroots_event = { workspace = true, default-features = false, features = ["std"] } +radroots_transport = { workspace = true, default-features = false } serde = { workspace = true, features = ["derive", "std"] } +serde_json = { workspace = true, features = ["std"] } sha2 = { workspace = true, features = ["std"] } -sqlx = { workspace = true, features = ["runtime-tokio", "sqlite-bundled"] } tempfile = { workspace = true } tokio = { workspace = true, features = ["macros", "rt"] } toml = { workspace = true } diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs @@ -0,0 +1,773 @@ +use radroots_event_codec::Codec; +use radroots_storage::{ + Error, EventStore, + event::{ + AdmissionDisposition, AdmissionReceipt, AdmissionStage, BoxFuture, EventAdmission, + EventCursor, EventId, EventPage, EventPosition, EventQuery, EventQueryBounds, + EventSequence, SourceGeneration, StoredEventProvenance, StoredRawEvent, + StoredVerifiedEvent, StoredVisibleEvent, + }, + status::{EventStoreHealth, EventStoreMode, EventStoreStatus}, +}; +use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool}; + +#[derive(Clone)] +pub struct SqliteStorage { + pool: SqlitePool, + generation: SourceGeneration, + mode: EventStoreMode, +} + +struct StoredEventRow { + position: EventPosition, + raw_json: String, + stage: AdmissionStage, +} + +impl SqliteStorage { + #[allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint. + pub(crate) const fn new( + pool: SqlitePool, + generation: SourceGeneration, + mode: EventStoreMode, + ) -> Self { + Self { + pool, + generation, + mode, + } + } + + async fn selected( + &self, + query: &EventQuery, + minimum_stage: AdmissionStage, + ) -> Result<(Vec<StoredEventRow>, Option<EventCursor>), Error> { + self.validate_cursor(query)?; + let after = query + .bounds() + .cursor() + .map_or(0, |cursor| cursor.sequence().get()); + let fetch_limit = u64::from(query.bounds().limit()) + 1; + let mut builder = QueryBuilder::<Sqlite>::new( + "SELECT source_generation, source_sequence, signed_event, admission_stage \ + FROM radroots_runtime_events WHERE source_sequence > ", + ); + builder.push_bind(i64_from_u64(after)?); + if minimum_stage == AdmissionStage::Verified { + builder.push(" AND admission_stage IN ('verified', 'visible')"); + } else if minimum_stage == AdmissionStage::Visible { + builder.push(" AND admission_stage = 'visible'"); + } + if !query.event_ids().is_empty() { + builder.push(" AND event_id IN ("); + let mut separated = builder.separated(", "); + for event_id in query.event_ids() { + separated.push_bind(event_id.as_bytes().to_vec()); + } + separated.push_unseparated(")"); + } + builder.push(" ORDER BY source_sequence LIMIT "); + builder.push_bind(i64_from_u64(fetch_limit)?); + + let rows = builder + .build() + .fetch_all(&self.pool) + .await + .map_err(map_backend)?; + let mut decoded = rows + .iter() + .map(|row| self.decode_event_row(row)) + .collect::<Result<Vec<_>, _>>()?; + let next = if decoded.len() > usize::from(query.bounds().limit()) { + decoded.truncate(usize::from(query.bounds().limit())); + decoded.last().map(|row| row.position) + } else { + None + }; + Ok((decoded, next)) + } + + fn validate_cursor(&self, query: &EventQuery) -> Result<(), Error> { + if query + .bounds() + .cursor() + .is_some_and(|cursor| cursor.generation() != self.generation) + { + return Err(Error::SourceGenerationChanged); + } + Ok(()) + } + + fn decode_event_row(&self, row: &sqlx::sqlite::SqliteRow) -> Result<StoredEventRow, Error> { + let generation = source_generation(row.try_get("source_generation").map_err(map_corrupt)?)?; + if generation != self.generation { + return Err(Error::CorruptStoredEvent); + } + let sequence = event_sequence(row.try_get("source_sequence").map_err(map_corrupt)?)?; + let raw_json = String::from_utf8(row.try_get("signed_event").map_err(map_corrupt)?) + .map_err(|_| Error::CorruptStoredEvent)?; + Codec::decode_signed_event(raw_json.as_str()).map_err(|_| Error::CorruptStoredEvent)?; + let stage = admission_stage(row.try_get("admission_stage").map_err(map_corrupt)?)?; + Ok(StoredEventRow { + position: EventPosition::new(generation, sequence), + raw_json, + stage, + }) + } + + async fn store_provenance( + transaction: &mut sqlx::Transaction<'_, Sqlite>, + admission: &EventAdmission, + ) -> Result<(), Error> { + let provenance = admission.provenance(); + let cursor = provenance.cursor().map_or("", |cursor| cursor.as_str()); + sqlx::query( + "INSERT OR IGNORE INTO radroots_runtime_event_provenance ( + event_id, transport_id, target_fingerprint, observed_at_unix_ms, cursor + ) VALUES (?, ?, ?, ?, ?)", + ) + .bind(admission.event_id().as_bytes().as_slice()) + .bind(provenance.transport_id().as_str()) + .bind(provenance.target().as_str()) + .bind(i64_from_u64(provenance.observed_at_unix_ms())?) + .bind(cursor) + .execute(&mut **transaction) + .await + .map_err(map_backend)?; + Ok(()) + } +} + +impl EventStore for SqliteStorage { + fn status(&self) -> BoxFuture<'_, Result<EventStoreStatus, Error>> { + Box::pin(async move { + let row = sqlx::query( + "SELECT + COUNT(*) AS raw_events, + COALESCE(SUM(CASE WHEN admission_stage IN ('verified', 'visible') THEN 1 ELSE 0 END), 0) + AS verified_events, + COALESCE(SUM(CASE WHEN admission_stage = 'visible' THEN 1 ELSE 0 END), 0) + AS visible_events + FROM radroots_runtime_events WHERE source_generation = ?", + ) + .bind(self.generation.as_bytes().as_slice()) + .fetch_one(&self.pool) + .await + .map_err(map_backend)?; + EventStoreStatus::new( + self.generation, + self.mode, + EventStoreHealth::Available, + u64_from_i64(row.try_get("raw_events").map_err(map_corrupt)?)?, + u64_from_i64(row.try_get("verified_events").map_err(map_corrupt)?)?, + u64_from_i64(row.try_get("visible_events").map_err(map_corrupt)?)?, + ) + }) + } + + fn admit(&self, admission: EventAdmission) -> BoxFuture<'_, Result<AdmissionReceipt, Error>> { + Box::pin(async move { + if self.mode == EventStoreMode::ReadOnly { + return Err(Error::BackendUnavailable); + } + let mut transaction = self + .pool + .begin_with("BEGIN IMMEDIATE") + .await + .map_err(map_backend)?; + let existing = sqlx::query( + "SELECT source_generation, source_sequence, signed_event, admission_stage + FROM radroots_runtime_events WHERE event_id = ?", + ) + .bind(admission.event_id().as_bytes().as_slice()) + .fetch_optional(&mut *transaction) + .await + .map_err(map_backend)?; + + let (position, disposition) = if let Some(row) = existing { + let stored = self.decode_event_row(&row)?; + if stored.raw_json.as_bytes() != admission.event().raw_json().as_bytes() { + return Err(Error::EventConflict); + } + if admission.stage() < stored.stage { + return Err(Error::AdmissionRegression); + } + let disposition = if admission.stage() == stored.stage { + AdmissionDisposition::Duplicate + } else { + sqlx::query( + "UPDATE radroots_runtime_events + SET admission_stage = ?, updated_at_unix_ms = MAX(updated_at_unix_ms, ?) + WHERE event_id = ?", + ) + .bind(stage_name(admission.stage())) + .bind(i64_from_u64(admission.provenance().observed_at_unix_ms())?) + .bind(admission.event_id().as_bytes().as_slice()) + .execute(&mut *transaction) + .await + .map_err(map_backend)?; + AdmissionDisposition::Advanced + }; + (stored.position, disposition) + } else { + let next = sqlx::query_scalar::<_, i64>( + "UPDATE radroots_runtime_source_generations + SET sequence_head = sequence_head + 1 + WHERE generation = ? AND state = 'active' + RETURNING sequence_head", + ) + .bind(self.generation.as_bytes().as_slice()) + .fetch_optional(&mut *transaction) + .await + .map_err(map_backend)? + .ok_or(Error::SourceGenerationChanged)?; + let sequence = event_sequence(next)?; + let observed_at = i64_from_u64(admission.provenance().observed_at_unix_ms())?; + sqlx::query( + "INSERT INTO radroots_runtime_events ( + source_generation, source_sequence, event_id, admission_stage, + signed_event, admitted_at_unix_ms, updated_at_unix_ms + ) VALUES (?, ?, ?, ?, ?, ?, ?)", + ) + .bind(self.generation.as_bytes().as_slice()) + .bind(next) + .bind(admission.event_id().as_bytes().as_slice()) + .bind(stage_name(admission.stage())) + .bind(admission.event().raw_json().as_bytes()) + .bind(observed_at) + .bind(observed_at) + .execute(&mut *transaction) + .await + .map_err(map_backend)?; + ( + EventPosition::new(self.generation, sequence), + AdmissionDisposition::Inserted, + ) + }; + Self::store_provenance(&mut transaction, &admission).await?; + transaction.commit().await.map_err(map_backend)?; + Ok(AdmissionReceipt::new( + *admission.event_id(), + position, + admission.stage(), + disposition, + )) + }) + } + + fn query_raw( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredRawEvent>, Error>> { + Box::pin(async move { + let (rows, next) = self.selected(&query, AdmissionStage::Raw).await?; + let items = rows + .into_iter() + .map(|row| { + let event = Codec::decode_signed_event(row.raw_json.as_str()) + .map_err(|_| Error::CorruptStoredEvent)?; + Ok(StoredRawEvent::new(row.position, event, row.stage)) + }) + .collect::<Result<Vec<_>, Error>>()?; + EventPage::new(self.generation, items, next, query.bounds()) + }) + } + + fn query_verified( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVerifiedEvent>, Error>> { + Box::pin(async move { + let (rows, next) = self.selected(&query, AdmissionStage::Verified).await?; + let items = rows + .into_iter() + .map(|row| { + let event = Codec::decode_signed_event(row.raw_json.as_str()) + .map_err(|_| Error::CorruptStoredEvent)?; + Ok(StoredVerifiedEvent::new(row.position, event)) + }) + .collect::<Result<Vec<_>, Error>>()?; + EventPage::new(self.generation, items, next, query.bounds()) + }) + } + + fn query_visible( + &self, + query: EventQuery, + ) -> BoxFuture<'_, Result<EventPage<StoredVisibleEvent>, Error>> { + Box::pin(async move { + let (rows, next) = self.selected(&query, AdmissionStage::Visible).await?; + let items = rows + .into_iter() + .map(|row| { + let event = Codec::decode_signed_event(row.raw_json.as_str()) + .map_err(|_| Error::CorruptStoredEvent)?; + Ok(StoredVisibleEvent::new(row.position, event)) + }) + .collect::<Result<Vec<_>, Error>>()?; + EventPage::new(self.generation, items, next, query.bounds()) + }) + } + + fn query_provenance( + &self, + event_id: EventId, + bounds: EventQueryBounds, + ) -> BoxFuture<'_, Result<EventPage<StoredEventProvenance>, Error>> { + Box::pin(async move { + if bounds + .cursor() + .is_some_and(|cursor| cursor.generation() != self.generation) + { + return Err(Error::SourceGenerationChanged); + } + let event_row = sqlx::query( + "SELECT source_generation, source_sequence + FROM radroots_runtime_events WHERE event_id = ?", + ) + .bind(event_id.as_bytes().as_slice()) + .fetch_optional(&self.pool) + .await + .map_err(map_backend)? + .ok_or(Error::EventNotFound)?; + let generation = source_generation( + event_row + .try_get("source_generation") + .map_err(map_corrupt)?, + )?; + if generation != self.generation { + return Err(Error::CorruptStoredEvent); + } + let position = EventPosition::new( + generation, + event_sequence(event_row.try_get("source_sequence").map_err(map_corrupt)?)?, + ); + let after = bounds.cursor().map_or(0, |cursor| cursor.sequence().get()); + let items = if position.sequence().get() <= after { + Vec::new() + } else { + let rows = sqlx::query( + "SELECT transport_id, target_fingerprint, observed_at_unix_ms, cursor + FROM radroots_runtime_event_provenance + WHERE event_id = ? + ORDER BY observed_at_unix_ms, transport_id, target_fingerprint, cursor + LIMIT ?", + ) + .bind(event_id.as_bytes().as_slice()) + .bind(i64::from(bounds.limit())) + .fetch_all(&self.pool) + .await + .map_err(map_backend)?; + rows.iter() + .map(|row| { + let cursor = row.try_get::<String, _>("cursor").map_err(map_corrupt)?; + StoredEventProvenance::from_stored_parts( + position, + row.try_get::<String, _>("transport_id") + .map_err(map_corrupt)? + .as_str(), + row.try_get::<String, _>("target_fingerprint") + .map_err(map_corrupt)? + .as_str(), + u64_from_i64(row.try_get("observed_at_unix_ms").map_err(map_corrupt)?)?, + (!cursor.is_empty()).then_some(cursor.as_str()), + ) + }) + .collect::<Result<Vec<_>, Error>>()? + }; + EventPage::new(self.generation, items, None, bounds) + }) + } +} + +const fn stage_name(stage: AdmissionStage) -> &'static str { + match stage { + AdmissionStage::Raw => "raw", + AdmissionStage::Verified => "verified", + AdmissionStage::Visible => "visible", + } +} + +fn admission_stage(value: String) -> Result<AdmissionStage, Error> { + match value.as_str() { + "raw" => Ok(AdmissionStage::Raw), + "verified" => Ok(AdmissionStage::Verified), + "visible" => Ok(AdmissionStage::Visible), + _ => Err(Error::CorruptStoredEvent), + } +} + +fn source_generation(value: Vec<u8>) -> Result<SourceGeneration, Error> { + SourceGeneration::new(value.try_into().map_err(|_| Error::CorruptStoredEvent)?) + .map_err(|_| Error::CorruptStoredEvent) +} + +fn event_sequence(value: i64) -> Result<EventSequence, Error> { + EventSequence::new(u64_from_i64(value)?).map_err(|_| Error::CorruptStoredEvent) +} + +fn i64_from_u64(value: u64) -> Result<i64, Error> { + i64::try_from(value).map_err(|_| Error::CorruptStoredEvent) +} + +fn u64_from_i64(value: i64) -> Result<u64, Error> { + u64::try_from(value).map_err(|_| Error::CorruptStoredEvent) +} + +fn map_backend(_: sqlx::Error) -> Error { + Error::BackendUnavailable +} + +fn map_corrupt(_: sqlx::Error) -> Error { + Error::CorruptStoredEvent +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::migration::runtime::{MIGRATIONS, migration_sql}; + use radroots_event::{ + SignedEvent, + admission::{AdmissionPolicy, RawEvent, VisibilityPolicy, VisibleEvent}, + wire::Nip01EventWire, + }; + use radroots_storage::event::EventQueryBounds; + use radroots_transport::{ + Target, TransportId, + source::{EventProvenance, FetchCursor, ObservedEvent}, + }; + use sqlx::sqlite::SqlitePoolOptions; + + 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-sqlite.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-sqlite.visibility.v1" + } + + fn make_visible( + &self, + _event: &radroots_event::admission::AdmittedEvent, + ) -> Result<(), Self::Error> { + Ok(()) + } + } + + async fn store(generation: SourceGeneration) -> SqliteStorage { + let pool = SqlitePoolOptions::new() + .max_connections(1) + .connect("sqlite::memory:") + .await + .expect("memory SQLite"); + sqlx::query("PRAGMA foreign_keys = ON") + .execute(&pool) + .await + .expect("foreign keys"); + for migration in MIGRATIONS { + sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) + .execute(&pool) + .await + .expect("runtime migration"); + } + sqlx::query( + "INSERT INTO radroots_runtime_source_generations ( + generation, sequence_head, state, created_at_unix_ms, retired_at_unix_ms + ) VALUES (?, 0, 'active', 1, NULL)", + ) + .bind(generation.as_bytes().as_slice()) + .execute(&pool) + .await + .expect("source generation"); + SqliteStorage::new(pool, generation, EventStoreMode::ReadWrite) + } + + fn signed_event(content: &str, pretty: bool) -> SignedEvent { + let mut wire = Nip01EventWire { + id: "0".repeat(64), + pubkey: "585591529da0bab31b3b1b1f986611cf5f435dca84f978c89ee8a40cca7103df".to_owned(), + created_at: 1_800_000_100, + kind: 0, + tags: vec![], + content: content.to_owned(), + sig: "42".repeat(64), + extra: Default::default(), + }; + wire.id = wire + .computed_event_id() + .expect("canonical event id") + .to_hex(); + let value = 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, + }); + let raw_json = if pretty { + serde_json::to_string_pretty(&value).expect("pretty event JSON") + } else { + value.to_string() + }; + SignedEvent::from_wire_verified_id(wire, raw_json).expect("signed event") + } + + fn observed(event: SignedEvent, at: u64, cursor: Option<&str>) -> ObservedEvent { + let target = Target::new(TransportId::NOSTR, "wss://relay.example").expect("target"); + let mut provenance = + EventProvenance::new(TransportId::NOSTR, target.fingerprint().clone(), at) + .expect("provenance"); + if let Some(cursor) = cursor { + provenance = provenance.with_cursor(FetchCursor::parse(cursor).expect("cursor")); + } + ObservedEvent::new(event, provenance) + } + + fn verified(event: &SignedEvent) -> radroots_event::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") + } + + #[tokio::test] + async fn admission_is_idempotent_monotonic_and_conflict_safe() { + let generation = SourceGeneration::new([7; 32]).expect("generation"); + let store = store(generation).await; + let event = signed_event( + "{\"display_name\":\"Moss Street Farm\",\"bot\":false}", + false, + ); + + let inserted = store + .admit(EventAdmission::raw(observed(event.clone(), 10, None))) + .await + .expect("insert raw"); + assert_eq!(inserted.disposition(), AdmissionDisposition::Inserted); + assert_eq!(inserted.position().sequence().get(), 1); + + let advanced = store + .admit( + EventAdmission::visible( + observed(event.clone(), 11, Some("relay-page-1")), + visible(&event), + ) + .expect("visible admission"), + ) + .await + .expect("advance visible"); + assert_eq!(advanced.disposition(), AdmissionDisposition::Advanced); + assert_eq!(advanced.position(), inserted.position()); + + let duplicate = store + .admit( + EventAdmission::visible( + observed(event.clone(), 11, Some("relay-page-1")), + visible(&event), + ) + .expect("visible admission"), + ) + .await + .expect("duplicate visible"); + assert_eq!(duplicate.disposition(), AdmissionDisposition::Duplicate); + assert_eq!( + store + .admit(EventAdmission::raw(observed(event.clone(), 12, None))) + .await, + Err(Error::AdmissionRegression) + ); + + let same_id_different_bytes = signed_event( + "{\"display_name\":\"Moss Street Farm\",\"bot\":false}", + true, + ); + assert_eq!(same_id_different_bytes.id(), event.id()); + assert_eq!( + store + .admit(EventAdmission::raw(observed( + same_id_different_bytes, + 13, + None, + ))) + .await, + Err(Error::EventConflict) + ); + } + + #[tokio::test] + async fn queries_preserve_stage_bounds_cursors_and_exact_provenance() { + let generation = SourceGeneration::new([8; 32]).expect("generation"); + let store = store(generation).await; + let empty_status = store.status().await.expect("empty status"); + assert_eq!(empty_status.raw_events(), 0); + assert_eq!(empty_status.verified_events(), 0); + assert_eq!(empty_status.visible_events(), 0); + let raw_event = signed_event("raw", false); + let visible_event = + signed_event("{\"display_name\":\"Visible Farm\",\"bot\":false}", false); + store + .admit(EventAdmission::raw(observed(raw_event.clone(), 20, None))) + .await + .expect("raw event"); + store + .admit( + EventAdmission::visible( + observed(visible_event.clone(), 21, Some("relay-page-2")), + visible(&visible_event), + ) + .expect("visible admission"), + ) + .await + .expect("visible event"); + + let first = store + .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"))) + .await + .expect("first page"); + assert_eq!(first.items().len(), 1); + assert_eq!(first.items()[0].event(), &raw_event); + let next = first.next_cursor().expect("continuation cursor"); + let second = store + .query_raw(EventQuery::all( + EventQueryBounds::first(1).expect("bounds").after(next), + )) + .await + .expect("second page"); + assert_eq!(second.items()[0].event(), &visible_event); + assert!(second.next_cursor().is_none()); + + let verified_page = store + .query_verified(EventQuery::all( + EventQueryBounds::first(10).expect("bounds"), + )) + .await + .expect("verified page"); + assert_eq!(verified_page.items().len(), 1); + assert_eq!(verified_page.items()[0].event(), &visible_event); + let visible_page = store + .query_visible( + EventQuery::for_ids( + EventQueryBounds::first(10).expect("bounds"), + vec![*visible_event.id()], + ) + .expect("id query"), + ) + .await + .expect("visible page"); + assert_eq!(visible_page.items().len(), 1); + + let provenance = store + .query_provenance( + *visible_event.id(), + EventQueryBounds::first(10).expect("bounds"), + ) + .await + .expect("provenance"); + assert_eq!(provenance.items().len(), 1); + assert_eq!(provenance.items()[0].provenance().observed_at_unix_ms(), 21); + assert_eq!( + provenance.items()[0] + .provenance() + .cursor() + .expect("cursor") + .as_str(), + "relay-page-2" + ); + + let status = store.status().await.expect("status"); + assert_eq!(status.raw_events(), 2); + assert_eq!(status.verified_events(), 1); + assert_eq!(status.visible_events(), 1); + let foreign_cursor = EventPosition::new( + SourceGeneration::new([9; 32]).expect("foreign generation"), + EventSequence::new(1).expect("sequence"), + ); + assert_eq!( + store + .query_raw(EventQuery::all( + EventQueryBounds::first(1) + .expect("bounds") + .after(foreign_cursor), + )) + .await, + Err(Error::SourceGenerationChanged) + ); + } + + #[tokio::test] + async fn corrupt_rows_fail_closed_and_source_history_is_immutable() { + let generation = SourceGeneration::new([10; 32]).expect("generation"); + let store = store(generation).await; + let event = signed_event("corruption", false); + store + .admit(EventAdmission::raw(observed(event, 30, None))) + .await + .expect("event"); + + assert!( + sqlx::query("DELETE FROM radroots_runtime_events") + .execute(&store.pool) + .await + .is_err() + ); + assert!( + sqlx::query("DELETE FROM radroots_runtime_source_generations") + .execute(&store.pool) + .await + .is_err() + ); + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(&store.pool) + .await + .expect("disable checks for corruption probe"); + sqlx::query("UPDATE radroots_runtime_events SET admission_stage = 'corrupt'") + .execute(&store.pool) + .await + .expect("forge corrupt stage"); + assert_eq!( + store + .query_raw(EventQuery::all(EventQueryBounds::first(1).expect("bounds"),)) + .await, + Err(Error::CorruptStoredEvent) + ); + } +} diff --git a/crates/storage_sqlite/src/lib.rs b/crates/storage_sqlite/src/lib.rs @@ -8,5 +8,8 @@ pub mod migration; pub mod open; pub mod status; +mod event; + pub use config::OpenOptions; +pub use event::SqliteStorage; pub use open::{Error, OpenMode, Paths}; diff --git a/crates/storage_sqlite/src/migration/runtime/0002_canonical_event_storage.up.sql b/crates/storage_sqlite/src/migration/runtime/0002_canonical_event_storage.up.sql @@ -0,0 +1,44 @@ +DROP INDEX radroots_runtime_event_provenance_observed_idx; + +ALTER TABLE radroots_runtime_event_provenance +RENAME TO radroots_runtime_event_provenance_v1; + +CREATE TABLE radroots_runtime_event_provenance ( + event_id BLOB NOT NULL REFERENCES radroots_runtime_events(event_id), + transport_id TEXT NOT NULL CHECK (length(transport_id) BETWEEN 1 AND 64), + target_fingerprint TEXT NOT NULL CHECK (length(target_fingerprint) = 64), + observed_at_unix_ms INTEGER NOT NULL CHECK (observed_at_unix_ms > 0), + cursor TEXT NOT NULL CHECK (length(cursor) <= 2048), + PRIMARY KEY (event_id, transport_id, target_fingerprint, observed_at_unix_ms, cursor) +) STRICT, WITHOUT ROWID; + +INSERT OR IGNORE INTO radroots_runtime_event_provenance ( + event_id, + transport_id, + target_fingerprint, + observed_at_unix_ms, + cursor +) +SELECT + event_id, + transport_kind, + CAST(endpoint_fingerprint AS TEXT), + last_observed_at_unix_ms, + '' +FROM radroots_runtime_event_provenance_v1; + +DROP TABLE radroots_runtime_event_provenance_v1; + +CREATE INDEX radroots_runtime_event_provenance_observed_idx +ON radroots_runtime_event_provenance(observed_at_unix_ms, event_id); + +CREATE UNIQUE INDEX radroots_runtime_source_generations_active_idx +ON radroots_runtime_source_generations(state) +WHERE state = 'active'; + +CREATE TRIGGER radroots_runtime_source_generations_sequence_guard +BEFORE UPDATE OF sequence_head ON radroots_runtime_source_generations +WHEN NEW.sequence_head < OLD.sequence_head +BEGIN + SELECT RAISE(ABORT, 'runtime source sequence cannot regress'); +END; diff --git a/crates/storage_sqlite/src/migration/runtime/mod.rs b/crates/storage_sqlite/src/migration/runtime/mod.rs @@ -6,10 +6,12 @@ /// Lowest runtime schema version this package can recognize. pub const MINIMUM_VERSION: u32 = 1; /// Current runtime schema version created by this package. -pub const CURRENT_VERSION: u32 = 1; +pub const CURRENT_VERSION: u32 = 2; #[allow(dead_code)] // Consumed by the migration executor introduced in its ordered RCL step. const RUNTIME_V1_SQL: &str = include_str!("0001_runtime.up.sql"); +#[allow(dead_code)] // Consumed by the migration executor introduced in its ordered RCL step. +const CANONICAL_EVENT_STORAGE_V2_SQL: &str = include_str!("0002_canonical_event_storage.up.sql"); /// Stable, non-SQL description of one forward runtime migration. #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -71,18 +73,58 @@ const RUNTIME_V1_OBJECTS: &[&str] = &[ "radroots_runtime_source_generations_identity_guard", ]; +const RUNTIME_V2_OBJECTS: &[&str] = &[ + "radroots_runtime_atomic_commits", + "radroots_runtime_delivery_evidence", + "radroots_runtime_delivery_evidence_item_idx", + "radroots_runtime_event_index_checkpoints", + "radroots_runtime_event_index_manifests", + "radroots_runtime_event_index_shards", + "radroots_runtime_event_provenance", + "radroots_runtime_event_provenance_observed_idx", + "radroots_runtime_events", + "radroots_runtime_events_admission_idx", + "radroots_runtime_events_delete_guard", + "radroots_runtime_events_event_id_idx", + "radroots_runtime_events_raw_update_guard", + "radroots_runtime_journal_idempotency_idx", + "radroots_runtime_journal_operations", + "radroots_runtime_journal_recovery_idx", + "radroots_runtime_outbox_items", + "radroots_runtime_outbox_ready_idx", + "radroots_runtime_outbox_targets", + "radroots_runtime_projection_checkpoints", + "radroots_runtime_projection_invalidations", + "radroots_runtime_projection_rebuilds", + "radroots_runtime_projection_rebuilds_stage_idx", + "radroots_runtime_source_generations", + "radroots_runtime_source_generations_active_idx", + "radroots_runtime_source_generations_delete_guard", + "radroots_runtime_source_generations_identity_guard", + "radroots_runtime_source_generations_sequence_guard", +]; + /// Ordered, immutable runtime migration plan. -pub const MIGRATIONS: &[MigrationDescriptor] = &[MigrationDescriptor { - version: 1, - name: "runtime_authority", - up_sha256: "3b869122dd5bd58f4a15e7a71fd1377879640b36cd496d7b3f15278ef1e128c9", - owned_objects: RUNTIME_V1_OBJECTS, -}]; +pub const MIGRATIONS: &[MigrationDescriptor] = &[ + MigrationDescriptor { + version: 1, + name: "runtime_authority", + up_sha256: "3b869122dd5bd58f4a15e7a71fd1377879640b36cd496d7b3f15278ef1e128c9", + owned_objects: RUNTIME_V1_OBJECTS, + }, + MigrationDescriptor { + version: 2, + name: "canonical_event_storage", + up_sha256: "35b036ba84eff7135665c4ae42fa8232d8bacd8115805b03011a3eb76423f4b8", + owned_objects: RUNTIME_V2_OBJECTS, + }, +]; #[allow(dead_code)] // Keeps raw SQL crate-private until the migration executor is installed. pub(crate) const fn migration_sql(version: u32) -> Option<&'static str> { match version { 1 => Some(RUNTIME_V1_SQL), + 2 => Some(CANONICAL_EVENT_STORAGE_V2_SQL), _ => None, } } @@ -124,9 +166,9 @@ mod tests { fn migration_plan_matches_governed_snapshot() { let snapshot = toml::from_str::<PlanSnapshot>(PLAN_SNAPSHOT).expect("valid snapshot"); assert_eq!(MINIMUM_VERSION, 1); - assert_eq!(CURRENT_VERSION, 1); - assert_eq!(MIGRATIONS.len(), 1); - let migration = MIGRATIONS[0]; + assert_eq!(CURRENT_VERSION, 2); + assert_eq!(MIGRATIONS.len(), 2); + let migration = MIGRATIONS[1]; assert_eq!(snapshot.schema_version, 1); assert_eq!(snapshot.database, "runtime.sqlite"); assert_eq!(snapshot.minimum_version, MINIMUM_VERSION); @@ -137,21 +179,22 @@ mod tests { assert!(!snapshot.raw_sql_public); assert_eq!(snapshot.authorities.len(), 8); assert_eq!(snapshot.source_invariants.len(), 5); - assert_eq!(snapshot.migrations.len(), 1); - assert_eq!(snapshot.migrations[0].version, migration.version()); - assert_eq!(snapshot.migrations[0].name, migration.name()); - assert_eq!(snapshot.migrations[0].sha256, migration.up_sha256()); - assert_eq!( - snapshot.migrations[0].owned_objects, - migration.owned_objects() - ); + assert_eq!(snapshot.migrations.len(), MIGRATIONS.len()); + for (expected, actual) in snapshot.migrations.iter().zip(MIGRATIONS) { + assert_eq!(expected.version, actual.version()); + assert_eq!(expected.name, actual.name()); + assert_eq!(expected.sha256, actual.up_sha256()); + assert_eq!(expected.owned_objects, actual.owned_objects()); + } } #[test] fn embedded_migration_checksum_is_pinned() { - let actual = format!("{:x}", Sha256::digest(migration_sql(1).expect("v1 SQL"))); - assert_eq!(actual, MIGRATIONS[0].up_sha256()); - assert_eq!(migration_sql(2), None); + for migration in MIGRATIONS { + let sql = migration_sql(migration.version()).expect("registered SQL"); + assert_eq!(format!("{:x}", Sha256::digest(sql)), migration.up_sha256()); + } + assert_eq!(migration_sql(3), None); } #[tokio::test] @@ -159,10 +202,12 @@ mod tests { let mut connection = SqliteConnection::connect("sqlite::memory:") .await .expect("open memory SQLite"); - sqlx::raw_sql(migration_sql(1).expect("v1 SQL")) - .execute(&mut connection) - .await - .expect("apply runtime schema"); + for migration in MIGRATIONS { + sqlx::raw_sql(migration_sql(migration.version()).expect("registered SQL")) + .execute(&mut connection) + .await + .expect("apply runtime schema"); + } let rows = sqlx::query( "SELECT name FROM sqlite_schema \ @@ -176,7 +221,13 @@ mod tests { .iter() .map(|row| row.get::<String, _>("name")) .collect::<Vec<_>>(); - assert_eq!(actual, MIGRATIONS[0].owned_objects()); + assert_eq!( + actual, + MIGRATIONS + .last() + .expect("current migration") + .owned_objects() + ); let foreign_key_violations = sqlx::query("PRAGMA foreign_key_check") .fetch_all(&mut connection) diff --git a/crates/storage_sqlite/tests/package_boundary.rs b/crates/storage_sqlite/tests/package_boundary.rs @@ -23,8 +23,11 @@ fn sqlite_storage_declares_the_final_backend_boundaries() { "radroots_event_codec", "radroots_secrets", "radroots_storage", + "sqlx", ]) ); + assert!(!ROOT.contains("SqlitePool")); + assert!(!ROOT.contains("sqlx")); for forbidden in [ "radroots_event_store", "radroots_outbox",