lib

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

commit 4a2285ea896cfc361f539c74c467a3fd5591ec74
parent b6667541b08b4e578199c6612cf51b7402071f5e
Author: triesap <tyson@radroots.org>
Date:   Sun, 20 Sep 2026 06:59:41 +0000

storage: inventory all authored payload schemas

- Bind all-schema traversal to explicit author and continuation authority
- Preserve exact queries and bounded metadata and payload reads
- Verify memory and SQLite parity through updates, reopen and byte limits
- Qualify public APIs, coverage, workspace and standalone release checks

Diffstat:
Mcontracts/api_baselines/radroots_storage.txt | 4+++-
Acontracts/architecture/decisions/authored_draft_inventory.v2.json | 13+++++++++++++
Mcrates/storage/src/authored_draft_query.rs | 83++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Acrates/storage/tests/authored_draft_all_schemas.rs | 210+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage/tests/authored_draft_query.rs | 4++--
Mcrates/storage_sqlite/src/authored_draft.rs | 5+++++
Acrates/storage_sqlite/src/authored_draft_all_schemas_tests.rs | 171+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/src/authored_draft_query.rs | 4++--
Mcrates/storage_sqlite/src/authored_draft_query_tests.rs | 4++--
9 files changed, 480 insertions(+), 18 deletions(-)

diff --git a/contracts/api_baselines/radroots_storage.txt b/contracts/api_baselines/radroots_storage.txt @@ -527,11 +527,12 @@ pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::after(& pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::author(&self) -> &[u8; 32] pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::cursor_after(&self, [u8; 16]) -> radroots_storage::authored_draft_query::AuthoredDraftCursor pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::for_author([u8; 32], impl core::convert::AsRef<str>, u16) -> core::result::Result<Self, radroots_storage::Error> +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::for_author_all_schemas([u8; 32], u16) -> core::result::Result<Self, radroots_storage::Error> pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::is_author_wide(&self) -> bool pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::limit(&self) -> u16 pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::matches(&self, &radroots_storage::authored_draft::AuthoredDraft) -> bool pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::new([u8; 32], impl core::convert::AsRef<str>, core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope>, u16) -> core::result::Result<Self, radroots_storage::Error> -pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::payload_schema(&self) -> &str +pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::payload_schema(&self) -> core::option::Option<&str> pub const fn radroots_storage::authored_draft_query::AuthoredDraftQuery::scope(&self) -> core::option::Option<radroots_storage::authored_draft_query::AuthoredDraftScope> pub fn radroots_storage::authored_draft_query::AuthoredDraftQuery::with_cursor(self, &radroots_storage::authored_draft_query::AuthoredDraftCursor) -> core::result::Result<Self, radroots_storage::Error> pub struct radroots_storage::authored_draft_query::AuthoredDraftScope(_) @@ -543,6 +544,7 @@ pub fn [u8; 32]::from(radroots_storage::authored_draft_query::AuthoredDraftScope impl core::convert::TryFrom<[u8; 32]> for radroots_storage::authored_draft_query::AuthoredDraftScope pub type radroots_storage::authored_draft_query::AuthoredDraftScope::Error = radroots_storage::Error pub fn radroots_storage::authored_draft_query::AuthoredDraftScope::try_from([u8; 32]) -> core::result::Result<Self, radroots_storage::Error> +pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION: u16 pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION: u16 pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION: u16 pub const radroots_storage::authored_draft_query::AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES: usize diff --git a/contracts/architecture/decisions/authored_draft_inventory.v2.json b/contracts/architecture/decisions/authored_draft_inventory.v2.json @@ -0,0 +1,13 @@ +{ + "schema": "radroots.authored-draft-inventory.v2", + "status": "implemented", + "owners": ["radroots_storage", "radroots_storage_sqlite"], + "authority": "Extend authored_draft_inventory.v1 without changing its exact optional-scope selection, persisted records or version-1/2 continuation bytes.", + "selection": "AuthoredDraftQuery::for_author_all_schemas explicitly selects all current payload schemas for one independently supplied author across scoped and unscoped heads. payload_schema() returns an optional exact schema; None denotes only this explicit all-schema selection, never an empty-string wildcard. Existing new and for_author constructors still require an exact nonempty schema.", + "cursor": "All-schema continuations require version 3, explicit author_all_schemas selection, null payload_schema and null scope. Required fields, unknown fields, unsupported versions, selection mismatches and author mismatches fail closed. Every call supplies independent query authority; a cursor cannot broaden an exact query or switch authors.", + "bounds": {"page_records": 256, "metadata_lookahead": 1, "decoded_page_payload_bytes": 4194304, "serialized_page_snapshot_bytes": 16777216}, + "backend": "Stable ascending draft IDs and latest revisions. Unknown nonempty application schemas remain opaque valid records; malformed snapshots and unknown empty schema yield bounded corruption locators. Metadata and byte budgets are checked before SQLite snapshot loading. Foreign authors remain outside selection. Existing exact-schema behavior is unchanged.", + "consistency": "One released SQLite read snapshot per page, bounded retained memory, and explicit continuation after record or byte exhaustion. This is live traversal, not a frozen multi-page snapshot or deletion authority. Callers fence their own mutations and retain resources on incomplete or unknown ownership.", + "compatibility": "No migrations, stored draft bytes, signed identities, dependencies or feature changes. The reviewed public schema accessor now represents explicit absence with Option. Owning SQLite selection and exact consumers must adopt that type directly; no alias or fallback API.", + "non_goals": ["application schema interpretation", "garbage collection policy", "author enumeration", "filesystem access", "new database", "raw SQL escape hatch", "mutation fencing", "release qualification"] +} diff --git a/crates/storage/src/authored_draft_query.rs b/crates/storage/src/authored_draft_query.rs @@ -14,6 +14,8 @@ pub const AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024; pub const AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION: u16 = 1; /// Explicit author/schema traversal has a distinct continuation authority. pub const AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION: u16 = 2; +/// Explicit author-wide traversal of every payload schema has separate authority. +pub const AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION: u16 = 3; /// An opaque, stable application-selected scope digest; never a credential. #[derive(Clone, Copy, Debug, Eq, Hash, Ord, PartialEq, PartialOrd)] @@ -48,7 +50,7 @@ impl From<AuthoredDraftScope> for [u8; 32] { #[derive(Clone, Debug, Eq, PartialEq)] pub struct AuthoredDraftQuery { author: [u8; 32], - payload_schema: String, + payload_schema: Option<String>, scope: Option<AuthoredDraftScope>, author_wide: bool, limit: u16, @@ -74,7 +76,7 @@ impl AuthoredDraftQuery { } Ok(Self { author, - payload_schema: schema.to_owned(), + payload_schema: Some(schema.to_owned()), scope, author_wide: false, limit, @@ -93,13 +95,34 @@ impl AuthoredDraftQuery { Ok(query) } - /// Whether scope filtering was explicitly omitted by `for_author`. + /// Selects every current payload schema for this author across all scopes. + /// Unknown application schemas remain opaque records, never absent owners. + pub fn for_author_all_schemas(author: [u8; 32], limit: u16) -> Result<Self, Error> { + if author.iter().all(|byte| *byte == 0) + || limit == 0 + || limit > AUTHORED_DRAFT_QUERY_LIMIT_MAX + { + return Err(Error::InvalidAuthoredDraft); + } + Ok(Self { + author, + payload_schema: None, + scope: None, + author_wide: true, + limit, + after: None, + }) + } + + /// Whether a named author-wide constructor explicitly omitted scope filtering. pub const fn is_author_wide(&self) -> bool { self.author_wide } const fn cursor_version(&self) -> u16 { - if self.author_wide { + if self.payload_schema.is_none() { + AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION + } else if self.author_wide { AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION } else { AUTHORED_DRAFT_CURSOR_SCHEMA_VERSION @@ -120,8 +143,9 @@ impl AuthoredDraftQuery { pub const fn author(&self) -> &[u8; 32] { &self.author } - pub fn payload_schema(&self) -> &str { - &self.payload_schema + /// Exact schema, or explicit all-schema selection. An empty string is never a wildcard. + pub fn payload_schema(&self) -> Option<&str> { + self.payload_schema.as_deref() } /// Exact optional scope for ordinary queries. Author-wide queries return /// `None`; backends must also honor `is_author_wide` when selecting rows. @@ -136,7 +160,9 @@ impl AuthoredDraftQuery { } pub fn matches(&self, draft: &AuthoredDraft) -> bool { draft.author() == &self.author - && draft.payload_schema() == self.payload_schema + && self + .payload_schema() + .is_none_or(|schema| draft.payload_schema() == schema) && (self.author_wide || draft.scope() == self.scope) } pub fn cursor_after(&self, after_id: [u8; 16]) -> AuthoredDraftCursor { @@ -158,7 +184,7 @@ impl AuthoredDraftQuery { pub struct AuthoredDraftCursor { schema_version: u16, author: [u8; 32], - payload_schema: String, + payload_schema: Option<String>, scope: Option<AuthoredDraftScope>, after_id: [u8; 16], } @@ -168,6 +194,7 @@ pub struct AuthoredDraftCursor { enum CursorWire { Exact(ExactCursorWire), Author(AuthorCursorWire), + AllSchemas(AllSchemasCursorWire), } #[cfg(feature = "serde")] @@ -191,11 +218,25 @@ struct AuthorCursorWire { after_id: [u8; 16], selection: CursorSelection, } + +#[cfg(feature = "serde")] +#[derive(serde::Serialize, serde::Deserialize)] +#[serde(deny_unknown_fields)] +struct AllSchemasCursorWire { + schema_version: u16, + author: [u8; 32], + payload_schema: (), + scope: (), + after_id: [u8; 16], + selection: CursorSelection, +} #[cfg(feature = "serde")] #[derive(serde::Serialize, serde::Deserialize)] enum CursorSelection { #[serde(rename = "author_schema")] AuthorSchema, + #[serde(rename = "author_all_schemas")] + AuthorAllSchemas, } #[cfg(feature = "serde")] impl TryFrom<CursorWire> for AuthoredDraftCursor { @@ -211,13 +252,23 @@ impl TryFrom<CursorWire> for AuthoredDraftCursor { ) } CursorWire::Author(value) - if value.schema_version == AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION => + if value.schema_version == AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION + && matches!(value.selection, CursorSelection::AuthorSchema) => { ( AuthoredDraftQuery::for_author(value.author, value.payload_schema, 1)?, value.after_id, ) } + CursorWire::AllSchemas(value) + if value.schema_version == AUTHORED_DRAFT_ALL_SCHEMAS_CURSOR_SCHEMA_VERSION + && matches!(value.selection, CursorSelection::AuthorAllSchemas) => + { + ( + AuthoredDraftQuery::for_author_all_schemas(value.author, 1)?, + value.after_id, + ) + } _ => return Err(Error::InvalidAuthoredDraft), }; Ok(query.cursor_after(after)) @@ -226,11 +277,21 @@ impl TryFrom<CursorWire> for AuthoredDraftCursor { #[cfg(feature = "serde")] impl From<AuthoredDraftCursor> for CursorWire { fn from(value: AuthoredDraftCursor) -> Self { + let Some(payload_schema) = value.payload_schema else { + return Self::AllSchemas(AllSchemasCursorWire { + schema_version: value.schema_version, + author: value.author, + payload_schema: (), + scope: (), + after_id: value.after_id, + selection: CursorSelection::AuthorAllSchemas, + }); + }; if value.schema_version == AUTHORED_DRAFT_AUTHOR_CURSOR_SCHEMA_VERSION { return Self::Author(AuthorCursorWire { schema_version: value.schema_version, author: value.author, - payload_schema: value.payload_schema, + payload_schema, scope: (), after_id: value.after_id, selection: CursorSelection::AuthorSchema, @@ -239,7 +300,7 @@ impl From<AuthoredDraftCursor> for CursorWire { Self::Exact(ExactCursorWire { schema_version: value.schema_version, author: value.author, - payload_schema: value.payload_schema, + payload_schema, scope: value.scope, after_id: value.after_id, }) diff --git a/crates/storage/tests/authored_draft_all_schemas.rs b/crates/storage/tests/authored_draft_all_schemas.rs @@ -0,0 +1,210 @@ +use futures_executor::block_on; +use radroots_storage::{ + authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, + authored_draft_query::{ + AuthoredDraftCursor, AuthoredDraftPage, AuthoredDraftQuery, AuthoredDraftQueryRecord, + AuthoredDraftScope, + }, + memory::MemoryStorage, +}; + +fn draft(id: u128, author: u8, schema: &str, payload: Vec<u8>) -> AuthoredDraft { + let value = AuthoredDraft::initial( + AuthoredDraftId::new(id.to_be_bytes()).unwrap(), + [author; 32], + schema, + payload, + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + if id.is_multiple_of(2) { + value + .with_scope(AuthoredDraftScope::new([3; 32]).unwrap()) + .unwrap() + } else { + value + } +} + +#[test] +fn all_schema_cursor_has_explicit_independent_authority_and_no_wildcard_schema() { + let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 37).unwrap(); + assert!(q.is_author_wide()); + assert_eq!(q.scope(), None); + assert_eq!(q.payload_schema(), None); + for (author, limit) in [([0; 32], 1), ([7; 32], 0), ([7; 32], 257)] { + assert!(AuthoredDraftQuery::for_author_all_schemas(author, limit).is_err()); + } + let cursor = q.cursor_after([0; 16]); + let encoded = serde_json::to_string(&cursor).unwrap(); + let expected = format!( + "{{\"schema_version\":3,\"author\":{},\"payload_schema\":null,\"scope\":null,\"after_id\":{},\"selection\":\"author_all_schemas\"}}", + serde_json::to_string(&[7; 32]).unwrap(), + serde_json::to_string(&[0; 16]).unwrap() + ); + assert_eq!(encoded, expected); + let decoded: AuthoredDraftCursor = serde_json::from_str(&encoded).unwrap(); + assert_eq!(decoded, cursor); + assert_eq!( + q.clone().with_cursor(&decoded).unwrap().after(), + Some([0; 16]) + ); + assert_eq!(serde_json::to_string(&decoded).unwrap(), encoded); + for other in [ + AuthoredDraftQuery::new([7; 32], "known.v1", None, 1).unwrap(), + AuthoredDraftQuery::for_author([7; 32], "known.v1", 1).unwrap(), + AuthoredDraftQuery::for_author_all_schemas([8; 32], 1).unwrap(), + ] { + assert!(q.clone().with_cursor(&other.cursor_after([0; 16])).is_err()); + assert!(other.with_cursor(&cursor).is_err()); + } + let wire = serde_json::to_value(&cursor).unwrap(); + for (field, value) in [ + ("schema_version", serde_json::json!(1)), + ("schema_version", serde_json::json!(2)), + ("schema_version", serde_json::json!(4)), + ("selection", serde_json::json!("author_schema")), + ("selection", serde_json::Value::Null), + ("selection", serde_json::json!("future")), + ("payload_schema", serde_json::json!("known.v1")), + ("payload_schema", serde_json::json!("")), + ("scope", serde_json::json!([3; 32].to_vec())), + ("author", serde_json::json!([0; 32].to_vec())), + ("after_id", serde_json::json!([0; 15].to_vec())), + ("unexpected", serde_json::json!(true)), + ] { + let mut changed = wire.clone(); + changed[field] = value; + assert!( + serde_json::from_value::<AuthoredDraftCursor>(changed).is_err(), + "{field}" + ); + } + for field in [ + "schema_version", + "author", + "payload_schema", + "scope", + "after_id", + "selection", + ] { + let mut changed = wire.clone(); + changed.as_object_mut().unwrap().remove(field); + assert!( + serde_json::from_value::<AuthoredDraftCursor>(changed).is_err(), + "missing {field}" + ); + } + let duplicate = encoded.replace("\"scope\":null", "\"scope\":null,\"scope\":null"); + assert!(serde_json::from_str::<AuthoredDraftCursor>(&duplicate).is_err()); + let exact = AuthoredDraftQuery::for_author([7; 32], "known.v1", 1).unwrap(); + let mut wrong_marker = serde_json::to_value(exact.cursor_after([0; 16])).unwrap(); + wrong_marker["selection"] = serde_json::json!("author_all_schemas"); + assert!(serde_json::from_value::<AuthoredDraftCursor>(wrong_marker).is_err()); + assert_eq!(exact.payload_schema(), Some("known.v1")); + assert!(AuthoredDraftQuery::new([7; 32], "", None, 1).is_err()); +} + +#[test] +fn all_schema_memory_inventory_covers_unknown_schemas_and_scopes_beyond_one_thousand() { + let store = MemoryStorage::default(); + for id in (1..=1001).rev() { + let schema = if id % 3 == 0 { + "future.unknown.v999" + } else { + "known.v1" + }; + block_on(store.append_authored_draft(draft(id, 7, schema, vec![1]), None)).unwrap(); + } + block_on(store.append_authored_draft(draft(1002, 8, "future.unknown.v999", vec![1]), None)) + .unwrap(); + let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 37).unwrap(); + let mut next = None; + let mut found = Vec::new(); + loop { + let query = next.as_ref().map_or_else( + || q.clone(), + |cursor| q.clone().with_cursor(cursor).unwrap(), + ); + let page = block_on(store.query_authored_drafts(query)).unwrap(); + assert!(page.records().len() <= 37); + for record in page.records() { + let AuthoredDraftQueryRecord::Draft(value) = record else { + panic!("valid fixture") + }; + let id = u128::from_be_bytes(record.draft_key()); + found.push(id); + assert_eq!( + value.payload_schema(), + if id % 3 == 0 { + "future.unknown.v999" + } else { + "known.v1" + } + ); + if id == 1001 { + assert_eq!(value.revision().get(), 2); + } + } + next = page.next_cursor().cloned(); + if found.len() == 37 { + for id in [1_u128, 1001] { + let value = block_on( + store.authored_draft_head(AuthoredDraftId::new(id.to_be_bytes()).unwrap()), + ) + .unwrap() + .unwrap(); + let revised = value + .successor(vec![2], AuthoredDraftStage::Draft, None, 11) + .unwrap(); + block_on(store.append_authored_draft(revised, Some(value.revision()))).unwrap(); + } + } + if next.is_none() { + break; + } + } + assert_eq!(found, (1..=1001).collect::<Vec<_>>()); + let fresh = block_on(store.query_authored_drafts(q)).unwrap(); + assert_eq!(fresh.records()[0].revision().get(), 2); +} + +#[test] +fn all_schema_memory_pages_keep_payload_bounds_and_validate_authority() { + let store = MemoryStorage::default(); + let first = draft(1, 7, "known.v1", vec![1; 3 * 1024 * 1024]); + let second = draft(2, 7, "future.v999", vec![2; 3 * 1024 * 1024]); + for value in [first.clone(), second.clone()] { + block_on(store.append_authored_draft(value, None)).unwrap(); + } + let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 256).unwrap(); + let page = block_on(store.query_authored_drafts(q.clone())).unwrap(); + assert_eq!( + page.records(), + [AuthoredDraftQueryRecord::Draft(first.clone())] + ); + let next = block_on( + store.query_authored_drafts(q.clone().with_cursor(page.next_cursor().unwrap()).unwrap()), + ) + .unwrap(); + assert_eq!( + next.records(), + [AuthoredDraftQueryRecord::Draft(second.clone())] + ); + assert!(next.next_cursor().is_none()); + assert!( + AuthoredDraftPage::new( + &q, + vec![ + AuthoredDraftQueryRecord::Draft(first), + AuthoredDraftQueryRecord::Draft(second) + ], + false + ) + .is_err() + ); + let foreign = AuthoredDraftQueryRecord::Draft(draft(3, 8, "known.v1", vec![1])); + assert!(AuthoredDraftPage::new(&q, vec![foreign], false).is_err()); +} diff --git a/crates/storage/tests/authored_draft_query.rs b/crates/storage/tests/authored_draft_query.rs @@ -57,7 +57,7 @@ fn author_wide_cursor_requires_explicit_selection_and_independent_authority() { for changed in [ query(None, 1), query(Some(AuthoredDraftScope::new([7; 32]).unwrap()), 1), - AuthoredDraftQuery::for_author([8; 32], q.payload_schema(), 1).unwrap(), + AuthoredDraftQuery::for_author([8; 32], q.payload_schema().unwrap(), 1).unwrap(), AuthoredDraftQuery::for_author([9; 32], "fixture.other.v1", 1).unwrap(), ] { assert!(changed.with_cursor(&decoded).is_err()); @@ -262,7 +262,7 @@ fn query_and_cursor_reject_invalid_or_changed_authority() { let q = query(None, 256); let cursor = q.cursor_after([0; 16]); for changed in [ - AuthoredDraftQuery::new([8; 32], q.payload_schema(), None, 1).unwrap(), + AuthoredDraftQuery::new([8; 32], q.payload_schema().unwrap(), None, 1).unwrap(), AuthoredDraftQuery::new([9; 32], "fixture.profile.v1", None, 1).unwrap(), query(Some(AuthoredDraftScope::new([7; 32]).unwrap()), 1), ] { diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs @@ -639,6 +639,11 @@ mod query_tests; #[cfg(test)] #[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_draft_all_schemas_tests.rs"] +mod all_schema_query_tests; + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] #[path = "authored_durability_tests.rs"] mod durability_tests; diff --git a/crates/storage_sqlite/src/authored_draft_all_schemas_tests.rs b/crates/storage_sqlite/src/authored_draft_all_schemas_tests.rs @@ -0,0 +1,171 @@ +use super::{ + query_tests::{corrupt, draft}, + tests::open_store, +}; +use radroots_storage::{ + authored_draft::{AuthoredDraft, AuthoredDraftId, AuthoredDraftStage, AuthoredDraftStore}, + authored_draft_query::{ + AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES, AuthoredDraftQuery, AuthoredDraftQueryRecord, + AuthoredDraftScope, + }, +}; +use tempfile::TempDir; + +#[tokio::test] +async fn all_schema_sqlite_pages_include_unknown_and_corrupt_owners_across_reopen() { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let scope = AuthoredDraftScope::new([3; 32]).unwrap(); + for id in 1_u128..=1001 { + let value = AuthoredDraft::initial( + AuthoredDraftId::new(id.to_be_bytes()).unwrap(), + [7; 32], + if id % 3 == 0 { + "future.unknown.v999" + } else { + "known.v1" + }, + vec![1], + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap(); + let value = if id % 2 == 0 { + value.with_scope(scope).unwrap() + } else { + value + }; + store.append_authored_draft(value, None).await.unwrap(); + } + corrupt(&store, 0, 7, "future.unknown.v999", Some(scope)).await; + corrupt(&store, 2, 8, "future.unknown.v999", None).await; + corrupt(&store, 3, 7, "", None).await; + corrupt(&store, 4, 7, "future.oversized.v999", None).await; + let update_guard: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE type = 'trigger' AND name = 'radroots_runtime_authored_draft_revisions_update_guard'" + ).fetch_one(store.pool()).await.unwrap(); + sqlx::query("DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard") + .execute(store.pool()) + .await + .unwrap(); + // The production CHECK rejects oversized writes. Simulate a corrupted file + // on one held fixture connection, then restore enforcement before querying. + let mut connection = store.pool().acquire().await.unwrap(); + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(&mut *connection) + .await + .unwrap(); + sqlx::query("UPDATE radroots_runtime_authored_draft_revisions SET snapshot = zeroblob(?) WHERE draft_id = ?") + .bind((AUTHORED_DRAFT_PAGE_SNAPSHOT_MAX_BYTES + 1) as i64).bind([4_u8; 16].as_slice()).execute(&mut *connection).await.unwrap(); + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(&mut *connection) + .await + .unwrap(); + // Exact SQL captured from this fresh fixture's governed migration, with no + // external input or interpolation. Restore its immutable catalog verbatim. + sqlx::query(sqlx::AssertSqlSafe(update_guard.as_str())) + .execute(&mut *connection) + .await + .unwrap(); + drop(connection); + let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 37).unwrap(); + let first = store.query_authored_drafts(q.clone()).await.unwrap(); + assert_eq!(first.records().len(), 37); + assert!(matches!( + first.records()[0], + AuthoredDraftQueryRecord::Corrupt { .. } + )); + assert!(first.records()[0].draft_id().is_err()); + let mut keys: Vec<_> = first + .records() + .iter() + .map(AuthoredDraftQueryRecord::draft_key) + .collect(); + let mut next = first.next_cursor().cloned(); + for id in [1_u128, 1001] { + let old = store + .authored_draft_head(AuthoredDraftId::new(id.to_be_bytes()).unwrap()) + .await + .unwrap() + .unwrap(); + let revised = old + .successor(vec![2], AuthoredDraftStage::Draft, None, 11) + .unwrap(); + store + .append_authored_draft(revised, Some(old.revision())) + .await + .unwrap(); + } + store.close().await.unwrap(); + let store = open_store(&temp).await; + while let Some(cursor) = next { + let cursor = serde_json::from_slice(&serde_json::to_vec(&cursor).unwrap()).unwrap(); + let page = store + .query_authored_drafts(q.clone().with_cursor(&cursor).unwrap()) + .await + .unwrap(); + assert!(page.records().len() <= 37); + for row in page.records() { + keys.push(row.draft_key()); + if row.draft_key() == 1001_u128.to_be_bytes() { + assert_eq!(row.revision().get(), 2); + } + if row.draft_key() == [3; 16] || row.draft_key() == [4; 16] { + assert!(matches!(row, AuthoredDraftQueryRecord::Corrupt { .. })); + } + } + next = page.next_cursor().cloned(); + } + let mut expected = vec![[0; 16]]; + expected.extend((1_u128..=1001).map(u128::to_be_bytes)); + expected.extend([[3; 16], [4; 16]]); + assert_eq!(keys, expected); + let fresh = store.query_authored_drafts(q).await.unwrap(); + assert_eq!(fresh.records()[1].revision().get(), 2); + let exact = store + .query_authored_drafts(AuthoredDraftQuery::for_author([7; 32], "known.v1", 37).unwrap()) + .await + .unwrap(); + assert!(exact.records().iter().all( + |r| matches!(r, AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "known.v1") + )); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn all_schema_sqlite_continuation_preserves_snapshot_and_payload_budgets() { + for byte in [0, 99] { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + for (id, schema, scope) in [ + (1, "known.v1", None), + ( + 2, + "future.v999", + Some(AuthoredDraftScope::new([3; 32]).unwrap()), + ), + ] { + store + .append_authored_draft(draft(id, schema, scope, vec![byte; 3 * 1024 * 1024]), None) + .await + .unwrap(); + } + let q = AuthoredDraftQuery::for_author_all_schemas([7; 32], 256).unwrap(); + let first = store.query_authored_drafts(q.clone()).await.unwrap(); + assert_eq!(first.records().len(), 1); + assert!( + matches!(&first.records()[0], AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "known.v1") + ); + let next = store + .query_authored_drafts(q.with_cursor(first.next_cursor().unwrap()).unwrap()) + .await + .unwrap(); + assert_eq!(next.records().len(), 1); + assert!( + matches!(&next.records()[0], AuthoredDraftQueryRecord::Draft(d) if d.payload_schema() == "future.v999") + ); + assert!(next.next_cursor().is_none()); + store.close().await.unwrap(); + } +} diff --git a/crates/storage_sqlite/src/authored_draft_query.rs b/crates/storage_sqlite/src/authored_draft_query.rs @@ -36,13 +36,13 @@ async fn read_page( revisions.payload_schema = '' AS unknown_schema FROM radroots_runtime_authored_draft_revisions AS revisions WHERE revisions.author = ? - AND ((revisions.payload_schema = ? AND (? OR revisions.payload_scope IS ?)) + AND (? OR (revisions.payload_schema = ? AND (? OR revisions.payload_scope IS ?)) OR revisions.payload_schema = '') AND (? IS NULL OR revisions.draft_id > ?) AND revisions.revision = (SELECT MAX(head.revision) FROM radroots_runtime_authored_draft_revisions AS head WHERE head.draft_id = revisions.draft_id) ORDER BY revisions.draft_id LIMIT ?" - ).bind(query.author().as_slice()).bind(query.payload_schema()) + ).bind(query.author().as_slice()).bind(query.payload_schema().is_none()).bind(query.payload_schema()) .bind(query.is_author_wide()) .bind(query.scope().map(|value| value.as_bytes().to_vec())) .bind(&after).bind(after).bind(i64::from(query.limit()) + 1) diff --git a/crates/storage_sqlite/src/authored_draft_query_tests.rs b/crates/storage_sqlite/src/authored_draft_query_tests.rs @@ -6,7 +6,7 @@ use radroots_storage::{ }; use tempfile::TempDir; -fn draft( +pub(super) fn draft( id: u8, schema: &str, scope: Option<AuthoredDraftScope>, @@ -156,7 +156,7 @@ async fn author_wide_sqlite_pages_preserve_snapshot_and_payload_budgets() { store.close().await.unwrap(); } } -async fn corrupt( +pub(super) async fn corrupt( store: &SqliteStorage, id: u8, author: u8,