lib

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

commit 91006304a47ff335c6f36fbbcfcb68ecfc76f6c7
parent 4a2285ea896cfc361f539c74c467a3fd5591ec74
Author: triesap <tyson@radroots.org>
Date:   Sun, 20 Sep 2026 09:00:42 +0000

storage_sqlite: bound authored row reads before materialization

- Bound snapshot and redundant metadata in closed SQLite projections
- Preserve corruption, exact history, replay and indexed head reads
- Test oversized fields, nullable corruption and query plans
- Verify unchanged APIs, coverage gates and full producer checks

Diffstat:
Acontracts/architecture/decisions/authored_draft_bounded_rows.v1.json | 13+++++++++++++
Mcrates/storage_sqlite/src/authored_draft.rs | 73++++++++++++++++++++++++++++---------------------------------------------
Mcrates/storage_sqlite/src/authored_draft_query.rs | 15+++++++++------
Acrates/storage_sqlite/src/authored_draft_row.rs | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/src/authored_draft_row_tests.rs | 258+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
5 files changed, 366 insertions(+), 51 deletions(-)

diff --git a/contracts/architecture/decisions/authored_draft_bounded_rows.v1.json b/contracts/architecture/decisions/authored_draft_bounded_rows.v1.json @@ -0,0 +1,13 @@ +{ + "schema": "radroots.authored-draft-bounded-rows.v1", + "status": "implemented", + "owner": "radroots_storage_sqlite", + "authority": "Preserve authored_draft_inventory.v1/v2, existing exact historical reads, append replay and transactional source validation while bounding selected SQLite result columns before SQLx materializes them.", + "reads": ["authored_draft_head", "authored_draft_revision", "append replay and current head", "transactional source head", "selected paged inventory revision"], + "bounds": {"snapshot_bytes": 16777216, "payload_schema_bytes": 128, "draft_id_bytes": 16, "author_bytes": 32, "operation_id_bytes": 16, "payload_sha256_bytes": 32, "payload_scope_bytes": 32}, + "projection": "Closed SQL CASE expressions validate SQLite storage type and octet_length before returning variable-width columns. Scalar metadata must be integer. Invalid mandatory columns project NULL and fail decoding; invalid nullable operation/scope values project a noncanonical one-byte sentinel so corruption cannot become valid absence. A selected malformed record never disappears through a validity WHERE filter. Paged locator keys and schema classification are bounded before loading selected revisions.", + "sqlite": "The bundled SQLite supports octet_length, whose column-byte-length optimization avoids materializing the complete TEXT/BLOB merely to inspect its byte count. This bounds result materialization; it is not an arbitrary corrupt-file engine memory guarantee or a full integrity check.", + "compatibility": "No public API, migration, stored bytes, dependency, feature, cursor, page limit, snapshot limit or payload budget changes. Exact revision and descending head lookup retain the compound primary key. Legacy bulk-head aggregate behavior is unchanged.", + "verification": ["actual projected oversized snapshot and redundant columns", "invalid nullable metadata remains corruption", "inclusive snapshot boundary", "malformed inventory key fails closed", "point read, replay, successor and transactional rejection", "historical identity and genuine absence", "compound-key query plans", "existing inventory, append and transaction tests"], + "non_goals": ["application ownership interpretation", "deletion policy", "cross-process mutation fencing", "author enumeration", "automatic integrity repair", "legacy bulk-head aggregate redesign"] +} diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs @@ -14,6 +14,9 @@ const SNAPSHOT_MAX_BYTES: usize = 16 * 1024 * 1024; #[path = "authored_draft_query.rs"] mod query; +#[path = "authored_draft_row.rs"] +mod bounded_row; + impl AuthoredDraftStore for SqliteStorage { fn query_authored_drafts( &self, @@ -38,15 +41,12 @@ impl AuthoredDraftStore for SqliteStorage { .begin_with("BEGIN IMMEDIATE") .await .map_err(map_backend)?; - if let Some(row) = sqlx::query( - "SELECT * FROM radroots_runtime_authored_draft_revisions - WHERE draft_id = ? AND revision = ?", + if let Some(row) = bounded_row::load( + &mut *transaction, + draft.draft_id().as_bytes(), + Some(draft.revision()), ) - .bind(draft.draft_id().as_bytes().as_slice()) - .bind(i64_from_u64(draft.revision().get())?) - .fetch_optional(&mut *transaction) - .await - .map_err(map_backend)? + .await? { let existing = decode_row(&row)?; transaction.rollback().await.map_err(map_backend)?; @@ -60,17 +60,11 @@ impl AuthoredDraftStore for SqliteStorage { }; } - let head = sqlx::query( - "SELECT * FROM radroots_runtime_authored_draft_revisions - WHERE draft_id = ? ORDER BY revision DESC LIMIT 1", - ) - .bind(draft.draft_id().as_bytes().as_slice()) - .fetch_optional(&mut *transaction) - .await - .map_err(map_backend)? - .as_ref() - .map(decode_row) - .transpose()?; + let head = bounded_row::load(&mut *transaction, draft.draft_id().as_bytes(), None) + .await? + .as_ref() + .map(decode_row) + .transpose()?; match (head.as_ref(), expected_head) { (None, None) if draft.revision() == AuthoredDraftRevision::INITIAL => {} (Some(previous), Some(expected)) if previous.revision() == expected => { @@ -96,17 +90,11 @@ impl AuthoredDraftStore for SqliteStorage { draft_id: AuthoredDraftId, ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { Box::pin(async move { - sqlx::query( - "SELECT * FROM radroots_runtime_authored_draft_revisions - WHERE draft_id = ? ORDER BY revision DESC LIMIT 1", - ) - .bind(draft_id.as_bytes().as_slice()) - .fetch_optional(self.pool()) - .await - .map_err(map_backend)? - .as_ref() - .map(decode_row) - .transpose() + bounded_row::load(self.pool(), draft_id.as_bytes(), None) + .await? + .as_ref() + .map(decode_row) + .transpose() }) } @@ -116,18 +104,11 @@ impl AuthoredDraftStore for SqliteStorage { revision: AuthoredDraftRevision, ) -> BoxFuture<'_, Result<Option<AuthoredDraft>, Error>> { Box::pin(async move { - sqlx::query( - "SELECT * FROM radroots_runtime_authored_draft_revisions - WHERE draft_id = ? AND revision = ?", - ) - .bind(draft_id.as_bytes().as_slice()) - .bind(i64_from_u64(revision.get())?) - .fetch_optional(self.pool()) - .await - .map_err(map_backend)? - .as_ref() - .map(decode_row) - .transpose() + bounded_row::load(self.pool(), draft_id.as_bytes(), Some(revision)) + .await? + .as_ref() + .map(decode_row) + .transpose() }) } @@ -651,9 +632,11 @@ pub(crate) async fn load_head_tx( transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, id: AuthoredDraftId, ) -> Result<Option<AuthoredDraft>, Error> { - sqlx::query("SELECT * FROM radroots_runtime_authored_draft_revisions WHERE draft_id = ? ORDER BY revision DESC LIMIT 1") - .bind(id.as_bytes().as_slice()).fetch_optional(&mut **transaction).await.map_err(map_backend)? - .as_ref().map(decode_row).transpose() + bounded_row::load(&mut **transaction, id.as_bytes(), None) + .await? + .as_ref() + .map(decode_row) + .transpose() } pub(crate) async fn insert_draft_tx( transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, diff --git a/crates/storage_sqlite/src/authored_draft_query.rs b/crates/storage_sqlite/src/authored_draft_query.rs @@ -1,4 +1,4 @@ -use super::{SqliteStorage, decode_row, map_backend}; +use super::{SqliteStorage, bounded_row, decode_row, map_backend}; use radroots_storage::{ Error, authored_draft::AuthoredDraftRevision, @@ -32,8 +32,11 @@ async fn read_page( ) -> Result<AuthoredDraftPage, Error> { let after = query.after().map(|value| value.to_vec()); let rows = sqlx::query( - "SELECT revisions.draft_id, revisions.revision, length(revisions.snapshot) AS snapshot_bytes, - revisions.payload_schema = '' AS unknown_schema + "SELECT CASE WHEN typeof(revisions.draft_id) = 'blob' AND octet_length(revisions.draft_id) = 16 THEN revisions.draft_id END AS draft_id, + CASE WHEN typeof(revisions.revision) = 'integer' THEN revisions.revision END AS revision, + octet_length(revisions.snapshot) AS snapshot_bytes, + CASE WHEN typeof(revisions.payload_schema) = 'text' AND octet_length(revisions.payload_schema) <= 128 + THEN revisions.payload_schema = '' ELSE 1 END AS unknown_schema FROM radroots_runtime_authored_draft_revisions AS revisions WHERE revisions.author = ? AND (? OR (revisions.payload_schema = ? AND (? OR revisions.payload_scope IS ?)) @@ -83,9 +86,9 @@ async fn read_page( break; } snapshot_bytes += size; - let row = sqlx::query("SELECT * FROM radroots_runtime_authored_draft_revisions WHERE draft_id = ? AND revision = ?") - .bind(key.as_slice()).bind(revision.get() as i64) - .fetch_one(&mut **transaction).await.map_err(map_backend)?; + let row = bounded_row::load(&mut **transaction, &key, Some(revision)) + .await? + .ok_or(Error::CorruptAuthoredDraft)?; match decode_row(&row) { Ok(draft) if query.matches(&draft) => { if payload_bytes + draft.payload().len() > AUTHORED_DRAFT_PAGE_PAYLOAD_MAX_BYTES { diff --git a/crates/storage_sqlite/src/authored_draft_row.rs b/crates/storage_sqlite/src/authored_draft_row.rs @@ -0,0 +1,58 @@ +//! Bound result columns before SQLx materializes a stored authored revision. +//! Invalid nullable metadata uses a noncanonical one-byte sentinel, never NULL: +//! corruption must not turn into an accepted missing operation or scope. + +use radroots_storage::{Error, authored_draft::AuthoredDraftRevision}; +use sqlx::{Executor, Sqlite, sqlite::SqliteRow}; + +macro_rules! select { + ($suffix:literal) => { + concat!( + "SELECT + CASE WHEN typeof(draft_id) = 'blob' AND octet_length(draft_id) = 16 THEN draft_id END AS draft_id, + CASE WHEN typeof(revision) = 'integer' THEN revision END AS revision, + CASE WHEN typeof(author) = 'blob' AND octet_length(author) = 32 THEN author END AS author, + CASE WHEN typeof(stage) = 'integer' THEN stage END AS stage, + CASE WHEN operation_id IS NULL OR (typeof(operation_id) = 'blob' AND octet_length(operation_id) = 16) + THEN operation_id ELSE X'00' END AS operation_id, + CASE WHEN typeof(payload_sha256) = 'blob' AND octet_length(payload_sha256) = 32 THEN payload_sha256 END AS payload_sha256, + CASE WHEN typeof(created_at_unix_ms) = 'integer' THEN created_at_unix_ms END AS created_at_unix_ms, + CASE WHEN typeof(updated_at_unix_ms) = 'integer' THEN updated_at_unix_ms END AS updated_at_unix_ms, + CASE WHEN typeof(snapshot) = 'blob' AND octet_length(snapshot) BETWEEN 1 AND 16777216 THEN snapshot END AS snapshot, + CASE WHEN typeof(payload_schema) = 'text' AND octet_length(payload_schema) <= 128 THEN payload_schema END AS payload_schema, + CASE WHEN payload_scope IS NULL OR (typeof(payload_scope) = 'blob' AND octet_length(payload_scope) = 32) + THEN payload_scope ELSE X'00' END AS payload_scope + FROM radroots_runtime_authored_draft_revisions ", + $suffix + ) + }; +} + +const HEAD: &str = select!( + "WHERE draft_id = ? AND ? IS NULL ORDER BY radroots_runtime_authored_draft_revisions.revision DESC LIMIT 1" +); +const REVISION: &str = select!("WHERE draft_id = ? AND revision = ?"); + +pub(super) async fn load<'e>( + executor: impl Executor<'e, Database = Sqlite>, + key: &[u8; 16], + revision: Option<AuthoredDraftRevision>, +) -> Result<Option<SqliteRow>, Error> { + let selected = if revision.is_some() { REVISION } else { HEAD }; + let revision = revision + .map(|value| super::i64_from_u64(value.get())) + .transpose()?; + // Both alternatives are closed compile-time SQL above. Every data value is + // bound separately; this private helper accepts no caller-supplied SQL. + sqlx::query(sqlx::AssertSqlSafe(selected)) + .bind(key.as_slice()) + .bind(revision) + .fetch_optional(executor) + .await + .map_err(super::map_backend) +} + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_draft_row_tests.rs"] +mod tests; diff --git a/crates/storage_sqlite/src/authored_draft_row_tests.rs b/crates/storage_sqlite/src/authored_draft_row_tests.rs @@ -0,0 +1,258 @@ +use super::*; +use crate::authored_draft::{decode_row, query_tests::draft, tests::open_store}; +use radroots_storage::{ + authored_draft::{AuthoredDraft, AuthoredDraftStore}, + authored_draft_query::{AuthoredDraftQuery, AuthoredDraftQueryRecord}, +}; +use sqlx::{Row, ValueRef}; +use tempfile::TempDir; + +// Fixture-only corruption, on one connection, with the exact migration guard +// restored before the owning reads run. No production schema is changed. +async fn corrupt_field(store: &crate::SqliteStorage, assignment: &'static str) { + let guard: String = sqlx::query_scalar( + "SELECT sql FROM sqlite_schema WHERE name = 'radroots_runtime_authored_draft_revisions_update_guard'", + ).fetch_one(store.pool()).await.unwrap(); + let mut connection = store.pool().acquire().await.unwrap(); + sqlx::query("DROP TRIGGER radroots_runtime_authored_draft_revisions_update_guard") + .execute(&mut *connection) + .await + .unwrap(); + sqlx::query("PRAGMA ignore_check_constraints = ON") + .execute(&mut *connection) + .await + .unwrap(); + // Only closed test literals below supply this assignment. + let statement = format!("UPDATE radroots_runtime_authored_draft_revisions SET {assignment}"); + sqlx::query(sqlx::AssertSqlSafe(statement.as_str())) + .execute(&mut *connection) + .await + .unwrap(); + sqlx::query("PRAGMA ignore_check_constraints = OFF") + .execute(&mut *connection) + .await + .unwrap(); + sqlx::query(sqlx::AssertSqlSafe(guard.as_str())) + .execute(&mut *connection) + .await + .unwrap(); +} + +async fn rejects_point_reads_and_replay(store: &crate::SqliteStorage, value: &AuthoredDraft) { + assert_eq!( + store.authored_draft_head(value.draft_id()).await, + Err(Error::CorruptAuthoredDraft) + ); + assert_eq!( + store + .authored_draft_revision(value.draft_id(), value.revision()) + .await, + Err(Error::CorruptAuthoredDraft) + ); + assert_eq!( + store.append_authored_draft(value.clone(), None).await, + Err(Error::CorruptAuthoredDraft) + ); + let successor = value.successor(vec![2], value.stage(), None, 11).unwrap(); + assert_eq!( + store + .append_authored_draft(successor, Some(value.revision())) + .await, + Err(Error::CorruptAuthoredDraft) + ); + let mut transaction = store.pool().begin().await.unwrap(); + assert_eq!( + crate::authored_draft::load_head_tx(&mut transaction, value.draft_id()).await, + Err(Error::CorruptAuthoredDraft) + ); + transaction.rollback().await.unwrap(); +} + +#[tokio::test] +async fn oversized_columns_are_bounded_in_sql_and_remain_corrupt() { + for (assignment, column, sentinel) in [ + ("snapshot = zeroblob(16777217)", "snapshot", false), + ("author = zeroblob(1048576)", "author", false), + ( + "payload_sha256 = zeroblob(1048576)", + "payload_sha256", + false, + ), + ("operation_id = zeroblob(1048576)", "operation_id", true), + ("payload_scope = zeroblob(1048576)", "payload_scope", true), + ( + "payload_schema = replace(hex(zeroblob(65536)), '0', 'é')", + "payload_schema", + false, + ), + ] { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let value = draft(1, "known.v1", None, vec![1]); + store + .append_authored_draft(value.clone(), None) + .await + .unwrap(); + corrupt_field(&store, assignment).await; + let row = load( + store.pool(), + value.draft_id().as_bytes(), + Some(value.revision()), + ) + .await + .unwrap() + .unwrap(); + if sentinel { + assert_eq!( + row.try_get::<Option<Vec<u8>>, _>(column).unwrap(), + Some(vec![0]) + ); + } else { + assert!(row.try_get_raw(column).unwrap().is_null(), "{column}"); + } + assert_eq!(decode_row(&row), Err(Error::CorruptAuthoredDraft)); + rejects_point_reads_and_replay(&store, &value).await; + // Foreign/corrupt authors remain outside this independently supplied + // author selection; every other malformed head remains a locator. + if column != "author" { + let page = store + .query_authored_drafts( + AuthoredDraftQuery::for_author_all_schemas([7; 32], 1).unwrap(), + ) + .await + .unwrap(); + assert!(matches!( + page.records(), + [AuthoredDraftQueryRecord::Corrupt { .. }] + )); + } + store.close().await.unwrap(); + } +} + +#[tokio::test] +async fn snapshot_limit_is_inclusive_and_empty_snapshot_is_rejected_before_decode() { + for (assignment, expected) in [ + ("snapshot = zeroblob(16777216)", Some(16777216)), + ("snapshot = X''", None), + ] { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let value = draft(1, "known.v1", None, vec![1]); + store + .append_authored_draft(value.clone(), None) + .await + .unwrap(); + corrupt_field(&store, assignment).await; + let row = load(store.pool(), value.draft_id().as_bytes(), None) + .await + .unwrap() + .unwrap(); + assert_eq!( + row.try_get::<Option<Vec<u8>>, _>("snapshot") + .unwrap() + .map(|bytes| bytes.len()), + expected + ); + assert_eq!(decode_row(&row), Err(Error::CorruptAuthoredDraft)); + store.close().await.unwrap(); + } +} + +#[tokio::test] +async fn malformed_key_fails_inventory_instead_of_becoming_absence() { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + store + .append_authored_draft(draft(1, "known.v1", None, vec![1]), None) + .await + .unwrap(); + corrupt_field(&store, "draft_id = zeroblob(1048576)").await; + assert!( + store + .query_authored_drafts(AuthoredDraftQuery::for_author_all_schemas([7; 32], 1).unwrap()) + .await + .is_err() + ); + store.close().await.unwrap(); +} + +#[tokio::test] +async fn bounded_load_preserves_history_absence_and_compound_index() { + let temp = TempDir::new().unwrap(); + let store = open_store(&temp).await; + let first = draft(1, "known.v1", None, vec![1]); + let second = first.successor(vec![2], first.stage(), None, 11).unwrap(); + store + .append_authored_draft(first.clone(), None) + .await + .unwrap(); + store + .append_authored_draft(second.clone(), Some(first.revision())) + .await + .unwrap(); + for (revision, expected) in [(None, &second), (Some(first.revision()), &first)] { + let row = load(store.pool(), first.draft_id().as_bytes(), revision) + .await + .unwrap() + .unwrap(); + assert_eq!(&decode_row(&row).unwrap(), expected); + assert!( + row.try_get::<Option<Vec<u8>>, _>("operation_id") + .unwrap() + .is_none() + ); + assert!( + row.try_get::<Option<Vec<u8>>, _>("payload_scope") + .unwrap() + .is_none() + ); + } + assert!(load(store.pool(), &[2; 16], None).await.unwrap().is_none()); + assert!( + load( + store.pool(), + first.draft_id().as_bytes(), + Some(AuthoredDraftRevision::new(3).unwrap()) + ) + .await + .unwrap() + .is_none() + ); + assert!(matches!( + load( + store.pool(), + &[1; 16], + Some(AuthoredDraftRevision::new(u64::MAX).unwrap()) + ) + .await, + Err(Error::InvalidAuthoredDraft) + )); + for statement in [HEAD, REVISION] { + let explain = format!("EXPLAIN QUERY PLAN {statement}"); + let rows = sqlx::query(sqlx::AssertSqlSafe(explain.as_str())) + .bind([1_u8; 16].as_slice()) + .bind(1_i64) + .fetch_all(store.pool()) + .await + .unwrap(); + let details: Vec<String> = rows.iter().map(|row| row.get("detail")).collect(); + assert!( + details.iter().any(|detail| detail.contains("SEARCH") + && detail.contains("PRIMARY KEY") + && detail.contains("draft_id=?")), + "{details:?}" + ); + assert!( + !details + .iter() + .any(|detail| detail.contains("SCAN ") || detail.contains("TEMP B-TREE")), + "{details:?}" + ); + } + store.close().await.unwrap(); + assert!(matches!( + load(store.pool(), &[1; 16], None).await, + Err(Error::BackendUnavailable) + )); +}