lib

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

commit 3e09cecc8afe7c312e4b9b2d62091728454aea14
parent e686fed05fd696eb0b41a945bb86c6f48b41ac9f
Author: triesap <tyson@radroots.org>
Date:   Mon, 27 Jul 2026 02:45:37 +0000

outbox: authenticate versioned migration history

- freeze the existing 0001 SQL identity and schema fingerprint
- derive supported bounds and inventories from the ordered registry
- adopt only the exact unledgered baseline in one transaction
- exercise successor rollback and hostile history mutations

Diffstat:
Mcrates/outbox/src/error.rs | 120+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/lib.rs | 6+++++-
Mcrates/outbox/src/migrations.rs | 622++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-
Acrates/outbox/src/schema.rs | 1420+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/outbox/src/store.rs | 52+++++++++++++++++++++++++---------------------------
5 files changed, 2190 insertions(+), 30 deletions(-)

diff --git a/crates/outbox/src/error.rs b/crates/outbox/src/error.rs @@ -4,6 +4,7 @@ use radroots_transport::RadrootsTransportError; use thiserror::Error; #[derive(Debug, Error)] +#[non_exhaustive] pub enum RadrootsOutboxError { #[cfg(feature = "sqlite")] #[error("SQLx error: {0}")] @@ -48,6 +49,125 @@ pub enum RadrootsOutboxError { #[error("SQLite outbox file connection did not enter WAL journal mode; reported `{actual}`")] SqliteFileJournalModeNotWal { actual: String }, + #[cfg(feature = "sqlite")] + #[error( + "temporary schema object `{name}` ({object_type}, table `{table_name}`) collides with outbox authority" + )] + TemporarySchemaCollision { + object_type: String, + name: String, + table_name: String, + }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration registry defect: {reason}")] + MigrationRegistryDefect { reason: String }, + + #[cfg(feature = "sqlite")] + #[error( + "embedded outbox migration {version} {direction} length mismatch: expected {expected}, found {actual}" + )] + EmbeddedMigrationLengthMismatch { + version: u32, + direction: &'static str, + expected: usize, + actual: usize, + }, + + #[cfg(feature = "sqlite")] + #[error( + "embedded outbox migration {version} {direction} checksum mismatch: expected {expected}, found {actual}" + )] + EmbeddedMigrationChecksumMismatch { + version: u32, + direction: &'static str, + expected: &'static str, + actual: String, + }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration {version} {direction} catalog delta mismatch: {reason}")] + MigrationCatalogDeltaMismatch { + version: u32, + direction: &'static str, + reason: String, + }, + + #[cfg(feature = "sqlite")] + #[error("outbox governed catalog exceeds the supported {max} rows")] + GovernedCatalogCapacityExceeded { max: usize }, + + #[cfg(feature = "sqlite")] + #[error("unmanaged outbox schema has fingerprint {actual_schema_sha256}")] + UnmanagedSchema { actual_schema_sha256: String }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration ledger catalog is invalid: {reason}")] + MigrationLedgerDrift { reason: String }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration history gap: expected version {expected}, found {actual:?}")] + MigrationHistoryGap { expected: u32, actual: Option<u32> }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration history references unknown version {version}")] + UnknownMigration { version: u32 }, + + #[cfg(feature = "sqlite")] + #[error("outbox schema version {database} is newer than supported version {current}")] + SchemaTooNew { current: u32, database: i64 }, + + #[cfg(feature = "sqlite")] + #[error("outbox migration {version} name drift: expected `{expected}`, found `{actual}`")] + MigrationHistoryNameDrift { + version: u32, + expected: &'static str, + actual: String, + }, + + #[cfg(feature = "sqlite")] + #[error( + "outbox migration {version} {field} checksum drift: expected {expected}, found {actual}" + )] + MigrationHistoryChecksumDrift { + version: u32, + field: &'static str, + expected: &'static str, + actual: String, + }, + + #[cfg(feature = "sqlite")] + #[error( + "outbox schema fingerprint mismatch at version {version}: expected {expected}, found {actual}" + )] + SchemaFingerprintMismatch { + version: u32, + expected: &'static str, + actual: String, + }, + + #[cfg(feature = "sqlite")] + #[error("outbox rollback target {target} is below the supported version floor {floor}")] + RollbackBelowVersionFloor { floor: u32, target: u32 }, + + #[cfg(feature = "sqlite")] + #[error("outbox rollback target {target} is ahead of managed version {current}")] + RollbackAhead { current: u32, target: u32 }, + + #[cfg(feature = "sqlite")] + #[error("outbox rollback requires a managed schema")] + RollbackUnmanaged, + + #[cfg(feature = "sqlite")] + #[error( + "outbox schema operation failed: {primary}; transaction rollback also failed: {rollback}" + )] + MigrationTransactionRollbackFailed { + #[source] + primary: Box<RadrootsOutboxError>, + rollback: sqlx::Error, + }, + #[error( "trade mutation outbox metadata does not match the canonical mutation content: {field}" )] diff --git a/crates/outbox/src/lib.rs b/crates/outbox/src/lib.rs @@ -6,11 +6,13 @@ mod error; mod migrations; mod model; #[cfg(feature = "sqlite")] +mod schema; +#[cfg(feature = "sqlite")] mod store; pub use error::RadrootsOutboxError; #[cfg(feature = "sqlite")] -pub use migrations::{OUTBOX_MIGRATION_DOWN, OUTBOX_MIGRATION_UP}; +pub use migrations::{RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, RADROOTS_OUTBOX_SCHEMA_VERSION_MIN}; pub use model::{ RadrootsOutboxClaimedEvent, RadrootsOutboxDeliveryAttemptRecord, RadrootsOutboxDeliveryPlanInput, RadrootsOutboxDeliveryPlanRecord, @@ -24,4 +26,6 @@ pub use model::{ RadrootsOutboxTradeMutationInput, }; #[cfg(feature = "sqlite")] +pub use schema::{RadrootsOutboxSchemaStatus, inspect_outbox_schema_status}; +#[cfg(feature = "sqlite")] pub use store::RadrootsOutbox; diff --git a/crates/outbox/src/migrations.rs b/crates/outbox/src/migrations.rs @@ -1,4 +1,622 @@ #![forbid(unsafe_code)] -pub const OUTBOX_MIGRATION_UP: &str = include_str!("../migrations/0001_outbox.up.sql"); -pub const OUTBOX_MIGRATION_DOWN: &str = include_str!("../migrations/0001_outbox.down.sql"); +use crate::RadrootsOutboxError; +use sha2::{Digest, Sha256}; +use std::collections::BTreeSet; + +pub(crate) const OUTBOX_LEDGER_NAME: &str = "radroots_outbox_schema_migrations"; +pub(crate) const OUTBOX_RESERVED_PREFIX: &str = "outbox_"; + +/// Oldest managed schema version that the runtime can preserve or target. +pub const RADROOTS_OUTBOX_SCHEMA_VERSION_MIN: u32 = OUTBOX_MIGRATIONS[0].version; +/// Latest managed schema version understood by this runtime. +pub const RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT: u32 = + OUTBOX_MIGRATIONS[OUTBOX_MIGRATIONS.len() - 1].version; + +pub(crate) const OUTBOX_LEDGER_DDL: &str = "CREATE TABLE radroots_outbox_schema_migrations ( + version INTEGER PRIMARY KEY NOT NULL CHECK (version > 0), + name TEXT NOT NULL UNIQUE CHECK (length(name) > 0), + up_sha256 TEXT NOT NULL CHECK (length(up_sha256) = 64 AND up_sha256 NOT GLOB '*[^0-9a-f]*'), + down_sha256 TEXT NOT NULL CHECK (length(down_sha256) = 64 AND down_sha256 NOT GLOB '*[^0-9a-f]*'), + schema_sha256 TEXT NOT NULL CHECK (length(schema_sha256) = 64 AND schema_sha256 NOT GLOB '*[^0-9a-f]*') +) STRICT, WITHOUT ROWID"; + +pub(crate) const OUTBOX_LEDGER_CREATE_DDL: &str = + "CREATE TABLE main.radroots_outbox_schema_migrations ( + version INTEGER PRIMARY KEY NOT NULL CHECK (version > 0), + name TEXT NOT NULL UNIQUE CHECK (length(name) > 0), + up_sha256 TEXT NOT NULL CHECK (length(up_sha256) = 64 AND up_sha256 NOT GLOB '*[^0-9a-f]*'), + down_sha256 TEXT NOT NULL CHECK (length(down_sha256) = 64 AND down_sha256 NOT GLOB '*[^0-9a-f]*'), + schema_sha256 TEXT NOT NULL CHECK (length(schema_sha256) = 64 AND schema_sha256 NOT GLOB '*[^0-9a-f]*') +) STRICT, WITHOUT ROWID"; + +pub(crate) const OUTBOX_BASELINE_OBJECT_NAMES: &[&str] = &[ + "outbox_delivery_attempt", + "outbox_delivery_attempt_target_idx", + "outbox_delivery_plan", + "outbox_delivery_plan_event_idx", + "outbox_delivery_target", + "outbox_delivery_target_ready_idx", + "outbox_event", + "outbox_event_event_id_idx", + "outbox_event_ready_idx", + "outbox_operation_idempotency_idx", + "outbox_operation_status_idx", + "outbox_operation_trade_mutation_idx", + "outbox_operations", +]; + +pub(crate) const OUTBOX_BASELINE_TABLE_NAMES: &[&str] = &[ + "outbox_delivery_attempt", + "outbox_delivery_plan", + "outbox_delivery_target", + "outbox_event", + "outbox_operations", +]; + +#[derive(Clone, Copy)] +pub(crate) struct OutboxMigration { + pub(crate) version: u32, + pub(crate) name: &'static str, + pub(crate) up_sql: &'static str, + pub(crate) down_sql: &'static str, + pub(crate) up_len: usize, + pub(crate) down_len: usize, + pub(crate) up_sha256: &'static str, + pub(crate) down_sha256: &'static str, + pub(crate) schema_sha256: &'static str, + pub(crate) owned_object_names: &'static [&'static str], + pub(crate) owned_table_names: &'static [&'static str], +} + +pub(crate) const OUTBOX_MIGRATIONS: &[OutboxMigration] = &[OutboxMigration { + version: 1, + name: "outbox", + up_sql: include_str!("../migrations/0001_outbox.up.sql"), + down_sql: include_str!("../migrations/0001_outbox.down.sql"), + up_len: 5_470, + down_len: 159, + up_sha256: "a7ee775d32c2b9f845961425362e1b1e558ce0d025f7d22dd58f118ba4dab4fa", + down_sha256: "5d56f978f9172dc5ecbc5043a6c286c75926974d8a2a9e44fffa7c134829af61", + schema_sha256: "e7eeba00de78ec6d990c620e7c056018166e8a00bb703e472ef6f67a00870293", + owned_object_names: OUTBOX_BASELINE_OBJECT_NAMES, + owned_table_names: OUTBOX_BASELINE_TABLE_NAMES, +}]; + +pub(crate) fn migration_for_version( + registry: &[OutboxMigration], + version: u32, +) -> Option<&OutboxMigration> { + registry + .iter() + .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 + .iter() + .flat_map(|migration| migration.owned_table_names) + .any(|owned| name.eq_ignore_ascii_case(owned)) +} + +pub(crate) fn is_outbox_governed_schema_name(registry: &[OutboxMigration], name: &str) -> bool { + name.eq_ignore_ascii_case(OUTBOX_LEDGER_NAME) + || sqlite_identifier_starts_with(name, OUTBOX_RESERVED_PREFIX) + || registry + .iter() + .flat_map(|migration| migration.owned_object_names) + .any(|owned| name.eq_ignore_ascii_case(owned)) +} + +pub(crate) fn sqlite_identifier_starts_with(name: &str, prefix: &str) -> bool { + name.get(..prefix.len()) + .is_some_and(|candidate| candidate.eq_ignore_ascii_case(prefix)) +} + +pub(crate) fn validate_embedded_migration_registry() -> Result<(), RadrootsOutboxError> { + validate_migration_registry( + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) +} + +pub(crate) fn validate_migration_registry( + registry: &[OutboxMigration], + minimum: u32, + current: u32, +) -> Result<(), RadrootsOutboxError> { + validate_ledger_ddl_identity(OUTBOX_LEDGER_DDL, OUTBOX_LEDGER_CREATE_DDL)?; + if minimum == 0 || current < minimum || registry.is_empty() { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "migration version range {minimum}..={current} requires a non-empty positive registry" + ), + }); + } + let mut expected_version = minimum; + let mut object_names = BTreeSet::new(); + let mut table_names = BTreeSet::new(); + for (index, migration) in registry.iter().enumerate() { + if migration.version != expected_version { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "expected migration version {expected_version}, found {}", + migration.version + ), + }); + } + if migration.name.is_empty() + || registry[..index] + .iter() + .any(|prior| prior.name == migration.name) + { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "migration version {} has an invalid or duplicate name", + migration.version + ), + }); + } + if migration.owned_object_names.is_empty() || migration.owned_table_names.is_empty() { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "migration version {} must own schema objects and tables", + migration.version + ), + }); + } + for name in migration.owned_object_names { + validate_owned_schema_name(migration.version, "object", name)?; + if !object_names.insert(*name) { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!("owned schema object `{name}` is declared more than once"), + }); + } + } + for name in migration.owned_table_names { + validate_owned_schema_name(migration.version, "table", name)?; + if !migration.owned_object_names.contains(name) || !table_names.insert(*name) { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "owned table `{name}` is missing from the object inventory or is duplicated" + ), + }); + } + } + validate_embedded_migration_input( + migration.version, + "up", + migration.up_sql, + migration.up_len, + migration.up_sha256, + )?; + validate_embedded_migration_input( + migration.version, + "down", + migration.down_sql, + migration.down_len, + migration.down_sha256, + )?; + validate_sha256_literal(migration.version, "schema", migration.schema_sha256)?; + expected_version = expected_version.checked_add(1).ok_or_else(|| { + RadrootsOutboxError::MigrationRegistryDefect { + reason: "migration version overflow".to_owned(), + } + })?; + } + if expected_version - 1 != current { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "migration registry ends at {}, expected {current}", + expected_version - 1 + ), + }); + } + Ok(()) +} + +fn validate_ledger_ddl_identity( + catalog_ddl: &str, + create_ddl: &str, +) -> Result<(), RadrootsOutboxError> { + let catalog_ddl = catalog_ddl.strip_prefix("CREATE TABLE "); + let create_ddl = create_ddl.strip_prefix("CREATE TABLE main."); + if catalog_ddl.is_none() || create_ddl.is_none() || create_ddl != catalog_ddl { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: "main-qualified ledger creation DDL does not match canonical catalog DDL" + .to_owned(), + }); + } + Ok(()) +} + +fn validate_owned_schema_name( + version: u32, + object_kind: &'static str, + name: &str, +) -> Result<(), RadrootsOutboxError> { + if name.is_empty() + || name == OUTBOX_LEDGER_NAME + || !sqlite_identifier_starts_with(name, OUTBOX_RESERVED_PREFIX) + || !name + .bytes() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') + { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!( + "migration version {version} has invalid owned {object_kind} name `{name}`" + ), + }); + } + Ok(()) +} + +fn validate_embedded_migration_input( + version: u32, + direction: &'static str, + sql: &str, + expected_len: usize, + expected_sha256: &'static str, +) -> Result<(), RadrootsOutboxError> { + if sql.len() != expected_len { + return Err(RadrootsOutboxError::EmbeddedMigrationLengthMismatch { + version, + direction, + expected: expected_len, + actual: sql.len(), + }); + } + validate_sha256_literal(version, direction, expected_sha256)?; + let actual = sha256_hex(sql.as_bytes()); + if actual != expected_sha256 { + return Err(RadrootsOutboxError::EmbeddedMigrationChecksumMismatch { + version, + direction, + expected: expected_sha256, + actual, + }); + } + Ok(()) +} + +fn validate_sha256_literal( + version: u32, + field: &'static str, + value: &str, +) -> Result<(), RadrootsOutboxError> { + if value.len() != 64 + || !value + .bytes() + .all(|byte| byte.is_ascii_digit() || (b'a'..=b'f').contains(&byte)) + { + return Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!("migration version {version} has an invalid {field} SHA-256 literal"), + }); + } + Ok(()) +} + +pub(crate) fn sha256_hex(bytes: &[u8]) -> String { + hex::encode(Sha256::digest(bytes)) +} + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +mod tests { + use super::*; + use std::fs; + use std::path::{Path, PathBuf}; + + #[derive(Debug, PartialEq, Eq)] + struct DiscoveredMigration { + version: u32, + name: String, + up: PathBuf, + down: PathBuf, + } + + fn discover(root: &Path) -> Result<Vec<DiscoveredMigration>, String> { + let directory = root.join("migrations"); + let mut entries = fs::read_dir(&directory) + .map_err(|error| format!("read {}: {error}", directory.display()))? + .collect::<Result<Vec<_>, _>>() + .map_err(|error| format!("read migration entry: {error}"))?; + entries.sort_by_key(|entry| entry.file_name()); + let mut pairs = + std::collections::BTreeMap::<(u32, String), (Option<PathBuf>, Option<PathBuf>)>::new(); + for entry in entries { + let file_type = entry.file_type().map_err(|error| error.to_string())?; + if !file_type.is_file() || file_type.is_symlink() { + return Err(format!( + "migration input must be a regular file: {}", + entry.path().display() + )); + } + let name = entry + .file_name() + .into_string() + .map_err(|_| "migration filename is not UTF-8".to_owned())?; + let (stem, direction) = if let Some(stem) = name.strip_suffix(".up.sql") { + (stem, "up") + } else if let Some(stem) = name.strip_suffix(".down.sql") { + (stem, "down") + } else { + return Err(format!("unknown migration file `{name}`")); + }; + let (version, migration_name) = stem + .split_once('_') + .ok_or_else(|| format!("invalid migration filename `{name}`"))?; + if version.len() != 4 + || !version.bytes().all(|byte| byte.is_ascii_digit()) + || migration_name.is_empty() + || !migration_name + .bytes() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || byte == b'_') + { + return Err(format!("invalid migration filename `{name}`")); + } + let version = version.parse::<u32>().map_err(|error| error.to_string())?; + let pair = pairs + .entry((version, migration_name.to_owned())) + .or_default(); + let slot = if direction == "up" { + &mut pair.0 + } else { + &mut pair.1 + }; + if slot.replace(entry.path()).is_some() { + return Err(format!( + "duplicate {direction} migration for version {version}" + )); + } + } + pairs + .into_iter() + .map(|((version, name), (up, down))| { + Ok(DiscoveredMigration { + version, + name, + up: up.ok_or_else(|| format!("migration {version} is missing up SQL"))?, + down: down.ok_or_else(|| format!("migration {version} is missing down SQL"))?, + }) + }) + .collect() + } + + #[test] + fn migration_source_discovery_is_exact_and_fail_closed() { + let root = Path::new(env!("CARGO_MANIFEST_DIR")); + let discovered = discover(root).expect("canonical migration discovery"); + assert_eq!(discovered.len(), OUTBOX_MIGRATIONS.len()); + assert_eq!(discovered[0].version, 1); + assert_eq!(discovered[0].name, "outbox"); + assert_eq!( + fs::read(&discovered[0].up).expect("up bytes"), + OUTBOX_MIGRATIONS[0].up_sql.as_bytes() + ); + assert_eq!( + fs::read(&discovered[0].down).expect("down bytes"), + OUTBOX_MIGRATIONS[0].down_sql.as_bytes() + ); + + let temp = tempfile::tempdir().expect("tempdir"); + fs::create_dir(temp.path().join("migrations")).expect("migrations"); + for (name, bytes) in [ + ("0001_outbox.up.sql", OUTBOX_MIGRATIONS[0].up_sql.as_bytes()), + ( + "0001_outbox.down.sql", + OUTBOX_MIGRATIONS[0].down_sql.as_bytes(), + ), + ] { + fs::write(temp.path().join("migrations").join(name), bytes).expect("fixture"); + } + fs::write( + temp.path().join("migrations/0002_unknown.txt"), + b"SELECT 1;", + ) + .expect("unknown"); + assert!( + discover(temp.path()) + .expect_err("unknown file") + .contains("unknown migration file") + ); + fs::remove_file(temp.path().join("migrations/0002_unknown.txt")).expect("remove"); + fs::remove_file(temp.path().join("migrations/0001_outbox.down.sql")).expect("remove down"); + assert!( + discover(temp.path()) + .expect_err("missing pair") + .contains("missing down SQL") + ); + } + + #[test] + fn embedded_registry_and_frozen_baseline_are_exact() { + validate_embedded_migration_registry().expect("registry"); + assert_eq!(OUTBOX_MIGRATIONS[0].up_len, 5_470); + assert_eq!(OUTBOX_MIGRATIONS[0].down_len, 159); + assert_eq!( + OUTBOX_MIGRATIONS[0].up_sha256, + "a7ee775d32c2b9f845961425362e1b1e558ce0d025f7d22dd58f118ba4dab4fa" + ); + assert_eq!( + OUTBOX_MIGRATIONS[0].down_sha256, + "5d56f978f9172dc5ecbc5043a6c286c75926974d8a2a9e44fffa7c134829af61" + ); + assert_eq!(OUTBOX_BASELINE_OBJECT_NAMES.len(), 13); + assert_eq!(OUTBOX_BASELINE_TABLE_NAMES.len(), 5); + } + + fn assert_registry_defect(result: Result<(), RadrootsOutboxError>) -> String { + match result.expect_err("registry defect") { + RadrootsOutboxError::MigrationRegistryDefect { reason } => reason, + other => panic!("unexpected error: {other}"), + } + } + + fn future_migration() -> OutboxMigration { + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.version = 2; + migration.name = "future"; + migration.owned_object_names = &["outbox_future"]; + migration.owned_table_names = &["outbox_future"]; + migration + } + + #[test] + fn namespace_predicates_cover_reserved_ledger_registry_and_unrelated_names() { + let mut legacy = OUTBOX_MIGRATIONS[0]; + legacy.owned_object_names = &["legacy_object"]; + legacy.owned_table_names = &["legacy_table"]; + let registry = [legacy]; + + assert!(migration_for_version(OUTBOX_MIGRATIONS, 1).is_some()); + assert!(migration_for_version(OUTBOX_MIGRATIONS, 2).is_none()); + assert!(sqlite_identifier_starts_with("OUTBOX_EVENT", "outbox_")); + assert!(!sqlite_identifier_starts_with("short", "outbox_")); + assert!(is_outbox_owned_table_name(&registry, "outbox_new")); + assert!(is_outbox_owned_table_name(&registry, "LEGACY_TABLE")); + assert!(!is_outbox_owned_table_name(&registry, "caller_table")); + assert!(is_outbox_governed_schema_name( + &registry, + "RADROOTS_OUTBOX_SCHEMA_MIGRATIONS" + )); + assert!(is_outbox_governed_schema_name(&registry, "outbox_new")); + assert!(is_outbox_governed_schema_name(&registry, "LEGACY_OBJECT")); + assert!(!is_outbox_governed_schema_name(&registry, "caller_object")); + } + + #[test] + fn ledger_identifiers_fail_closed() { + validate_ledger_ddl_identity(OUTBOX_LEDGER_DDL, OUTBOX_LEDGER_CREATE_DDL) + .expect("ledger DDL"); + for (catalog, create) in [ + (OUTBOX_LEDGER_DDL, "CREATE TABLE main.counterfeit"), + ("counterfeit", OUTBOX_LEDGER_CREATE_DDL), + (OUTBOX_LEDGER_DDL, "counterfeit"), + ("counterfeit", "counterfeit"), + ] { + assert_registry_defect(validate_ledger_ddl_identity(catalog, create)); + } + } + + #[test] + fn registry_shape_validation_rejects_every_structural_defect() { + assert_registry_defect(validate_migration_registry(OUTBOX_MIGRATIONS, 0, 1)); + assert_registry_defect(validate_migration_registry(OUTBOX_MIGRATIONS, 2, 1)); + assert_registry_defect(validate_migration_registry(&[], 1, 1)); + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.version = 2; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.name = ""; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + + let mut duplicate_name = future_migration(); + duplicate_name.name = OUTBOX_MIGRATIONS[0].name; + assert_registry_defect(validate_migration_registry( + &[OUTBOX_MIGRATIONS[0], duplicate_name], + 1, + 2, + )); + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.owned_object_names = &[]; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.owned_table_names = &[]; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + + for invalid_name in ["", OUTBOX_LEDGER_NAME, "caller_object", "outbox_UPPER"] { + assert_registry_defect(validate_owned_schema_name(1, "object", invalid_name)); + } + validate_owned_schema_name(1, "object", "outbox_123").expect("numeric identifier"); + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.owned_object_names = &["outbox_valid"]; + migration.owned_table_names = &["outbox_UPPER"]; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + + let mut duplicate_object = future_migration(); + duplicate_object.owned_object_names = &["outbox_event"]; + duplicate_object.owned_table_names = &["outbox_event"]; + assert_registry_defect(validate_migration_registry( + &[OUTBOX_MIGRATIONS[0], duplicate_object], + 1, + 2, + )); + + let mut missing_table = future_migration(); + missing_table.owned_object_names = &["outbox_future_index"]; + assert_registry_defect(validate_migration_registry( + &[OUTBOX_MIGRATIONS[0], missing_table], + 1, + 2, + )); + + let mut duplicate_table = future_migration(); + duplicate_table.owned_table_names = &["outbox_future", "outbox_future"]; + assert_registry_defect(validate_migration_registry( + &[OUTBOX_MIGRATIONS[0], duplicate_table], + 1, + 2, + )); + + assert_registry_defect(validate_migration_registry(OUTBOX_MIGRATIONS, 1, 2)); + validate_migration_registry(&[OUTBOX_MIGRATIONS[0], future_migration()], 1, 2) + .expect("contiguous synthetic registry"); + + let mut overflow = OUTBOX_MIGRATIONS[0]; + overflow.version = u32::MAX; + validate_migration_registry(&[overflow], u32::MAX, u32::MAX).expect_err("version overflow"); + } + + #[test] + fn registry_checksum_validation_rejects_lengths_literals_and_bytes() { + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.up_len += 1; + assert!(matches!( + validate_migration_registry(&[migration], 1, 1), + Err(RadrootsOutboxError::EmbeddedMigrationLengthMismatch { + direction: "up", + .. + }) + )); + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.down_len += 1; + assert!(matches!( + validate_migration_registry(&[migration], 1, 1), + Err(RadrootsOutboxError::EmbeddedMigrationLengthMismatch { + direction: "down", + .. + }) + )); + + for invalid_sha in [ + "short", + "gggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggggg", + ] { + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.up_sha256 = invalid_sha; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + } + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.up_sha256 = "bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb"; + assert!(matches!( + validate_migration_registry(&[migration], 1, 1), + Err(RadrootsOutboxError::EmbeddedMigrationChecksumMismatch { + direction: "up", + .. + }) + )); + + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.schema_sha256 = "short"; + assert_registry_defect(validate_migration_registry(&[migration], 1, 1)); + } +} diff --git a/crates/outbox/src/schema.rs b/crates/outbox/src/schema.rs @@ -0,0 +1,1420 @@ +#![forbid(unsafe_code)] + +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, +}; +use sha2::{Digest, Sha256}; +use sqlx::{Row, Sqlite, SqliteConnection, SqlitePool, Transaction}; +use std::collections::{BTreeMap, BTreeSet}; + +#[cfg(test)] +const EMPTY_SCHEMA_SHA256: &str = + "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855"; + +#[derive(Clone, Debug, PartialEq, Eq)] +#[non_exhaustive] +/// Authenticated lifecycle state of the governed outbox schema. +pub enum RadrootsOutboxSchemaStatus { + /// No governed outbox schema objects exist. + Uninitialized, + /// The exact frozen baseline exists without its migration ledger. + UnledgeredBaseline, + /// The schema and ledger match a supported managed version. + Managed { version: u32 }, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +struct CatalogRow { + object_type: String, + name: String, + table_name: String, + sql: Option<String>, +} + +#[derive(Clone, Debug, PartialEq, Eq)] +struct AppliedMigration { + version: i64, + name: String, + up_sha256: String, + down_sha256: String, + schema_sha256: String, +} + +/// Inspects and authenticates schema state without adopting or migrating it. +pub async fn inspect_outbox_schema_status( + pool: &SqlitePool, +) -> Result<RadrootsOutboxSchemaStatus, RadrootsOutboxError> { + validate_embedded_migration_registry()?; + let mut transaction = pool.begin().await?; + let result = inspect_schema_on_connection( + &mut transaction, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) + .await; + finish_schema_transaction(transaction, result).await +} + +pub(crate) async fn migrate_outbox_schema(pool: &SqlitePool) -> Result<(), RadrootsOutboxError> { + validate_embedded_migration_registry()?; + migrate_outbox_schema_with_registry( + pool, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) + .await +} + +async fn migrate_outbox_schema_with_registry( + pool: &SqlitePool, + registry: &[OutboxMigration], + minimum: u32, + supported_current: u32, +) -> Result<(), RadrootsOutboxError> { + validate_migration_registry(registry, minimum, supported_current)?; + let mut transaction = pool.begin_with("BEGIN IMMEDIATE").await?; + let result = migrate_schema_on_connection(&mut transaction, registry, supported_current).await; + finish_schema_transaction(transaction, result).await +} + +#[cfg(test)] +async fn rollback_outbox_schema_offline( + pool: &SqlitePool, + target: u32, +) -> Result<(), RadrootsOutboxError> { + if target < RADROOTS_OUTBOX_SCHEMA_VERSION_MIN { + return Err(RadrootsOutboxError::RollbackBelowVersionFloor { + floor: RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + target, + }); + } + validate_embedded_migration_registry()?; + let mut transaction = pool.begin_with("BEGIN EXCLUSIVE").await?; + let result = rollback_schema_on_connection( + &mut transaction, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + target, + ) + .await; + finish_schema_transaction(transaction, result).await +} + +#[cfg(test)] +pub(crate) async fn destroy_outbox_schema_for_migration_test( + pool: &SqlitePool, +) -> Result<(), RadrootsOutboxError> { + validate_embedded_migration_registry()?; + let mut transaction = pool.begin_with("BEGIN EXCLUSIVE").await?; + let result = destroy_schema_on_connection(&mut transaction).await; + finish_schema_transaction(transaction, result).await +} + +async fn finish_schema_transaction<T>( + transaction: Transaction<'static, Sqlite>, + result: Result<T, RadrootsOutboxError>, +) -> Result<T, RadrootsOutboxError> { + match result { + Ok(value) => { + transaction.commit().await?; + Ok(value) + } + Err(primary) => match transaction.rollback().await { + Ok(()) => Err(primary), + Err(rollback) => Err(RadrootsOutboxError::MigrationTransactionRollbackFailed { + primary: Box::new(primary), + rollback, + }), + }, + } +} + +async fn migrate_schema_on_connection( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + supported_current: u32, +) -> Result<(), RadrootsOutboxError> { + let status = inspect_schema_on_connection(connection, registry, supported_current).await?; + let current_version = match status { + RadrootsOutboxSchemaStatus::Uninitialized => { + apply_migration_up(connection, registry, &registry[0]).await?; + create_ledger(connection, registry).await?; + insert_ledger_row(connection, &registry[0]).await?; + registry[0].version + } + RadrootsOutboxSchemaStatus::UnledgeredBaseline => { + create_ledger(connection, registry).await?; + insert_ledger_row(connection, &registry[0]).await?; + registry[0].version + } + RadrootsOutboxSchemaStatus::Managed { version } if version == supported_current => version, + RadrootsOutboxSchemaStatus::Managed { version } => version, + }; + + for migration in registry + .iter() + .filter(|migration| migration.version > current_version) + { + apply_migration_up(connection, registry, migration).await?; + insert_ledger_row(connection, migration).await?; + } + + match inspect_schema_on_connection(connection, registry, supported_current).await? { + RadrootsOutboxSchemaStatus::Managed { version } if version == supported_current => Ok(()), + status => Err(RadrootsOutboxError::MigrationRegistryDefect { + reason: format!("migration completed in unexpected state {status:?}"), + }), + } +} + +#[cfg(test)] +async fn rollback_schema_on_connection( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + supported_current: u32, + target: u32, +) -> Result<(), RadrootsOutboxError> { + let RadrootsOutboxSchemaStatus::Managed { + version: current_version, + } = inspect_schema_on_connection(connection, registry, supported_current).await? + else { + return Err(RadrootsOutboxError::RollbackUnmanaged); + }; + if target > current_version { + return Err(RadrootsOutboxError::RollbackAhead { + current: current_version, + target, + }); + } + + for version in ((target + 1)..=current_version).rev() { + let migration = migration_for_version(registry, version) + .ok_or(RadrootsOutboxError::UnknownMigration { version })?; + apply_migration_down(connection, registry, migration).await?; + let prior = migration_for_version(registry, version - 1).ok_or( + RadrootsOutboxError::MigrationHistoryGap { + expected: version - 1, + actual: None, + }, + )?; + validate_schema_fingerprint(connection, registry, prior).await?; + let deleted = + sqlx::query("DELETE FROM main.radroots_outbox_schema_migrations WHERE version = ?") + .bind(i64::from(version)) + .execute(&mut *connection) + .await?; + if deleted.rows_affected() != 1 { + return Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: format!( + "rollback expected one ledger row for version {version}, deleted {}", + deleted.rows_affected() + ), + }); + } + } + + match inspect_schema_on_connection(connection, registry, supported_current).await? { + RadrootsOutboxSchemaStatus::Managed { version } if version == target => Ok(()), + status => Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: format!("rollback completed in unexpected state {status:?}"), + }), + } +} + +#[cfg(test)] +async fn destroy_schema_on_connection( + connection: &mut SqliteConnection, +) -> Result<(), RadrootsOutboxError> { + match inspect_schema_on_connection( + connection, + OUTBOX_MIGRATIONS, + RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + ) + .await? + { + RadrootsOutboxSchemaStatus::Managed { version } => { + for version in (RADROOTS_OUTBOX_SCHEMA_VERSION_MIN..=version).rev() { + let migration = migration_for_version(OUTBOX_MIGRATIONS, version) + .ok_or(RadrootsOutboxError::UnknownMigration { version })?; + apply_migration_down(connection, OUTBOX_MIGRATIONS, migration).await?; + let deleted = sqlx::query( + "DELETE FROM main.radroots_outbox_schema_migrations WHERE version = ?", + ) + .bind(i64::from(version)) + .execute(&mut *connection) + .await?; + if deleted.rows_affected() != 1 { + return Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: format!( + "test destruction expected one ledger row for version {version}, deleted {}", + deleted.rows_affected() + ), + }); + } + } + validate_empty_governed_catalog(connection, OUTBOX_MIGRATIONS).await?; + sqlx::query("DROP TABLE main.radroots_outbox_schema_migrations") + .execute(&mut *connection) + .await?; + } + RadrootsOutboxSchemaStatus::UnledgeredBaseline => { + apply_migration_down(connection, OUTBOX_MIGRATIONS, &OUTBOX_MIGRATIONS[0]).await?; + validate_empty_governed_catalog(connection, OUTBOX_MIGRATIONS).await?; + } + RadrootsOutboxSchemaStatus::Uninitialized => {} + } + Ok(()) +} + +async fn apply_migration_up( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + migration: &OutboxMigration, +) -> Result<(), RadrootsOutboxError> { + let before = read_catalog_bounded(connection, registry).await?; + sqlx::raw_sql(migration.up_sql) + .execute(&mut *connection) + .await?; + let after = read_catalog_bounded(connection, registry).await?; + validate_catalog_delta(&before, &after, migration, "up")?; + validate_schema_fingerprint(connection, registry, migration).await +} + +#[cfg(test)] +async fn apply_migration_down( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + migration: &OutboxMigration, +) -> Result<(), RadrootsOutboxError> { + let before = read_catalog_bounded(connection, registry).await?; + sqlx::raw_sql(migration.down_sql) + .execute(&mut *connection) + .await?; + let after = read_catalog_bounded(connection, registry).await?; + validate_catalog_delta(&before, &after, migration, "down") +} + +fn validate_catalog_delta( + before: &[CatalogRow], + after: &[CatalogRow], + migration: &OutboxMigration, + direction: &'static str, +) -> Result<(), RadrootsOutboxError> { + let before = before + .iter() + .map(|row| (row.name.as_str(), row)) + .collect::<BTreeMap<_, _>>(); + let after = after + .iter() + .map(|row| (row.name.as_str(), row)) + .collect::<BTreeMap<_, _>>(); + let added = after + .keys() + .filter(|name| !before.contains_key(**name)) + .copied() + .collect::<BTreeSet<_>>(); + let removed = before + .keys() + .filter(|name| !after.contains_key(**name)) + .copied() + .collect::<BTreeSet<_>>(); + let changed = before + .iter() + .filter_map(|(name, row)| { + after + .get(name) + .filter(|after_row| *after_row != row) + .map(|_| *name) + }) + .collect::<BTreeSet<_>>(); + let expected = migration + .owned_object_names + .iter() + .copied() + .collect::<BTreeSet<_>>(); + let valid = match direction { + "up" => added == expected && removed.is_empty() && changed.is_empty(), + "down" => removed == expected && added.is_empty() && changed.is_empty(), + _ => false, + }; + if !valid { + return Err(RadrootsOutboxError::MigrationCatalogDeltaMismatch { + version: migration.version, + direction, + reason: format!( + "expected {expected:?}; added {added:?}, removed {removed:?}, changed {changed:?}" + ), + }); + } + Ok(()) +} + +async fn create_ledger( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], +) -> Result<(), RadrootsOutboxError> { + sqlx::query(OUTBOX_LEDGER_CREATE_DDL) + .execute(&mut *connection) + .await?; + validate_ledger_catalog(&read_catalog_bounded(connection, registry).await?)?; + Ok(()) +} + +async fn insert_ledger_row( + connection: &mut SqliteConnection, + migration: &OutboxMigration, +) -> Result<(), RadrootsOutboxError> { + sqlx::query( + "INSERT INTO main.radroots_outbox_schema_migrations(version, name, up_sha256, down_sha256, schema_sha256) VALUES (?, ?, ?, ?, ?)", + ) + .bind(i64::from(migration.version)) + .bind(migration.name) + .bind(migration.up_sha256) + .bind(migration.down_sha256) + .bind(migration.schema_sha256) + .execute(&mut *connection) + .await?; + Ok(()) +} + +async fn inspect_schema_on_connection( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + supported_current: u32, +) -> Result<RadrootsOutboxSchemaStatus, RadrootsOutboxError> { + validate_outbox_temp_schema_with_registry(connection, registry).await?; + let catalog = read_catalog_bounded(connection, registry).await?; + let has_ledger = validate_ledger_catalog(&catalog)?; + let governed = governed_catalog(&catalog, registry); + let actual_schema_sha256 = catalog_fingerprint(&governed); + + if !has_ledger { + if governed.is_empty() { + return Ok(RadrootsOutboxSchemaStatus::Uninitialized); + } + let baseline = &registry[0]; + if governed.len() == baseline.owned_object_names.len() + && actual_schema_sha256 == baseline.schema_sha256 + { + return Ok(RadrootsOutboxSchemaStatus::UnledgeredBaseline); + } + return Err(RadrootsOutboxError::UnmanagedSchema { + actual_schema_sha256, + }); + } + + let history = read_history_bounded(connection, supported_current).await?; + let current = validate_history_against_registry(&history, registry, supported_current)?; + let expected = migration_for_version(registry, current) + .ok_or(RadrootsOutboxError::UnknownMigration { version: current })?; + if actual_schema_sha256 != expected.schema_sha256 { + return Err(RadrootsOutboxError::SchemaFingerprintMismatch { + version: current, + expected: expected.schema_sha256, + actual: actual_schema_sha256, + }); + } + Ok(RadrootsOutboxSchemaStatus::Managed { version: current }) +} + +async fn validate_outbox_temp_schema_with_registry( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], +) -> Result<(), RadrootsOutboxError> { + let collision = sqlx::query( + "SELECT type, name, tbl_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(?) + OR name = ? COLLATE NOCASE + OR tbl_name = ? COLLATE NOCASE + ORDER BY type, name, tbl_name LIMIT 1", + ) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_LEDGER_NAME) + .bind(OUTBOX_LEDGER_NAME) + .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")?; + if matches!(object_type.as_str(), "trigger" | "view") + || is_outbox_governed_schema_name(registry, &name) + || is_outbox_governed_schema_name(registry, &table_name) + { + return Err(RadrootsOutboxError::TemporarySchemaCollision { + object_type, + name, + table_name, + }); + } + } + Ok(()) +} + +fn catalog_row_limit(registry: &[OutboxMigration]) -> Result<i64, RadrootsOutboxError> { + let max = registry + .iter() + .flat_map(|migration| migration.owned_object_names.iter().copied()) + .collect::<BTreeSet<_>>() + .len() + .checked_add(1) + .ok_or_else(|| RadrootsOutboxError::MigrationRegistryDefect { + reason: "governed catalog object limit overflow".to_owned(), + })?; + i64::try_from(max.checked_add(1).ok_or_else(|| { + RadrootsOutboxError::MigrationRegistryDefect { + reason: "governed catalog collision limit overflow".to_owned(), + } + })?) + .map_err(|_| RadrootsOutboxError::MigrationRegistryDefect { + reason: "governed catalog object limit is outside SQLite range".to_owned(), + }) +} + +async fn read_catalog_bounded( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], +) -> 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 + WHERE lower(substr(name, 1, 7)) != 'sqlite_' + AND (lower(substr(name, 1, length(?))) = lower(?) + OR lower(substr(tbl_name, 1, length(?))) = lower(?) + OR name = ? COLLATE NOCASE + OR tbl_name = ? COLLATE NOCASE) + ORDER BY type, name, tbl_name LIMIT ?", + ) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_LEDGER_NAME) + .bind(OUTBOX_LEDGER_NAME) + .bind(row_limit) + .fetch_all(&mut *connection) + .await?; + if i64::try_from(rows.len()).ok() == Some(row_limit) { + return Err(RadrootsOutboxError::GovernedCatalogCapacityExceeded { + max: usize::try_from(row_limit - 1).unwrap_or(usize::MAX), + }); + } + 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")?, + }) + }) + .collect() +} + +fn validate_ledger_catalog(catalog: &[CatalogRow]) -> Result<bool, RadrootsOutboxError> { + let rows = catalog + .iter() + .filter(|row| { + row.name.eq_ignore_ascii_case(OUTBOX_LEDGER_NAME) + || row.table_name.eq_ignore_ascii_case(OUTBOX_LEDGER_NAME) + }) + .collect::<Vec<_>>(); + if rows.is_empty() { + return Ok(false); + } + if rows.len() != 1 { + return Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: format!( + "expected exactly one non-internal ledger catalog object, found {}", + rows.len() + ), + }); + } + let row = rows[0]; + if row.object_type != "table" + || row.name != OUTBOX_LEDGER_NAME + || row.table_name != OUTBOX_LEDGER_NAME + || row.sql.as_deref() != Some(OUTBOX_LEDGER_DDL) + { + return Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: "ledger table definition does not match canonical catalog SQL".to_owned(), + }); + } + Ok(true) +} + +fn governed_catalog(catalog: &[CatalogRow], registry: &[OutboxMigration]) -> Vec<CatalogRow> { + catalog + .iter() + .filter(|row| !row.name.eq_ignore_ascii_case(OUTBOX_LEDGER_NAME)) + .filter(|row| { + is_outbox_governed_schema_name(registry, &row.name) + || is_outbox_governed_schema_name(registry, &row.table_name) + }) + .cloned() + .collect() +} + +fn catalog_fingerprint(catalog: &[CatalogRow]) -> String { + let mut rows = catalog.to_vec(); + rows.sort_by(|left, right| { + ( + left.object_type.as_bytes(), + left.name.as_bytes(), + left.table_name.as_bytes(), + left.sql.as_deref().unwrap_or("").as_bytes(), + ) + .cmp(&( + right.object_type.as_bytes(), + right.name.as_bytes(), + right.table_name.as_bytes(), + right.sql.as_deref().unwrap_or("").as_bytes(), + )) + }); + let mut digest = Sha256::new(); + for row in rows { + for field in [ + row.object_type.as_str(), + row.name.as_str(), + row.table_name.as_str(), + row.sql.as_deref().unwrap_or(""), + ] { + digest.update(field.as_bytes()); + digest.update([0]); + } + } + hex::encode(digest.finalize()) +} + +async fn read_history_bounded( + connection: &mut SqliteConnection, + supported_current: u32, +) -> Result<Vec<AppliedMigration>, RadrootsOutboxError> { + let row_limit = i64::from(supported_current).checked_add(1).ok_or_else(|| { + RadrootsOutboxError::MigrationRegistryDefect { + reason: "migration history row limit overflow".to_owned(), + } + })?; + let rows = sqlx::query( + "SELECT version, name, up_sha256, down_sha256, schema_sha256 + FROM main.radroots_outbox_schema_migrations + ORDER BY version LIMIT ?", + ) + .bind(row_limit) + .fetch_all(&mut *connection) + .await?; + rows.into_iter() + .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")?, + }) + }) + .collect() +} + +fn validate_history_against_registry( + history: &[AppliedMigration], + registry: &[OutboxMigration], + supported_current: u32, +) -> Result<u32, RadrootsOutboxError> { + if history.is_empty() { + return Err(RadrootsOutboxError::MigrationLedgerDrift { + reason: "ledger exists without migration history".to_owned(), + }); + } + let database_version = history + .iter() + .map(|row| row.version) + .max() + .unwrap_or_default(); + if database_version > i64::from(supported_current) { + return Err(RadrootsOutboxError::SchemaTooNew { + current: supported_current, + database: database_version, + }); + } + let mut expected_version = registry[0].version; + for row in history { + let version = + u32::try_from(row.version).map_err(|_| RadrootsOutboxError::MigrationLedgerDrift { + reason: format!( + "ledger version {} is outside the positive range", + row.version + ), + })?; + if version != expected_version { + return Err(RadrootsOutboxError::MigrationHistoryGap { + expected: expected_version, + actual: Some(version), + }); + } + let migration = migration_for_version(registry, version) + .ok_or(RadrootsOutboxError::UnknownMigration { version })?; + if row.name != migration.name { + return Err(RadrootsOutboxError::MigrationHistoryNameDrift { + version, + expected: migration.name, + actual: row.name.clone(), + }); + } + validate_history_checksum(version, "up_sha256", &row.up_sha256, migration.up_sha256)?; + validate_history_checksum( + version, + "down_sha256", + &row.down_sha256, + migration.down_sha256, + )?; + validate_history_checksum( + version, + "schema_sha256", + &row.schema_sha256, + migration.schema_sha256, + )?; + expected_version = expected_version.checked_add(1).ok_or_else(|| { + RadrootsOutboxError::MigrationLedgerDrift { + reason: "migration history version overflow".to_owned(), + } + })?; + } + Ok(expected_version - 1) +} + +fn validate_history_checksum( + version: u32, + field: &'static str, + actual: &str, + expected: &'static str, +) -> Result<(), RadrootsOutboxError> { + if actual != expected { + return Err(RadrootsOutboxError::MigrationHistoryChecksumDrift { + version, + field, + expected, + actual: actual.to_owned(), + }); + } + Ok(()) +} + +async fn validate_schema_fingerprint( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], + migration: &OutboxMigration, +) -> Result<(), RadrootsOutboxError> { + let catalog = read_catalog_bounded(connection, registry).await?; + let actual = catalog_fingerprint(&governed_catalog(&catalog, registry)); + if actual != migration.schema_sha256 { + return Err(RadrootsOutboxError::SchemaFingerprintMismatch { + version: migration.version, + expected: migration.schema_sha256, + actual, + }); + } + Ok(()) +} + +#[cfg(test)] +async fn validate_empty_governed_catalog( + connection: &mut SqliteConnection, + registry: &[OutboxMigration], +) -> Result<(), RadrootsOutboxError> { + let catalog = read_catalog_bounded(connection, registry).await?; + let actual = catalog_fingerprint(&governed_catalog(&catalog, registry)); + if actual != EMPTY_SCHEMA_SHA256 { + return Err(RadrootsOutboxError::SchemaFingerprintMismatch { + version: 0, + expected: EMPTY_SCHEMA_SHA256, + actual, + }); + } + Ok(()) +} + +#[cfg(test)] +#[cfg_attr(coverage_nightly, coverage(off))] +mod tests { + use super::*; + use sqlx::sqlite::{SqliteConnectOptions, SqlitePoolOptions}; + use std::str::FromStr; + + async fn memory_pool() -> SqlitePool { + SqlitePoolOptions::new() + .max_connections(1) + .connect_with(SqliteConnectOptions::from_str("sqlite::memory:").expect("options")) + .await + .expect("pool") + } + + #[tokio::test] + async fn fresh_migration_reaches_the_authenticated_schema_fingerprint() { + let pool = memory_pool().await; + migrate_outbox_schema(&pool).await.expect("fresh migration"); + assert_eq!( + inspect_outbox_schema_status(&pool) + .await + .expect("managed schema"), + RadrootsOutboxSchemaStatus::Managed { version: 1 } + ); + } + + async fn apply_unledgered_baseline(pool: &SqlitePool) { + sqlx::raw_sql(OUTBOX_MIGRATIONS[0].up_sql) + .execute(pool) + .await + .expect("unledgered baseline"); + } + + async fn ledger_count(pool: &SqlitePool) -> i64 { + sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE type = 'table' AND name = ?", + ) + .bind(OUTBOX_LEDGER_NAME) + .fetch_one(pool) + .await + .expect("ledger count") + } + + async fn synthetic_v2_registry(pool: &SqlitePool) -> [OutboxMigration; 2] { + const UP_SQL: &str = "CREATE TABLE outbox_future (value TEXT NOT NULL) STRICT;\n"; + const DOWN_SQL: &str = "DROP TABLE outbox_future;\n"; + + migrate_outbox_schema(pool).await.expect("version 1"); + sqlx::raw_sql(UP_SQL) + .execute(pool) + .await + .expect("synthetic version 2 schema"); + let catalog = sqlx::query( + "SELECT type, name, tbl_name, 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(?))", + ) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .bind(OUTBOX_RESERVED_PREFIX) + .fetch_all(pool) + .await + .expect("synthetic version 2 catalog") + .into_iter() + .map(|row| CatalogRow { + object_type: row.try_get("type").expect("catalog type"), + name: row.try_get("name").expect("catalog name"), + table_name: row.try_get("tbl_name").expect("catalog table name"), + sql: row.try_get("sql").expect("catalog SQL"), + }) + .collect::<Vec<_>>(); + let schema_sha256 = Box::leak(catalog_fingerprint(&catalog).into_boxed_str()); + sqlx::raw_sql(DOWN_SQL) + .execute(pool) + .await + .expect("remove synthetic version 2 schema"); + + let up_sha256 = + Box::leak(crate::migrations::sha256_hex(UP_SQL.as_bytes()).into_boxed_str()); + let down_sha256 = + Box::leak(crate::migrations::sha256_hex(DOWN_SQL.as_bytes()).into_boxed_str()); + [ + OUTBOX_MIGRATIONS[0], + OutboxMigration { + version: 2, + name: "future", + up_sql: UP_SQL, + down_sql: DOWN_SQL, + up_len: UP_SQL.len(), + down_len: DOWN_SQL.len(), + up_sha256, + down_sha256, + schema_sha256, + owned_object_names: &["outbox_future"], + owned_table_names: &["outbox_future"], + }, + ] + } + + async fn inspect_with_registry( + pool: &SqlitePool, + registry: &[OutboxMigration], + supported_current: u32, + ) -> Result<RadrootsOutboxSchemaStatus, RadrootsOutboxError> { + let mut transaction = pool.begin().await?; + let result = + inspect_schema_on_connection(&mut transaction, registry, supported_current).await; + finish_schema_transaction(transaction, result).await + } + + #[tokio::test] + async fn additive_future_migration_advances_and_rolls_back_to_an_explicit_target() { + let pool = memory_pool().await; + sqlx::raw_sql( + "CREATE TABLE caller_state (value TEXT NOT NULL); + INSERT INTO caller_state(value) VALUES ('preserved');", + ) + .execute(&pool) + .await + .expect("caller state"); + let registry = synthetic_v2_registry(&pool).await; + + migrate_outbox_schema_with_registry(&pool, &registry, 1, 2) + .await + .expect("advance to version 2"); + assert_eq!( + inspect_with_registry(&pool, &registry, 2) + .await + .expect("managed version 2"), + RadrootsOutboxSchemaStatus::Managed { version: 2 } + ); + let future_objects: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", + ) + .fetch_one(&pool) + .await + .expect("future object count"); + assert_eq!(future_objects, 1); + + let mut transaction = pool + .begin_with("BEGIN EXCLUSIVE") + .await + .expect("rollback transaction"); + let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 1).await; + finish_schema_transaction(transaction, result) + .await + .expect("rollback to version 1"); + assert_eq!( + inspect_with_registry(&pool, &registry, 2) + .await + .expect("managed version 1"), + RadrootsOutboxSchemaStatus::Managed { version: 1 } + ); + let future_objects: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", + ) + .fetch_one(&pool) + .await + .expect("rolled-back object count"); + assert_eq!(future_objects, 0); + let caller_value: String = sqlx::query_scalar("SELECT value FROM caller_state") + .fetch_one(&pool) + .await + .expect("preserved caller state"); + assert_eq!(caller_value, "preserved"); + + let mut transaction = pool + .begin_with("BEGIN EXCLUSIVE") + .await + .expect("ahead transaction"); + let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 2).await; + assert!(matches!( + result, + Err(RadrootsOutboxError::RollbackAhead { .. }) + )); + transaction + .rollback() + .await + .expect("rollback ahead fixture"); + + let unmanaged = memory_pool().await; + let mut transaction = unmanaged + .begin_with("BEGIN EXCLUSIVE") + .await + .expect("unmanaged transaction"); + let result = rollback_schema_on_connection(&mut transaction, &registry, 2, 1).await; + assert!(matches!( + result, + Err(RadrootsOutboxError::RollbackUnmanaged) + )); + transaction + .rollback() + .await + .expect("rollback unmanaged fixture"); + } + + #[tokio::test] + async fn migration_post_up_authority_failure_rolls_back_catalog_ledger_and_caller_state() { + let pool = memory_pool().await; + sqlx::raw_sql( + "CREATE TABLE caller_state (key TEXT PRIMARY KEY, value TEXT NOT NULL); + INSERT INTO caller_state(key, value) VALUES ('victoria', 'preserved');", + ) + .execute(&pool) + .await + .expect("caller state"); + let mut registry = synthetic_v2_registry(&pool).await; + + let mut connection = pool.acquire().await.expect("catalog connection"); + let before_catalog = read_catalog_bounded(&mut connection, &registry) + .await + .expect("before catalog"); + let before_fingerprint = catalog_fingerprint(&governed_catalog(&before_catalog, &registry)); + let before_history = read_history_bounded(&mut connection, 2) + .await + .expect("before history"); + drop(connection); + + registry[1].owned_object_names = &["outbox_expected"]; + registry[1].owned_table_names = &["outbox_expected"]; + let error = migrate_outbox_schema_with_registry(&pool, &registry, 1, 2) + .await + .expect_err("post-UP catalog delta must fail"); + assert!(matches!( + error, + RadrootsOutboxError::MigrationCatalogDeltaMismatch { + version: 2, + direction: "up", + .. + } + )); + + let mut connection = pool.acquire().await.expect("verification connection"); + let after_catalog = read_catalog_bounded(&mut connection, &registry) + .await + .expect("after catalog"); + let after_fingerprint = catalog_fingerprint(&governed_catalog(&after_catalog, &registry)); + let after_history = read_history_bounded(&mut connection, 2) + .await + .expect("after history"); + drop(connection); + assert_eq!(after_fingerprint, before_fingerprint); + assert_eq!(after_history, before_history); + assert_eq!( + inspect_outbox_schema_status(&pool) + .await + .expect("restored managed schema"), + RadrootsOutboxSchemaStatus::Managed { version: 1 } + ); + let future_objects: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'outbox_future'", + ) + .fetch_one(&pool) + .await + .expect("rolled-back future object count"); + assert_eq!(future_objects, 0); + let caller_value: String = + sqlx::query_scalar("SELECT value FROM caller_state WHERE key = 'victoria'") + .fetch_one(&pool) + .await + .expect("preserved caller row"); + assert_eq!(caller_value, "preserved"); + } + + #[tokio::test] + async fn migration_exact_unledgered_baseline_is_adopted_without_replaying_or_losing_caller_state() + { + let pool = memory_pool().await; + sqlx::query("CREATE TABLE caller_state (key TEXT PRIMARY KEY, value TEXT NOT NULL)") + .execute(&pool) + .await + .expect("caller table"); + sqlx::query("INSERT INTO caller_state(key, value) VALUES ('victoria', 'preserved')") + .execute(&pool) + .await + .expect("caller row"); + apply_unledgered_baseline(&pool).await; + sqlx::query( + "INSERT INTO outbox_operations(operation_kind, expected_pubkey, semantic_scope, trade_id, mutation_id, canonical_payload_sha256, idempotency_key, operation_idempotency_digest, status, created_at_ms, updated_at_ms) + VALUES ('post', 'author', 'generic_event', NULL, NULL, NULL, NULL, ?, 'queued', 1, 1)", + ) + .bind("a".repeat(64)) + .execute(&pool) + .await + .expect("legacy row"); + + assert_eq!( + inspect_outbox_schema_status(&pool) + .await + .expect("unledgered status"), + RadrootsOutboxSchemaStatus::UnledgeredBaseline + ); + migrate_outbox_schema(&pool).await.expect("adoption"); + assert_eq!(ledger_count(&pool).await, 1); + let caller: String = + sqlx::query_scalar("SELECT value FROM caller_state WHERE key = 'victoria'") + .fetch_one(&pool) + .await + .expect("caller row preserved"); + assert_eq!(caller, "preserved"); + let operations: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM outbox_operations") + .fetch_one(&pool) + .await + .expect("legacy row count"); + assert_eq!(operations, 1); + } + + #[tokio::test] + async fn migration_test_only_destruction_handles_unledgered_and_uninitialized_schemas() { + let pool = memory_pool().await; + apply_unledgered_baseline(&pool).await; + destroy_outbox_schema_for_migration_test(&pool) + .await + .expect("destroy unledgered baseline"); + assert_eq!( + inspect_outbox_schema_status(&pool) + .await + .expect("uninitialized after destruction"), + RadrootsOutboxSchemaStatus::Uninitialized + ); + destroy_outbox_schema_for_migration_test(&pool) + .await + .expect("destroy uninitialized schema"); + } + + #[tokio::test] + async fn migration_fresh_initialization_preserves_unrelated_caller_schema_and_rows() { + let pool = memory_pool().await; + sqlx::raw_sql( + "CREATE TABLE caller_table (value TEXT NOT NULL); + INSERT INTO caller_table(value) VALUES ('keep'); + CREATE INDEX caller_table_value_idx ON caller_table(value);", + ) + .execute(&pool) + .await + .expect("caller schema"); + migrate_outbox_schema(&pool).await.expect("migration"); + let value: String = sqlx::query_scalar("SELECT value FROM caller_table") + .fetch_one(&pool) + .await + .expect("caller row"); + assert_eq!(value, "keep"); + let index_count: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE name = 'caller_table_value_idx'", + ) + .fetch_one(&pool) + .await + .expect("caller index"); + assert_eq!(index_count, 1); + } + + #[tokio::test] + async fn migration_partial_changed_and_unknown_unledgered_catalogs_fail_before_adoption() { + for mutation in [ + "DROP INDEX outbox_event_event_id_idx", + "DROP INDEX outbox_event_event_id_idx; CREATE INDEX outbox_event_event_id_idx ON outbox_event(expected_pubkey)", + "CREATE TABLE outbox_counterfeit (value TEXT NOT NULL)", + ] { + let pool = memory_pool().await; + apply_unledgered_baseline(&pool).await; + sqlx::raw_sql(mutation) + .execute(&pool) + .await + .expect("catalog mutation"); + assert!(matches!( + migrate_outbox_schema(&pool).await, + Err(RadrootsOutboxError::UnmanagedSchema { .. }) + )); + assert_eq!(ledger_count(&pool).await, 0); + } + } + + #[tokio::test] + async fn migration_counterfeit_ledger_shape_fails_before_any_schema_mutation() { + let pool = memory_pool().await; + sqlx::query("CREATE TABLE radroots_outbox_schema_migrations (version INTEGER)") + .execute(&pool) + .await + .expect("counterfeit ledger"); + assert!(matches!( + migrate_outbox_schema(&pool).await, + Err(RadrootsOutboxError::MigrationLedgerDrift { .. }) + )); + let outbox_objects: i64 = sqlx::query_scalar( + "SELECT COUNT(*) FROM main.sqlite_schema WHERE lower(substr(name, 1, 7)) = 'outbox_'", + ) + .fetch_one(&pool) + .await + .expect("outbox object count"); + assert_eq!(outbox_objects, 0); + } + + #[tokio::test] + async fn migration_ledger_name_checksum_and_catalog_mutations_fail_closed() { + for statement in [ + "UPDATE radroots_outbox_schema_migrations SET name = 'counterfeit' WHERE version = 1", + "UPDATE radroots_outbox_schema_migrations SET up_sha256 = 'bbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbbb' WHERE version = 1", + "UPDATE radroots_outbox_schema_migrations SET schema_sha256 = 'cccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccccc' WHERE version = 1", + "DROP INDEX outbox_event_event_id_idx; CREATE INDEX outbox_event_event_id_idx ON outbox_event(expected_pubkey)", + ] { + let pool = memory_pool().await; + migrate_outbox_schema(&pool).await.expect("migration"); + sqlx::raw_sql(statement) + .execute(&pool) + .await + .expect("managed mutation"); + assert!(migrate_outbox_schema(&pool).await.is_err(), "{statement}"); + let rows: i64 = + sqlx::query_scalar("SELECT COUNT(*) FROM radroots_outbox_schema_migrations") + .fetch_one(&pool) + .await + .expect("ledger rows"); + assert_eq!(rows, 1); + } + } + + #[tokio::test] + async fn migration_newer_history_and_governed_catalog_overflow_are_bounded_and_rejected() { + let newer = memory_pool().await; + migrate_outbox_schema(&newer).await.expect("migration"); + sqlx::query( + "INSERT INTO radroots_outbox_schema_migrations(version, name, up_sha256, down_sha256, schema_sha256) + VALUES (2, 'future', ?, ?, ?)", + ) + .bind("a".repeat(64)) + .bind("b".repeat(64)) + .bind("c".repeat(64)) + .execute(&newer) + .await + .expect("future history"); + assert!(matches!( + inspect_outbox_schema_status(&newer).await, + Err(RadrootsOutboxError::SchemaTooNew { + current: 1, + database: 2 + }) + )); + + let overflow = memory_pool().await; + migrate_outbox_schema(&overflow).await.expect("migration"); + sqlx::query("CREATE TABLE outbox_unknown (value TEXT)") + .execute(&overflow) + .await + .expect("unknown governed object"); + assert!(matches!( + inspect_outbox_schema_status(&overflow).await, + Err(RadrootsOutboxError::GovernedCatalogCapacityExceeded { max: 14 }) + )); + } + + #[test] + fn migration_history_validation_rejects_gaps_and_unknown_versions() { + let row = |version, migration: &OutboxMigration| AppliedMigration { + version, + name: migration.name.to_owned(), + up_sha256: migration.up_sha256.to_owned(), + down_sha256: migration.down_sha256.to_owned(), + schema_sha256: migration.schema_sha256.to_owned(), + }; + assert!(matches!( + validate_history_against_registry( + &[row(2, &OUTBOX_MIGRATIONS[0])], + OUTBOX_MIGRATIONS, + 3 + ), + Err(RadrootsOutboxError::MigrationHistoryGap { + expected: 1, + actual: Some(2) + }) + )); + assert!(matches!( + validate_history_against_registry( + &[row(1, &OUTBOX_MIGRATIONS[0]), row(2, &OUTBOX_MIGRATIONS[0])], + OUTBOX_MIGRATIONS, + 3, + ), + Err(RadrootsOutboxError::UnknownMigration { version: 2 }) + )); + + assert!(matches!( + validate_history_against_registry(&[], OUTBOX_MIGRATIONS, 1), + Err(RadrootsOutboxError::MigrationLedgerDrift { .. }) + )); + assert!(matches!( + validate_history_against_registry( + &[row(-1, &OUTBOX_MIGRATIONS[0])], + OUTBOX_MIGRATIONS, + 1 + ), + Err(RadrootsOutboxError::MigrationLedgerDrift { .. }) + )); + assert_eq!( + validate_history_against_registry( + &[row(1, &OUTBOX_MIGRATIONS[0])], + OUTBOX_MIGRATIONS, + 1 + ) + .expect("canonical history"), + 1 + ); + + let mut name_drift = row(1, &OUTBOX_MIGRATIONS[0]); + name_drift.name = "counterfeit".to_owned(); + assert!(matches!( + validate_history_against_registry(&[name_drift], OUTBOX_MIGRATIONS, 1), + Err(RadrootsOutboxError::MigrationHistoryNameDrift { .. }) + )); + for field in ["up_sha256", "down_sha256", "schema_sha256"] { + let mut checksum_drift = row(1, &OUTBOX_MIGRATIONS[0]); + match field { + "up_sha256" => checksum_drift.up_sha256 = "a".repeat(64), + "down_sha256" => checksum_drift.down_sha256 = "a".repeat(64), + "schema_sha256" => checksum_drift.schema_sha256 = "a".repeat(64), + _ => unreachable!(), + } + assert!(matches!( + validate_history_against_registry(&[checksum_drift], OUTBOX_MIGRATIONS, 1), + Err(RadrootsOutboxError::MigrationHistoryChecksumDrift { + field: actual_field, + .. + }) if actual_field == field + )); + } + } + + fn catalog_row( + object_type: &str, + name: &str, + table_name: &str, + sql: Option<&str>, + ) -> CatalogRow { + CatalogRow { + object_type: object_type.to_owned(), + name: name.to_owned(), + table_name: table_name.to_owned(), + sql: sql.map(str::to_owned), + } + } + + #[test] + fn migration_catalog_delta_and_ledger_validators_fail_closed() { + let mut migration = OUTBOX_MIGRATIONS[0]; + migration.version = 2; + migration.owned_object_names = &["outbox_future"]; + migration.owned_table_names = &["outbox_future"]; + let future = catalog_row( + "table", + "outbox_future", + "outbox_future", + Some("CREATE TABLE outbox_future (value TEXT)"), + ); + validate_catalog_delta(&[], std::slice::from_ref(&future), &migration, "up") + .expect("additive delta"); + validate_catalog_delta(std::slice::from_ref(&future), &[], &migration, "down") + .expect("rollback delta"); + assert!(matches!( + validate_catalog_delta(&[], std::slice::from_ref(&future), &migration, "sideways"), + Err(RadrootsOutboxError::MigrationCatalogDeltaMismatch { .. }) + )); + let changed = catalog_row( + "table", + "outbox_future", + "outbox_future", + Some("CREATE TABLE outbox_future (changed TEXT)"), + ); + assert!(matches!( + validate_catalog_delta( + std::slice::from_ref(&future), + std::slice::from_ref(&changed), + &migration, + "up" + ), + Err(RadrootsOutboxError::MigrationCatalogDeltaMismatch { .. }) + )); + + assert!(!validate_ledger_catalog(&[]).expect("absent ledger")); + let canonical = catalog_row( + "table", + OUTBOX_LEDGER_NAME, + OUTBOX_LEDGER_NAME, + Some(OUTBOX_LEDGER_DDL), + ); + assert!(validate_ledger_catalog(std::slice::from_ref(&canonical)).expect("ledger")); + assert!(matches!( + validate_ledger_catalog(&[canonical.clone(), canonical.clone()]), + Err(RadrootsOutboxError::MigrationLedgerDrift { .. }) + )); + for counterfeit in [ + catalog_row( + "view", + OUTBOX_LEDGER_NAME, + OUTBOX_LEDGER_NAME, + Some(OUTBOX_LEDGER_DDL), + ), + catalog_row( + "table", + "counterfeit", + OUTBOX_LEDGER_NAME, + Some(OUTBOX_LEDGER_DDL), + ), + catalog_row( + "table", + OUTBOX_LEDGER_NAME, + "counterfeit", + Some(OUTBOX_LEDGER_DDL), + ), + catalog_row( + "table", + OUTBOX_LEDGER_NAME, + OUTBOX_LEDGER_NAME, + Some("counterfeit"), + ), + ] { + assert!(matches!( + validate_ledger_catalog(&[counterfeit]), + Err(RadrootsOutboxError::MigrationLedgerDrift { .. }) + )); + } + } + + #[tokio::test] + async fn temporary_authority_collisions_are_rejected_before_migration() { + for sql in [ + "CREATE TEMP TABLE outbox_event (value TEXT)", + "CREATE TEMP TABLE radroots_outbox_schema_migrations (value TEXT)", + "CREATE TEMP VIEW caller_view AS SELECT 1 AS value", + ] { + let pool = memory_pool().await; + sqlx::raw_sql(sql) + .execute(&pool) + .await + .expect("temp fixture"); + assert!(matches!( + migrate_outbox_schema(&pool).await, + Err(RadrootsOutboxError::TemporarySchemaCollision { .. }) + )); + assert_eq!(ledger_count(&pool).await, 0); + } + } + + #[tokio::test] + async fn migration_current_schema_reopen_is_idempotent_and_does_not_rewrite_history() { + let pool = memory_pool().await; + migrate_outbox_schema(&pool).await.expect("first migration"); + let before_changes: i64 = sqlx::query_scalar("SELECT total_changes()") + .fetch_one(&pool) + .await + .expect("before total changes"); + migrate_outbox_schema(&pool) + .await + .expect("current migration"); + let after_changes: i64 = sqlx::query_scalar("SELECT total_changes()") + .fetch_one(&pool) + .await + .expect("after total changes"); + assert_eq!(before_changes, after_changes); + assert_eq!(ledger_count(&pool).await, 1); + } + + #[tokio::test] + async fn migration_rollback_wrapper_rejects_exactly_below_the_registry_floor() { + let pool = memory_pool().await; + assert!(matches!( + rollback_outbox_schema_offline(&pool, RADROOTS_OUTBOX_SCHEMA_VERSION_MIN - 1).await, + Err(RadrootsOutboxError::RollbackBelowVersionFloor { + floor: RADROOTS_OUTBOX_SCHEMA_VERSION_MIN, + target: 0, + }) + )); + } +} diff --git a/crates/outbox/src/store.rs b/crates/outbox/src/store.rs @@ -1,7 +1,6 @@ #![forbid(unsafe_code)] use crate::RadrootsOutboxError; -use crate::migrations::{OUTBOX_MIGRATION_DOWN, OUTBOX_MIGRATION_UP}; #[cfg(feature = "event-store-adapter")] use crate::model::RadrootsOutboxEventStoreIngestReceipt; use crate::model::{ @@ -15,6 +14,9 @@ use crate::model::{ RadrootsOutboxSignedOperationInput, RadrootsOutboxSignedTradeMutationInput, RadrootsOutboxStatusSummary, RadrootsOutboxTradeMutationInput, }; +use crate::schema::{ + RadrootsOutboxSchemaStatus, inspect_outbox_schema_status, migrate_outbox_schema, +}; use radroots_event::RadrootsEventKindClass; use radroots_event::draft::{ RadrootsEventDraft, RadrootsSignedEvent, validate_signed_nostr_event_matches_draft, @@ -58,7 +60,7 @@ impl RadrootsOutbox { .connect_with(options) .await?; configure_pool(&pool, false).await?; - apply_up(&pool).await?; + migrate_outbox_schema(&pool).await?; Ok(Self { pool }) } @@ -71,7 +73,7 @@ impl RadrootsOutbox { .connect_with(options) .await?; configure_pool(&pool, true).await?; - apply_up(&pool).await?; + migrate_outbox_schema(&pool).await?; Ok(Self { pool }) } @@ -80,7 +82,7 @@ impl RadrootsOutbox { file_backed: bool, ) -> Result<Self, RadrootsOutboxError> { configure_pool(&pool, file_backed).await?; - apply_up(&pool).await?; + migrate_outbox_schema(&pool).await?; Ok(Self { pool }) } @@ -88,8 +90,14 @@ impl RadrootsOutbox { &self.pool } - pub async fn migrate_down(&self) -> Result<(), RadrootsOutboxError> { - apply_down(&self.pool).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 } pub async fn pragma_foreign_keys(&self) -> Result<i64, RadrootsOutboxError> { @@ -1884,18 +1892,6 @@ async fn file_pool_with_immutable_delete_journal(path: &Path) -> SqlitePool { } #[cfg_attr(coverage_nightly, coverage(off))] -async fn apply_up(pool: &SqlitePool) -> Result<(), RadrootsOutboxError> { - sqlx::raw_sql(OUTBOX_MIGRATION_UP).execute(pool).await?; - Ok(()) -} - -#[cfg_attr(coverage_nightly, coverage(off))] -async fn apply_down(pool: &SqlitePool) -> Result<(), RadrootsOutboxError> { - sqlx::raw_sql(OUTBOX_MIGRATION_DOWN).execute(pool).await?; - Ok(()) -} - -#[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?; Ok(row.try_get(0)?) @@ -4184,7 +4180,7 @@ mod tests { } #[tokio::test] - async fn migration_applies_delivery_plan_schema_and_migrates_down() { + async fn migration_applies_and_authenticates_delivery_plan_schema() { let outbox = RadrootsOutbox::open_memory().await.expect("open"); assert_eq!(outbox.pragma_foreign_keys().await.expect("foreign keys"), 1); @@ -4241,14 +4237,16 @@ mod tests { "outcome_kind" ); - outbox.migrate_down().await.expect("migrate down"); - let row = sqlx::query( - "SELECT name FROM sqlite_master WHERE type = 'table' AND name = 'outbox_event'", - ) - .fetch_optional(outbox.pool()) - .await - .expect("table query"); - assert!(row.is_none()); + assert_eq!( + outbox.schema_status().await.expect("schema status"), + RadrootsOutboxSchemaStatus::Managed { + version: crate::RADROOTS_OUTBOX_SCHEMA_VERSION_CURRENT, + } + ); + outbox + .migrate_to_current_schema() + .await + .expect("idempotent migration"); } #[tokio::test]