commit b6667541b08b4e578199c6612cf51b7402071f5e parent 963dd8dc99fc5f16ef7d7eea11f8d964770a696e Author: triesap <tyson@radroots.org> Date: Sun, 20 Sep 2026 01:28:20 +0000 storage: add bounded projection document inventory - Bind continuations to explicit projection and generation selection - Bound metadata and payload reads while retaining corrupt record evidence - Verify memory and SQLite parity through updates, reopen and byte limits - Qualify public APIs, coverage, workspace and standalone release checks Diffstat:
12 files changed, 816 insertions(+), 0 deletions(-)
diff --git a/contracts/api_baselines/radroots_storage.txt b/contracts/api_baselines/radroots_storage.txt @@ -1012,6 +1012,7 @@ pub fn radroots_storage::memory::MemoryStorage::put_event_index_checkpoint(&self pub fn radroots_storage::memory::MemoryStorage::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>> @@ -1373,6 +1374,37 @@ pub fn radroots_storage::memory::MemoryStorage::tombstone(&self, radroots_storag pub mod radroots_storage::projection pub use radroots_storage::projection::BoxFuture pub use radroots_storage::projection::EventId +pub mod radroots_storage::projection::document_query +pub enum radroots_storage::projection::document_query::ProjectionDocumentGenerations +pub radroots_storage::projection::document_query::ProjectionDocumentGenerations::All +pub radroots_storage::projection::document_query::ProjectionDocumentGenerations::Exact(radroots_storage::projection::ProjectionGeneration) +pub struct radroots_storage::projection::document_query::ProjectionDocumentCursor +pub struct radroots_storage::projection::document_query::ProjectionDocumentPage +impl radroots_storage::projection::document_query::ProjectionDocumentPage +pub fn radroots_storage::projection::document_query::ProjectionDocumentPage::into_records(self) -> alloc::vec::Vec<radroots_storage::projection::document_query::ProjectionDocumentRecord> +pub fn radroots_storage::projection::document_query::ProjectionDocumentPage::new(&radroots_storage::projection::document_query::ProjectionDocumentQuery, alloc::vec::Vec<radroots_storage::projection::document_query::ProjectionDocumentRecord>, bool) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::projection::document_query::ProjectionDocumentPage::next_cursor(&self) -> core::option::Option<&radroots_storage::projection::document_query::ProjectionDocumentCursor> +pub fn radroots_storage::projection::document_query::ProjectionDocumentPage::records(&self) -> &[radroots_storage::projection::document_query::ProjectionDocumentRecord] +pub struct radroots_storage::projection::document_query::ProjectionDocumentQuery +impl radroots_storage::projection::document_query::ProjectionDocumentQuery +pub fn radroots_storage::projection::document_query::ProjectionDocumentQuery::after(&self) -> core::option::Option<(radroots_storage::projection::ProjectionGeneration, &str)> +pub const fn radroots_storage::projection::document_query::ProjectionDocumentQuery::generations(&self) -> radroots_storage::projection::document_query::ProjectionDocumentGenerations +pub const fn radroots_storage::projection::document_query::ProjectionDocumentQuery::limit(&self) -> u16 +pub fn radroots_storage::projection::document_query::ProjectionDocumentQuery::matches(&self, radroots_storage::projection::ProjectionGeneration, &str) -> bool +pub fn radroots_storage::projection::document_query::ProjectionDocumentQuery::new(radroots_storage::projection::ProjectionId, radroots_storage::projection::document_query::ProjectionDocumentGenerations, u16) -> core::result::Result<Self, radroots_storage::Error> +pub fn radroots_storage::projection::document_query::ProjectionDocumentQuery::projection_id(&self) -> &radroots_storage::projection::ProjectionId +pub fn radroots_storage::projection::document_query::ProjectionDocumentQuery::with_cursor(self, &radroots_storage::projection::document_query::ProjectionDocumentCursor) -> core::result::Result<Self, radroots_storage::Error> +pub struct radroots_storage::projection::document_query::ProjectionDocumentRecord +impl radroots_storage::projection::document_query::ProjectionDocumentRecord +pub fn radroots_storage::projection::document_query::ProjectionDocumentRecord::corrupt(radroots_storage::projection::ProjectionGeneration, alloc::string::String) -> core::result::Result<Self, radroots_storage::Error> +pub const fn radroots_storage::projection::document_query::ProjectionDocumentRecord::document(&self) -> core::option::Option<&radroots_storage::projection::ProjectionDocument> +pub const fn radroots_storage::projection::document_query::ProjectionDocumentRecord::generation(&self) -> radroots_storage::projection::ProjectionGeneration +pub fn radroots_storage::projection::document_query::ProjectionDocumentRecord::key(&self) -> &str +pub fn radroots_storage::projection::document_query::ProjectionDocumentRecord::new(radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> core::result::Result<Self, radroots_storage::Error> +impl core::fmt::Debug for radroots_storage::projection::document_query::ProjectionDocumentRecord +pub fn radroots_storage::projection::document_query::ProjectionDocumentRecord::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub const radroots_storage::projection::document_query::PROJECTION_DOCUMENT_PAGE_BYTES_MAX: usize +pub const radroots_storage::projection::document_query::PROJECTION_DOCUMENT_QUERY_LIMIT_MAX: u16 pub enum radroots_storage::projection::InvalidationReason pub radroots_storage::projection::InvalidationReason::EventIndexManifestChanged pub radroots_storage::projection::InvalidationReason::IntegrityFailure @@ -1545,6 +1577,7 @@ pub fn radroots_storage::projection::ProjectionStore::put_event_index_checkpoint pub fn radroots_storage::projection::ProjectionStore::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::projection::ProjectionStore::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::projection::ProjectionStore::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> +pub fn radroots_storage::projection::ProjectionStore::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::Error>> pub fn radroots_storage::projection::ProjectionStore::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>> pub fn radroots_storage::projection::ProjectionStore::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>> pub fn radroots_storage::projection::ProjectionStore::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>> @@ -1561,6 +1594,7 @@ pub fn radroots_storage::memory::MemoryStorage::put_event_index_checkpoint(&self pub fn radroots_storage::memory::MemoryStorage::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>> @@ -1840,6 +1874,7 @@ pub fn radroots_storage::ProjectionStore::put_event_index_checkpoint(&self, radr pub fn radroots_storage::ProjectionStore::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::ProjectionStore::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::ProjectionStore::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> +pub fn radroots_storage::ProjectionStore::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::Error>> pub fn radroots_storage::ProjectionStore::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>> pub fn radroots_storage::ProjectionStore::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>> pub fn radroots_storage::ProjectionStore::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>> @@ -1856,6 +1891,7 @@ pub fn radroots_storage::memory::MemoryStorage::put_event_index_checkpoint(&self pub fn radroots_storage::memory::MemoryStorage::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::Error>> +pub fn radroots_storage::memory::MemoryStorage::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::Error>> pub fn radroots_storage::memory::MemoryStorage::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::Error>> diff --git a/contracts/api_baselines/radroots_storage_sqlite.txt b/contracts/api_baselines/radroots_storage_sqlite.txt @@ -611,6 +611,7 @@ pub fn radroots_storage_sqlite::SqliteStorage::put_event_index_checkpoint(&self, pub fn radroots_storage_sqlite::SqliteStorage::put_event_index_manifest(&self, radroots_storage::projection::EventIndexManifest) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::put_projection_document(&self, radroots_storage::projection::ProjectionId, radroots_storage::projection::ProjectionGeneration, radroots_storage::projection::ProjectionDocument) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::put_projection_snapshot(&self, radroots_storage::projection::ProjectionSnapshot) -> radroots_transport::source::BoxFuture<'_, core::result::Result<(), radroots_storage::error::Error>> +pub fn radroots_storage_sqlite::SqliteStorage::query_projection_documents(&self, radroots_storage::projection::document_query::ProjectionDocumentQuery) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::document_query::ProjectionDocumentPage, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::rebuild(&self, radroots_storage::projection::RebuildTicketId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::RebuildTicket>, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::request_rebuild(&self, radroots_storage::projection::RebuildTicket) -> radroots_transport::source::BoxFuture<'_, core::result::Result<radroots_storage::projection::RebuildTicket, radroots_storage::error::Error>> pub fn radroots_storage_sqlite::SqliteStorage::status(&self, radroots_storage::projection::ProjectionId) -> radroots_transport::source::BoxFuture<'_, core::result::Result<core::option::Option<radroots_storage::projection::ProjectionStatus>, radroots_storage::error::Error>> diff --git a/contracts/architecture/decisions/projection_document_inventory.v1.json b/contracts/architecture/decisions/projection_document_inventory.v1.json @@ -0,0 +1,11 @@ +{ + "schema": "radroots.projection-document-inventory.v1", + "status": "implemented", + "owners": ["radroots_storage", "radroots_storage_sqlite"], + "selection": "An independently supplied projection identity and explicit exact-generation or all-generation selection bind every continuation. Ascending generation bytes and UTF-8 key bytes define traversal. Typed cursors cannot broaden selection.", + "bounds": {"page_records": 256, "metadata_lookahead": 1, "key_bytes": 512, "generation_bytes": 32, "digest_bytes": 32, "page_payload_bytes": 16777216}, + "corruption": "Invalid key or generation prevents a safe continuation and fails the query. Invalid value length or digest produces a bounded corrupt-record locator. No malformed record is silently omitted or reported as an empty complete inventory. SQLite bounds metadata before decoding and checks the remaining byte budget before reading each payload.", + "consistency": "One released SQLite read snapshot per page; bounded memory selection under its existing mutex. This is a live scan, not an immutable cross-page snapshot. Updates retain ordering. Callers fence their own mutations or retain resources when reference completeness cannot be established.", + "compatibility": "Keep exact document lookup, persisted bytes, migrations, dependencies and features unchanged. A backend that has not implemented inventory fails explicitly with BackendUnavailable.", + "non_goals": ["application schema interpretation", "garbage collection policy", "author enumeration", "filesystem access", "new database", "raw SQL escape hatch", "release qualification"] +} diff --git a/crates/storage/README.md b/crates/storage/README.md @@ -5,6 +5,13 @@ Radroots hosts. It owns canonical event persistence, durable operation journal state, outbox and delivery evidence, projection coordination, protected-record metadata, reliability operations, and high-level atomic workflow commits. +Projection document inventory selects one projection and explicitly selects +one generation or all generations. Pages contain at most 256 records and +16 MiB of opaque values, with independently scoped continuations and explicit +corrupt-record locators. This is a live scan, not a frozen snapshot. Callers +must fence their own mutations before treating a traversal as complete. +Storage neither interprets document payloads nor grants deletion authority. + The package does not expose SQL, filesystem handles, database pools, raw transactions, encryption keys, transport clients, schedulers, or application state. Concrete backends implement these contracts; `radroots_storage_sqlite` diff --git a/crates/storage/src/memory.rs b/crates/storage/src/memory.rs @@ -1038,6 +1038,42 @@ impl ProjectionStore for MemoryStorage { }) } + fn query_projection_documents( + &self, + query: crate::projection::document_query::ProjectionDocumentQuery, + ) -> BoxFuture<'_, Result<crate::projection::document_query::ProjectionDocumentPage, Error>> + { + use crate::projection::document_query::{ + PROJECTION_DOCUMENT_PAGE_BYTES_MAX, ProjectionDocumentPage, ProjectionDocumentRecord, + }; + Box::pin(async move { + let state = self.state()?; + let mut selected = std::collections::BTreeMap::new(); + let maximum = usize::from(query.limit()) + 1; + for (id, generation, document) in &state.projection_documents { + if id == query.projection_id() && query.matches(*generation, document.key()) { + selected.insert((*generation, document.key()), document); + if selected.len() > maximum { + selected.pop_last(); + } + } + } + let mut has_more = selected.len() > usize::from(query.limit()); + let mut records = Vec::new(); + let mut bytes = 0; + for ((generation, _), document) in selected.into_iter().take(usize::from(query.limit())) + { + if bytes + document.value().len() > PROJECTION_DOCUMENT_PAGE_BYTES_MAX { + has_more = true; + break; + } + bytes += document.value().len(); + records.push(ProjectionDocumentRecord::new(generation, document.clone())?); + } + ProjectionDocumentPage::new(&query, records, has_more) + }) + } + fn projection_document( &self, projection_id: ProjectionId, diff --git a/crates/storage/src/projection.rs b/crates/storage/src/projection.rs @@ -3,6 +3,8 @@ //! Storage owns durable coordination metadata. Domain reducers and projected //! row representations remain in their domain packages. +pub mod document_query; + pub use radroots_event::EventId; pub use radroots_transport::BoxFuture; use sha2::{Digest, Sha256}; @@ -1084,6 +1086,13 @@ pub trait ProjectionStore: Send + Sync { generation: ProjectionGeneration, key: String, ) -> BoxFuture<'_, Result<Option<ProjectionDocument>, Error>>; + /// Inventories are live bounded scans. Unsupported backends fail explicitly. + fn query_projection_documents( + &self, + _query: document_query::ProjectionDocumentQuery, + ) -> BoxFuture<'_, Result<document_query::ProjectionDocumentPage, Error>> { + Box::pin(async { Err(Error::BackendUnavailable) }) + } /// Persists one immutable frozen-query snapshot idempotently. fn put_projection_snapshot( &self, diff --git a/crates/storage/src/projection/document_query.rs b/crates/storage/src/projection/document_query.rs @@ -0,0 +1,186 @@ +//! Bounded opaque document inventory. A continuation is not a frozen snapshot. +#[cfg(test)] +mod tests; +use super::{Error, ProjectionDocument, ProjectionGeneration, ProjectionId, valid_document_key}; + +pub const PROJECTION_DOCUMENT_QUERY_LIMIT_MAX: u16 = 256; +pub const PROJECTION_DOCUMENT_PAGE_BYTES_MAX: usize = super::PROJECTION_DOCUMENT_VALUE_MAX_BYTES; + +#[derive(Clone, Copy, Debug, Eq, PartialEq)] +pub enum ProjectionDocumentGenerations { + Exact(ProjectionGeneration), + All, +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProjectionDocumentQuery { + projection_id: ProjectionId, + generations: ProjectionDocumentGenerations, + limit: u16, + after: Option<(ProjectionGeneration, String)>, +} + +impl ProjectionDocumentQuery { + pub fn new( + projection_id: ProjectionId, + generations: ProjectionDocumentGenerations, + limit: u16, + ) -> Result<Self, Error> { + ProjectionId::parse(projection_id.as_str())?; + if let ProjectionDocumentGenerations::Exact(generation) = generations { + ProjectionGeneration::new(*generation.as_bytes())?; + } + if limit == 0 || limit > PROJECTION_DOCUMENT_QUERY_LIMIT_MAX { + return Err(Error::InvalidProjectionDocument); + } + Ok(Self { + projection_id, + generations, + limit, + after: None, + }) + } + pub fn with_cursor(mut self, cursor: &ProjectionDocumentCursor) -> Result<Self, Error> { + if self.projection_id != cursor.projection_id || self.generations != cursor.generations { + return Err(Error::InvalidProjectionDocument); + } + self.after = Some(cursor.after.clone()); + Ok(self) + } + pub fn projection_id(&self) -> &ProjectionId { + &self.projection_id + } + pub const fn generations(&self) -> ProjectionDocumentGenerations { + self.generations + } + pub const fn limit(&self) -> u16 { + self.limit + } + pub fn after(&self) -> Option<(ProjectionGeneration, &str)> { + self.after + .as_ref() + .map(|(generation, key)| (*generation, key.as_str())) + } + pub fn matches(&self, generation: ProjectionGeneration, key: &str) -> bool { + (match self.generations { + ProjectionDocumentGenerations::Exact(expected) => generation == expected, + ProjectionDocumentGenerations::All => true, + }) && self.after().is_none_or(|after| (generation, key) > after) + } +} + +/// In-process continuation; the next call independently supplies its scope. +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProjectionDocumentCursor { + projection_id: ProjectionId, + generations: ProjectionDocumentGenerations, + after: (ProjectionGeneration, String), +} + +/// An absent document is explicit corruption requiring caller reconciliation. +#[derive(Clone, Eq, PartialEq)] +pub struct ProjectionDocumentRecord { + generation: ProjectionGeneration, + key: String, + document: Option<ProjectionDocument>, +} +impl ProjectionDocumentRecord { + pub fn new( + generation: ProjectionGeneration, + document: ProjectionDocument, + ) -> Result<Self, Error> { + ProjectionGeneration::new(*generation.as_bytes())?; + Ok(Self { + generation, + key: document.key().to_owned(), + document: Some(document), + }) + } + pub fn corrupt(generation: ProjectionGeneration, key: String) -> Result<Self, Error> { + ProjectionGeneration::new(*generation.as_bytes())?; + if !valid_document_key(&key) { + return Err(Error::CorruptProjectionDocument); + } + Ok(Self { + generation, + key, + document: None, + }) + } + pub const fn generation(&self) -> ProjectionGeneration { + self.generation + } + pub fn key(&self) -> &str { + &self.key + } + pub const fn document(&self) -> Option<&ProjectionDocument> { + self.document.as_ref() + } +} +impl std::fmt::Debug for ProjectionDocumentRecord { + fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + formatter + .debug_struct("ProjectionDocumentRecord") + .field("generation", &self.generation) + .field("key", &self.key) + .field("corrupt", &self.document.is_none()) + .finish() + } +} + +#[derive(Clone, Debug, Eq, PartialEq)] +pub struct ProjectionDocumentPage { + records: Vec<ProjectionDocumentRecord>, + next_cursor: Option<ProjectionDocumentCursor>, +} +impl ProjectionDocumentPage { + /// Backends must bound bytes read, including invalid payloads, before constructing a page. + pub fn new( + query: &ProjectionDocumentQuery, + records: Vec<ProjectionDocumentRecord>, + has_more: bool, + ) -> Result<Self, Error> { + if records.len() > usize::from(query.limit) || (has_more && records.is_empty()) { + return Err(Error::InvalidProjectionDocument); + } + let mut previous = query.after(); + let mut bytes = 0; + for record in &records { + let position = (record.generation, record.key()); + if !query.matches(position.0, position.1) + || previous.is_some_and(|last| position <= last) + { + return Err(Error::InvalidProjectionDocument); + } + bytes += record + .document() + .map_or(0, |document| document.value().len()); + if bytes > PROJECTION_DOCUMENT_PAGE_BYTES_MAX { + return Err(Error::InvalidProjectionDocument); + } + previous = Some(position); + } + let next_cursor = if has_more { + previous.map(|(generation, key)| ProjectionDocumentCursor { + projection_id: query.projection_id.clone(), + generations: query.generations, + after: (generation, key.to_owned()), + }) + } else { + None + }; + Ok(Self { + records, + next_cursor, + }) + } + pub fn records(&self) -> &[ProjectionDocumentRecord] { + &self.records + } + pub const fn next_cursor(&self) -> Option<&ProjectionDocumentCursor> { + self.next_cursor.as_ref() + } + pub fn into_records(self) -> Vec<ProjectionDocumentRecord> { + self.records + } +} diff --git a/crates/storage/src/projection/document_query/tests.rs b/crates/storage/src/projection/document_query/tests.rs @@ -0,0 +1,153 @@ +use super::*; + +fn generation() -> ProjectionGeneration { + ProjectionGeneration::new([1; 32]).unwrap() +} +fn query(limit: u16) -> ProjectionDocumentQuery { + ProjectionDocumentQuery::new( + ProjectionId::parse("fixture").unwrap(), + ProjectionDocumentGenerations::All, + limit, + ) + .unwrap() +} +fn record(key: &str, size: usize) -> ProjectionDocumentRecord { + ProjectionDocumentRecord::new( + generation(), + ProjectionDocument::new(key.into(), vec![123; size]).unwrap(), + ) + .unwrap() +} + +#[test] +fn independent_scope_and_bounded_page_construction() { + let q = query(2); + for limit in [0, 257] { + assert!( + ProjectionDocumentQuery::new(q.projection_id().clone(), q.generations(), limit) + .is_err() + ); + } + for key in ["", " bad", "bad\n", &"x".repeat(513)] { + assert!(ProjectionDocumentRecord::corrupt(generation(), key.into()).is_err()); + } + let corrupt = ProjectionDocumentRecord::corrupt(generation(), "b".into()).unwrap(); + let page = + ProjectionDocumentPage::new(&q, vec![record("a", 1), corrupt.clone()], true).unwrap(); + let cursor = page.next_cursor().unwrap(); + assert!(corrupt.document().is_none()); + assert_eq!(corrupt.generation(), generation()); + assert_eq!(corrupt.key(), "b"); + assert_eq!(page.clone().into_records(), page.records()); + assert!(!format!("{:?}", record("redacted", 100)).contains("123")); + assert_eq!( + q.clone().with_cursor(cursor).unwrap().after(), + Some((generation(), "b")) + ); + for changed in [ + ProjectionDocumentQuery::new(ProjectionId::parse("other").unwrap(), q.generations(), 1) + .unwrap(), + ProjectionDocumentQuery::new( + q.projection_id().clone(), + ProjectionDocumentGenerations::Exact(generation()), + 1, + ) + .unwrap(), + ] { + assert!(changed.with_cursor(cursor).is_err()); + } + assert!(ProjectionDocumentPage::new(&q, vec![], true).is_err()); + assert!(ProjectionDocumentPage::new(&q, vec![record("a", 1), record("a", 1)], false).is_err()); + assert!( + ProjectionDocumentPage::new(&query(1), vec![record("a", 1), record("b", 1)], false) + .is_err() + ); + let resumed = q.clone().with_cursor(cursor).unwrap(); + assert!(ProjectionDocumentPage::new(&resumed, vec![record("b", 1)], false).is_err()); + assert!( + ProjectionDocumentPage::new( + &q, + vec![ + record("a", PROJECTION_DOCUMENT_PAGE_BYTES_MAX), + record("b", 1) + ], + false + ) + .is_err() + ); + let exact = ProjectionDocumentQuery::new( + q.projection_id().clone(), + ProjectionDocumentGenerations::Exact(ProjectionGeneration::new([2; 32]).unwrap()), + 2, + ) + .unwrap(); + assert!(ProjectionDocumentPage::new(&exact, vec![record("a", 1)], false).is_err()); + assert!( + ProjectionDocumentPage::new(&q, vec![], false) + .unwrap() + .next_cursor() + .is_none() + ); +} + +#[cfg(feature = "serde")] +#[test] +fn unchecked_legacy_serde_values_cannot_enter_queries_or_records() { + let invalid_id: ProjectionId = serde_json::from_str("\"bad id\"").unwrap(); + let zero: ProjectionGeneration = + serde_json::from_str(&serde_json::to_string(&[0; 32]).unwrap()).unwrap(); + assert!( + ProjectionDocumentQuery::new(invalid_id, ProjectionDocumentGenerations::All, 1).is_err() + ); + assert!( + ProjectionDocumentQuery::new( + query(1).projection_id().clone(), + ProjectionDocumentGenerations::Exact(zero), + 1 + ) + .is_err() + ); + assert!(ProjectionDocumentRecord::corrupt(zero, "key".into()).is_err()); + assert!( + ProjectionDocumentRecord::new( + zero, + ProjectionDocument::new("key".into(), vec![1]).unwrap() + ) + .is_err() + ); +} + +#[cfg(feature = "memory")] +#[test] +fn memory_inventory_continues_after_byte_exhaustion_and_reports_closed_store() { + use crate::{ProjectionStore, event::SourceGeneration, memory::MemoryStorage}; + futures_executor::block_on(async { + let store = MemoryStorage::new(SourceGeneration::new([1; 32]).unwrap()); + let q = query(256); + for (key, size) in [("a", PROJECTION_DOCUMENT_PAGE_BYTES_MAX), ("b", 1)] { + store + .put_projection_document( + q.projection_id().clone(), + generation(), + ProjectionDocument::new(key.into(), vec![1; size]).unwrap(), + ) + .await + .unwrap(); + } + let first = store.query_projection_documents(q.clone()).await.unwrap(); + assert_eq!(first.records().len(), 1); + let second = store + .query_projection_documents(q.with_cursor(first.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert_eq!(second.records()[0].key(), "b"); + assert!(second.next_cursor().is_none()); + crate::backup::StorageReliability::close(&store) + .await + .unwrap(); + assert_eq!( + store.query_projection_documents(query(1)).await, + Err(Error::BackendUnavailable) + ); + }); +} diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md @@ -2,6 +2,13 @@ SQLite storage backend for Radroots. +Projection document inventory uses the existing canonical projection table. +Each page releases its read snapshot before returning a continuation. Metadata +is bounded before decoding; the 16 MiB payload budget is charged before each +value is loaded, including values that fail digest verification. Corrupt +payloads retain an explicit locator, while unpageable corrupt keys or +generations fail the query. Cross-page mutation fencing remains caller-owned. + Runtime schema v15 permits late verified authored signatures to remain durable alongside a signing stop. The logical snapshot preserves cancelled or terminal state; a separate SQL stop column preserves the original signed/raw constraint. diff --git a/crates/storage_sqlite/src/projection/document_query.rs b/crates/storage_sqlite/src/projection/document_query.rs @@ -0,0 +1,112 @@ +use super::{SqliteStorage, map_backend}; +use radroots_storage::{ + Error, + projection::{ + ProjectionDocument, ProjectionGeneration, + document_query::{ + PROJECTION_DOCUMENT_PAGE_BYTES_MAX, ProjectionDocumentGenerations, + ProjectionDocumentPage, ProjectionDocumentQuery, ProjectionDocumentRecord, + }, + }, +}; +use sqlx::Row; + +pub(super) async fn page( + store: &SqliteStorage, + query: ProjectionDocumentQuery, +) -> Result<ProjectionDocumentPage, Error> { + let mut transaction = store.pool().begin().await.map_err(map_backend)?; + let result = read_page(&mut transaction, &query).await; + let rollback = transaction.rollback().await.map_err(map_backend); + match result { + Ok(page) => { + rollback?; + Ok(page) + } + Err(error) => Err(error), + } +} + +async fn read_page( + transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, + query: &ProjectionDocumentQuery, +) -> Result<ProjectionDocumentPage, Error> { + let generation = match query.generations() { + ProjectionDocumentGenerations::All => None, + ProjectionDocumentGenerations::Exact(value) => Some(value.as_bytes().to_vec()), + }; + let after_generation = query + .after() + .map(|(generation, _)| generation.as_bytes().to_vec()); + let after_key = query.after().map(|(_, key)| key); + // CASE bounds every variable-width metadata value before SQLx allocates it. + // Payloads are loaded separately, only after charging this page's byte budget. + let rows = sqlx::query( + "SELECT CASE WHEN length(generation) = 32 THEN generation END AS generation, + CASE WHEN length(CAST(document_key AS BLOB)) BETWEEN 1 AND 512 THEN document_key END AS document_key, + CASE WHEN length(value_sha256) = 32 THEN value_sha256 END AS value_sha256, + length(value) AS value_bytes + FROM radroots_runtime_projection_documents AS documents + WHERE projection_id = ? AND (? IS NULL OR generation = ?) + AND (? IS NULL OR (generation, document_key) > (?, ?)) + ORDER BY documents.generation, documents.document_key LIMIT ?" + ).bind(query.projection_id().as_str()).bind(&generation).bind(generation) + .bind(&after_generation).bind(&after_generation).bind(after_key) + .bind(i64::from(query.limit()) + 1) + .fetch_all(&mut **transaction).await.map_err(map_backend)?; + let mut has_more = rows.len() > usize::from(query.limit()); + let mut records = Vec::new(); + let mut bytes = 0; + for row in rows.iter().take(usize::from(query.limit())) { + let generation: Vec<u8> = row + .try_get("generation") + .map_err(|_| Error::CorruptProjectionDocument)?; + let generation = ProjectionGeneration::new( + generation + .try_into() + .map_err(|_| Error::CorruptProjectionDocument)?, + ) + .map_err(|_| Error::CorruptProjectionDocument)?; + let key: String = row + .try_get("document_key") + .map_err(|_| Error::CorruptProjectionDocument)?; + let corrupt = ProjectionDocumentRecord::corrupt(generation, key.clone())?; + let digest: Option<Vec<u8>> = row + .try_get("value_sha256") + .map_err(|_| Error::CorruptProjectionDocument)?; + let size: i64 = row + .try_get("value_bytes") + .map_err(|_| Error::CorruptProjectionDocument)?; + let Some(digest) = digest else { + records.push(corrupt); + continue; + }; + let digest: [u8; 32] = digest + .try_into() + .map_err(|_| Error::CorruptProjectionDocument)?; + if size <= 0 || size > PROJECTION_DOCUMENT_PAGE_BYTES_MAX as i64 { + records.push(corrupt); + continue; + } + let size = size as usize; + if bytes + size > PROJECTION_DOCUMENT_PAGE_BYTES_MAX { + has_more = true; + break; + } + bytes += size; + let value: Vec<u8> = sqlx::query_scalar( + "SELECT value FROM radroots_runtime_projection_documents WHERE projection_id = ? AND generation = ? AND document_key = ?" + ).bind(query.projection_id().as_str()).bind(generation.as_bytes().as_slice()).bind(&key) + .fetch_one(&mut **transaction).await.map_err(map_backend)?; + records.push( + match ProjectionDocument::from_stored_parts(key, value, digest) { + Ok(document) => ProjectionDocumentRecord::new(generation, document)?, + Err(_) => corrupt, + }, + ); + } + ProjectionDocumentPage::new(query, records, has_more) +} + +#[cfg(test)] +mod tests; diff --git a/crates/storage_sqlite/src/projection/document_query/tests.rs b/crates/storage_sqlite/src/projection/document_query/tests.rs @@ -0,0 +1,246 @@ +use super::*; +use crate::{OpenMode, OpenOptions, Paths}; +use radroots_storage::{ + ProjectionStore, event::SourceGeneration, memory::MemoryStorage, projection::ProjectionId, +}; +use tempfile::TempDir; + +async fn open(temp: &TempDir) -> SqliteStorage { + SqliteStorage::open( + OpenOptions::new( + Paths::from_directory(temp.path()).unwrap(), + OpenMode::Create, + ) + .with_source_generation(SourceGeneration::new([9; 32]).unwrap(), 9) + .unwrap(), + ) + .await + .unwrap() +} +fn generation(byte: u8) -> ProjectionGeneration { + ProjectionGeneration::new([byte; 32]).unwrap() +} +fn query(selection: ProjectionDocumentGenerations, limit: u16) -> ProjectionDocumentQuery { + ProjectionDocumentQuery::new(ProjectionId::parse("fixture").unwrap(), selection, limit).unwrap() +} +async fn put(store: &dyn ProjectionStore, generation: u8, key: &str, value: Vec<u8>) { + store + .put_projection_document( + ProjectionId::parse("fixture").unwrap(), + self::generation(generation), + ProjectionDocument::new(key.into(), value).unwrap(), + ) + .await + .unwrap(); +} + +#[tokio::test] +async fn inventory_thousand_records_matches_memory_across_generations_updates_and_reopen() { + let temp = TempDir::new().unwrap(); + let store = open(&temp).await; + let memory = MemoryStorage::new(SourceGeneration::new([9; 32]).unwrap()); + for i in (0..1000).rev() { + for backend in [&store as &dyn ProjectionStore, &memory] { + put(backend, 1 + (i % 2) as u8, &format!("key.{i:04}"), vec![1]).await; + } + } + store + .put_projection_document( + ProjectionId::parse("foreign").unwrap(), + generation(1), + ProjectionDocument::new("key.0000".into(), vec![3]).unwrap(), + ) + .await + .unwrap(); + let q = query(ProjectionDocumentGenerations::All, 37); + let first = store.query_projection_documents(q.clone()).await.unwrap(); + assert_eq!( + first, + memory.query_projection_documents(q.clone()).await.unwrap() + ); + for backend in [&store as &dyn ProjectionStore, &memory] { + put(backend, 1, "key.0000", vec![2]).await; + put(backend, 2, "key.0999", vec![2]).await; + } + store.close().await.unwrap(); + let store = open(&temp).await; + let mut count = first.records().len(); + let mut cursor = first.next_cursor().cloned(); + while let Some(current) = cursor { + let query = q.clone().with_cursor(¤t).unwrap(); + let page = store + .query_projection_documents(query.clone()) + .await + .unwrap(); + assert_eq!( + page, + memory.query_projection_documents(query).await.unwrap() + ); + count += page.records().len(); + cursor = page.next_cursor().cloned(); + } + assert_eq!(count, 1000); + for byte in [1, 2, 3] { + let q = query(ProjectionDocumentGenerations::Exact(generation(byte)), 256); + let page = store.query_projection_documents(q.clone()).await.unwrap(); + assert_eq!( + page, + memory.query_projection_documents(q.clone()).await.unwrap() + ); + assert!( + page.records() + .iter() + .all(|r| r.generation() == generation(byte)) + ); + if byte == 3 { + assert!(page.records().is_empty()); + } else { + let last = store + .query_projection_documents(q.with_cursor(page.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert_eq!(last.records().len() + page.records().len(), 500); + assert!(last.next_cursor().is_none()); + } + } + store.close().await.unwrap(); + assert_eq!( + store.query_projection_documents(q).await, + Err(Error::BackendUnavailable) + ); +} + +#[tokio::test] +async fn inventory_readonly_preserves_identical_keys_and_byte_order() { + let temp = TempDir::new().unwrap(); + let store = open(&temp).await; + for byte in [1, 2] { + for key in ["same", "z", "é"] { + put(&store, byte, key, vec![byte]).await; + } + } + store.close().await.unwrap(); + let store = SqliteStorage::open(OpenOptions::new( + Paths::from_directory(temp.path()).unwrap(), + OpenMode::ReadOnly, + )) + .await + .unwrap(); + let q = query(ProjectionDocumentGenerations::All, 1); + let mut current = q.clone(); + let mut found = Vec::new(); + loop { + let page = store.query_projection_documents(current).await.unwrap(); + found.extend( + page.records() + .iter() + .map(|r| (r.generation(), r.key().to_owned())), + ); + let Some(cursor) = page.next_cursor() else { + break; + }; + current = q.clone().with_cursor(cursor).unwrap(); + } + assert_eq!( + found, + [1, 2] + .into_iter() + .flat_map(|byte| ["same", "z", "é"].map(|key| (generation(byte), key.to_owned()))) + .collect::<Vec<_>>() + ); + let exact = query(ProjectionDocumentGenerations::Exact(generation(1)), 1); + let page = store.query_projection_documents(exact).await.unwrap(); + assert!( + query(ProjectionDocumentGenerations::Exact(generation(2)), 1) + .with_cursor(page.next_cursor().unwrap()) + .is_err() + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn inventory_bounds_bytes_including_corrupt_payloads_without_losing_continuation() { + let temp = TempDir::new().unwrap(); + let store = open(&temp).await; + put(&store, 1, "a", vec![7; PROJECTION_DOCUMENT_PAGE_BYTES_MAX]).await; + put(&store, 1, "b", vec![8]).await; + let q = query(ProjectionDocumentGenerations::All, 256); + for corrupt in [false, true] { + if corrupt { + sqlx::query("UPDATE radroots_runtime_projection_documents SET value_sha256 = zeroblob(32) WHERE document_key = 'a'") + .execute(store.pool()).await.unwrap(); + } + let first = store.query_projection_documents(q.clone()).await.unwrap(); + assert_eq!(first.records().len(), 1); + assert_eq!(first.records()[0].document().is_none(), corrupt); + let next = store + .query_projection_documents( + q.clone().with_cursor(first.next_cursor().unwrap()).unwrap(), + ) + .await + .unwrap(); + assert_eq!(next.records()[0].key(), "b"); + assert!(next.next_cursor().is_none()); + } + store.close().await.unwrap(); +} + +#[tokio::test] +async fn inventory_exposes_bounded_corruption_and_rejects_unpageable_keys() { + let temp = TempDir::new().unwrap(); + let store = open(&temp).await; + // One held connection makes disabled-check injection local to this fixture. + let mut connection = store.pool().acquire().await.unwrap(); + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(&mut *connection) + .await + .unwrap(); + for (key, size, digest) in [("a", 0, 32), ("b", 16777217, 32), ("c", 1, 33)] { + sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', ?, ?, zeroblob(?), zeroblob(?))") + .bind(generation(1).as_bytes().as_slice()).bind(key).bind(size).bind(digest) + .execute(&mut *connection).await.unwrap(); + } + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(&mut *connection) + .await + .unwrap(); + drop(connection); + put(&store, 1, "d", vec![1]).await; + let q = query(ProjectionDocumentGenerations::All, 2); + let page = store.query_projection_documents(q.clone()).await.unwrap(); + assert!(page.records().iter().all(|r| r.document().is_none())); + let next = store + .query_projection_documents(q.with_cursor(page.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert!(next.records()[0].document().is_none()); + assert!(next.records()[1].document().is_some()); + assert!(next.next_cursor().is_none()); + for key in ["\n".to_owned(), "é".repeat(512)] { + sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', ?, ?, x'01', zeroblob(32))") + .bind(generation(2).as_bytes().as_slice()).bind(&key).execute(store.pool()).await.unwrap(); + assert_eq!( + store + .query_projection_documents(query( + ProjectionDocumentGenerations::Exact(generation(2)), + 1 + )) + .await, + Err(Error::CorruptProjectionDocument) + ); + sqlx::query("DELETE FROM radroots_runtime_projection_documents WHERE generation = ?") + .bind(generation(2).as_bytes().as_slice()) + .execute(store.pool()) + .await + .unwrap(); + } + sqlx::query("INSERT INTO radroots_runtime_projection_documents VALUES ('fixture', zeroblob(32), 'zero', x'01', zeroblob(32))") + .execute(store.pool()).await.unwrap(); + assert_eq!( + store + .query_projection_documents(query(ProjectionDocumentGenerations::All, 1)) + .await, + Err(Error::CorruptProjectionDocument) + ); + store.close().await.unwrap(); +} diff --git a/crates/storage_sqlite/src/projection/mod.rs b/crates/storage_sqlite/src/projection/mod.rs @@ -13,6 +13,8 @@ use radroots_storage::{ }; use sqlx::{Row, Sqlite, SqliteConnection}; +mod document_query; + #[cfg_attr(coverage_nightly, coverage(off))] impl ProjectionStore for SqliteStorage { fn status( @@ -505,6 +507,16 @@ impl ProjectionStore for SqliteStorage { }) } + fn query_projection_documents( + &self, + query: radroots_storage::projection::document_query::ProjectionDocumentQuery, + ) -> BoxFuture< + '_, + Result<radroots_storage::projection::document_query::ProjectionDocumentPage, Error>, + > { + Box::pin(document_query::page(self, query)) + } + fn put_projection_snapshot( &self, snapshot: ProjectionSnapshot,