lib

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

commit 8f9cd897c4d3c5395f2e569372950ad82694264d
parent adf4e3b820333242616351ecae46930b1cfa3470
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 04:01:08 +0000

outbox: seal authenticated SQLite lifecycle

- own bounded lazy pools under one five-second open deadline
- authenticate bounded schema, ledger, foreign-key, and integrity state before WAL
- move exact destructive rollback to a closed path with explicit acknowledgement
- bind lifecycle features, release notes, generated authority, and hostile tests

Diffstat:
MCHANGELOG.md | 12+++++++++++-
Mcontracts/outbox_feature_matrix.toml | 2+-
Mcontracts/releases/1.0.0-alpha.1.toml | 13+++++++++++++
Mcrates/outbox/Cargo.toml | 5+++--
Mcrates/outbox/README | 29++++++++++++++++++++---------
Mcrates/outbox/contracts/migration_authority_v1.manifest.json | 62++++++++++++++++++++++++++++++++++++--------------------------
Mcrates/outbox/contracts/migration_authority_v1.manifest.schema.json | 6+++---
Mcrates/outbox/contracts/migration_authority_v1.manifest.sha256 | 2+-
Mcrates/outbox/src/error.rs | 101+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Mcrates/outbox/src/lib.rs | 8++++++++
Mcrates/outbox/src/migrations.rs | 1-
Mcrates/outbox/src/schema.rs | 356+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------
Acrates/outbox/src/sqlite_lifecycle.rs | 831+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/store.rs | 276++++++++++---------------------------------------------------------------------
Mtools/xtask/src/contract/outbox_migration.rs | 10+++++++---
15 files changed, 1391 insertions(+), 323 deletions(-)

diff --git a/CHANGELOG.md b/CHANGELOG.md @@ -9,6 +9,16 @@ publish policy both pass for the same source revision. ### Changed +<!-- release-change: outbox-authenticated-sqlite-lifecycle --> +- Outbox SQLite open now owns a lazy pool capped at four file connections and + authenticates UTF-8, bounded catalog and ledger text, the managed schema, + owned foreign keys, and scoped integrity under one five-second deadline + before exposing the store. Persistent WAL configuration follows authority; + caller pools, raw pool access, and live-store migration were removed. + Destructive rollback is now an explicit data-loss-acknowledged offline-path + operation that rejects live owned handles, while full-database integrity is + a separate maintenance call and public diagnostics are capped at 4,096 + bytes. <!-- release-change: outbox-versioned-migration-authority --> - Outbox schema initialization now uses an ordered, generated migration registry with immutable source checksums, authenticated catalog fingerprints, @@ -17,7 +27,7 @@ publish policy both pass for the same source revision. identity; partial, changed, counterfeit, gapped, newer, or unknown outbox state fails before governed mutation while unrelated caller tables remain untouched. Raw migration SQL exports and live `migrate_down` were replaced by - authenticated schema status and migrate-to-current APIs. A machine-readable + authenticated schema status and owned open-time migration. A machine-readable matrix now governs no-default, SQLite, Tokio, event-store-adapter, and all-feature builds. - Event-store schema initialization now uses a transactional, checksummed diff --git a/contracts/outbox_feature_matrix.toml b/contracts/outbox_feature_matrix.toml @@ -4,7 +4,7 @@ package = "radroots_outbox" [feature_edges] default = ["event-store-adapter"] sqlite = ["dep:sqlx", "sqlx/sqlite-bundled"] -runtime-tokio = ["sqlite", "sqlx/runtime-tokio"] +runtime-tokio = ["sqlite", "sqlx/runtime-tokio", "dep:tokio"] event-store-adapter = [ "runtime-tokio", "dep:radroots_event_store", diff --git a/contracts/releases/1.0.0-alpha.1.toml b/contracts/releases/1.0.0-alpha.1.toml @@ -534,6 +534,19 @@ semver_impacts = [ summary = "Add transport-neutral BUD-02 plus BUD-01 publication-readiness evidence with bounded complete-byte verification, canonical persisted reload, exact sequential JPEG entropy accounting plus pinned strict zune-jpeg decoding, and internally authoritative full decoding for static JPEG, PNG, and WebP rasters." [[changes]] +id = "outbox-authenticated-sqlite-lifecycle" +classification = "breaking" +semver_impacts = [ + "add_exported_type", + "add_exported_function", + "add_exported_constant", + "add_enum_variant", + "remove_exported_function", + "change_exported_algorithm_behavior", +] +summary = "Replace caller-owned outbox pools and live schema mutation with a bounded lazy owned SQLite lifecycle: authenticate UTF-8, bounded ledger/catalog authority, owned foreign keys, and scoped integrity under one five-second deadline before WAL; cap diagnostics; expose explicit full integrity; and require a closed path plus typed data-loss acknowledgement for exclusive exact-target rollback." + +[[changes]] id = "outbox-versioned-migration-authority" classification = "breaking" semver_impacts = [ diff --git a/crates/outbox/Cargo.toml b/crates/outbox/Cargo.toml @@ -14,7 +14,7 @@ readme = "README" [features] default = ["event-store-adapter"] sqlite = ["dep:sqlx", "sqlx/sqlite-bundled"] -runtime-tokio = ["sqlite", "sqlx/runtime-tokio"] +runtime-tokio = ["sqlite", "sqlx/runtime-tokio", "dep:tokio"] event-store-adapter = [ "runtime-tokio", "dep:radroots_event_store", @@ -35,6 +35,7 @@ serde_json = { workspace = true, features = ["std"] } sha2 = { workspace = true } sqlx = { workspace = true, optional = true, features = ["derive"] } thiserror = { workspace = true } +tokio = { workspace = true, optional = true, features = ["time"] } [dev-dependencies] radroots_nostr = { workspace = true, default-features = false, features = [ @@ -42,7 +43,7 @@ radroots_nostr = { workspace = true, default-features = false, features = [ "events", ] } tempfile = { workspace = true } -tokio = { workspace = true, features = ["macros", "rt"] } +tokio = { workspace = true, features = ["macros", "rt", "time"] } [lints.rust] unexpected_cfgs = { level = "warn", check-cfg = ['cfg(coverage_nightly)'] } diff --git a/crates/outbox/README b/crates/outbox/README @@ -9,10 +9,14 @@ HTTP-auth signatures, must be constructed and sent inside their owning transport exchange. This keeps event-store insertion and duplicate outcomes unambiguous for every queued event. -`open_pool` validates the declared SQLite backing mode, rejects -multi-connection in-memory pools in every supported URL form, and configures -every file-pool connection with foreign-key enforcement and the required busy -timeout before migrations or writes. +`open_memory` and `open_file` own their SQLite pool and configuration. File +stores use a lazy pool capped at four connections. One five-second total +deadline covers connection acquisition, schema authentication, migration, +owned integrity, and WAL configuration; lock contention retries consume that +same deadline. UTF-8, connection-local foreign-key enforcement, the complete +ledger/catalog identity, owned-table foreign keys, and table-scoped integrity +are verified before the store is exposed. WAL is configured only after the +existing database is proven fresh, exactly adoptable, or managed. Schema identity is governed by an ordered, generated migration registry. The original `0001_outbox` up and down files are immutable contract inputs. An @@ -23,8 +27,15 @@ fingerprint; unknown objects, counterfeit ledgers, history gaps, and newer schema versions fail before governed mutation. Unrelated caller tables in a shared SQLite database are outside the outbox fingerprint and remain untouched. -`schema_status` authenticates the current catalog and history, while -`migrate_to_current_schema` serializes exact adoption or append-only migration. -Migration rollback descriptors remain executable in the migration-authority -test suite; the separately governed store-lifecycle checkpoint owns any future -offline production rollback surface. +Catalog identifiers, catalog SQL, migration-ledger text, file paths, and +integrity results are bounded before they enter public diagnostics. +`RadrootsOutboxError::public_diagnostic` enforces the 4,096-byte diagnostic +ceiling and redacts underlying SQLx detail from its display value. + +`schema_status` authenticates the current catalog and history. Migration and +exact adoption occur only during owned open. Production does not accept a +caller pool or expose its raw pool. Destructive rollback accepts only a file +path with no live owned handles, requires explicit data-loss acknowledgement, +uses an exclusive exact-target transaction, and leaves the database closed. +Full-database `integrity_check` is available only as an explicit maintenance +operation; normal open performs scoped owned-table checks. diff --git a/crates/outbox/contracts/migration_authority_v1.manifest.json b/crates/outbox/contracts/migration_authority_v1.manifest.json @@ -21,7 +21,8 @@ { "enables": [ "sqlite", - "sqlx/runtime-tokio" + "sqlx/runtime-tokio", + "dep:tokio" ], "feature": "runtime-tokio" }, @@ -78,10 +79,10 @@ } ], "source": { - "byte_length": 889, + "byte_length": 902, "hash_algorithm": "sha256_bytes_v1", "path": "contracts/outbox_feature_matrix.toml", - "sha256": "c18b8f1b3c17e732def304dd8273a006bef061ee8cf40223e6d65ddfa5ab8084" + "sha256": "846a9c0a41f538e5f55f352602ca898f4508e0b6d657ddb11b77eb0f76d19e55" } }, "generated_runtime": { @@ -96,13 +97,13 @@ "migration_transaction": "begin_immediate_v1", "name": "radroots_outbox_schema_migrations", "reserved_prefix": "outbox_", - "rollback_transaction": "begin_exclusive_test_executor_v1" + "rollback_transaction": "begin_exclusive_offline_path_v1" }, "manifest_schema": { - "byte_length": 7461, + "byte_length": 7460, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/contracts/migration_authority_v1.manifest.schema.json", - "sha256": "65e6a2e5309d7f31b05b5c80d6e1b548aea82f32fa95488a79a0660f73b866ac" + "sha256": "df7958b384b4105617dbca3176741344cb455567fe92d3ac567266ee46ae79c5" }, "migrations": [ { @@ -177,28 +178,28 @@ "source_files": [ { "file": { - "byte_length": 1444, + "byte_length": 1532, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/Cargo.toml", - "sha256": "77f3c8e51aa2f3676c679803b4d7388dba67987a146169ba05c8d715b713d67d" + "sha256": "144c3da533519c9e5cf830dfd9ca7e7c8309cd8b581245141cf869f6ae41a47b" }, "role": "outbox_package_manifest" }, { "file": { - "byte_length": 1377, + "byte_length": 1684, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/lib.rs", - "sha256": "ad57e5979036ba29daad640acf898af7544aa06320112a063ffbc642ad886130" + "sha256": "8a5ac45f4b88eee78078773d6400ae658f5180771a8f51c3ba1c6c6a95d2fe76" }, "role": "outbox_public_surface" }, { "file": { - "byte_length": 8502, + "byte_length": 11111, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/error.rs", - "sha256": "7d10bfbf43dfbccf2433a3e67faac1cbab1fae530b6acb5d5ec06e161d6efe5b" + "sha256": "6f7a918670700696c6e853e45c2b722a1f68551334ca709e77c148f4617809a9" }, "role": "outbox_error_surface" }, @@ -213,28 +214,37 @@ }, { "file": { - "byte_length": 22852, + "byte_length": 22839, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/migrations.rs", - "sha256": "70614a806a7443e7077ba585f638483f44137a6ba3a1ade8e5571e10dc2c9da5" + "sha256": "f207cab712cbc4bdf8b9fe1892e553ced3512adab6b7b6883c8cbe79c7205255" }, "role": "migration_registry_runtime" }, { "file": { - "byte_length": 51902, + "byte_length": 63715, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/schema.rs", - "sha256": "f2cb9a8f0f9416aa02b18eeb966d5dae54373c92ae621173e299a98c79b57191" + "sha256": "135341d8d47171bf587bd17bb5ec5085c00648290ba84cdfc9feb348e6e214dd" }, "role": "schema_runtime" }, { "file": { - "byte_length": 326431, + "byte_length": 29959, + "hash_algorithm": "sha256_bytes_v1", + "path": "crates/outbox/src/sqlite_lifecycle.rs", + "sha256": "8fdad57a1a22a905d0fc449962e5f46cc266ea6e9b2b8010b7f26392f56d5369" + }, + "role": "sqlite_lifecycle_runtime" + }, + { + "file": { + "byte_length": 318687, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/src/store.rs", - "sha256": "47824d14cea2d00f0310f0676f3cd0b5825d642bea668b99c8194c53b0d76107" + "sha256": "6aa76ba0e86554a92dd3e21b6b74853bd5c061edef69c5024e740dcbf777d193" }, "role": "store_integration" }, @@ -249,10 +259,10 @@ }, { "file": { - "byte_length": 46017, + "byte_length": 46145, "hash_algorithm": "sha256_bytes_v1", "path": "tools/xtask/src/contract/outbox_migration.rs", - "sha256": "f41a2ac2486210e9cd6733149378d89cb32278730e8c7188207e63bc2d949faa" + "sha256": "2ee39f355148e42a80ed73f7d292a7f03c42b12e9bdd03d849c01dda17254784" }, "role": "contract_governance" }, @@ -294,28 +304,28 @@ }, { "file": { - "byte_length": 1663, + "byte_length": 2465, "hash_algorithm": "sha256_bytes_v1", "path": "crates/outbox/README", - "sha256": "4bcb2c1860e3c6b75436a6dce722e91ddcb4de70c624f3141cb8167f79dca822" + "sha256": "5f7de28d9b2cd2a68b03e21b1849994a6881b409369e4695a2e1f919d4f09ae8" }, "role": "outbox_readme" }, { "file": { - "byte_length": 24634, + "byte_length": 25312, "hash_algorithm": "sha256_bytes_v1", "path": "contracts/releases/1.0.0-alpha.1.toml", - "sha256": "755db4ea8775c1ad4c7faa223ebb3c06b127d00ccff3b51034fc6bb42d13d272" + "sha256": "1ae79c394b13c64d785436e98a0743789eba7786adb2ec4edff7df1344c4f2c6" }, "role": "release_record" }, { "file": { - "byte_length": 35643, + "byte_length": 36323, "hash_algorithm": "sha256_bytes_v1", "path": "CHANGELOG.md", - "sha256": "0581d0f0e95a8a4f81b78b2810a231a9061bb4057ca699201355b3ae3257ff4f" + "sha256": "3ee8e0d37238113448ddcf44ecd7f60e38d74e87a45ba712a1a3730d157510a3" }, "role": "release_notes" } diff --git a/crates/outbox/contracts/migration_authority_v1.manifest.schema.json b/crates/outbox/contracts/migration_authority_v1.manifest.schema.json @@ -209,7 +209,7 @@ "const": "outbox_" }, "rollback_transaction": { - "const": "begin_exclusive_test_executor_v1" + "const": "begin_exclusive_offline_path_v1" } }, "required": [ @@ -291,8 +291,8 @@ "items": { "$ref": "#/$defs/source_file" }, - "maxItems": 16, - "minItems": 16, + "maxItems": 17, + "minItems": 17, "type": "array" }, "version_bounds": { diff --git a/crates/outbox/contracts/migration_authority_v1.manifest.sha256 b/crates/outbox/contracts/migration_authority_v1.manifest.sha256 @@ -1 +1 @@ -362f12c15cc5772cf4c31d70acd9911d767b3af734e3a99b721f3aff779a40ce +cbc5f03a84e21694878a6e3cd3173cbb0da222f8258a5f25a565fdebf00d3c49 diff --git a/crates/outbox/src/error.rs b/crates/outbox/src/error.rs @@ -7,8 +7,12 @@ use thiserror::Error; #[non_exhaustive] pub enum RadrootsOutboxError { #[cfg(feature = "sqlite")] - #[error("SQLx error: {0}")] - Sqlx(#[from] sqlx::Error), + #[error("SQLite operation failed")] + Sqlx( + #[source] + #[from] + sqlx::Error, + ), #[error("JSON error: {0}")] Json(#[from] serde_json::Error), @@ -38,18 +42,57 @@ pub enum RadrootsOutboxError { #[error("ephemeral event kind {kind} cannot enter the durable generic outbox")] EphemeralEventNotQueueable { kind: u32 }, - #[error("in-memory SQLite pools must have exactly one connection; configured {actual}")] - UnsafeInMemoryPoolConnectionCount { actual: u32 }, - - #[error( - "SQLite pool backing does not match file_backed={file_backed}: configured filename {filename}" - )] - SqlitePoolBackingMismatch { file_backed: bool, filename: String }, - #[error("SQLite outbox file connection did not enter WAL journal mode; reported `{actual}`")] SqliteFileJournalModeNotWal { actual: String }, #[cfg(feature = "sqlite")] + #[error("SQLite outbox main database must use UTF-8 encoding; reported `{actual}`")] + SqliteMainDatabaseEncodingNotUtf8 { actual: String }, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox connection must enforce foreign keys; reported {actual}")] + SqliteForeignKeysNotEnabled { actual: i64 }, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox open deadline of {limit_ms} ms was exhausted during {stage}")] + SqliteOpenDeadlineExceeded { stage: &'static str, limit_ms: u64 }, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox file path is {actual} bytes; maximum is {max}")] + SqliteFilePathTooLong { max: usize, actual: usize }, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox file path cannot identify one database file")] + SqliteFilePathInvalid, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox file path could not be resolved")] + SqliteFilePathResolutionFailed { + #[source] + source: std::io::Error, + }, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox offline rollback is already active for this database")] + SqliteOfflineRollbackInProgress, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox offline rollback requires all owned live handles to be closed")] + SqliteOfflineRollbackHasLiveHandles, + + #[cfg(feature = "sqlite")] + #[error("SQLite outbox lifecycle failed during {stage}")] + SqliteLifecycleFailure { stage: &'static str }, + + #[cfg(feature = "sqlite")] + #[error("SQLite {field} is {actual} bytes; maximum is {max}")] + SqliteTextLimitExceeded { + field: &'static str, + max: usize, + actual: usize, + }, + + #[cfg(feature = "sqlite")] #[error( "temporary schema object `{name}` ({object_type}, table `{table_name}`) collides with outbox authority" )] @@ -168,6 +211,21 @@ pub enum RadrootsOutboxError { rollback: sqlx::Error, }, + #[cfg(feature = "sqlite")] + #[error("outbox SQLite integrity check failed: {detail}")] + IntegrityCheckFailed { detail: String }, + + #[cfg(feature = "sqlite")] + #[error( + "outbox foreign-key violation in `{table}` row {rowid:?} against `{parent}` constraint {foreign_key_id}" + )] + ForeignKeyViolation { + table: String, + rowid: Option<i64>, + parent: String, + foreign_key_id: i64, + }, + #[error( "trade mutation outbox metadata does not match the canonical mutation content: {field}" )] @@ -250,3 +308,26 @@ impl From<RadrootsTransportError> for RadrootsOutboxError { Self::Transport(value) } } + +impl RadrootsOutboxError { + /// Returns a stable public diagnostic capped at the governed byte ceiling. + pub fn public_diagnostic(&self) -> String { + #[cfg(feature = "sqlite")] + const LIMIT: usize = crate::RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX; + #[cfg(not(feature = "sqlite"))] + const LIMIT: usize = 4_096; + + let diagnostic = self.to_string(); + if diagnostic.len() <= LIMIT { + return diagnostic; + } + let suffix = "…"; + let mut end = LIMIT.saturating_sub(suffix.len()); + while !diagnostic.is_char_boundary(end) { + end = end.saturating_sub(1); + } + let mut bounded = diagnostic[..end].to_owned(); + bounded.push_str(suffix); + bounded + } +} diff --git a/crates/outbox/src/lib.rs b/crates/outbox/src/lib.rs @@ -10,6 +10,8 @@ mod model; #[cfg(feature = "sqlite")] mod schema; #[cfg(feature = "sqlite")] +mod sqlite_lifecycle; +#[cfg(feature = "sqlite")] mod store; pub use error::RadrootsOutboxError; @@ -30,4 +32,10 @@ pub use model::{ #[cfg(feature = "sqlite")] pub use schema::{RadrootsOutboxSchemaStatus, inspect_outbox_schema_status}; #[cfg(feature = "sqlite")] +pub use sqlite_lifecycle::{ + RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX, RADROOTS_OUTBOX_FILE_CONNECTION_LIMIT, + RADROOTS_OUTBOX_FILE_PATH_BYTES_MAX, RADROOTS_OUTBOX_OPEN_DEADLINE_MILLIS, + RadrootsOutboxRollbackConfirmation, +}; +#[cfg(feature = "sqlite")] pub use store::RadrootsOutbox; diff --git a/crates/outbox/src/migrations.rs b/crates/outbox/src/migrations.rs @@ -55,7 +55,6 @@ pub(crate) fn migration_for_version( .find(|migration| migration.version == version) } -#[cfg(test)] pub(crate) fn is_outbox_owned_table_name(registry: &[OutboxMigration], name: &str) -> bool { sqlite_identifier_starts_with(name, OUTBOX_RESERVED_PREFIX) || registry diff --git a/crates/outbox/src/schema.rs b/crates/outbox/src/schema.rs @@ -4,13 +4,19 @@ use crate::RadrootsOutboxError; use crate::migrations::{ OUTBOX_LEDGER_CREATE_DDL, OUTBOX_LEDGER_DDL, OUTBOX_LEDGER_NAME, OUTBOX_MIGRATIONS, OUTBOX_RESERVED_PREFIX, OutboxMigration, RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, - RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, is_outbox_governed_schema_name, migration_for_version, - validate_embedded_migration_registry, validate_migration_registry, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, is_outbox_governed_schema_name, is_outbox_owned_table_name, + migration_for_version, validate_embedded_migration_registry, validate_migration_registry, }; use sha2::{Digest, Sha256}; -use sqlx::{Row, Sqlite, SqliteConnection, SqlitePool, Transaction}; +use sqlx::{Connection, Row, Sqlite, SqliteConnection, SqlitePool, Transaction}; use std::collections::{BTreeMap, BTreeSet}; +pub(crate) const SQLITE_IDENTIFIER_BYTES_MAX: usize = 255; +const SQLITE_CATALOG_SQL_BYTES_MAX: usize = 65_536; +pub(crate) const SQLITE_LEDGER_NAME_BYTES_MAX: usize = 128; +const SQLITE_LEDGER_DIGEST_BYTES_MAX: usize = 64; +const SQLITE_INTEGRITY_RESULT_ROWS_MAX: i64 = 2; + #[cfg(test)] const EMPTY_SCHEMA_SHA256: &str = "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"; @@ -59,6 +65,21 @@ pub async fn inspect_outbox_schema_status( finish_schema_transaction(transaction, result).await } +pub(crate) async fn inspect_outbox_schema_on_connection( + connection: &mut SqliteConnection, +) -> Result<RadrootsOutboxSchemaStatus, RadrootsOutboxError> { + validate_embedded_migration_registry()?; + let mut transaction = connection.begin().await?; + let result = inspect_schema_on_connection( + &mut transaction, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) + .await; + finish_schema_transaction(transaction, result).await +} + +#[cfg(test)] pub(crate) async fn migrate_outbox_schema(pool: &SqlitePool) -> Result<(), RadrootsOutboxError> { validate_embedded_migration_registry()?; migrate_outbox_schema_with_registry( @@ -70,6 +91,26 @@ pub(crate) async fn migrate_outbox_schema(pool: &SqlitePool) -> Result<(), Radro .await } +pub(crate) async fn migrate_outbox_schema_on_connection( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + validate_embedded_migration_registry()?; + validate_migration_registry( + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + )?; + let mut transaction = connection.begin_with("BEGIN IMMEDIATE").await?; + let result = migrate_schema_on_connection( + &mut transaction, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) + .await; + finish_schema_transaction(transaction, result).await +} + +#[cfg(test)] async fn migrate_outbox_schema_with_registry( pool: &SqlitePool, registry: &[OutboxMigration], @@ -82,9 +123,8 @@ async fn migrate_outbox_schema_with_registry( finish_schema_transaction(transaction, result).await } -#[cfg(test)] -async fn rollback_outbox_schema_offline( - pool: &SqlitePool, +pub(crate) async fn rollback_outbox_schema_on_connection( + connection: &mut SqliteConnection, target: u32, ) -> Result<(), RadrootsOutboxError> { if target < RADROOTS_OUTBOX_SCHEMA_VERSION_MIN { @@ -94,7 +134,7 @@ async fn rollback_outbox_schema_offline( }); } validate_embedded_migration_registry()?; - let mut transaction = pool.begin_with("BEGIN EXCLUSIVE").await?; + let mut transaction = connection.begin_with("BEGIN EXCLUSIVE").await?; let result = rollback_schema_on_connection( &mut transaction, OUTBOX_MIGRATIONS, @@ -116,7 +156,7 @@ pub(crate) async fn destroy_outbox_schema_for_migration_test( } async fn finish_schema_transaction<T>( - transaction: Transaction<'static, Sqlite>, + transaction: Transaction<'_, Sqlite>, result: Result<T, RadrootsOutboxError>, ) -> Result<T, RadrootsOutboxError> { match result { @@ -172,7 +212,6 @@ async fn migrate_schema_on_connection( } } -#[cfg(test)] async fn rollback_schema_on_connection( connection: &mut SqliteConnection, registry: &[OutboxMigration], @@ -285,7 +324,6 @@ async fn apply_migration_up( validate_schema_fingerprint(connection, registry, migration).await } -#[cfg(test)] async fn apply_migration_down( connection: &mut SqliteConnection, registry: &[OutboxMigration], @@ -427,7 +465,14 @@ async fn validate_outbox_temp_schema_with_registry( registry: &[OutboxMigration], ) -> Result<(), RadrootsOutboxError> { let collision = sqlx::query( - "SELECT type, name, tbl_name FROM temp.sqlite_schema + "SELECT + length(CAST(type AS BLOB)) AS type_bytes, + CASE WHEN length(CAST(type AS BLOB)) <= ? THEN type END AS bounded_type, + length(CAST(name AS BLOB)) AS name_bytes, + CASE WHEN length(CAST(name AS BLOB)) <= ? THEN name END AS bounded_name, + length(CAST(tbl_name AS BLOB)) AS table_name_bytes, + CASE WHEN length(CAST(tbl_name AS BLOB)) <= ? THEN tbl_name END AS bounded_table_name + FROM temp.sqlite_schema WHERE type IN ('trigger', 'view') OR lower(substr(name, 1, length(?))) = lower(?) OR lower(substr(tbl_name, 1, length(?))) = lower(?) @@ -435,6 +480,9 @@ async fn validate_outbox_temp_schema_with_registry( OR tbl_name = ? COLLATE NOCASE ORDER BY type, name, tbl_name LIMIT 1", ) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) .bind(OUTBOX_RESERVED_PREFIX) .bind(OUTBOX_RESERVED_PREFIX) .bind(OUTBOX_RESERVED_PREFIX) @@ -444,9 +492,27 @@ async fn validate_outbox_temp_schema_with_registry( .fetch_optional(&mut *connection) .await?; if let Some(row) = collision { - let object_type: String = row.try_get("type")?; - let name: String = row.try_get("name")?; - let table_name: String = row.try_get("tbl_name")?; + let object_type = bounded_required_text( + &row, + "bounded_type", + "type_bytes", + "temporary catalog object type", + SQLITE_IDENTIFIER_BYTES_MAX, + )?; + let name = bounded_required_text( + &row, + "bounded_name", + "name_bytes", + "temporary catalog object name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?; + let table_name = bounded_required_text( + &row, + "bounded_table_name", + "table_name_bytes", + "temporary catalog table name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?; if matches!(object_type.as_str(), "trigger" | "view") || is_outbox_governed_schema_name(registry, &name) || is_outbox_governed_schema_name(registry, &table_name) @@ -487,7 +553,16 @@ async fn read_catalog_bounded( ) -> Result<Vec<CatalogRow>, RadrootsOutboxError> { let row_limit = catalog_row_limit(registry)?; let rows = sqlx::query( - "SELECT type, name, tbl_name, sql FROM main.sqlite_schema + "SELECT + length(CAST(type AS BLOB)) AS type_bytes, + CASE WHEN length(CAST(type AS BLOB)) <= ? THEN type END AS bounded_type, + length(CAST(name AS BLOB)) AS name_bytes, + CASE WHEN length(CAST(name AS BLOB)) <= ? THEN name END AS bounded_name, + length(CAST(tbl_name AS BLOB)) AS table_name_bytes, + CASE WHEN length(CAST(tbl_name AS BLOB)) <= ? THEN tbl_name END AS bounded_table_name, + length(CAST(sql AS BLOB)) AS sql_bytes, + CASE WHEN sql IS NULL OR length(CAST(sql AS BLOB)) <= ? THEN sql END AS bounded_sql + FROM main.sqlite_schema WHERE lower(substr(name, 1, 7)) != 'sqlite_' AND (lower(substr(name, 1, length(?))) = lower(?) OR lower(substr(tbl_name, 1, length(?))) = lower(?) @@ -495,6 +570,10 @@ async fn read_catalog_bounded( OR tbl_name = ? COLLATE NOCASE) ORDER BY type, name, tbl_name LIMIT ?", ) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_CATALOG_SQL_BYTES_MAX).unwrap_or(i64::MAX)) .bind(OUTBOX_RESERVED_PREFIX) .bind(OUTBOX_RESERVED_PREFIX) .bind(OUTBOX_RESERVED_PREFIX) @@ -512,15 +591,86 @@ async fn read_catalog_bounded( rows.into_iter() .map(|row| { Ok(CatalogRow { - object_type: row.try_get("type")?, - name: row.try_get("name")?, - table_name: row.try_get("tbl_name")?, - sql: row.try_get("sql")?, + object_type: bounded_required_text( + &row, + "bounded_type", + "type_bytes", + "catalog object type", + SQLITE_IDENTIFIER_BYTES_MAX, + )?, + name: bounded_required_text( + &row, + "bounded_name", + "name_bytes", + "catalog object name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?, + table_name: bounded_required_text( + &row, + "bounded_table_name", + "table_name_bytes", + "catalog table name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?, + sql: bounded_optional_text( + &row, + "bounded_sql", + "sql_bytes", + "catalog SQL", + SQLITE_CATALOG_SQL_BYTES_MAX, + )?, }) }) .collect() } +fn bounded_required_text( + row: &sqlx::sqlite::SqliteRow, + value_column: &'static str, + length_column: &'static str, + field: &'static str, + max: usize, +) -> Result<String, RadrootsOutboxError> { + let actual = bounded_text_length(row.try_get(length_column)?, field, max)?; + if actual > max { + return Err(RadrootsOutboxError::SqliteTextLimitExceeded { field, max, actual }); + } + row.try_get::<Option<String>, _>(value_column)?.ok_or( + RadrootsOutboxError::SqliteLifecycleFailure { + stage: "bounded SQLite text decode", + }, + ) +} + +fn bounded_optional_text( + row: &sqlx::sqlite::SqliteRow, + value_column: &'static str, + length_column: &'static str, + field: &'static str, + max: usize, +) -> Result<Option<String>, RadrootsOutboxError> { + let Some(length) = row.try_get::<Option<i64>, _>(length_column)? else { + return Ok(None); + }; + let actual = bounded_text_length(length, field, max)?; + if actual > max { + return Err(RadrootsOutboxError::SqliteTextLimitExceeded { field, max, actual }); + } + row.try_get(value_column).map_err(Into::into) +} + +fn bounded_text_length( + length: i64, + field: &'static str, + max: usize, +) -> Result<usize, RadrootsOutboxError> { + usize::try_from(length).map_err(|_| RadrootsOutboxError::SqliteTextLimitExceeded { + field, + max, + actual: usize::MAX, + }) +} + fn validate_ledger_catalog(catalog: &[CatalogRow]) -> Result<bool, RadrootsOutboxError> { let rows = catalog .iter() @@ -606,10 +756,22 @@ async fn read_history_bounded( } })?; let rows = sqlx::query( - "SELECT version, name, up_sha256, down_sha256, schema_sha256 + "SELECT version, + length(CAST(name AS BLOB)) AS name_bytes, + CASE WHEN length(CAST(name AS BLOB)) <= ? THEN name END AS bounded_name, + length(CAST(up_sha256 AS BLOB)) AS up_sha256_bytes, + CASE WHEN length(CAST(up_sha256 AS BLOB)) <= ? THEN up_sha256 END AS bounded_up_sha256, + length(CAST(down_sha256 AS BLOB)) AS down_sha256_bytes, + CASE WHEN length(CAST(down_sha256 AS BLOB)) <= ? THEN down_sha256 END AS bounded_down_sha256, + length(CAST(schema_sha256 AS BLOB)) AS schema_sha256_bytes, + CASE WHEN length(CAST(schema_sha256 AS BLOB)) <= ? THEN schema_sha256 END AS bounded_schema_sha256 FROM main.radroots_outbox_schema_migrations ORDER BY version LIMIT ?", ) + .bind(i64::try_from(SQLITE_LEDGER_NAME_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_LEDGER_DIGEST_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_LEDGER_DIGEST_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_LEDGER_DIGEST_BYTES_MAX).unwrap_or(i64::MAX)) .bind(row_limit) .fetch_all(&mut *connection) .await?; @@ -617,10 +779,34 @@ async fn read_history_bounded( .map(|row| { Ok(AppliedMigration { version: row.try_get("version")?, - name: row.try_get("name")?, - up_sha256: row.try_get("up_sha256")?, - down_sha256: row.try_get("down_sha256")?, - schema_sha256: row.try_get("schema_sha256")?, + name: bounded_required_text( + &row, + "bounded_name", + "name_bytes", + "migration ledger name", + SQLITE_LEDGER_NAME_BYTES_MAX, + )?, + up_sha256: bounded_required_text( + &row, + "bounded_up_sha256", + "up_sha256_bytes", + "migration ledger up checksum", + SQLITE_LEDGER_DIGEST_BYTES_MAX, + )?, + down_sha256: bounded_required_text( + &row, + "bounded_down_sha256", + "down_sha256_bytes", + "migration ledger down checksum", + SQLITE_LEDGER_DIGEST_BYTES_MAX, + )?, + schema_sha256: bounded_required_text( + &row, + "bounded_schema_sha256", + "schema_sha256_bytes", + "migration ledger schema checksum", + SQLITE_LEDGER_DIGEST_BYTES_MAX, + )?, }) }) .collect() @@ -727,6 +913,121 @@ async fn validate_schema_fingerprint( Ok(()) } +pub(crate) async fn validate_outbox_owned_integrity( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + validate_embedded_migration_registry()?; + let tables = OUTBOX_MIGRATIONS + .iter() + .flat_map(|migration| migration.owned_table_names.iter().copied()) + .collect::<BTreeSet<_>>(); + for table in tables { + validate_owned_table_foreign_keys(connection, table).await?; + validate_integrity_results(connection, Some(table)).await?; + } + Ok(()) +} + +pub(crate) async fn validate_full_database_integrity( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + validate_integrity_results(connection, None).await +} + +async fn validate_owned_table_foreign_keys( + connection: &mut SqliteConnection, + table: &'static str, +) -> Result<(), RadrootsOutboxError> { + if !is_outbox_owned_table_name(OUTBOX_MIGRATIONS, table) { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: "scoped integrity table is outside outbox authority".to_owned(), + }); + } + let row = sqlx::query( + "SELECT + length(CAST(\"table\" AS BLOB)) AS table_bytes, + CASE WHEN length(CAST(\"table\" AS BLOB)) <= ? THEN \"table\" END AS bounded_table, + rowid, + length(CAST(parent AS BLOB)) AS parent_bytes, + CASE WHEN length(CAST(parent AS BLOB)) <= ? THEN parent END AS bounded_parent, + fkid + FROM pragma_foreign_key_check(?) LIMIT 1", + ) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(i64::try_from(SQLITE_IDENTIFIER_BYTES_MAX).unwrap_or(i64::MAX)) + .bind(table) + .fetch_optional(&mut *connection) + .await?; + if let Some(row) = row { + return Err(RadrootsOutboxError::ForeignKeyViolation { + table: bounded_required_text( + &row, + "bounded_table", + "table_bytes", + "foreign-key table name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?, + rowid: row.try_get("rowid")?, + parent: bounded_required_text( + &row, + "bounded_parent", + "parent_bytes", + "foreign-key parent name", + SQLITE_IDENTIFIER_BYTES_MAX, + )?, + foreign_key_id: row.try_get("fkid")?, + }); + } + Ok(()) +} + +async fn validate_integrity_results( + connection: &mut SqliteConnection, + table: Option<&'static str>, +) -> Result<(), RadrootsOutboxError> { + let max = i64::try_from(crate::RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX).unwrap_or(i64::MAX); + let rows = if let Some(table) = table { + sqlx::query( + "SELECT + length(CAST(integrity_check AS BLOB)) AS detail_bytes, + CASE WHEN length(CAST(integrity_check AS BLOB)) <= ? THEN integrity_check END AS bounded_detail + FROM pragma_integrity_check(?) LIMIT ?", + ) + .bind(max) + .bind(table) + .bind(SQLITE_INTEGRITY_RESULT_ROWS_MAX) + .fetch_all(&mut *connection) + .await? + } else { + sqlx::query( + "SELECT + length(CAST(integrity_check AS BLOB)) AS detail_bytes, + CASE WHEN length(CAST(integrity_check AS BLOB)) <= ? THEN integrity_check END AS bounded_detail + FROM pragma_integrity_check LIMIT ?", + ) + .bind(max) + .bind(SQLITE_INTEGRITY_RESULT_ROWS_MAX) + .fetch_all(&mut *connection) + .await? + }; + if rows.len() == 1 { + let detail = bounded_required_text( + &rows[0], + "bounded_detail", + "detail_bytes", + "integrity diagnostic", + crate::RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX, + )?; + if detail == "ok" { + return Ok(()); + } + return Err(RadrootsOutboxError::IntegrityCheckFailed { detail }); + } + Err(RadrootsOutboxError::IntegrityCheckFailed { + detail: "integrity check returned an invalid result cardinality".to_owned(), + }) +} + #[cfg(test)] async fn validate_empty_governed_catalog( connection: &mut SqliteConnection, @@ -1409,8 +1710,13 @@ mod tests { #[tokio::test] async fn migration_rollback_wrapper_rejects_exactly_below_the_registry_floor() { let pool = memory_pool().await; + let mut connection = pool.acquire().await.expect("connection"); assert!(matches!( - rollback_outbox_schema_offline(&pool, RADROOTS_OUTBOX_SCHEMA_VERSION_MIN - 1).await, + rollback_outbox_schema_on_connection( + &mut connection, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN - 1, + ) + .await, Err(RadrootsOutboxError::RollbackBelowVersionFloor { floor: RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, target: 0, diff --git a/crates/outbox/src/sqlite_lifecycle.rs b/crates/outbox/src/sqlite_lifecycle.rs @@ -0,0 +1,831 @@ +#![forbid(unsafe_code)] + +use crate::RadrootsOutboxError; +use crate::migrations::{ + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, +}; +use crate::schema::{ + RadrootsOutboxSchemaStatus, inspect_outbox_schema_on_connection, + migrate_outbox_schema_on_connection, rollback_outbox_schema_on_connection, + validate_full_database_integrity, validate_outbox_owned_integrity, +}; +#[cfg(test)] +use crate::schema::{SQLITE_IDENTIFIER_BYTES_MAX, SQLITE_LEDGER_NAME_BYTES_MAX}; +use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions}; +use sqlx::{SqliteConnection, SqlitePool}; +use std::collections::BTreeMap; +use std::future::Future; +use std::path::{Path, PathBuf}; +use std::str::FromStr; +use std::sync::{Arc, Mutex, MutexGuard, OnceLock}; +use std::time::{Duration, Instant}; + +/// Maximum public error diagnostic size in UTF-8 bytes. +pub const RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX: usize = 4_096; +/// Maximum accepted encoded database path size in bytes. +pub const RADROOTS_OUTBOX_FILE_PATH_BYTES_MAX: usize = 4_096; +/// Governed maximum number of connections in a file-backed outbox pool. +pub const RADROOTS_OUTBOX_FILE_CONNECTION_LIMIT: u32 = 4; +/// Total open and offline-maintenance deadline in milliseconds. +pub const RADROOTS_OUTBOX_OPEN_DEADLINE_MILLIS: u64 = 5_000; + +const MEMORY_CONNECTION_LIMIT: u32 = 1; +const OPEN_DEADLINE: Duration = Duration::from_millis(RADROOTS_OUTBOX_OPEN_DEADLINE_MILLIS); + +/// Explicit acknowledgement required by destructive offline rollback. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct RadrootsOutboxRollbackConfirmation(()); + +impl RadrootsOutboxRollbackConfirmation { + /// Acknowledges that migrations above the exact target may lose data. + pub const fn acknowledge_data_loss() -> Self { + Self(()) + } +} + +pub(crate) struct OpenedOutboxDatabase { + pub(crate) pool: SqlitePool, + pub(crate) file_lease: Option<Arc<OutboxFileLease>>, +} + +#[derive(Debug)] +pub(crate) struct OutboxFileLease { + path: PathBuf, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum FileLifecycleState { + Open(usize), + OfflineRollback, +} + +static FILE_LIFECYCLES: OnceLock<Mutex<BTreeMap<PathBuf, FileLifecycleState>>> = OnceLock::new(); + +fn file_lifecycles() -> MutexGuard<'static, BTreeMap<PathBuf, FileLifecycleState>> { + FILE_LIFECYCLES + .get_or_init(|| Mutex::new(BTreeMap::new())) + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()) +} + +impl OutboxFileLease { + fn acquire(path: PathBuf) -> Result<Arc<Self>, RadrootsOutboxError> { + let mut lifecycles = file_lifecycles(); + match lifecycles.get_mut(&path) { + Some(FileLifecycleState::Open(count)) => { + *count = + count + .checked_add(1) + .ok_or(RadrootsOutboxError::SqliteLifecycleFailure { + stage: "file lifecycle lease count", + })?; + } + Some(FileLifecycleState::OfflineRollback) => { + return Err(RadrootsOutboxError::SqliteOfflineRollbackInProgress); + } + None => { + lifecycles.insert(path.clone(), FileLifecycleState::Open(1)); + } + } + drop(lifecycles); + Ok(Arc::new(Self { path })) + } +} + +impl Drop for OutboxFileLease { + fn drop(&mut self) { + let mut lifecycles = file_lifecycles(); + match lifecycles.get_mut(&self.path) { + Some(FileLifecycleState::Open(1)) => { + lifecycles.remove(&self.path); + } + Some(FileLifecycleState::Open(count)) => *count -= 1, + Some(FileLifecycleState::OfflineRollback) | None => {} + } + } +} + +struct OfflineRollbackLease { + path: PathBuf, +} + +impl OfflineRollbackLease { + fn acquire(path: PathBuf) -> Result<Self, RadrootsOutboxError> { + let mut lifecycles = file_lifecycles(); + match lifecycles.get(&path) { + Some(FileLifecycleState::Open(_)) => { + return Err(RadrootsOutboxError::SqliteOfflineRollbackHasLiveHandles); + } + Some(FileLifecycleState::OfflineRollback) => { + return Err(RadrootsOutboxError::SqliteOfflineRollbackInProgress); + } + None => { + lifecycles.insert(path.clone(), FileLifecycleState::OfflineRollback); + } + } + drop(lifecycles); + Ok(Self { path }) + } +} + +impl Drop for OfflineRollbackLease { + fn drop(&mut self) { + let mut lifecycles = file_lifecycles(); + if matches!( + lifecycles.get(&self.path), + Some(FileLifecycleState::OfflineRollback) + ) { + lifecycles.remove(&self.path); + } + } +} + +#[derive(Clone, Copy)] +enum OpenFailureInjection { + None, + #[cfg(test)] + AfterAuthority, +} + +struct OpenDeadline { + started: Instant, + limit: Duration, +} + +impl OpenDeadline { + fn new(limit: Duration) -> Self { + Self { + started: Instant::now(), + limit, + } + } + + fn exceeded(&self, stage: &'static str) -> RadrootsOutboxError { + RadrootsOutboxError::SqliteOpenDeadlineExceeded { + stage, + limit_ms: u64::try_from(self.limit.as_millis()).unwrap_or(u64::MAX), + } + } + + fn remaining(&self, stage: &'static str) -> Result<Duration, RadrootsOutboxError> { + self.limit + .checked_sub(self.started.elapsed()) + .filter(|remaining| !remaining.is_zero()) + .ok_or_else(|| self.exceeded(stage)) + } + + async fn run<T, F>(&self, stage: &'static str, operation: F) -> Result<T, RadrootsOutboxError> + where + F: Future<Output = Result<T, RadrootsOutboxError>>, + { + let _remaining = self.remaining(stage)?; + #[cfg(feature = "runtime-tokio")] + { + tokio::time::timeout(_remaining, operation) + .await + .map_err(|_| self.exceeded(stage))? + } + #[cfg(not(feature = "runtime-tokio"))] + { + let result = operation.await; + if self.started.elapsed() >= self.limit { + return Err(self.exceeded(stage)); + } + result + } + } + + async fn wait_for_lock_retry(&self) -> Result<(), RadrootsOutboxError> { + const RETRY_INTERVAL: Duration = Duration::from_millis(10); + let remaining = self.remaining("SQLite lock contention")?; + let delay = remaining.min(RETRY_INTERVAL); + #[cfg(feature = "runtime-tokio")] + tokio::time::sleep(delay).await; + #[cfg(not(feature = "runtime-tokio"))] + std::thread::sleep(delay); + Ok(()) + } +} + +pub(crate) async fn open_memory() -> Result<OpenedOutboxDatabase, RadrootsOutboxError> { + let deadline = OpenDeadline::new(OPEN_DEADLINE); + let options = SqliteConnectOptions::from_str("sqlite::memory:")? + .foreign_keys(true) + .busy_timeout(OPEN_DEADLINE); + open_owned_pool( + options, + false, + MEMORY_CONNECTION_LIMIT, + deadline, + None, + OpenFailureInjection::None, + ) + .await +} + +pub(crate) async fn open_file(path: &Path) -> Result<OpenedOutboxDatabase, RadrootsOutboxError> { + open_file_with(path, OPEN_DEADLINE, OpenFailureInjection::None).await +} + +async fn open_file_with( + path: &Path, + limit: Duration, + injection: OpenFailureInjection, +) -> Result<OpenedOutboxDatabase, RadrootsOutboxError> { + let deadline = OpenDeadline::new(limit); + let lifecycle_path = canonical_lifecycle_path(path)?; + deadline.remaining("file path resolution")?; + let file_lease = OutboxFileLease::acquire(lifecycle_path)?; + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(true) + .foreign_keys(true) + .busy_timeout(limit); + open_owned_pool( + options, + true, + RADROOTS_OUTBOX_FILE_CONNECTION_LIMIT, + deadline, + Some(file_lease), + injection, + ) + .await +} + +async fn open_owned_pool( + options: SqliteConnectOptions, + file_backed: bool, + connection_limit: u32, + deadline: OpenDeadline, + file_lease: Option<Arc<OutboxFileLease>>, + injection: OpenFailureInjection, +) -> Result<OpenedOutboxDatabase, RadrootsOutboxError> { + let pool = SqlitePoolOptions::new() + .min_connections(0) + .max_connections(connection_limit) + .acquire_timeout(deadline.limit) + .connect_lazy_with(options.clone()); + let result = loop { + match authenticate_owned_pool(&pool, options.clone(), file_backed, &deadline, injection) + .await + { + Err(error) if is_sqlite_lock_contention(&error) => { + if let Err(deadline_error) = deadline.wait_for_lock_retry().await { + break Err(deadline_error); + } + } + result => break result, + } + }; + if let Err(error) = result { + pool.close().await; + return Err(error); + } + Ok(OpenedOutboxDatabase { pool, file_lease }) +} + +fn is_sqlite_lock_contention(error: &RadrootsOutboxError) -> bool { + let RadrootsOutboxError::Sqlx(sqlx::Error::Database(database)) = error else { + return false; + }; + database + .code() + .and_then(|code| code.parse::<u32>().ok()) + .is_some_and(|code| matches!(code & 0xff, 5 | 6)) +} + +async fn authenticate_owned_pool( + pool: &SqlitePool, + options: SqliteConnectOptions, + file_backed: bool, + deadline: &OpenDeadline, + injection: OpenFailureInjection, +) -> Result<(), RadrootsOutboxError> { + let mut connection = deadline + .run("connection", async { + pool.acquire().await.map_err(RadrootsOutboxError::from) + }) + .await?; + deadline + .run("connection-local safety settings", async { + validate_connection_local_authority(&mut connection).await + }) + .await?; + deadline + .run("schema authentication", async { + inspect_outbox_schema_on_connection(&mut connection) + .await + .map(|_| ()) + }) + .await?; + #[cfg(test)] + if matches!(injection, OpenFailureInjection::AfterAuthority) { + return Err(RadrootsOutboxError::SqliteLifecycleFailure { + stage: "injected post-authority failure", + }); + } + #[cfg(not(test))] + let _ = injection; + deadline + .run("schema migration", async { + migrate_outbox_schema_on_connection(&mut connection).await + }) + .await?; + deadline + .run("managed schema authentication", async { + match inspect_outbox_schema_on_connection(&mut connection).await? { + RadrootsOutboxSchemaStatus::Managed { version } + if version == RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT => + { + Ok(()) + } + _ => Err(RadrootsOutboxError::SqliteLifecycleFailure { + stage: "managed schema postcondition", + }), + } + }) + .await?; + deadline + .run("owned integrity", async { + validate_outbox_owned_integrity(&mut connection).await + }) + .await?; + if file_backed { + deadline + .run("file journal mode", async { + configure_file_journal_mode(&mut connection).await + }) + .await?; + pool.set_connect_options(options.journal_mode(SqliteJournalMode::Wal)); + } + Ok(()) +} + +async fn validate_connection_local_authority( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + let encoding: String = sqlx::query_scalar("PRAGMA main.encoding") + .fetch_one(&mut *connection) + .await?; + if encoding != "UTF-8" { + return Err(RadrootsOutboxError::SqliteMainDatabaseEncodingNotUtf8 { actual: encoding }); + } + let foreign_keys: i64 = sqlx::query_scalar("PRAGMA foreign_keys") + .fetch_one(&mut *connection) + .await?; + if foreign_keys != 1 { + return Err(RadrootsOutboxError::SqliteForeignKeysNotEnabled { + actual: foreign_keys, + }); + } + Ok(()) +} + +async fn configure_file_journal_mode( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + let actual: String = sqlx::query_scalar("PRAGMA main.journal_mode = WAL") + .fetch_one(&mut *connection) + .await?; + if actual != "wal" { + return Err(RadrootsOutboxError::SqliteFileJournalModeNotWal { actual }); + } + Ok(()) +} + +pub(crate) async fn verify_full_integrity(pool: &SqlitePool) -> Result<(), RadrootsOutboxError> { + let mut connection = pool.acquire().await?; + validate_full_database_integrity(&mut connection).await +} + +pub(crate) async fn rollback_file_offline( + path: &Path, + target: u32, + confirmation: RadrootsOutboxRollbackConfirmation, +) -> Result<(), RadrootsOutboxError> { + let deadline = OpenDeadline::new(OPEN_DEADLINE); + if target < RADROOTS_OUTBOX_SCHEMA_VERSION_MIN { + return Err(RadrootsOutboxError::RollbackBelowVersionFloor { + floor: RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + target, + }); + } + let lifecycle_path = canonical_lifecycle_path(path)?; + deadline.remaining("offline rollback path resolution")?; + let _rollback_lease = OfflineRollbackLease::acquire(lifecycle_path)?; + let _confirmation = confirmation; + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(false) + .foreign_keys(true) + .busy_timeout(OPEN_DEADLINE); + let pool = SqlitePoolOptions::new() + .min_connections(0) + .max_connections(1) + .acquire_timeout(OPEN_DEADLINE) + .connect_lazy_with(options); + let result = rollback_owned_pool(&pool, target, &deadline).await; + pool.close().await; + result +} + +async fn rollback_owned_pool( + pool: &SqlitePool, + target: u32, + deadline: &OpenDeadline, +) -> Result<(), RadrootsOutboxError> { + let mut connection = deadline + .run("offline rollback connection", async { + pool.acquire().await.map_err(RadrootsOutboxError::from) + }) + .await?; + deadline + .run("offline rollback connection authority", async { + validate_connection_local_authority(&mut connection).await + }) + .await?; + let current = deadline + .run("offline rollback preflight", async { + match inspect_outbox_schema_on_connection(&mut connection).await? { + RadrootsOutboxSchemaStatus::Managed { version } => Ok(version), + _ => Err(RadrootsOutboxError::RollbackUnmanaged), + } + }) + .await?; + if target > current { + return Err(RadrootsOutboxError::RollbackAhead { current, target }); + } + deadline + .run("offline rollback owned integrity preflight", async { + validate_outbox_owned_integrity(&mut connection).await + }) + .await?; + deadline + .run("offline rollback transaction", async { + rollback_outbox_schema_on_connection(&mut connection, target).await + }) + .await?; + deadline + .run("offline rollback postcondition", async { + match inspect_outbox_schema_on_connection(&mut connection).await? { + RadrootsOutboxSchemaStatus::Managed { version } if version == target => { + validate_outbox_owned_integrity(&mut connection).await + } + _ => Err(RadrootsOutboxError::SqliteLifecycleFailure { + stage: "offline rollback target postcondition", + }), + } + }) + .await +} + +fn canonical_lifecycle_path(path: &Path) -> Result<PathBuf, RadrootsOutboxError> { + validate_path_bytes(path)?; + let filename = path + .file_name() + .filter(|filename| !filename.is_empty()) + .ok_or(RadrootsOutboxError::SqliteFilePathInvalid)?; + let resolved = if path.exists() { + std::fs::canonicalize(path) + } else { + let parent = path + .parent() + .filter(|parent| !parent.as_os_str().is_empty()); + std::fs::canonicalize(parent.unwrap_or_else(|| Path::new("."))) + .map(|parent| parent.join(filename)) + } + .map_err(|source| RadrootsOutboxError::SqliteFilePathResolutionFailed { source })?; + validate_path_bytes(&resolved)?; + Ok(resolved) +} + +fn validate_path_bytes(path: &Path) -> Result<(), RadrootsOutboxError> { + let actual = path.as_os_str().as_encoded_bytes().len(); + if actual > RADROOTS_OUTBOX_FILE_PATH_BYTES_MAX { + return Err(RadrootsOutboxError::SqliteFilePathTooLong { + max: RADROOTS_OUTBOX_FILE_PATH_BYTES_MAX, + actual, + }); + } + Ok(()) +} + +#[cfg(test)] +mod tests { + use super::*; + use crate::{RadrootsOutbox, RadrootsOutboxSchemaStatus}; + use sqlx::Connection; + + async fn direct_connection(path: &Path) -> SqliteConnection { + SqliteConnection::connect_with( + &SqliteConnectOptions::new() + .filename(path) + .create_if_missing(true), + ) + .await + .expect("direct connection") + } + + async fn journal_mode(path: &Path) -> String { + let mut connection = direct_connection(path).await; + sqlx::query_scalar("PRAGMA main.journal_mode") + .fetch_one(&mut connection) + .await + .expect("journal mode") + } + + #[tokio::test] + async fn sqlite_lifecycle_owned_pool_is_bounded_lazy_and_reopen_is_authenticated() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("owned.sqlite"); + let store = RadrootsOutbox::open_file(&path).await.expect("open"); + assert_eq!( + store.pool().options().get_max_connections(), + RADROOTS_OUTBOX_FILE_CONNECTION_LIMIT + ); + assert_eq!(store.pool().size(), 1); + assert_eq!(store.pragma_foreign_keys().await.expect("foreign keys"), 1); + assert_eq!( + store.pragma_busy_timeout().await.expect("busy timeout"), + i64::try_from(RADROOTS_OUTBOX_OPEN_DEADLINE_MILLIS).expect("deadline") + ); + assert_eq!(store.pragma_journal_mode().await.expect("journal"), "wal"); + assert_eq!( + store.schema_status().await.expect("schema"), + RadrootsOutboxSchemaStatus::Managed { + version: RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + } + ); + store + .verify_full_database_integrity() + .await + .expect("explicit full integrity"); + store.close().await; + + let reopened = RadrootsOutbox::open_file(&path).await.expect("reopen"); + assert_eq!(reopened.pool().size(), 1); + reopened.close().await; + } + + #[tokio::test] + async fn sqlite_lifecycle_rejects_hostile_schema_before_persistent_mutation() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("hostile.sqlite"); + let mut connection = direct_connection(&path).await; + sqlx::raw_sql( + "CREATE TABLE caller_state(value TEXT NOT NULL); + INSERT INTO caller_state(value) VALUES ('preserved'); + CREATE TABLE outbox_counterfeit(value TEXT NOT NULL);", + ) + .execute(&mut connection) + .await + .expect("hostile schema"); + connection.close().await.expect("close"); + assert_eq!(journal_mode(&path).await, "delete"); + + assert!(matches!( + RadrootsOutbox::open_file(&path).await, + Err(RadrootsOutboxError::UnmanagedSchema { .. }) + )); + assert_eq!(journal_mode(&path).await, "delete"); + let mut connection = direct_connection(&path).await; + let caller: String = sqlx::query_scalar("SELECT value FROM caller_state") + .fetch_one(&mut connection) + .await + .expect("caller state"); + assert_eq!(caller, "preserved"); + let ledger: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM sqlite_schema WHERE name = 'radroots_outbox_schema_migrations'", + ) + .fetch_one(&mut connection) + .await + .expect("ledger count"); + assert_eq!(ledger, 0); + } + + #[tokio::test] + async fn sqlite_lifecycle_bounds_catalog_ledger_path_and_diagnostics() { + let directory = tempfile::tempdir().expect("tempdir"); + let oversized_path = directory + .path() + .join("x".repeat(RADROOTS_OUTBOX_FILE_PATH_BYTES_MAX + 1)); + assert!(matches!( + RadrootsOutbox::open_file(&oversized_path).await, + Err(RadrootsOutboxError::SqliteFilePathTooLong { .. }) + )); + assert!(!oversized_path.exists()); + + let catalog_path = directory.path().join("catalog.sqlite"); + let mut connection = direct_connection(&catalog_path).await; + let oversized_name = format!("outbox_{}", "n".repeat(SQLITE_IDENTIFIER_BYTES_MAX)); + let sql = format!("CREATE TABLE \"{oversized_name}\"(value TEXT)"); + sqlx::query(sqlx::AssertSqlSafe(sql)) + .execute(&mut connection) + .await + .expect("oversized catalog identifier"); + connection.close().await.expect("close catalog"); + assert!(matches!( + RadrootsOutbox::open_file(&catalog_path).await, + Err(RadrootsOutboxError::SqliteTextLimitExceeded { + field: "catalog object name", + .. + }) + )); + + let ledger_path = directory.path().join("ledger.sqlite"); + let store = RadrootsOutbox::open_file(&ledger_path) + .await + .expect("ledger store"); + store.close().await; + let mut connection = direct_connection(&ledger_path).await; + sqlx::query("UPDATE radroots_outbox_schema_migrations SET name = ? WHERE version = 1") + .bind("n".repeat(SQLITE_LEDGER_NAME_BYTES_MAX + 1)) + .execute(&mut connection) + .await + .expect("oversized ledger name"); + connection.close().await.expect("close ledger"); + assert!(matches!( + RadrootsOutbox::open_file(&ledger_path).await, + Err(RadrootsOutboxError::SqliteTextLimitExceeded { + field: "migration ledger name", + .. + }) + )); + + let diagnostic = RadrootsOutboxError::IntegrityCheckFailed { + detail: "é".repeat(RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX), + } + .public_diagnostic(); + assert!(diagnostic.len() <= RADROOTS_OUTBOX_DIAGNOSTIC_BYTES_MAX); + assert!(diagnostic.is_char_boundary(diagnostic.len())); + } + + #[tokio::test] + async fn sqlite_lifecycle_rejects_non_utf8_and_owned_foreign_key_drift() { + let directory = tempfile::tempdir().expect("tempdir"); + let encoding_path = directory.path().join("utf16.sqlite"); + let mut encoding = direct_connection(&encoding_path).await; + sqlx::query("PRAGMA main.encoding = 'UTF-16le'") + .execute(&mut encoding) + .await + .expect("UTF-16 encoding"); + sqlx::query("CREATE TABLE encoding_anchor(value TEXT)") + .execute(&mut encoding) + .await + .expect("encoding anchor"); + encoding.close().await.expect("close encoding"); + assert!(matches!( + RadrootsOutbox::open_file(&encoding_path).await, + Err(RadrootsOutboxError::SqliteMainDatabaseEncodingNotUtf8 { .. }) + )); + assert_eq!(journal_mode(&encoding_path).await, "delete"); + + let foreign_key_path = directory.path().join("foreign-key.sqlite"); + let store = RadrootsOutbox::open_file(&foreign_key_path) + .await + .expect("foreign-key store"); + store.close().await; + let mut connection = direct_connection(&foreign_key_path).await; + sqlx::query("PRAGMA foreign_keys = OFF") + .execute(&mut connection) + .await + .expect("disable foreign keys"); + sqlx::query( + "INSERT INTO outbox_event(outbox_event_id, operation_id, event_id, expected_pubkey, draft_json, signed_event_json, raw_event_json, state, attempt_count, claim_token, claim_owner, claim_expires_at_ms, active_delivery_plan_id, next_attempt_after_ms, last_error, event_store_ingested, event_store_inserted, event_store_ingested_at_ms, created_at_ms, updated_at_ms) + VALUES (1, 999, 'event', 'author', '{}', NULL, NULL, 'draft_queued', 0, NULL, NULL, NULL, NULL, 0, NULL, 0, 0, NULL, 0, 0)", + ) + .execute(&mut connection) + .await + .expect("foreign-key drift"); + connection.close().await.expect("close foreign-key fixture"); + assert!(matches!( + RadrootsOutbox::open_file(&foreign_key_path).await, + Err(RadrootsOutboxError::ForeignKeyViolation { + table, + .. + }) if table == "outbox_event" + )); + } + + #[tokio::test] + async fn sqlite_lifecycle_deadline_concurrency_and_failure_injection_are_transactional() { + let directory = tempfile::tempdir().expect("tempdir"); + let concurrent_path = directory.path().join("concurrent.sqlite"); + let (left, right) = tokio::join!( + RadrootsOutbox::open_file(&concurrent_path), + RadrootsOutbox::open_file(&concurrent_path), + ); + left.expect("left concurrent open").close().await; + right.expect("right concurrent open").close().await; + + let locked_path = directory.path().join("locked.sqlite"); + let store = RadrootsOutbox::open_file(&locked_path) + .await + .expect("locked store"); + store.close().await; + let mut locker = direct_connection(&locked_path).await; + let transaction = locker + .begin_with("BEGIN EXCLUSIVE") + .await + .expect("exclusive lock"); + let lock_error = match open_file_with( + &locked_path, + Duration::from_millis(50), + OpenFailureInjection::None, + ) + .await + { + Ok(_) => panic!("locked open must exhaust its deadline"), + Err(error) => error, + }; + assert!( + matches!( + lock_error, + RadrootsOutboxError::SqliteOpenDeadlineExceeded { + stage: "schema migration" | "SQLite lock contention", + limit_ms: 50, + } + ), + "{lock_error:?}" + ); + transaction.rollback().await.expect("release lock"); + + let injected_path = directory.path().join("injected.sqlite"); + let mut connection = direct_connection(&injected_path).await; + sqlx::raw_sql( + "CREATE TABLE caller_state(value TEXT NOT NULL); + INSERT INTO caller_state(value) VALUES ('preserved');", + ) + .execute(&mut connection) + .await + .expect("caller fixture"); + connection.close().await.expect("close fixture"); + assert!(matches!( + open_file_with( + &injected_path, + OPEN_DEADLINE, + OpenFailureInjection::AfterAuthority, + ) + .await, + Err(RadrootsOutboxError::SqliteLifecycleFailure { + stage: "injected post-authority failure", + }) + )); + assert_eq!(journal_mode(&injected_path).await, "delete"); + let mut connection = direct_connection(&injected_path).await; + let outbox_objects: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM sqlite_schema WHERE name LIKE 'outbox_%' OR name = 'radroots_outbox_schema_migrations'", + ) + .fetch_one(&mut connection) + .await + .expect("outbox object count"); + assert_eq!(outbox_objects, 0); + let caller: String = sqlx::query_scalar("SELECT value FROM caller_state") + .fetch_one(&mut connection) + .await + .expect("caller state"); + assert_eq!(caller, "preserved"); + } + + #[tokio::test] + async fn sqlite_lifecycle_offline_rollback_rejects_live_handles_and_exactly_targets() { + let directory = tempfile::tempdir().expect("tempdir"); + let path = directory.path().join("rollback.sqlite"); + let store = RadrootsOutbox::open_file(&path).await.expect("store"); + assert!(matches!( + RadrootsOutbox::rollback_file_schema_offline( + &path, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await, + Err(RadrootsOutboxError::SqliteOfflineRollbackHasLiveHandles) + )); + store.close().await; + RadrootsOutbox::rollback_file_schema_offline( + &path, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await + .expect("exact target no-op"); + assert!(matches!( + RadrootsOutbox::rollback_file_schema_offline( + &path, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT + 1, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await, + Err(RadrootsOutboxError::RollbackAhead { .. }) + )); + assert!(matches!( + RadrootsOutbox::rollback_file_schema_offline( + &path, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN - 1, + RadrootsOutboxRollbackConfirmation::acknowledge_data_loss(), + ) + .await, + Err(RadrootsOutboxError::RollbackBelowVersionFloor { .. }) + )); + } +} diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -14,8 +14,10 @@ use crate::model::{ RadrootsOutboxSignedOperationInput, RadrootsOutboxSignedTradeMutationInput, RadrootsOutboxStatusSummary, RadrootsOutboxTradeMutationInput, }; -use crate::schema::{ - RadrootsOutboxSchemaStatus, inspect_outbox_schema_status, migrate_outbox_schema, +use crate::schema::{RadrootsOutboxSchemaStatus, inspect_outbox_schema_status}; +use crate::sqlite_lifecycle::{ + OutboxFileLease, RadrootsOutboxRollbackConfirmation, open_file, open_memory, + rollback_file_offline, verify_full_integrity, }; use radroots_event::RadrootsEventKindClass; use radroots_event::draft::{ @@ -38,66 +40,62 @@ use radroots_transport::{ }; use serde::Serialize; use sha2::{Digest, Sha256}; -use sqlx::sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteQueryResult}; -#[cfg(test)] -use sqlx::{Connection, SqliteConnection}; +use sqlx::sqlite::SqliteQueryResult; use sqlx::{Row, SqlitePool}; use std::collections::BTreeSet; use std::path::Path; -use std::str::FromStr; -use std::time::Duration; +use std::sync::Arc; #[derive(Clone)] pub struct RadrootsOutbox { pool: SqlitePool, + _file_lease: Option<Arc<OutboxFileLease>>, } impl RadrootsOutbox { pub async fn open_memory() -> Result<Self, RadrootsOutboxError> { - let options = SqliteConnectOptions::from_str("sqlite::memory:")?; - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) - .await?; - configure_pool(&pool, false).await?; - migrate_outbox_schema(&pool).await?; - Ok(Self { pool }) + let opened = open_memory().await?; + Ok(Self { + pool: opened.pool, + _file_lease: opened.file_lease, + }) } pub async fn open_file(path: impl AsRef<Path>) -> Result<Self, RadrootsOutboxError> { - let options = SqliteConnectOptions::new() - .filename(path) - .create_if_missing(true); - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) - .await?; - configure_pool(&pool, true).await?; - migrate_outbox_schema(&pool).await?; - Ok(Self { pool }) + let opened = open_file(path.as_ref()).await?; + Ok(Self { + pool: opened.pool, + _file_lease: opened.file_lease, + }) } - pub async fn open_pool( - pool: SqlitePool, - file_backed: bool, - ) -> Result<Self, RadrootsOutboxError> { - configure_pool(&pool, file_backed).await?; - migrate_outbox_schema(&pool).await?; - Ok(Self { pool }) + /// Performs an exact destructive migration rollback against a closed file. + pub async fn rollback_file_schema_offline( + path: impl AsRef<Path>, + target: u32, + confirmation: RadrootsOutboxRollbackConfirmation, + ) -> Result<(), RadrootsOutboxError> { + rollback_file_offline(path.as_ref(), target, confirmation).await } - pub fn pool(&self) -> &SqlitePool { + #[cfg(test)] + pub(crate) fn pool(&self) -> &SqlitePool { &self.pool } + /// Closes this owned pool. Existing clones are closed by the same action. + pub async fn close(self) { + self.pool.close().await; + } + /// Inspects and authenticates the governed outbox schema. pub async fn schema_status(&self) -> Result<RadrootsOutboxSchemaStatus, RadrootsOutboxError> { inspect_outbox_schema_status(&self.pool).await } - /// Serializes schema adoption or migration to the current supported version. - pub async fn migrate_to_current_schema(&self) -> Result<(), RadrootsOutboxError> { - migrate_outbox_schema(&self.pool).await + /// Runs an explicit full-database integrity maintenance check. + pub async fn verify_full_database_integrity(&self) -> Result<(), RadrootsOutboxError> { + verify_full_integrity(&self.pool).await } pub async fn pragma_foreign_keys(&self) -> Result<i64, RadrootsOutboxError> { @@ -1815,82 +1813,6 @@ fn publish_lifecycle_from_plan_evaluation<'a>( } } -async fn configure_pool(pool: &SqlitePool, file_backed: bool) -> Result<(), RadrootsOutboxError> { - let max_connections = pool.options().get_max_connections(); - let existing_options = pool.connect_options(); - let main_filename: String = - sqlx::query_scalar("SELECT file FROM pragma_database_list WHERE name = 'main'") - .fetch_one(pool) - .await?; - let database_is_memory = main_filename.is_empty(); - if file_backed == database_is_memory { - return Err(RadrootsOutboxError::SqlitePoolBackingMismatch { - file_backed, - filename: main_filename, - }); - } - if !file_backed && max_connections != 1 { - return Err(RadrootsOutboxError::UnsafeInMemoryPoolConnectionCount { - actual: max_connections, - }); - } - - let mut connect_options = existing_options - .as_ref() - .clone() - .foreign_keys(true) - .busy_timeout(Duration::from_millis(5_000)); - if file_backed { - connect_options = connect_options.journal_mode(SqliteJournalMode::Wal); - } - pool.set_connect_options(connect_options); - - let mut connections = Vec::with_capacity(max_connections as usize); - for _ in 0..max_connections { - connections.push(pool.acquire().await?); - } - for connection in &mut connections { - sqlx::query("PRAGMA foreign_keys = ON") - .execute(&mut **connection) - .await?; - sqlx::query("PRAGMA busy_timeout = 5000") - .execute(&mut **connection) - .await?; - if file_backed { - let actual = sqlx::query_scalar::<_, String>("PRAGMA main.journal_mode = WAL") - .fetch_one(&mut **connection) - .await?; - if actual != "wal" { - return Err(RadrootsOutboxError::SqliteFileJournalModeNotWal { actual }); - } - } - } - Ok(()) -} - -#[cfg(test)] -async fn file_pool_with_immutable_delete_journal(path: &Path) -> SqlitePool { - let mut writer = SqliteConnection::connect_with( - &SqliteConnectOptions::new() - .filename(path) - .create_if_missing(true), - ) - .await - .expect("writer connection"); - let mode: String = sqlx::query_scalar("PRAGMA main.journal_mode = DELETE") - .fetch_one(&mut writer) - .await - .expect("delete journal mode"); - assert_eq!(mode, "delete"); - writer.close().await.expect("close writer"); - - SqlitePoolOptions::new() - .max_connections(1) - .connect_with(SqliteConnectOptions::new().filename(path).immutable(true)) - .await - .expect("immutable pool") -} - #[cfg_attr(coverage_nightly, coverage(off))] async fn query_i64(pool: &SqlitePool, sql: &'static str) -> Result<i64, RadrootsOutboxError> { let row = sqlx::query(sql).fetch_one(pool).await?; @@ -4243,23 +4165,6 @@ mod tests { version: crate::RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, } ); - outbox - .migrate_to_current_schema() - .await - .expect("idempotent migration"); - } - - #[tokio::test] - async fn file_pool_rejects_successful_non_wal_journal_result() { - let directory = tempfile::tempdir().expect("tempdir"); - let path = directory.path().join("immutable-delete.sqlite"); - let pool = file_pool_with_immutable_delete_journal(&path).await; - - assert!(matches!( - RadrootsOutbox::open_pool(pool, true).await, - Err(RadrootsOutboxError::SqliteFileJournalModeNotWal { actual }) - if actual == "delete" - )); } #[tokio::test] @@ -4276,15 +4181,7 @@ mod tests { "wal" ); - let options = SqliteConnectOptions::from_str("sqlite::memory:").expect("memory options"); - let pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with(options) - .await - .expect("memory pool"); - let outbox = RadrootsOutbox::open_pool(pool, false) - .await - .expect("pool outbox"); + let outbox = RadrootsOutbox::open_memory().await.expect("memory outbox"); let draft = generic_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "preflight"); let signed_event = radroots_nostr_sign_frozen_draft(&fixture_keys(), &draft) @@ -4360,109 +4257,6 @@ mod tests { } #[tokio::test] - async fn open_pool_configures_every_file_connection_and_rejects_unsafe_memory_pools() { - let memory_options = SqliteConnectOptions::from_str("sqlite::memory:") - .expect("memory options") - .foreign_keys(false); - let memory_pool = SqlitePoolOptions::new() - .max_connections(2) - .connect_with(memory_options) - .await - .expect("memory pool"); - let memory_error = match RadrootsOutbox::open_pool(memory_pool, false).await { - Ok(_) => panic!("multi-connection memory pool must be rejected"), - Err(error) => error, - }; - assert!( - matches!( - memory_error, - RadrootsOutboxError::UnsafeInMemoryPoolConnectionCount { actual: 2 } - ), - "{memory_error:?}" - ); - - let mislabeled_memory_pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with( - SqliteConnectOptions::from_str("sqlite::memory:") - .expect("memory options") - .foreign_keys(false), - ) - .await - .expect("mislabeled memory pool"); - assert!(matches!( - RadrootsOutbox::open_pool(mislabeled_memory_pool, true).await, - Err(RadrootsOutboxError::SqlitePoolBackingMismatch { - file_backed: true, - .. - }) - )); - for memory_url in ["sqlite://?mode=memory", "sqlite://named?mode=memory"] { - let mode_memory_pool = SqlitePoolOptions::new() - .max_connections(2) - .connect_with( - SqliteConnectOptions::from_str(memory_url) - .expect("mode-memory options") - .foreign_keys(false), - ) - .await - .expect("mode-memory pool"); - assert!(matches!( - RadrootsOutbox::open_pool(mode_memory_pool, false).await, - Err(RadrootsOutboxError::UnsafeInMemoryPoolConnectionCount { actual: 2 }) - )); - - let mislabeled_mode_memory_pool = SqlitePoolOptions::new() - .max_connections(1) - .connect_with( - SqliteConnectOptions::from_str(memory_url) - .expect("mode-memory options") - .foreign_keys(false), - ) - .await - .expect("mislabeled mode-memory pool"); - assert!(matches!( - RadrootsOutbox::open_pool(mislabeled_mode_memory_pool, true).await, - Err(RadrootsOutboxError::SqlitePoolBackingMismatch { - file_backed: true, - .. - }) - )); - } - - let directory = tempfile::tempdir().expect("tempdir"); - let file_path = directory.path().join("multi-connection-outbox.sqlite"); - let file_options = SqliteConnectOptions::new() - .filename(&file_path) - .create_if_missing(true) - .foreign_keys(false); - let file_pool = SqlitePoolOptions::new() - .max_connections(3) - .connect_with(file_options) - .await - .expect("file pool"); - let outbox = RadrootsOutbox::open_pool(file_pool, true) - .await - .expect("file outbox"); - let mut connections = Vec::new(); - for _ in 0..3 { - connections.push(outbox.pool().acquire().await.expect("connection")); - } - for connection in &mut connections { - let foreign_keys: i64 = sqlx::query_scalar("PRAGMA foreign_keys") - .fetch_one(&mut **connection) - .await - .expect("foreign keys"); - let busy_timeout: i64 = sqlx::query_scalar("PRAGMA busy_timeout") - .fetch_one(&mut **connection) - .await - .expect("busy timeout"); - assert_eq!(foreign_keys, 1); - assert_eq!(busy_timeout, 5_000); - } - } - - #[tokio::test] async fn defensive_storage_decoding_and_remaining_idempotency_edges_are_explicit() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); let draft = generic_draft(FIXTURE_ALICE_PUBLIC_KEY_HEX, "defensive storage"); diff --git a/tools/xtask/src/contract/outbox_migration.rs b/tools/xtask/src/contract/outbox_migration.rs @@ -78,6 +78,10 @@ const SOURCE_FILES: &[(&str, &str)] = &[ "crates/outbox/src/migrations.rs", ), ("schema_runtime", "crates/outbox/src/schema.rs"), + ( + "sqlite_lifecycle_runtime", + "crates/outbox/src/sqlite_lifecycle.rs", + ), ("store_integration", "crates/outbox/src/store.rs"), ("vector_executor", VECTOR_EXECUTOR_RELATIVE), ( @@ -893,7 +897,7 @@ fn expected_manifest( "reserved_prefix": RESERVED_PREFIX, "catalog_fingerprint": "sha256_type_nul_name_nul_table_name_nul_sql_nul_sorted_v1", "migration_transaction": "begin_immediate_v1", - "rollback_transaction": "begin_exclusive_test_executor_v1", + "rollback_transaction": "begin_exclusive_offline_path_v1", "adoption": "exact_unledgered_0001_catalog_only_v1", }, "migrations": migrations, @@ -959,7 +963,7 @@ fn manifest_schema() -> Value { "reserved_prefix": { "const": RESERVED_PREFIX }, "catalog_fingerprint": { "const": "sha256_type_nul_name_nul_table_name_nul_sql_nul_sorted_v1" }, "migration_transaction": { "const": "begin_immediate_v1" }, - "rollback_transaction": { "const": "begin_exclusive_test_executor_v1" }, + "rollback_transaction": { "const": "begin_exclusive_offline_path_v1" }, "adoption": { "const": "exact_unledgered_0001_catalog_only_v1" } } }, @@ -1000,7 +1004,7 @@ fn manifest_schema() -> Value { } }, "source_files": { - "type": "array", "minItems": 16, "maxItems": 16, + "type": "array", "minItems": SOURCE_FILES.len(), "maxItems": SOURCE_FILES.len(), "items": { "$ref": "#/$defs/source_file" } }, "release": {