lib

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

commit d3332ff1c68245a232f093434ee35fe2a342ed38
parent 363b9be18831be2fd6ef505d6154fb878ed8efd7
Author: triesap <tyson@radroots.org>
Date:   Thu, 10 Sep 2026 19:44:20 +0000

storage: Qualify durable authored writes

- Request full filesystem sync while preserving WAL and FULL
- Verify commit capacity lock and write-denial failures
- Prove acknowledged and uncommitted process-crash recovery
- Preserve APIs and pass current coverage and workspace gates

Diffstat:
Mcontracts/storage/connection_policy_v1.toml | 3+++
Mcontracts/storage/failure_injection_policy_v1.toml | 10++++++++++
Mcrates/storage_sqlite/README.md | 15+++++++++++++++
Mcrates/storage_sqlite/src/authored_draft.rs | 5+++++
Acrates/storage_sqlite/src/authored_durability_crash_tests.rs | 149+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/src/authored_durability_policy_tests.rs | 49+++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/storage_sqlite/src/authored_durability_tests.rs | 345+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/storage_sqlite/src/open.rs | 13+++++++++++++
8 files changed, 589 insertions(+), 0 deletions(-)

diff --git a/contracts/storage/connection_policy_v1.toml b/contracts/storage/connection_policy_v1.toml @@ -4,6 +4,9 @@ max_connections_per_database = 4 foreign_keys = true journal_mode = "wal" synchronous = "full" +# Request the VFS full flush where supported; SQLite may fall back to fsync. +# This flag alone does not establish physical power-loss durability. +fullfsync = true busy_timeout_min_ms = 1 busy_timeout_default_ms = 5000 busy_timeout_max_ms = 60000 diff --git a/contracts/storage/failure_injection_policy_v1.toml b/contracts/storage/failure_injection_policy_v1.toml @@ -56,3 +56,13 @@ lock_close_points = [ ] restore_recovery = "idempotent_forward_completion_before_connection_open" failure_reporting = "stable_typed_error_without_backend_details" + +[authored_write] +acknowledgment = "only_after_successful_commit_or_exact_committed_replay" +commit_fault = "deferred_foreign_key_at_actual_commit" +capacity_fault = "bounded_sqlite_max_page_count" +write_denied_fault = "owned_connection_query_only" +busy_fault = "owned_begin_immediate_with_bounded_timeout" +process_termination_points = ["after_acknowledgment", "during_uncommitted_insert"] +power_loss_qualified = false +protected_data_policy_owner = "native_host" diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md @@ -22,6 +22,21 @@ source-head CAS, including after reopen or later editing. Existing atomic receipt snapshots retain their 4 MiB limit; oversized submissions fail without partial writes. No additional public database or connection API is introduced. +Every connection in both owned pools uses `synchronous=FULL` and requests +`fullfsync=ON`; writable stores retain WAL. Authored append acknowledges only +after a successful SQLite COMMIT, or an exact replay of an already committed +revision. Connection-policy validation rejects a missing full-sync request. +The existing owner supplies this policy without a second database or journal. + +Qualification covers actual deferred-constraint COMMIT failure, bounded SQLite +page-capacity exhaustion, lock contention, denied writes, closed/read-only +stores, and child process termination before and after acknowledgment. Failed writes retain the +previous revision and return a typed failure without a success receipt. These +tests do not fill the host disk or establish physical power-loss survival. +SQLite requests `F_FULLFSYNC` where its VFS supports it and may fall back to +`fsync`; actual filesystem, device and power-cut behavior remains a separate +host qualification. Native protected-data policy remains the host's concern. + Backend status and the last integrity result are passive. Hosts invoke `check_integrity` explicitly with their own positive timestamp when they want full SQLite and foreign-key validation across both owned files. `close` drains diff --git a/crates/storage_sqlite/src/authored_draft.rs b/crates/storage_sqlite/src/authored_draft.rs @@ -637,6 +637,11 @@ mod tests { #[path = "authored_draft_query_tests.rs"] mod query_tests; +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_durability_tests.rs"] +mod durability_tests; + pub(crate) async fn load_head_tx( transaction: &mut sqlx::Transaction<'_, sqlx::Sqlite>, id: AuthoredDraftId, diff --git a/crates/storage_sqlite/src/authored_durability_crash_tests.rs b/crates/storage_sqlite/src/authored_durability_crash_tests.rs @@ -0,0 +1,149 @@ +use super::*; +use std::{ + fs, + process::{Child, Command, Stdio}, + time::{Duration, Instant}, +}; + +const CHILD_TEST: &str = "authored_draft::durability_tests::crash::authored_durability_child"; +const CHILD_ROOT: &str = "RADROOTS_AUTHORED_CRASH_FIXTURE_ROOT"; +const CHILD_PHASE: &str = "RADROOTS_AUTHORED_CRASH_FIXTURE_PHASE"; + +struct ChildGuard(Child); + +impl Drop for ChildGuard { + fn drop(&mut self) { + let _ = self.0.kill(); + let _ = self.0.wait(); + } +} + +#[tokio::test] +async fn authored_durability_child() { + let Some(directory) = std::env::var_os(CHILD_ROOT) else { + return; + }; + let directory = Path::new(&directory); + assert!(directory.is_absolute() && directory.is_dir()); + let phase = std::env::var(CHILD_PHASE).unwrap(); + assert!(matches!( + phase.as_str(), + "after_acknowledgment" | "during_uncommitted_insert" + )); + let store = open(directory, OpenMode::ReadWriteExisting).await; + let baseline = first(); + let pending = next(&baseline, b"complete child revision".to_vec()); + if phase == "after_acknowledgment" { + let receipt = store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await + .unwrap(); + assert_eq!(receipt.disposition(), DraftAppendDisposition::Inserted); + assert_head(&store, &pending).await; + fs::write(directory.join("child-ready"), phase.as_bytes()).unwrap(); + std::future::pending::<()>().await; + } else { + let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); + insert_draft_tx(&mut transaction, &pending).await.unwrap(); + assert_eq!( + load_head_tx(&mut transaction, baseline.draft_id()) + .await + .unwrap(), + Some(pending) + ); + fs::write(directory.join("child-ready"), phase.as_bytes()).unwrap(); + std::future::pending::<()>().await; + // Keep the actual uncommitted transaction alive until process termination. + transaction.rollback().await.unwrap(); + } +} + +async fn kill_and_reopen(phase: &str) { + let temp = TempDir::new().unwrap(); + let store = open(temp.path(), OpenMode::Create).await; + let baseline = first(); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + store.close().await.unwrap(); + let output = fs::File::create(temp.path().join("child-output")).unwrap(); + let mut child = ChildGuard( + Command::new(std::env::current_exe().unwrap()) + .args(["--exact", CHILD_TEST, "--nocapture", "--test-threads=1"]) + .env(CHILD_ROOT, temp.path()) + .env(CHILD_PHASE, phase) + .stdin(Stdio::null()) + .stdout(output.try_clone().unwrap()) + .stderr(output) + .spawn() + .unwrap(), + ); + let deadline = Instant::now() + Duration::from_secs(30); + loop { + assert!( + child.0.try_wait().unwrap().is_none(), + "child exited before ready: {}", + fs::read_to_string(temp.path().join("child-output")).unwrap() + ); + if fs::read(temp.path().join("child-ready")).ok().as_deref() == Some(phase.as_bytes()) { + break; + } + assert!( + Instant::now() < deadline, + "child did not reach exact write boundary" + ); + std::thread::sleep(Duration::from_millis(10)); + } + assert!(temp.path().join("runtime.sqlite-wal").is_file()); + child.0.kill().unwrap(); + assert!(!child.0.wait().unwrap().success()); + let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; + let pending = next(&baseline, b"complete child revision".to_vec()); + let expected = if phase == "after_acknowledgment" { + &pending + } else { + &baseline + }; + assert_head(&reopened, expected).await; + assert_eq!( + reopened + .authored_draft_revision(baseline.draft_id(), baseline.revision()) + .await + .unwrap(), + Some(baseline.clone()) + ); + if phase == "during_uncommitted_insert" { + assert!( + reopened + .authored_draft_revision(pending.draft_id(), pending.revision()) + .await + .unwrap() + .is_none() + ); + } + let receipt = reopened + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await + .unwrap(); + assert_eq!( + receipt.disposition(), + if phase == "after_acknowledgment" { + DraftAppendDisposition::Replay + } else { + DraftAppendDisposition::Inserted + } + ); + assert_head(&reopened, &pending).await; + reopened.close().await.unwrap(); +} + +#[tokio::test] +async fn authored_durability_acknowledged_revision_survives_process_kill() { + kill_and_reopen("after_acknowledgment").await; +} + +#[tokio::test] +async fn authored_durability_uncommitted_revision_is_absent_after_process_kill() { + kill_and_reopen("during_uncommitted_insert").await; +} diff --git a/crates/storage_sqlite/src/authored_durability_policy_tests.rs b/crates/storage_sqlite/src/authored_durability_policy_tests.rs @@ -0,0 +1,49 @@ +use super::*; +use radroots_storage::event::SourceGeneration; +use tempfile::TempDir; + +#[tokio::test] +async fn authored_durability_policy_holds_for_every_owned_pool_connection() { + let temp = TempDir::new().unwrap(); + let store = SqliteStorage::open( + OpenOptions::new( + Paths::from_directory(temp.path()).unwrap(), + OpenMode::Create, + ) + .with_source_generation(SourceGeneration::new([93; 32]).unwrap(), 9) + .unwrap(), + ) + .await + .unwrap(); + for (pool, database) in [ + (store.pool(), RUNTIME_DATABASE_NAME), + (store.private_pool(), PRIVATE_DATABASE_NAME), + ] { + let mut connections = Vec::new(); + for _ in 0..MAX_CONNECTIONS_PER_DATABASE { + connections.push(pool.acquire().await.unwrap()); + } + assert_eq!(connections.len(), 4); + for connection in &mut connections { + verify_connection(connection, database, Duration::from_millis(5_000)) + .await + .unwrap(); + sqlx::query("PRAGMA fullfsync = OFF") + .execute(&mut **connection) + .await + .unwrap(); + assert!(matches!( + verify_connection(connection, database, Duration::from_millis(5_000)).await, + Err(Error::ConnectionPolicyMismatch { database: actual }) if actual == database + )); + sqlx::query("PRAGMA fullfsync = ON") + .execute(&mut **connection) + .await + .unwrap(); + verify_connection(connection, database, Duration::from_millis(5_000)) + .await + .unwrap(); + } + } + store.close().await.unwrap(); +} diff --git a/crates/storage_sqlite/src/authored_durability_tests.rs b/crates/storage_sqlite/src/authored_durability_tests.rs @@ -0,0 +1,345 @@ +use super::*; +use crate::{OpenMode, OpenOptions, Paths}; +use radroots_storage::event::SourceGeneration; +use std::path::Path; +use tempfile::TempDir; + +#[path = "authored_durability_crash_tests.rs"] +mod crash; + +async fn open(directory: &Path, mode: OpenMode) -> SqliteStorage { + let mut options = OpenOptions::new(Paths::from_directory(directory).unwrap(), mode); + if matches!(mode, OpenMode::Create) { + options = options + .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9) + .unwrap(); + } + SqliteStorage::open(options).await.unwrap() +} + +fn first() -> AuthoredDraft { + AuthoredDraft::initial( + AuthoredDraftId::new([41; 16]).unwrap(), + [7; 32], + "radroots.durability-fixture.v1", + b"acknowledged baseline".to_vec(), + AuthoredDraftStage::Draft, + None, + 10, + ) + .unwrap() +} + +fn next(previous: &AuthoredDraft, payload: Vec<u8>) -> AuthoredDraft { + previous + .successor(payload, AuthoredDraftStage::Draft, None, 11) + .unwrap() +} + +async fn assert_head(store: &SqliteStorage, expected: &AuthoredDraft) { + assert_eq!( + store + .authored_draft_head(expected.draft_id()) + .await + .unwrap(), + Some(expected.clone()) + ); +} + +#[tokio::test] +async fn authored_durability_commit_fault_never_acknowledges_and_exact_retry_recovers() { + let temp = TempDir::new().unwrap(); + let store = open(temp.path(), OpenMode::Create).await; + let baseline = first(); + let pending = next(&baseline, b"pending complete revision".to_vec()); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + sqlx::query("CREATE TABLE authored_commit_parent (id INTEGER PRIMARY KEY)") + .execute(store.pool()) + .await + .unwrap(); + sqlx::query("CREATE TABLE authored_commit_fault (id INTEGER REFERENCES authored_commit_parent(id) DEFERRABLE INITIALLY DEFERRED)") + .execute(store.pool()).await.unwrap(); + sqlx::query("CREATE TRIGGER authored_commit_fault_trigger AFTER INSERT ON radroots_runtime_authored_draft_revisions BEGIN INSERT INTO authored_commit_fault VALUES (99); END") + .execute(store.pool()).await.unwrap(); + + // Establish that the insertion succeeds and the actual COMMIT is the fault. + let mut transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); + insert_draft_tx(&mut transaction, &pending).await.unwrap(); + let failure = transaction.commit().await.unwrap_err(); + assert_eq!( + failure.as_database_error().unwrap().kind(), + sqlx::error::ErrorKind::ForeignKeyViolation + ); + assert_eq!( + store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + assert_head(&store, &baseline).await; + assert!( + store + .authored_draft_revision(pending.draft_id(), pending.revision()) + .await + .unwrap() + .is_none() + ); + + for statement in [ + "DROP TRIGGER authored_commit_fault_trigger", + "DROP TABLE authored_commit_fault", + "DROP TABLE authored_commit_parent", + ] { + sqlx::query(statement).execute(store.pool()).await.unwrap(); + } + let receipt = store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await + .unwrap(); + assert_eq!(receipt.disposition(), DraftAppendDisposition::Inserted); + store.close().await.unwrap(); + let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; + assert_head(&reopened, &pending).await; + assert_eq!( + reopened + .append_authored_draft(pending, Some(baseline.revision())) + .await + .unwrap() + .disposition(), + DraftAppendDisposition::Replay + ); + reopened.close().await.unwrap(); +} + +#[tokio::test] +async fn authored_durability_sqlite_capacity_failure_preserves_the_acknowledged_head() { + let temp = TempDir::new().unwrap(); + let store = open(temp.path(), OpenMode::Create).await; + let baseline = first(); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + let mut connections = Vec::new(); + for _ in 0..4 { + connections.push(store.pool().acquire().await.unwrap()); + } + let mut original_limits = Vec::new(); + for connection in &mut connections { + original_limits.push( + sqlx::query_scalar::<_, i64>("PRAGMA max_page_count") + .fetch_one(&mut **connection) + .await + .unwrap(), + ); + let pages = sqlx::query_scalar::<_, i64>("PRAGMA page_count") + .fetch_one(&mut **connection) + .await + .unwrap(); + // PRAGMA assignments do not accept bind parameters; only an i64 is interpolated. + let limit = sqlx::query_scalar::<_, i64>(sqlx::AssertSqlSafe(format!( + "PRAGMA max_page_count = {}", + pages + 2 + ))) + .fetch_one(&mut **connection) + .await + .unwrap(); + assert_eq!(limit, pages + 2); + } + drop(connections); + // This bounded allocation proves SQLITE_FULL without filling the host disk. + let failure = + sqlx::query("CREATE TABLE authored_capacity_probe AS SELECT zeroblob(1048576) AS data") + .execute(store.pool()) + .await + .unwrap_err(); + assert_eq!( + failure.as_database_error().unwrap().code().as_deref(), + Some("13") + ); + let pending = next(&baseline, vec![42; 256 * 1024]); + assert_eq!( + store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + assert_head(&store, &baseline).await; + + let mut connections = Vec::new(); + for _ in 0..4 { + connections.push(store.pool().acquire().await.unwrap()); + } + for (connection, limit) in connections.iter_mut().zip(original_limits) { + sqlx::query(sqlx::AssertSqlSafe(format!( + "PRAGMA max_page_count = {limit}" + ))) + .execute(&mut **connection) + .await + .unwrap(); + } + drop(connections); + store.close().await.unwrap(); + let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; + assert_head(&reopened, &baseline).await; + assert_eq!( + reopened + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await + .unwrap() + .disposition(), + DraftAppendDisposition::Inserted + ); + assert_head(&reopened, &pending).await; + reopened.close().await.unwrap(); +} + +#[tokio::test] +async fn authored_durability_denied_writes_have_no_receipt_or_head_advance() { + let temp = TempDir::new().unwrap(); + let store = open(temp.path(), OpenMode::Create).await; + let baseline = first(); + let pending = next(&baseline, b"denied revision".to_vec()); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + let mut connections = Vec::new(); + for _ in 0..4 { + connections.push(store.pool().acquire().await.unwrap()); + } + for connection in &mut connections { + sqlx::query("PRAGMA query_only = ON") + .execute(&mut **connection) + .await + .unwrap(); + } + drop(connections); + assert_eq!( + store + .append_authored_draft(pending, Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + assert_head(&store, &baseline).await; + store.close().await.unwrap(); + let reopened = open(temp.path(), OpenMode::ReadWriteExisting).await; + assert_head(&reopened, &baseline).await; + reopened.close().await.unwrap(); +} + +#[tokio::test] +async fn authored_durability_busy_writer_never_returns_a_success_receipt() { + let temp = TempDir::new().unwrap(); + let store = SqliteStorage::open( + OpenOptions::new( + Paths::from_directory(temp.path()).unwrap(), + OpenMode::Create, + ) + .with_source_generation(SourceGeneration::new([91; 32]).unwrap(), 9) + .unwrap() + .with_busy_timeout(std::time::Duration::from_millis(10)) + .unwrap(), + ) + .await + .unwrap(); + let baseline = first(); + let pending = next(&baseline, b"blocked complete revision".to_vec()); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + let transaction = store.pool().begin_with("BEGIN IMMEDIATE").await.unwrap(); + assert_eq!( + store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + assert_head(&store, &baseline).await; + transaction.rollback().await.unwrap(); + assert_eq!( + store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await + .unwrap() + .disposition(), + DraftAppendDisposition::Inserted + ); + assert_head(&store, &pending).await; + store.close().await.unwrap(); +} + +#[tokio::test] +async fn authored_durability_read_only_and_closed_stores_cannot_acknowledge() { + let temp = TempDir::new().unwrap(); + let store = open(temp.path(), OpenMode::Create).await; + let baseline = first(); + let pending = next(&baseline, b"not writable".to_vec()); + store + .append_authored_draft(baseline.clone(), None) + .await + .unwrap(); + store.close().await.unwrap(); + assert_eq!( + store + .append_authored_draft(pending.clone(), Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + let read_only = open(temp.path(), OpenMode::ReadOnly).await; + assert_head(&read_only, &baseline).await; + assert_eq!( + read_only + .append_authored_draft(pending, Some(baseline.revision())) + .await, + Err(Error::BackendUnavailable) + ); + read_only.close().await.unwrap(); +} + +#[test] +fn authored_durability_contract_distinguishes_crash_and_power_loss() { + let policy: toml::Value = toml::from_str(include_str!( + "../../../contracts/storage/failure_injection_policy_v1.toml" + )) + .unwrap(); + let authored = &policy["authored_write"]; + assert_eq!( + authored["acknowledgment"].as_str(), + Some("only_after_successful_commit_or_exact_committed_replay") + ); + assert_eq!( + authored["commit_fault"].as_str(), + Some("deferred_foreign_key_at_actual_commit") + ); + assert_eq!( + authored["capacity_fault"].as_str(), + Some("bounded_sqlite_max_page_count") + ); + assert_eq!( + authored["write_denied_fault"].as_str(), + Some("owned_connection_query_only") + ); + assert_eq!( + authored["busy_fault"].as_str(), + Some("owned_begin_immediate_with_bounded_timeout") + ); + assert_eq!( + authored["process_termination_points"] + .as_array() + .unwrap() + .iter() + .map(|value| value.as_str().unwrap()) + .collect::<Vec<_>>(), + ["after_acknowledgment", "during_uncommitted_insert"] + ); + assert_eq!(authored["power_loss_qualified"].as_bool(), Some(false)); + assert_eq!( + authored["protected_data_policy_owner"].as_str(), + Some("native_host") + ); +} diff --git a/crates/storage_sqlite/src/open.rs b/crates/storage_sqlite/src/open.rs @@ -715,6 +715,7 @@ fn connect_options(path: &Path, mode: OpenMode, busy_timeout: Duration) -> Sqlit .foreign_keys(true) .busy_timeout(busy_timeout) .synchronous(SqliteSynchronous::Full) + .pragma("fullfsync", "ON") .disable_statement_logging(); if mode.is_writable() { options = options.journal_mode(SqliteJournalMode::Wal); @@ -880,10 +881,15 @@ async fn verify_connection( .map_err(|_| Error::ConnectionPolicyMismatch { database })?; let expected_busy_timeout = i64::try_from(busy_timeout.as_millis()) .map_err(|_| Error::ConnectionPolicyMismatch { database })?; + let fullfsync = sqlx::query_scalar::<_, i64>("PRAGMA fullfsync") + .fetch_one(&mut *connection) + .await + .map_err(|_| Error::ConnectionPolicyMismatch { database })?; if foreign_keys == 1 && journal_mode.eq_ignore_ascii_case("wal") && configured_busy_timeout == expected_busy_timeout && synchronous == 2 + && fullfsync == 1 { Ok(()) } else { @@ -1026,6 +1032,7 @@ mod policy_tests { foreign_keys: bool, journal_mode: String, synchronous: String, + fullfsync: bool, busy_timeout_min_ms: u64, busy_timeout_default_ms: u64, busy_timeout_max_ms: u64, @@ -1049,6 +1056,7 @@ mod policy_tests { assert!(policy.foreign_keys); assert_eq!(policy.journal_mode, "wal"); assert_eq!(policy.synchronous, "full"); + assert!(policy.fullfsync); assert_eq!(policy.busy_timeout_min_ms, 1); assert_eq!(policy.busy_timeout_default_ms, 5_000); assert_eq!(policy.busy_timeout_max_ms, 60_000); @@ -1060,3 +1068,8 @@ mod policy_tests { assert!(!policy.raw_handles_public); } } + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +#[path = "authored_durability_policy_tests.rs"] +mod authored_durability_policy_tests;