lib

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

commit 90dbe592dee3fcab7dfafc9c285491dafc94d87f
parent 0029b3ce8ea5cc9c4a09a200ebc4a8fc98caa841
Author: triesap <tyson@radroots.org>
Date:   Tue, 11 Aug 2026 15:26:39 +0000

service-sqlite: verify exact schema catalogs

- define bounded table, index, trigger, snapshot, and catalog identities
- bind exact schema versions to the governed migration catalog digest
- verify initialization, migration, open, connection, and checkout schemas
- reject transactional and runtime drift with redacted integrity evidence

Diffstat:
Mcrates/service_sqlite/src/initialize.rs | 147++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----------
Acrates/service_sqlite/src/integrity/catalog.rs | 889+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Acrates/service_sqlite/src/integrity/mod.rs | 415+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/service_sqlite/src/lib.rs | 5+++++
Mcrates/service_sqlite/src/metadata.rs | 132+++++++++++++++++++++++++++++++++++++++++++++++++++----------------------------
Mcrates/service_sqlite/src/migration.rs | 277+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------------
Mcrates/service_sqlite/src/open.rs | 295+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------
Mcrates/service_sqlite/tests/package_boundary.rs | 86+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------------
8 files changed, 2085 insertions(+), 161 deletions(-)

diff --git a/crates/service_sqlite/src/initialize.rs b/crates/service_sqlite/src/initialize.rs @@ -4,7 +4,7 @@ use core::{fmt, future::Future}; use std::{error::Error, path::PathBuf}; use crate::{ - OpenMode, ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind, + OpenMode, SchemaCatalog, ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqlitePaths, WriterAuthority, }; @@ -20,6 +20,7 @@ pub async fn initialize_database<F, Fut, E>( paths: &ServiceSqlitePaths, mode: OpenMode, metadata: &ServiceDatabaseMetadata, + schema_catalog: &SchemaCatalog, initialize_schema: F, ) -> Result<WriterAuthority, ServiceSqliteError> where @@ -48,6 +49,7 @@ where paths, authority, metadata, + schema_catalog, initialize_schema, &SystemInitializationOperations, ) @@ -56,7 +58,7 @@ where #[cfg(not(any(target_os = "linux", target_os = "macos")))] { - drop((authority, metadata, initialize_schema)); + drop((authority, metadata, schema_catalog, initialize_schema)); Err(initialization_error(InitializationCause::new( InitializationFailureKind::CreateUnavailable, ))) @@ -439,6 +441,7 @@ mod supported { paths: &ServiceSqlitePaths, authority: WriterAuthority, metadata: &ServiceDatabaseMetadata, + schema_catalog: &SchemaCatalog, initialize_schema: F, operations: &O, ) -> Result<WriterAuthority, ServiceSqliteError> @@ -486,7 +489,8 @@ mod supported { ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source) })?; let write_result = - crate::metadata::write_database_metadata(&mut connection, metadata).await; + crate::metadata::write_database_metadata(&mut connection, metadata, schema_catalog) + .await; let close_result = connection.close().await.map_err(|source| { ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source) }); @@ -637,6 +641,39 @@ mod supported { .expect("database metadata") } + fn schema_catalog(service_objects: Vec<crate::SchemaObject>) -> crate::SchemaCatalog { + let migrations = crate::MigrationCatalog::new([]).expect("empty migrations"); + let digest = + crate::SchemaVersionCatalog::computed_digest(1, service_objects.iter().cloned()) + .expect("schema digest"); + let version = crate::SchemaVersionCatalog::new(1, service_objects, digest) + .expect("schema version"); + crate::SchemaCatalog::new(&migrations, [version]).expect("schema catalog") + } + + fn base_schema_catalog() -> crate::SchemaCatalog { + schema_catalog(Vec::new()) + } + + fn service_schema_catalog() -> crate::SchemaCatalog { + const SQL: &str = "CREATE TABLE service_schema (id INTEGER PRIMARY KEY)"; + let object = crate::SchemaObject::new( + crate::SchemaObjectKind::Table, + "service_schema", + "service_schema", + SQL, + crate::SchemaObject::computed_digest( + crate::SchemaObjectKind::Table, + "service_schema", + "service_schema", + SQL, + ) + .expect("service schema digest"), + ) + .expect("service schema object"); + schema_catalog(vec![object]) + } + #[tokio::test(flavor = "current_thread")] async fn successful_initialization_is_create_new_mode_0600_and_retains_authority() { let root = tempfile::tempdir().expect("root"); @@ -644,8 +681,13 @@ mod supported { prepare(&paths); let expected_path = paths.state_database().to_path_buf(); let metadata = metadata(&paths); - let mut authority = - initialize_database(&paths, OpenMode::Initialize, &metadata, |path| async move { + let schema_catalog = service_schema_catalog(); + let mut authority = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + |path| async move { assert_eq!(path, expected_path); use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; @@ -659,9 +701,10 @@ mod supported { .await?; connection.close().await?; Ok::<(), sqlx::Error>(()) - }) - .await - .expect("initialize"); + }, + ) + .await + .expect("initialize"); let filesystem_metadata = fs::metadata(paths.state_database()).unwrap(); assert_eq!(filesystem_metadata.permissions().mode() & 0o777, 0o600); @@ -693,10 +736,17 @@ mod supported { } let called = Cell::new(false); let metadata = metadata(&paths); - let error = initialize_database(&paths, OpenMode::Initialize, &metadata, |_| { - called.set(true); - ready(Ok::<(), CallbackFailure>(())) - }) + let schema_catalog = base_schema_catalog(); + let error = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + |_| { + called.set(true); + ready(Ok::<(), CallbackFailure>(())) + }, + ) .await .expect_err("existing state must fail"); assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); @@ -712,7 +762,8 @@ mod supported { let paths = paths(root.path(), "wrong-mode"); let called = Cell::new(false); let metadata = metadata(&paths); - let error = initialize_database(&paths, mode, &metadata, |_| { + let schema_catalog = base_schema_catalog(); + let error = initialize_database(&paths, mode, &metadata, &schema_catalog, |_| { called.set(true); ready(Ok::<(), CallbackFailure>(())) }) @@ -732,9 +783,14 @@ mod supported { let paths = paths(root.path(), "callback-failure"); prepare(&paths); let metadata = metadata(&paths); - let error = initialize_database(&paths, OpenMode::Initialize, &metadata, |_| { - ready(Err::<(), _>(CallbackFailure)) - }) + let schema_catalog = base_schema_catalog(); + let error = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + |_| ready(Err::<(), _>(CallbackFailure)), + ) .await .expect_err("callback failure"); assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); @@ -756,9 +812,13 @@ mod supported { Some("secret callback path=/private/state.sqlite") ); - let retry = initialize_database(&paths, OpenMode::Initialize, &metadata, |_| { - ready(Ok::<(), CallbackFailure>(())) - }) + let retry = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + |_| ready(Ok::<(), CallbackFailure>(())), + ) .await .expect("retry after cleanup"); assert!(retry.is_held()); @@ -783,10 +843,12 @@ mod supported { let paths = paths(root.path(), instance); prepare(&paths); let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); let error = initialize_database( &paths, OpenMode::Initialize, &metadata, + &schema_catalog, move |path| async move { use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; @@ -814,15 +876,55 @@ mod supported { } #[tokio::test(flavor = "current_thread")] + async fn schema_catalog_mismatch_cleans_the_exact_reserved_database() { + let root = tempfile::tempdir().expect("root"); + let paths = paths(root.path(), "schema-mismatch"); + prepare(&paths); + let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); + let error = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + |path| async move { + use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; + + let options = SqliteConnectOptions::new() + .filename(path) + .create_if_missing(false) + .disable_statement_logging(); + let mut connection = sqlx::SqliteConnection::connect_with(&options).await?; + sqlx::query("CREATE TABLE unexpected (value INTEGER)") + .execute(&mut connection) + .await?; + connection.close().await?; + Ok::<(), sqlx::Error>(()) + }, + ) + .await + .expect_err("unlisted schema object must fail initialization"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert!(!paths.state_database().exists()); + assert!( + WriterAuthority::acquire(&paths, OpenMode::Initialize) + .unwrap() + .is_some() + ); + } + + #[tokio::test(flavor = "current_thread")] async fn cancellation_rolls_back_and_releases_authority() { let root = tempfile::tempdir().expect("root"); let paths = paths(root.path(), "cancelled"); prepare(&paths); let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); let (poll, future) = poll_once(initialize_database( &paths, OpenMode::Initialize, &metadata, + &schema_catalog, |_| pending::<Result<(), CallbackFailure>>(), )); assert!(poll.is_pending()); @@ -843,11 +945,13 @@ mod supported { let paths = paths(root.path(), "replacement"); prepare(&paths); let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); let replacement_path = paths.state_database().to_path_buf(); let error = initialize_database( &paths, OpenMode::Initialize, &metadata, + &schema_catalog, move |path| async move { fs::remove_file(&path)?; fs::write(&replacement_path, b"replacement")?; @@ -866,6 +970,7 @@ mod supported { let paths = paths(root.path(), "parent-replacement"); prepare(&paths); let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); let state_directory = paths.state_database().parent().unwrap().to_path_buf(); let displaced_directory = state_directory.with_file_name("parent-replacement-old"); let displaced_for_callback = displaced_directory.clone(); @@ -875,6 +980,7 @@ mod supported { &paths, OpenMode::Initialize, &metadata, + &schema_catalog, move |_| async move { fs::rename(&state_directory, &displaced_for_callback)?; fs::create_dir(&state_directory)?; @@ -904,6 +1010,7 @@ mod supported { let paths = paths(root.path(), scenario); prepare(&paths); let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); let authority = WriterAuthority::acquire(&paths, OpenMode::Initialize) .unwrap() .unwrap(); @@ -924,6 +1031,7 @@ mod supported { &paths, authority, &metadata, + &schema_catalog, |_| ready(Err::<(), _>(CallbackFailure)), &operations, ) @@ -933,6 +1041,7 @@ mod supported { &paths, authority, &metadata, + &schema_catalog, |_| ready(Ok::<(), CallbackFailure>(())), &operations, ) diff --git a/crates/service_sqlite/src/integrity/catalog.rs b/crates/service_sqlite/src/integrity/catalog.rs @@ -0,0 +1,889 @@ +//! Immutable expected SQLite schema-object catalogs. + +use core::fmt; +use std::{collections::BTreeSet, error::Error}; + +use sha2::{Digest, Sha256}; + +use crate::{MigrationCatalog, MigrationChecksum}; + +pub(crate) const MAX_SCHEMA_OBJECT_COUNT: usize = 4096; +pub(crate) const MAX_SCHEMA_SQL_UTF8_BYTES: usize = 1024 * 1024; +pub(crate) const MAX_SCHEMA_CATALOG_UTF8_BYTES: usize = 16 * 1024 * 1024; +const MAX_SCHEMA_NAME_UTF8_BYTES: usize = 128; +const MAX_SCHEMA_VERSION_COUNT: usize = 4097; + +const OBJECT_DOMAIN: &[u8] = b"radroots.service_sqlite.schema_object.v1\0"; +const SNAPSHOT_DOMAIN: &[u8] = b"radroots.service_sqlite.schema_snapshot.v1\0"; +const CATALOG_DOMAIN: &[u8] = b"radroots.service_sqlite.schema_catalog.v1\0"; + +pub(crate) const CREATE_METADATA_TABLE_SQL: &str = r#"CREATE TABLE radroots_service_metadata ( + singleton INTEGER NOT NULL PRIMARY KEY CHECK (singleton = 1), + service_id TEXT NOT NULL, + instance_id TEXT NOT NULL, + source_generation BLOB NOT NULL CHECK (length(source_generation) = 32), + state_schema_version INTEGER NOT NULL + CHECK (state_schema_version BETWEEN 1 AND 4294967295), + created_at_unix_ms INTEGER NOT NULL CHECK (created_at_unix_ms > 0) +) STRICT"#; +pub(crate) const CREATE_METADATA_GUARD_TRIGGER_SQL: &str = r#"CREATE TRIGGER radroots_service_metadata_guard_update +BEFORE UPDATE ON radroots_service_metadata +WHEN NEW.singleton != OLD.singleton + OR NEW.service_id != OLD.service_id + OR NEW.instance_id != OLD.instance_id + OR NEW.source_generation != OLD.source_generation + OR NEW.created_at_unix_ms != OLD.created_at_unix_ms + OR NEW.state_schema_version <= OLD.state_schema_version +BEGIN + SELECT RAISE(ABORT, 'service metadata identity is immutable'); +END"#; +pub(crate) const CREATE_METADATA_NO_DELETE_TRIGGER_SQL: &str = r#"CREATE TRIGGER radroots_service_metadata_no_delete +BEFORE DELETE ON radroots_service_metadata +BEGIN + SELECT RAISE(ABORT, 'service metadata is immutable'); +END"#; +pub(crate) const CREATE_MIGRATION_LEDGER_TABLE_SQL: &str = r#"CREATE TABLE schema_migrations ( + version INTEGER NOT NULL PRIMARY KEY + CHECK (version BETWEEN 2 AND 4294967295), + name TEXT NOT NULL UNIQUE + CHECK (length(CAST(name AS BLOB)) BETWEEN 1 AND 128), + checksum BLOB NOT NULL CHECK (length(checksum) = 32), + applied_at_unix_s INTEGER NOT NULL CHECK (applied_at_unix_s BETWEEN 0 AND 9223372036854775807), + service_version TEXT NOT NULL CHECK (length(CAST(service_version AS BLOB)) BETWEEN 1 AND 128), + service_commit TEXT NOT NULL CHECK (length(CAST(service_commit AS BLOB)) = 40), + lib_revision TEXT NOT NULL CHECK (length(CAST(lib_revision AS BLOB)) = 40), + rust_version TEXT NOT NULL CHECK (length(CAST(rust_version AS BLOB)) BETWEEN 1 AND 128), + target TEXT NOT NULL CHECK (length(CAST(target AS BLOB)) BETWEEN 1 AND 128), + feature_profile TEXT NOT NULL CHECK (length(CAST(feature_profile AS BLOB)) BETWEEN 1 AND 128), + config_contract_version INTEGER NOT NULL CHECK (config_contract_version BETWEEN 1 AND 4294967295), + state_contract_version INTEGER NOT NULL CHECK (state_contract_version BETWEEN 1 AND 4294967295), + admin_contract_version INTEGER NOT NULL CHECK (admin_contract_version BETWEEN 1 AND 4294967295), + status_contract_version INTEGER NOT NULL CHECK (status_contract_version BETWEEN 1 AND 4294967295), + provider_contract_version INTEGER NOT NULL CHECK (provider_contract_version BETWEEN 1 AND 4294967295) +) STRICT"#; +pub(crate) const CREATE_MIGRATION_NO_UPDATE_TRIGGER_SQL: &str = r#"CREATE TRIGGER schema_migrations_no_update +BEFORE UPDATE ON schema_migrations +BEGIN + SELECT RAISE(ABORT, 'migration history is immutable'); +END"#; +pub(crate) const CREATE_MIGRATION_NO_DELETE_TRIGGER_SQL: &str = r#"CREATE TRIGGER schema_migrations_no_delete +BEFORE DELETE ON schema_migrations +BEGIN + SELECT RAISE(ABORT, 'migration history is immutable'); +END"#; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) const METADATA_SCHEMA_SQL: [&str; 3] = [ + CREATE_METADATA_TABLE_SQL, + CREATE_METADATA_GUARD_TRIGGER_SQL, + CREATE_METADATA_NO_DELETE_TRIGGER_SQL, +]; +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) const MIGRATION_LEDGER_SCHEMA_SQL: [&str; 3] = [ + CREATE_MIGRATION_LEDGER_TABLE_SQL, + CREATE_MIGRATION_NO_UPDATE_TRIGGER_SQL, + CREATE_MIGRATION_NO_DELETE_TRIGGER_SQL, +]; + +/// Supported persistent object kinds in a governed service schema. +#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Hash)] +pub enum SchemaObjectKind { + Table, + Index, + Trigger, +} + +impl SchemaObjectKind { + pub(crate) const fn tag(self) -> u8 { + match self { + Self::Table => 0, + Self::Index => 1, + Self::Trigger => 2, + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + pub(crate) fn from_sqlite(value: &str) -> Option<Self> { + match value { + "table" => Some(Self::Table), + "index" => Some(Self::Index), + "trigger" => Some(Self::Trigger), + _ => None, + } + } +} + +/// A SHA-256 digest over an object, version snapshot, or bound schema catalog. +#[derive(Clone, Copy, PartialEq, Eq, Hash)] +pub struct SchemaDigest([u8; 32]); + +impl SchemaDigest { + /// Constructs an independently reviewed digest from exact bytes. + #[must_use] + pub const fn from_bytes(bytes: [u8; 32]) -> Self { + Self(bytes) + } + + /// Returns the exact digest bytes. + #[must_use] + pub const fn as_bytes(&self) -> &[u8; 32] { + &self.0 + } +} + +impl fmt::Debug for SchemaDigest { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("SchemaDigest([redacted])") + } +} + +/// One immutable service-owned SQLite schema object definition. +#[derive(Clone, PartialEq, Eq)] +pub struct SchemaObject { + kind: SchemaObjectKind, + name: &'static str, + table_name: &'static str, + sql: &'static str, + digest: SchemaDigest, +} + +impl SchemaObject { + /// Validates an embedded object definition and its independently pinned digest. + pub fn new( + kind: SchemaObjectKind, + name: &'static str, + table_name: &'static str, + sql: &'static str, + expected_digest: SchemaDigest, + ) -> Result<Self, SchemaCatalogContractError> { + validate_service_object(kind, name, table_name, sql)?; + let actual_digest = object_digest(kind, name, table_name, sql); + if actual_digest != expected_digest { + return Err(SchemaCatalogContractError::ObjectDigestMismatch); + } + Ok(Self { + kind, + name, + table_name, + sql, + digest: actual_digest, + }) + } + + #[must_use] + pub const fn kind(&self) -> SchemaObjectKind { + self.kind + } + + #[must_use] + pub const fn name(&self) -> &'static str { + self.name + } + + #[must_use] + pub const fn table_name(&self) -> &'static str { + self.table_name + } + + #[must_use] + pub const fn digest(&self) -> SchemaDigest { + self.digest + } + + /// Computes the frozen object digest for independent pin generation. + pub fn computed_digest( + kind: SchemaObjectKind, + name: &str, + table_name: &str, + sql: &str, + ) -> Result<SchemaDigest, SchemaCatalogContractError> { + validate_service_object(kind, name, table_name, sql)?; + Ok(object_digest(kind, name, table_name, sql)) + } +} + +impl fmt::Debug for SchemaObject { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("SchemaObject") + .field("kind", &self.kind) + .field("name", &self.name) + .field("table_name", &self.table_name) + .field("digest", &self.digest) + .field("sql", &"[redacted]") + .finish() + } +} + +/// Exact expected non-internal object snapshot for one schema version. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub struct SchemaVersionCatalog { + version: u32, + object_count: u32, + digest: SchemaDigest, +} + +impl SchemaVersionCatalog { + /// Validates service-owned objects, adds the shared objects, and pins the snapshot. + pub fn new<I>( + version: u32, + service_objects: I, + expected_digest: SchemaDigest, + ) -> Result<Self, SchemaCatalogContractError> + where + I: IntoIterator<Item = SchemaObject>, + { + if version == 0 { + return Err(SchemaCatalogContractError::InvalidVersionSequence); + } + let service_objects: Vec<_> = service_objects + .into_iter() + .take(MAX_SCHEMA_OBJECT_COUNT + 1) + .collect(); + if service_objects.len() + shared_objects().len() > MAX_SCHEMA_OBJECT_COUNT { + return Err(SchemaCatalogContractError::TooManyObjects); + } + validate_object_set(&service_objects)?; + let mut objects = shared_objects(); + objects.extend(service_objects.iter().map(ObjectRef::from)); + let actual_digest = snapshot_digest(version, &objects); + if actual_digest != expected_digest { + return Err(SchemaCatalogContractError::SnapshotDigestMismatch); + } + Ok(Self { + version, + object_count: u32::try_from(objects.len()).expect("schema object bound fits in u32"), + digest: actual_digest, + }) + } + + /// Computes a snapshot digest for pin generation after validating the object set. + pub fn computed_digest<I>( + version: u32, + service_objects: I, + ) -> Result<SchemaDigest, SchemaCatalogContractError> + where + I: IntoIterator<Item = SchemaObject>, + { + if version == 0 { + return Err(SchemaCatalogContractError::InvalidVersionSequence); + } + let service_objects: Vec<_> = service_objects + .into_iter() + .take(MAX_SCHEMA_OBJECT_COUNT + 1) + .collect(); + if service_objects.len() + shared_objects().len() > MAX_SCHEMA_OBJECT_COUNT { + return Err(SchemaCatalogContractError::TooManyObjects); + } + validate_object_set(&service_objects)?; + let mut objects = shared_objects(); + objects.extend(service_objects.iter().map(ObjectRef::from)); + Ok(snapshot_digest(version, &objects)) + } + + #[must_use] + pub const fn version(self) -> u32 { + self.version + } + + #[must_use] + pub const fn object_count(self) -> u32 { + self.object_count + } + + #[must_use] + pub const fn digest(self) -> SchemaDigest { + self.digest + } +} + +/// Ordered exact schema snapshots bound to one migration catalog. +#[derive(Clone, PartialEq, Eq)] +pub struct SchemaCatalog { + versions: Box<[SchemaVersionCatalog]>, + migration_catalog_digest: MigrationChecksum, + digest: SchemaDigest, +} + +impl SchemaCatalog { + /// Validates one exact snapshot for every migration-catalog schema version. + pub fn new<I>( + migrations: &MigrationCatalog, + versions: I, + ) -> Result<Self, SchemaCatalogContractError> + where + I: IntoIterator<Item = SchemaVersionCatalog>, + { + let versions: Vec<_> = versions + .into_iter() + .take(MAX_SCHEMA_VERSION_COUNT + 1) + .collect(); + if versions.len() > MAX_SCHEMA_VERSION_COUNT { + return Err(SchemaCatalogContractError::TooManyVersions); + } + let expected_len = usize::try_from(migrations.current_version()) + .map_err(|_| SchemaCatalogContractError::MigrationCatalogMismatch)?; + if versions.len() != expected_len + || versions + .iter() + .enumerate() + .any(|(index, entry)| entry.version != u32::try_from(index + 1).unwrap_or(0)) + { + return Err(SchemaCatalogContractError::InvalidVersionSequence); + } + let migration_catalog_digest = migrations.digest(); + let digest = catalog_digest(migration_catalog_digest, &versions); + Ok(Self { + versions: versions.into_boxed_slice(), + migration_catalog_digest, + digest, + }) + } + + #[must_use] + pub fn versions(&self) -> &[SchemaVersionCatalog] { + &self.versions + } + + #[must_use] + pub const fn migration_catalog_digest(&self) -> MigrationChecksum { + self.migration_catalog_digest + } + + #[must_use] + pub const fn digest(&self) -> SchemaDigest { + self.digest + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + pub(crate) fn matches_migrations(&self, migrations: &MigrationCatalog) -> bool { + self.migration_catalog_digest == migrations.digest() + && self.versions.len() + == usize::try_from(migrations.current_version()).unwrap_or(usize::MAX) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + pub(crate) fn version(&self, version: u32) -> Option<SchemaVersionCatalog> { + let index = usize::try_from(version.checked_sub(1)?).ok()?; + self.versions.get(index).copied() + } +} + +impl fmt::Debug for SchemaCatalog { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("SchemaCatalog") + .field("version_count", &self.versions.len()) + .field("migration_catalog_digest", &self.migration_catalog_digest) + .field("digest", &self.digest) + .finish() + } +} + +/// Invalid immutable schema-catalog construction. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum SchemaCatalogContractError { + InvalidName, + ReservedName, + InvalidBinding, + InvalidSql, + ObjectDigestMismatch, + DuplicateObject, + TooManyObjects, + SnapshotDigestMismatch, + InvalidVersionSequence, + TooManyVersions, + MigrationCatalogMismatch, +} + +impl fmt::Display for SchemaCatalogContractError { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self { + Self::InvalidName => "SQLite schema object name is invalid", + Self::ReservedName => "SQLite schema object name is reserved", + Self::InvalidBinding => "SQLite schema object table binding is invalid", + Self::InvalidSql => "SQLite schema object definition is invalid", + Self::ObjectDigestMismatch => "SQLite schema object digest does not match", + Self::DuplicateObject => "SQLite schema object identity is duplicated", + Self::TooManyObjects => "SQLite schema snapshot has too many objects", + Self::SnapshotDigestMismatch => "SQLite schema snapshot digest does not match", + Self::InvalidVersionSequence => "SQLite schema catalog version sequence is invalid", + Self::TooManyVersions => "SQLite schema catalog has too many versions", + Self::MigrationCatalogMismatch => { + "SQLite schema catalog does not match the migration catalog" + } + }) + } +} + +impl Error for SchemaCatalogContractError {} + +#[derive(Clone, Copy)] +pub(crate) struct ObjectRef<'a> { + pub(crate) kind: SchemaObjectKind, + pub(crate) name: &'a str, + pub(crate) table_name: &'a str, + pub(crate) sql: &'a str, + pub(crate) digest: SchemaDigest, +} + +impl<'a> From<&'a SchemaObject> for ObjectRef<'a> { + fn from(value: &'a SchemaObject) -> Self { + Self { + kind: value.kind, + name: value.name, + table_name: value.table_name, + sql: value.sql, + digest: value.digest, + } + } +} + +pub(crate) fn object_digest( + kind: SchemaObjectKind, + name: &str, + table_name: &str, + sql: &str, +) -> SchemaDigest { + let mut hasher = Sha256::new(); + hasher.update(OBJECT_DOMAIN); + hasher.update([kind.tag()]); + update_length_prefixed(&mut hasher, name.as_bytes()); + update_length_prefixed(&mut hasher, table_name.as_bytes()); + update_length_prefixed(&mut hasher, sql.as_bytes()); + SchemaDigest(hasher.finalize().into()) +} + +pub(crate) fn snapshot_digest(version: u32, objects: &[ObjectRef<'_>]) -> SchemaDigest { + let mut objects = objects.to_vec(); + objects.sort_by(|left, right| { + ( + left.kind.tag(), + left.name.as_bytes(), + left.table_name.as_bytes(), + ) + .cmp(&( + right.kind.tag(), + right.name.as_bytes(), + right.table_name.as_bytes(), + )) + }); + let mut hasher = Sha256::new(); + hasher.update(SNAPSHOT_DOMAIN); + hasher.update(version.to_be_bytes()); + hasher.update( + u32::try_from(objects.len()) + .expect("schema object bound fits in u32") + .to_be_bytes(), + ); + for object in objects { + hasher.update([object.kind.tag()]); + update_length_prefixed(&mut hasher, object.name.as_bytes()); + update_length_prefixed(&mut hasher, object.table_name.as_bytes()); + hasher.update(object.digest.as_bytes()); + } + SchemaDigest(hasher.finalize().into()) +} + +fn catalog_digest( + migration_digest: MigrationChecksum, + versions: &[SchemaVersionCatalog], +) -> SchemaDigest { + let mut hasher = Sha256::new(); + hasher.update(CATALOG_DOMAIN); + hasher.update(migration_digest.as_bytes()); + hasher.update( + u32::try_from(versions.len()) + .expect("schema version bound fits in u32") + .to_be_bytes(), + ); + for version in versions { + hasher.update(version.version.to_be_bytes()); + hasher.update(version.object_count.to_be_bytes()); + hasher.update(version.digest.as_bytes()); + } + SchemaDigest(hasher.finalize().into()) +} + +fn update_length_prefixed(hasher: &mut Sha256, value: &[u8]) { + hasher.update( + u64::try_from(value.len()) + .expect("bounded schema field fits in u64") + .to_be_bytes(), + ); + hasher.update(value); +} + +fn validate_service_object( + kind: SchemaObjectKind, + name: &str, + table_name: &str, + sql: &str, +) -> Result<(), SchemaCatalogContractError> { + if !valid_name(name) || !valid_name(table_name) { + return Err(SchemaCatalogContractError::InvalidName); + } + if is_reserved(name) || is_reserved(table_name) { + return Err(SchemaCatalogContractError::ReservedName); + } + if (kind == SchemaObjectKind::Table) != (name == table_name) { + return Err(SchemaCatalogContractError::InvalidBinding); + } + if sql.is_empty() || sql.len() > MAX_SCHEMA_SQL_UTF8_BYTES || sql.as_bytes().contains(&0) { + return Err(SchemaCatalogContractError::InvalidSql); + } + Ok(()) +} + +fn validate_object_set(objects: &[SchemaObject]) -> Result<(), SchemaCatalogContractError> { + let mut identities = BTreeSet::new(); + let tables = objects + .iter() + .filter(|object| object.kind == SchemaObjectKind::Table) + .map(|object| object.name) + .collect::<BTreeSet<_>>(); + let mut total_sql_bytes = shared_objects() + .iter() + .map(|object| object.sql.len()) + .sum::<usize>(); + for object in objects { + if !identities.insert((object.kind, object.name)) { + return Err(SchemaCatalogContractError::DuplicateObject); + } + if object.kind != SchemaObjectKind::Table && !tables.contains(object.table_name) { + return Err(SchemaCatalogContractError::InvalidBinding); + } + total_sql_bytes = total_sql_bytes + .checked_add(object.sql.len()) + .ok_or(SchemaCatalogContractError::TooManyObjects)?; + } + if total_sql_bytes > MAX_SCHEMA_CATALOG_UTF8_BYTES { + return Err(SchemaCatalogContractError::TooManyObjects); + } + Ok(()) +} + +fn valid_name(value: &str) -> bool { + let bytes = value.as_bytes(); + !bytes.is_empty() + && bytes.len() <= MAX_SCHEMA_NAME_UTF8_BYTES + && bytes[0].is_ascii_lowercase() + && bytes[bytes.len() - 1].is_ascii_alphanumeric() + && !bytes.windows(2).any(|pair| pair == b"__") + && bytes + .iter() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') +} + +fn is_reserved(value: &str) -> bool { + value.starts_with("sqlite_") || shared_objects().iter().any(|object| object.name == value) +} + +fn shared_objects() -> Vec<ObjectRef<'static>> { + [ + ( + SchemaObjectKind::Table, + "radroots_service_metadata", + "radroots_service_metadata", + CREATE_METADATA_TABLE_SQL, + ), + ( + SchemaObjectKind::Trigger, + "radroots_service_metadata_guard_update", + "radroots_service_metadata", + CREATE_METADATA_GUARD_TRIGGER_SQL, + ), + ( + SchemaObjectKind::Trigger, + "radroots_service_metadata_no_delete", + "radroots_service_metadata", + CREATE_METADATA_NO_DELETE_TRIGGER_SQL, + ), + ( + SchemaObjectKind::Table, + "schema_migrations", + "schema_migrations", + CREATE_MIGRATION_LEDGER_TABLE_SQL, + ), + ( + SchemaObjectKind::Trigger, + "schema_migrations_no_update", + "schema_migrations", + CREATE_MIGRATION_NO_UPDATE_TRIGGER_SQL, + ), + ( + SchemaObjectKind::Trigger, + "schema_migrations_no_delete", + "schema_migrations", + CREATE_MIGRATION_NO_DELETE_TRIGGER_SQL, + ), + ] + .into_iter() + .map(|(kind, name, table_name, sql)| ObjectRef { + kind, + name, + table_name, + sql, + digest: object_digest(kind, name, table_name, sql), + }) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + + const TABLE_SQL: &str = "CREATE TABLE alpha (value INTEGER NOT NULL) STRICT"; + + fn table() -> SchemaObject { + SchemaObject::new( + SchemaObjectKind::Table, + "alpha", + "alpha", + TABLE_SQL, + SchemaObject::computed_digest(SchemaObjectKind::Table, "alpha", "alpha", TABLE_SQL) + .unwrap(), + ) + .unwrap() + } + + fn empty_migrations() -> MigrationCatalog { + MigrationCatalog::new([]).unwrap() + } + + #[test] + fn exact_object_snapshot_and_catalog_vectors_are_stable() { + let object = table(); + assert_eq!( + object.digest().as_bytes(), + &[ + 0xf1, 0xa0, 0x6b, 0x60, 0x76, 0x0f, 0x73, 0xae, 0x0b, 0x43, 0x44, 0x1b, 0xb8, 0xfc, + 0x12, 0x56, 0x48, 0xa9, 0xcb, 0xf5, 0x63, 0xdb, 0x59, 0xd6, 0xc1, 0xfc, 0x60, 0x3b, + 0x5e, 0x92, 0x0b, 0x7c, + ] + ); + let snapshot_digest = SchemaVersionCatalog::computed_digest(1, [object.clone()]).unwrap(); + assert_eq!( + snapshot_digest.as_bytes(), + &[ + 0x9e, 0x50, 0xd2, 0x25, 0xfe, 0xdc, 0xe1, 0x4f, 0x0d, 0x40, 0x41, 0x84, 0xb2, 0x00, + 0xbc, 0xd5, 0xcf, 0xe9, 0xb3, 0x56, 0x98, 0x16, 0xf6, 0x29, 0xc2, 0x40, 0x86, 0xdf, + 0x93, 0x2f, 0x52, 0x4a, + ] + ); + let version = SchemaVersionCatalog::new(1, [object], snapshot_digest).unwrap(); + let catalog = SchemaCatalog::new(&empty_migrations(), [version]).unwrap(); + assert_eq!(version.version(), 1); + assert_eq!(version.object_count(), 7); + assert_eq!(catalog.versions(), &[version]); + assert_eq!( + catalog.digest().as_bytes(), + &[ + 0xff, 0x9d, 0xbe, 0x4f, 0x32, 0x42, 0xb3, 0x3f, 0x6f, 0xd8, 0x76, 0x9d, 0x3e, 0x11, + 0x21, 0x7b, 0x38, 0xb2, 0x77, 0x3e, 0xa6, 0xa2, 0x98, 0x6b, 0x85, 0xb6, 0xed, 0x8d, + 0xe2, 0x4b, 0x52, 0x51, + ] + ); + } + + #[test] + fn object_validation_is_closed_and_redacted() { + let digest = SchemaDigest::from_bytes([0; 32]); + for result in [ + SchemaObject::new(SchemaObjectKind::Table, "", "", "x", digest), + SchemaObject::new(SchemaObjectKind::Table, "Bad", "Bad", "x", digest), + SchemaObject::new( + SchemaObjectKind::Table, + "sqlite_bad", + "sqlite_bad", + "x", + digest, + ), + SchemaObject::new( + SchemaObjectKind::Table, + "schema_migrations", + "schema_migrations", + "x", + digest, + ), + SchemaObject::new( + SchemaObjectKind::Index, + "alpha_idx", + "alpha_idx", + "x", + digest, + ), + SchemaObject::new(SchemaObjectKind::Table, "alpha", "alpha", "", digest), + ] { + assert!(result.is_err()); + } + let debug = format!("{:?}", table()); + assert!(!debug.contains(TABLE_SQL)); + assert!(debug.contains("[redacted]")); + + let maximum_name = Box::leak("a".repeat(MAX_SCHEMA_NAME_UTF8_BYTES).into_boxed_str()); + let maximum_digest = + SchemaObject::computed_digest(SchemaObjectKind::Table, maximum_name, maximum_name, "x") + .unwrap(); + assert!( + SchemaObject::new( + SchemaObjectKind::Table, + maximum_name, + maximum_name, + "x", + maximum_digest, + ) + .is_ok() + ); + let excessive_name = Box::leak("a".repeat(MAX_SCHEMA_NAME_UTF8_BYTES + 1).into_boxed_str()); + assert_eq!( + SchemaObject::computed_digest( + SchemaObjectKind::Table, + excessive_name, + excessive_name, + "x", + ), + Err(SchemaCatalogContractError::InvalidName) + ); + assert_eq!( + SchemaObject::computed_digest( + SchemaObjectKind::Table, + "alpha__beta", + "alpha__beta", + "x", + ), + Err(SchemaCatalogContractError::InvalidName) + ); + + let maximum_sql = Box::leak("x".repeat(MAX_SCHEMA_SQL_UTF8_BYTES).into_boxed_str()); + assert!( + SchemaObject::computed_digest( + SchemaObjectKind::Table, + "maximum_sql", + "maximum_sql", + maximum_sql, + ) + .is_ok() + ); + let excessive_sql = Box::leak("x".repeat(MAX_SCHEMA_SQL_UTF8_BYTES + 1).into_boxed_str()); + assert_eq!( + SchemaObject::computed_digest( + SchemaObjectKind::Table, + "excessive_sql", + "excessive_sql", + excessive_sql, + ), + Err(SchemaCatalogContractError::InvalidSql) + ); + } + + #[test] + fn object_sets_reject_duplicates_missing_tables_and_bounds() { + let duplicate_digest = SchemaVersionCatalog::computed_digest(1, [table(), table()]); + assert_eq!( + duplicate_digest, + Err(SchemaCatalogContractError::DuplicateObject) + ); + + const INDEX_SQL: &str = "CREATE INDEX alpha_idx ON missing(value)"; + let index = SchemaObject::new( + SchemaObjectKind::Index, + "alpha_idx", + "missing", + INDEX_SQL, + SchemaObject::computed_digest( + SchemaObjectKind::Index, + "alpha_idx", + "missing", + INDEX_SQL, + ) + .unwrap(), + ) + .unwrap(); + assert_eq!( + SchemaVersionCatalog::computed_digest(1, [index]), + Err(SchemaCatalogContractError::InvalidBinding) + ); + + let excessive = std::iter::repeat_with(table).take(MAX_SCHEMA_OBJECT_COUNT + 1); + assert_eq!( + SchemaVersionCatalog::computed_digest(1, excessive), + Err(SchemaCatalogContractError::TooManyObjects) + ); + let infinite = std::iter::repeat_with(table); + assert_eq!( + SchemaVersionCatalog::computed_digest(1, infinite), + Err(SchemaCatalogContractError::TooManyObjects) + ); + + let maximum = (0..(MAX_SCHEMA_OBJECT_COUNT - shared_objects().len())) + .map(|index| { + let name = Box::leak(format!("table_{index}").into_boxed_str()); + let sql = + Box::leak(format!("CREATE TABLE {name} (value INTEGER)").into_boxed_str()); + SchemaObject::new( + SchemaObjectKind::Table, + name, + name, + sql, + SchemaObject::computed_digest(SchemaObjectKind::Table, name, name, sql) + .unwrap(), + ) + .unwrap() + }) + .collect::<Vec<_>>(); + assert!(SchemaVersionCatalog::computed_digest(1, maximum.iter().cloned()).is_ok()); + let mut excessive = maximum; + excessive.push(table()); + assert_eq!( + SchemaVersionCatalog::computed_digest(1, excessive), + Err(SchemaCatalogContractError::TooManyObjects) + ); + } + + #[test] + fn catalog_requires_every_exact_migration_version_and_terminates() { + let migrations = empty_migrations(); + let digest = SchemaVersionCatalog::computed_digest(1, []).unwrap(); + let v1 = SchemaVersionCatalog::new(1, [], digest).unwrap(); + assert!(SchemaCatalog::new(&migrations, [v1]).is_ok()); + assert_eq!( + SchemaCatalog::new(&migrations, []), + Err(SchemaCatalogContractError::InvalidVersionSequence) + ); + let infinite = std::iter::repeat(v1); + assert_eq!( + SchemaCatalog::new(&migrations, infinite), + Err(SchemaCatalogContractError::TooManyVersions) + ); + + let descriptors = (0..4096_u32) + .map(|index| { + let name = Box::leak(format!("migration_{index}").into_boxed_str()); + crate::MigrationDescriptor::callback( + index + 2, + name, + b"x", + MigrationChecksum::for_callback(b"x"), + ) + .unwrap() + }) + .collect::<Vec<_>>(); + let maximum_migrations = MigrationCatalog::new(descriptors).unwrap(); + let maximum_versions = (1..=4097_u32) + .map(|version| { + let digest = SchemaVersionCatalog::computed_digest(version, []).unwrap(); + SchemaVersionCatalog::new(version, [], digest).unwrap() + }) + .collect::<Vec<_>>(); + assert!(SchemaCatalog::new(&maximum_migrations, maximum_versions).is_ok()); + + let digest = SchemaVersionCatalog::computed_digest(1, []).unwrap(); + let v1 = SchemaVersionCatalog::new(1, [], digest).unwrap(); + let duplicate = SchemaVersionCatalog { version: 1, ..v1 }; + assert_eq!( + SchemaCatalog::new(&maximum_migrations, [v1, duplicate]), + Err(SchemaCatalogContractError::InvalidVersionSequence) + ); + } +} diff --git a/crates/service_sqlite/src/integrity/mod.rs b/crates/service_sqlite/src/integrity/mod.rs @@ -0,0 +1,415 @@ +//! Exact schema-object catalog verification. + +pub(crate) mod catalog; + +pub use catalog::{ + SchemaCatalog, SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind, + SchemaVersionCatalog, +}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use core::fmt; +#[cfg(any(target_os = "linux", target_os = "macos"))] +use std::error::Error; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use sqlx::{Row, SqliteConnection}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use crate::{ServiceSqliteError, ServiceSqliteErrorKind}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use self::catalog::{ + MAX_SCHEMA_CATALOG_UTF8_BYTES, MAX_SCHEMA_OBJECT_COUNT, MAX_SCHEMA_SQL_UTF8_BYTES, ObjectRef, + object_digest, snapshot_digest, +}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +const READ_SCHEMA_CATALOG_SQL: &str = r#" +WITH raw_catalog AS ( + SELECT type, name, tbl_name, sql + FROM main.sqlite_schema + WHERE NOT (typeof(name) = 'text' AND substr(name, 1, 7) = 'sqlite_') +), bounded_catalog AS ( + SELECT + type, + name, + tbl_name, + sql, + COUNT(*) OVER () AS object_count, + COALESCE(SUM(length(CAST(sql AS BLOB))) OVER (), 0) AS total_sql_bytes + FROM raw_catalog +) +SELECT + object_count, + total_sql_bytes, + CASE + WHEN object_count <= 4096 + AND typeof(type) = 'text' + AND length(CAST(type AS BLOB)) BETWEEN 1 AND 16 + THEN type + END AS bounded_type, + CASE + WHEN object_count <= 4096 + AND typeof(name) = 'text' + AND length(CAST(name AS BLOB)) BETWEEN 1 AND 128 + THEN name + END AS bounded_name, + CASE + WHEN object_count <= 4096 + AND typeof(tbl_name) = 'text' + AND length(CAST(tbl_name AS BLOB)) BETWEEN 1 AND 128 + THEN tbl_name + END AS bounded_table_name, + CASE + WHEN object_count <= 4096 + AND total_sql_bytes <= 16777216 + AND typeof(sql) = 'text' + AND length(CAST(sql AS BLOB)) BETWEEN 1 AND 1048576 + THEN sql + END AS bounded_sql +FROM bounded_catalog +LIMIT 4097 +"#; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[derive(Clone, Debug, PartialEq, Eq)] +pub(crate) struct SchemaVerificationReport { + version: u32, + expected_count: u32, + actual_count: u32, + expected_digest: SchemaDigest, + actual_digest: SchemaDigest, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum SchemaIntegrityFailureKind { + CatalogMismatch, + CatalogCorrupt, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +struct SchemaIntegrityFailure { + kind: SchemaIntegrityFailureKind, + report: Option<SchemaVerificationReport>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl fmt::Debug for SchemaIntegrityFailure { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("SchemaIntegrityFailure") + .field("kind", &self.kind) + .field("report", &self.report) + .finish() + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl fmt::Display for SchemaIntegrityFailure { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + SchemaIntegrityFailureKind::CatalogMismatch => { + "SQLite schema object catalog does not match" + } + SchemaIntegrityFailureKind::CatalogCorrupt => "SQLite schema object catalog is invalid", + }) + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl Error for SchemaIntegrityFailure {} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn integrity_error(kind: SchemaIntegrityFailureKind) -> ServiceSqliteError { + ServiceSqliteError::with_source( + ServiceSqliteErrorKind::Integrity, + SchemaIntegrityFailure { kind, report: None }, + ) +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn mismatch_error(report: SchemaVerificationReport) -> ServiceSqliteError { + ServiceSqliteError::with_source( + ServiceSqliteErrorKind::Integrity, + SchemaIntegrityFailure { + kind: SchemaIntegrityFailureKind::CatalogMismatch, + report: Some(report), + }, + ) +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +struct RuntimeSchemaObject { + kind: SchemaObjectKind, + name: String, + table_name: String, + sql: String, + digest: SchemaDigest, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn verify_schema_catalog( + connection: &mut SqliteConnection, + catalog: &SchemaCatalog, + version: u32, +) -> Result<SchemaVerificationReport, ServiceSqliteError> { + let expected = catalog + .version(version) + .ok_or_else(|| integrity_error(SchemaIntegrityFailureKind::CatalogMismatch))?; + let rows = sqlx::query(READ_SCHEMA_CATALOG_SQL) + .fetch_all(connection) + .await + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + if rows.len() > MAX_SCHEMA_OBJECT_COUNT { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + let mut objects = Vec::with_capacity(rows.len()); + let mut reported_count = None; + let mut reported_total = None; + for row in rows { + let count = row + .try_get::<i64, _>("object_count") + .ok() + .and_then(|value| u32::try_from(value).ok()) + .ok_or_else(|| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + let total = row + .try_get::<i64, _>("total_sql_bytes") + .ok() + .and_then(|value| usize::try_from(value).ok()) + .ok_or_else(|| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + if count > u32::try_from(MAX_SCHEMA_OBJECT_COUNT).unwrap_or(u32::MAX) + || total > MAX_SCHEMA_CATALOG_UTF8_BYTES + || reported_count.is_some_and(|observed| observed != count) + || reported_total.is_some_and(|observed| observed != total) + { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + reported_count = Some(count); + reported_total = Some(total); + let object_type = row + .try_get::<String, _>("bounded_type") + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + let name = row + .try_get::<String, _>("bounded_name") + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + let table_name = row + .try_get::<String, _>("bounded_table_name") + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + let sql = row + .try_get::<String, _>("bounded_sql") + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + if sql.len() > MAX_SCHEMA_SQL_UTF8_BYTES { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + let kind = SchemaObjectKind::from_sqlite(&object_type) + .ok_or_else(|| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + if !runtime_name_is_valid(&name) + || !runtime_name_is_valid(&table_name) + || (kind == SchemaObjectKind::Table) != (name == table_name) + { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + let digest = object_digest(kind, &name, &table_name, &sql); + objects.push(RuntimeSchemaObject { + kind, + name, + table_name, + sql, + digest, + }); + } + let actual_count = u32::try_from(objects.len()) + .map_err(|_| integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt))?; + if reported_count.unwrap_or(0) != actual_count { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + let mut identities = std::collections::BTreeSet::new(); + if objects + .iter() + .any(|object| !identities.insert((object.kind, object.name.as_str()))) + { + return Err(integrity_error(SchemaIntegrityFailureKind::CatalogCorrupt)); + } + let refs = objects + .iter() + .map(|object| ObjectRef { + kind: object.kind, + name: &object.name, + table_name: &object.table_name, + sql: &object.sql, + digest: object.digest, + }) + .collect::<Vec<_>>(); + let actual_digest = snapshot_digest(version, &refs); + let report = SchemaVerificationReport { + version, + expected_count: expected.object_count(), + actual_count, + expected_digest: expected.digest(), + actual_digest, + }; + if report.expected_count != report.actual_count + || report.expected_digest != report.actual_digest + { + return Err(mismatch_error(report)); + } + Ok(report) +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn runtime_name_is_valid(value: &str) -> bool { + let bytes = value.as_bytes(); + !bytes.is_empty() + && bytes.len() <= 128 + && bytes[0].is_ascii_lowercase() + && bytes[bytes.len() - 1].is_ascii_alphanumeric() + && !bytes.windows(2).any(|pair| pair == b"__") + && bytes + .iter() + .all(|byte| byte.is_ascii_lowercase() || byte.is_ascii_digit() || *byte == b'_') +} + +#[cfg(all(test, any(target_os = "linux", target_os = "macos")))] +mod tests { + use super::*; + use sqlx::{Connection, Executor, sqlite::SqliteConnectOptions}; + + const TABLE_SQL: &str = + "CREATE TABLE alpha (id INTEGER PRIMARY KEY, value INTEGER NOT NULL) STRICT"; + const INDEX_SQL: &str = "CREATE INDEX alpha_value_idx ON alpha(value)"; + const TRIGGER_SQL: &str = "CREATE TRIGGER alpha_guard BEFORE UPDATE ON alpha BEGIN SELECT RAISE(ABORT, 'blocked'); END"; + + fn object( + kind: SchemaObjectKind, + name: &'static str, + table_name: &'static str, + sql: &'static str, + ) -> SchemaObject { + SchemaObject::new( + kind, + name, + table_name, + sql, + SchemaObject::computed_digest(kind, name, table_name, sql).unwrap(), + ) + .unwrap() + } + + fn expected(objects: Vec<SchemaObject>) -> SchemaCatalog { + let migrations = crate::MigrationCatalog::new([]).unwrap(); + let digest = SchemaVersionCatalog::computed_digest(1, objects.iter().cloned()).unwrap(); + let version = SchemaVersionCatalog::new(1, objects, digest).unwrap(); + SchemaCatalog::new(&migrations, [version]).unwrap() + } + + fn full_objects() -> Vec<SchemaObject> { + vec![ + object( + SchemaObjectKind::Trigger, + "alpha_guard", + "alpha", + TRIGGER_SQL, + ), + object( + SchemaObjectKind::Index, + "alpha_value_idx", + "alpha", + INDEX_SQL, + ), + object(SchemaObjectKind::Table, "alpha", "alpha", TABLE_SQL), + ] + } + + async fn shared_database() -> SqliteConnection { + let mut connection = + SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:")) + .await + .unwrap(); + for statement in catalog::METADATA_SCHEMA_SQL + .into_iter() + .chain(catalog::MIGRATION_LEDGER_SCHEMA_SQL) + { + connection.execute(statement).await.unwrap(); + } + connection + } + + async fn full_database() -> SqliteConnection { + let mut connection = shared_database().await; + for statement in [TABLE_SQL, INDEX_SQL, TRIGGER_SQL] { + connection.execute(statement).await.unwrap(); + } + connection + } + + #[tokio::test(flavor = "current_thread")] + async fn exact_catalog_is_order_independent_and_report_is_bounded() { + let mut connection = full_database().await; + let catalog = expected(full_objects()); + let first = verify_schema_catalog(&mut connection, &catalog, 1) + .await + .unwrap(); + let second = verify_schema_catalog(&mut connection, &catalog, 1) + .await + .unwrap(); + assert_eq!(first, second); + assert_eq!(first.version, 1); + assert_eq!(first.expected_count, 9); + assert_eq!(first.actual_count, 9); + assert_eq!(first.expected_digest, first.actual_digest); + assert!(!format!("{first:?}").contains(TABLE_SQL)); + } + + #[tokio::test(flavor = "current_thread")] + async fn missing_extra_replaced_index_trigger_column_and_view_fail_closed() { + let cases = [ + "DROP INDEX alpha_value_idx", + "CREATE TABLE extra (value INTEGER)", + "DROP INDEX alpha_value_idx; CREATE INDEX alpha_value_idx ON alpha(value DESC)", + "DROP TRIGGER alpha_guard; CREATE TRIGGER alpha_guard BEFORE UPDATE ON alpha BEGIN SELECT RAISE(ABORT, 'changed'); END", + "ALTER TABLE alpha ADD COLUMN changed TEXT", + "CREATE VIEW alpha_view AS SELECT value FROM alpha", + ]; + for mutation in cases { + let mut connection = full_database().await; + sqlx::raw_sql(mutation) + .execute(&mut connection) + .await + .unwrap(); + let error = verify_schema_catalog(&mut connection, &expected(full_objects()), 1) + .await + .expect_err("schema drift must fail"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert!(!error.to_string().contains("alpha")); + assert!(!format!("{error:?}").contains("alpha")); + } + + let mut missing = shared_database().await; + assert_eq!( + verify_schema_catalog(&mut missing, &expected(full_objects()), 1) + .await + .expect_err("missing table") + .kind(), + ServiceSqliteErrorKind::Integrity + ); + } + + #[tokio::test(flavor = "current_thread")] + async fn oversized_persisted_sql_is_rejected_before_rust_decode() { + let mut connection = shared_database().await; + let oversized = "x".repeat(catalog::MAX_SCHEMA_SQL_UTF8_BYTES + 1); + let statement = format!("CREATE TABLE oversized (value TEXT DEFAULT '{oversized}')"); + sqlx::query(sqlx::AssertSqlSafe(statement)) + .execute(&mut connection) + .await + .unwrap(); + let error = verify_schema_catalog(&mut connection, &expected(Vec::new()), 1) + .await + .expect_err("oversized SQL must fail"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert!(!error.to_string().contains(&oversized)); + assert!(!format!("{error:?}").contains(&oversized)); + } +} diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs @@ -6,6 +6,7 @@ mod authority; mod config; mod error; mod initialize; +mod integrity; mod metadata; mod migration; mod open; @@ -17,6 +18,10 @@ pub use error::{ SafeServiceSqliteError, ServiceSqliteError, ServiceSqliteErrorCode, ServiceSqliteErrorKind, }; pub use initialize::initialize_database; +pub use integrity::{ + SchemaCatalog, SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind, + SchemaVersionCatalog, +}; pub use metadata::{ ServiceDatabaseIdentity, ServiceDatabaseMetadata, ServiceSqliteApplicationId, ServiceSqliteMetadataValueError, diff --git a/crates/service_sqlite/src/metadata.rs b/crates/service_sqlite/src/metadata.rs @@ -17,35 +17,6 @@ use sqlx::{Connection, Row, SqliteConnection}; const MAX_APPLICATION_ID: u32 = i32::MAX as u32; const MAX_CREATED_AT_UNIX_MS: u64 = i64::MAX as u64; -#[cfg(any(target_os = "linux", target_os = "macos"))] -const CREATE_METADATA_SQL: &str = r#" -CREATE TABLE radroots_service_metadata ( - singleton INTEGER NOT NULL PRIMARY KEY CHECK (singleton = 1), - service_id TEXT NOT NULL, - instance_id TEXT NOT NULL, - source_generation BLOB NOT NULL CHECK (length(source_generation) = 32), - state_schema_version INTEGER NOT NULL - CHECK (state_schema_version BETWEEN 1 AND 4294967295), - created_at_unix_ms INTEGER NOT NULL CHECK (created_at_unix_ms > 0) -) STRICT; -CREATE TRIGGER radroots_service_metadata_guard_update -BEFORE UPDATE ON radroots_service_metadata -WHEN NEW.singleton != OLD.singleton - OR NEW.service_id != OLD.service_id - OR NEW.instance_id != OLD.instance_id - OR NEW.source_generation != OLD.source_generation - OR NEW.created_at_unix_ms != OLD.created_at_unix_ms - OR NEW.state_schema_version <= OLD.state_schema_version -BEGIN - SELECT RAISE(ABORT, 'service metadata identity is immutable'); -END; -CREATE TRIGGER radroots_service_metadata_no_delete -BEFORE DELETE ON radroots_service_metadata -BEGIN - SELECT RAISE(ABORT, 'service metadata is immutable'); -END; -"#; - /// A validated nonzero SQLite application identifier. #[derive(Clone, Copy, Debug, PartialEq, Eq)] pub struct ServiceSqliteApplicationId(u32); @@ -317,6 +288,7 @@ impl Error for MigrationLedgerInitializationFailure {} pub(crate) async fn write_database_metadata( connection: &mut SqliteConnection, expected: &ServiceDatabaseMetadata, + schema_catalog: &crate::SchemaCatalog, ) -> Result<(), ServiceSqliteError> { if expected.state_schema_version().get() != 1 { return Err(metadata_error(MetadataFailureKind::Mismatch)); @@ -329,19 +301,23 @@ pub(crate) async fn write_database_metadata( .begin() .await .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; - sqlx::raw_sql(CREATE_METADATA_SQL) - .execute(&mut *transaction) - .await - .map_err(|_| metadata_error(MetadataFailureKind::AlreadyPresent))?; - sqlx::raw_sql(crate::migration::CREATE_MIGRATION_LEDGER_SQL) - .execute(&mut *transaction) - .await - .map_err(|_source| { - ServiceSqliteError::with_source( - ServiceSqliteErrorKind::Migration, - MigrationLedgerInitializationFailure, - ) - })?; + for statement in crate::integrity::catalog::METADATA_SCHEMA_SQL { + sqlx::query(statement) + .execute(&mut *transaction) + .await + .map_err(|_| metadata_error(MetadataFailureKind::AlreadyPresent))?; + } + for statement in crate::integrity::catalog::MIGRATION_LEDGER_SCHEMA_SQL { + sqlx::query(statement) + .execute(&mut *transaction) + .await + .map_err(|_source| { + ServiceSqliteError::with_source( + ServiceSqliteErrorKind::Migration, + MigrationLedgerInitializationFailure, + ) + })?; + } sqlx::query( "INSERT INTO radroots_service_metadata ( singleton, service_id, instance_id, source_generation, @@ -368,6 +344,12 @@ pub(crate) async fn write_database_metadata( .execute(&mut *transaction) .await .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; + crate::integrity::verify_schema_catalog( + &mut transaction, + schema_catalog, + expected.state_schema_version().get(), + ) + .await?; transaction .commit() .await @@ -533,6 +515,7 @@ mod tests { ServiceSqlitePaths::from_runtime_context(&context).expect("SQLite paths") } + #[cfg(any(target_os = "linux", target_os = "macos"))] fn metadata( paths: &ServiceSqlitePaths, generation_byte: u8, @@ -557,6 +540,16 @@ mod tests { .expect("memory SQLite") } + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn base_schema_catalog() -> crate::SchemaCatalog { + let migrations = crate::MigrationCatalog::new([]).expect("empty migration catalog"); + let digest = crate::SchemaVersionCatalog::computed_digest(1, []) + .expect("base schema snapshot digest"); + let version = + crate::SchemaVersionCatalog::new(1, [], digest).expect("base schema version catalog"); + crate::SchemaCatalog::new(&migrations, [version]).expect("base schema catalog") + } + #[test] fn application_id_and_creation_time_bounds_are_exact() { assert_eq!( @@ -613,7 +606,7 @@ mod tests { let expected = metadata(&paths, 7, 1, 1_700_000_000_000, 0x5244_5351); let mut connection = memory_connection().await; - write_database_metadata(&mut connection, &expected) + write_database_metadata(&mut connection, &expected, &base_schema_catalog()) .await .expect("write metadata"); verify_database_metadata(&mut connection, &expected.identity()) @@ -671,7 +664,7 @@ mod tests { ); } assert_eq!( - write_database_metadata(&mut connection, &expected) + write_database_metadata(&mut connection, &expected, &base_schema_catalog()) .await .expect_err("second write must fail") .kind(), @@ -687,11 +680,56 @@ mod tests { #[cfg(any(target_os = "linux", target_os = "macos"))] #[tokio::test(flavor = "current_thread")] + async fn schema_mismatch_rolls_back_shared_objects_metadata_and_application_id() { + let paths = sqlite_paths("myc", "primary"); + let expected = metadata(&paths, 7, 1, 1_700_000_000_000, 0x5244_5351); + let mut connection = memory_connection().await; + sqlx::query("CREATE TABLE unlisted_service_object (id INTEGER PRIMARY KEY)") + .execute(&mut connection) + .await + .expect("pre-existing service object"); + + let error = write_database_metadata(&mut connection, &expected, &base_schema_catalog()) + .await + .expect_err("schema mismatch must roll back initialization transaction"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert_eq!(read_application_id(&mut connection).await.unwrap(), 0); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM sqlite_schema + WHERE name IN ( + 'radroots_service_metadata', + 'radroots_service_metadata_guard_update', + 'radroots_service_metadata_no_delete', + 'schema_migrations', + 'schema_migrations_no_update', + 'schema_migrations_no_delete' + )", + ) + .fetch_one(&mut connection) + .await + .expect("shared schema rollback evidence"), + 0 + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM sqlite_schema + WHERE type = 'table' AND name = 'unlisted_service_object'", + ) + .fetch_one(&mut connection) + .await + .expect("pre-existing object evidence"), + 1 + ); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] async fn every_identity_dimension_must_match() { let paths = sqlite_paths("myc", "primary"); let expected = metadata(&paths, 7, 1, 1_700_000_000_000, 0x5244_5351); let mut connection = memory_connection().await; - write_database_metadata(&mut connection, &expected) + write_database_metadata(&mut connection, &expected, &base_schema_catalog()) .await .expect("write metadata"); @@ -746,14 +784,14 @@ mod tests { let mut newer_state = memory_connection().await; let version_two = metadata(&paths, 7, 2, 1_700_000_000_001, 0x5244_5351); assert_eq!( - write_database_metadata(&mut newer_state, &version_two) + write_database_metadata(&mut newer_state, &version_two, &base_schema_catalog()) .await .expect_err("fresh metadata must start at v1") .kind(), ServiceSqliteErrorKind::Metadata ); let version_one = metadata(&paths, 7, 1, 1_700_000_000_001, 0x5244_5351); - write_database_metadata(&mut newer_state, &version_one) + write_database_metadata(&mut newer_state, &version_one, &base_schema_catalog()) .await .expect("write v1 metadata"); sqlx::query( diff --git a/crates/service_sqlite/src/migration.rs b/crates/service_sqlite/src/migration.rs @@ -20,7 +20,7 @@ use sha2::{Digest, Sha256}; use sqlx::{Connection, Row, SqliteConnection}; #[cfg(any(target_os = "linux", target_os = "macos"))] -use crate::{ServiceSqliteError, ServiceSqliteErrorKind}; +use crate::{SchemaCatalog, ServiceSqliteError, ServiceSqliteErrorKind}; const MIGRATION_CONTENT_DOMAIN: &[u8] = b"radroots.service_sqlite.migration_content.v1\0"; const MIGRATION_CATALOG_DOMAIN: &[u8] = b"radroots.service_sqlite.migration_catalog.v1\0"; @@ -30,39 +30,6 @@ const MAX_MIGRATION_CONTENT_BYTES: usize = 1024 * 1024; const MAX_MIGRATION_COUNT: usize = 4096; const MAX_MIGRATION_BUILD_ID_UTF8_BYTES: usize = 128; -#[cfg(any(target_os = "linux", target_os = "macos"))] -pub(crate) const CREATE_MIGRATION_LEDGER_SQL: &str = r#" -CREATE TABLE schema_migrations ( - version INTEGER NOT NULL PRIMARY KEY - CHECK (version BETWEEN 2 AND 4294967295), - name TEXT NOT NULL UNIQUE - CHECK (length(CAST(name AS BLOB)) BETWEEN 1 AND 128), - checksum BLOB NOT NULL CHECK (length(checksum) = 32), - applied_at_unix_s INTEGER NOT NULL CHECK (applied_at_unix_s BETWEEN 0 AND 9223372036854775807), - service_version TEXT NOT NULL CHECK (length(CAST(service_version AS BLOB)) BETWEEN 1 AND 128), - service_commit TEXT NOT NULL CHECK (length(CAST(service_commit AS BLOB)) = 40), - lib_revision TEXT NOT NULL CHECK (length(CAST(lib_revision AS BLOB)) = 40), - rust_version TEXT NOT NULL CHECK (length(CAST(rust_version AS BLOB)) BETWEEN 1 AND 128), - target TEXT NOT NULL CHECK (length(CAST(target AS BLOB)) BETWEEN 1 AND 128), - feature_profile TEXT NOT NULL CHECK (length(CAST(feature_profile AS BLOB)) BETWEEN 1 AND 128), - config_contract_version INTEGER NOT NULL CHECK (config_contract_version BETWEEN 1 AND 4294967295), - state_contract_version INTEGER NOT NULL CHECK (state_contract_version BETWEEN 1 AND 4294967295), - admin_contract_version INTEGER NOT NULL CHECK (admin_contract_version BETWEEN 1 AND 4294967295), - status_contract_version INTEGER NOT NULL CHECK (status_contract_version BETWEEN 1 AND 4294967295), - provider_contract_version INTEGER NOT NULL CHECK (provider_contract_version BETWEEN 1 AND 4294967295) -) STRICT; -CREATE TRIGGER schema_migrations_no_update -BEFORE UPDATE ON schema_migrations -BEGIN - SELECT RAISE(ABORT, 'migration history is immutable'); -END; -CREATE TRIGGER schema_migrations_no_delete -BEFORE DELETE ON schema_migrations -BEGIN - SELECT RAISE(ABORT, 'migration history is immutable'); -END; -"#; - /// A bounded stable lower-snake migration name. #[derive(Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Hash)] pub struct MigrationName(&'static str); @@ -926,14 +893,23 @@ struct AppliedMigration { pub(crate) async fn verify_migration_history( connection: &mut SqliteConnection, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, require_current: bool, ) -> Result<u32, ServiceSqliteError> { + if !schema_catalog.matches_migrations(catalog) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); + } let mut transaction = connection .begin() .await .map_err(|source| migration_source(MigrationFailureKind::HistoryCorrupt, source))?; - let result = - verify_migration_history_snapshot(&mut transaction, catalog, require_current).await; + let result = verify_migration_history_snapshot( + &mut transaction, + catalog, + schema_catalog, + require_current, + ) + .await; let rollback = transaction.rollback().await; match (result, rollback) { (Err(error), _) => Err(error), @@ -949,11 +925,13 @@ pub(crate) async fn verify_migration_history( async fn verify_migration_history_snapshot( connection: &mut SqliteConnection, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, require_current: bool, ) -> Result<u32, ServiceSqliteError> { let version = read_state_schema_version(connection).await?; let history = read_migration_history(connection).await?; validate_migration_prefix(catalog, version, &history)?; + crate::integrity::verify_schema_catalog(connection, schema_catalog, version).await?; if require_current && version != catalog.current_version() { return Err(migration_error(MigrationFailureKind::CatalogMismatch)); } @@ -964,6 +942,7 @@ async fn verify_migration_history_snapshot( pub(crate) async fn apply_governed_migrations<V>( connection: &mut SqliteConnection, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, applied_at: MigrationAppliedAtUnixSeconds, build: &MigrationBuildIdentity, callback_bindings: &[MigrationCallbackBinding], @@ -976,6 +955,7 @@ where apply_governed_migrations_with_observer( connection, catalog, + schema_catalog, applied_at, build, callback_bindings, @@ -986,9 +966,11 @@ where } #[cfg(any(target_os = "linux", target_os = "macos"))] +#[allow(clippy::too_many_arguments)] async fn apply_governed_migrations_with_observer<V, O>( connection: &mut SqliteConnection, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, applied_at: MigrationAppliedAtUnixSeconds, build: &MigrationBuildIdentity, callback_bindings: &[MigrationCallbackBinding], @@ -999,6 +981,9 @@ where V: FnMut() -> Result<(), ServiceSqliteError>, O: FnMut() -> Result<(), ServiceSqliteError>, { + if !schema_catalog.matches_migrations(catalog) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); + } let callbacks = validate_callback_bindings(catalog, callback_bindings)?; let mut initial_version = None; let mut applied_count = 0_u32; @@ -1017,7 +1002,8 @@ where let mut transaction = transaction_result .map_err(|source| migration_source(MigrationFailureKind::Execution, source))?; let transactional_result = - verify_migration_history_snapshot(&mut transaction, catalog, false).await; + verify_migration_history_snapshot(&mut transaction, catalog, schema_catalog, false) + .await; validate_authority()?; let current = transactional_result?; initial_version.get_or_insert(current); @@ -1047,6 +1033,22 @@ where let transaction_result = assert_governed_transaction(&mut transaction).await; validate_authority()?; transaction_result?; + let policy_result = read_connection_policy(&mut transaction).await; + validate_authority()?; + if policy_result? != initial_policy { + return Err(migration_error(MigrationFailureKind::Execution)); + } + let transaction_result = assert_governed_transaction(&mut transaction).await; + validate_authority()?; + transaction_result?; + let schema_result = crate::integrity::verify_schema_catalog( + &mut transaction, + schema_catalog, + descriptor.target_version(), + ) + .await; + validate_authority()?; + schema_result?; let insert_result = insert_migration_row(&mut transaction, descriptor, applied_at, build).await; validate_authority()?; @@ -1086,7 +1088,7 @@ where } validate_authority()?; - let final_result = verify_migration_history(connection, catalog, true).await; + let final_result = verify_migration_history(connection, catalog, schema_catalog, true).await; validate_authority()?; let final_version = final_result?; Ok(MigrationApplicationOutcome { @@ -1538,6 +1540,72 @@ mod tests { .expect("valid SQL descriptor") } + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn table_object(name: &'static str, sql: &'static str) -> crate::SchemaObject { + crate::SchemaObject::new( + crate::SchemaObjectKind::Table, + name, + name, + sql, + crate::SchemaObject::computed_digest(crate::SchemaObjectKind::Table, name, name, sql) + .expect("schema table digest"), + ) + .expect("schema table") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn schema_catalog( + migrations: &MigrationCatalog, + versions: Vec<Vec<crate::SchemaObject>>, + ) -> crate::SchemaCatalog { + let versions = versions + .into_iter() + .enumerate() + .map(|(index, objects)| { + let version = u32::try_from(index + 1).expect("schema version"); + let digest = + crate::SchemaVersionCatalog::computed_digest(version, objects.iter().cloned()) + .expect("schema digest"); + crate::SchemaVersionCatalog::new(version, objects, digest) + .expect("schema version catalog") + }) + .collect::<Vec<_>>(); + crate::SchemaCatalog::new(migrations, versions).expect("schema catalog") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn unchanged_schema_catalog(migrations: &MigrationCatalog) -> crate::SchemaCatalog { + schema_catalog( + migrations, + (0..migrations.current_version()) + .map(|_| Vec::new()) + .collect(), + ) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn alpha_schema_catalog(migrations: &MigrationCatalog) -> crate::SchemaCatalog { + const ALPHA_SQL: &str = "CREATE TABLE alpha (id INTEGER PRIMARY KEY)"; + let alpha = table_object("alpha", ALPHA_SQL); + let mut versions = vec![Vec::new()]; + versions.extend((1..migrations.current_version()).map(|_| vec![alpha.clone()])); + schema_catalog(migrations, versions) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn alpha_beta_schema_catalog( + migrations: &MigrationCatalog, + beta_sql: &'static str, + ) -> crate::SchemaCatalog { + const ALPHA_SQL: &str = "CREATE TABLE alpha (id INTEGER PRIMARY KEY)"; + let alpha = table_object("alpha", ALPHA_SQL); + let beta = table_object("beta", beta_sql); + schema_catalog( + migrations, + vec![Vec::new(), vec![alpha.clone()], vec![alpha, beta]], + ) + } + fn build_identity() -> MigrationBuildIdentity { MigrationBuildIdentity::new( "0.1.0-alpha", @@ -1598,7 +1666,9 @@ mod tests { let mut connection = SqliteConnection::connect_with(&options) .await .expect("test SQLite"); - crate::metadata::write_database_metadata(&mut connection, &metadata) + let migrations = MigrationCatalog::new([]).expect("empty migration catalog"); + let schema_catalog = unchanged_schema_catalog(&migrations); + crate::metadata::write_database_metadata(&mut connection, &metadata, &schema_catalog) .await .expect("initialize metadata and ledger"); connection @@ -1878,6 +1948,7 @@ mod tests { ); let catalog = MigrationCatalog::new([sql_descriptor, callback_descriptor]).expect("catalog"); + let schema_catalog = alpha_schema_catalog(&catalog); let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(); let build = build_identity(); let mut connection = initialized_memory_database().await; @@ -1886,6 +1957,7 @@ mod tests { let outcome = apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &[callback], @@ -1981,6 +2053,7 @@ mod tests { let reopened = apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, MigrationAppliedAtUnixSeconds::new(1_900_000_000).unwrap(), &build, &[callback], @@ -2009,6 +2082,7 @@ mod tests { let first = sql(2, "create_alpha", SQL_TWO); let invalid = sql(3, "create_beta", INVALID_SQL); let invalid_catalog = MigrationCatalog::new([first.clone(), invalid]).unwrap(); + let invalid_schema_catalog = alpha_schema_catalog(&invalid_catalog); let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(); let build = build_identity(); let mut connection = initialized_memory_database().await; @@ -2017,6 +2091,7 @@ mod tests { let error = apply_governed_migrations( &mut connection, &invalid_catalog, + &invalid_schema_catalog, applied_at, &build, &[], @@ -2042,9 +2117,14 @@ mod tests { let recovered_catalog = MigrationCatalog::new([first, sql(3, "create_beta", RECOVERY_SQL)]).unwrap(); + let recovered_schema_catalog = alpha_beta_schema_catalog( + &recovered_catalog, + "CREATE TABLE beta (id INTEGER PRIMARY KEY)", + ); let recovered = apply_governed_migrations( &mut connection, &recovered_catalog, + &recovered_schema_catalog, applied_at, &build, &[], @@ -2059,6 +2139,91 @@ mod tests { #[cfg(any(target_os = "linux", target_os = "macos"))] #[tokio::test(flavor = "current_thread")] + async fn schema_mismatch_before_or_after_execution_never_commits_a_step() { + let catalog = MigrationCatalog::new([sql(2, "create_alpha", SQL_TWO)]).unwrap(); + let wrong_target_catalog = unchanged_schema_catalog(&catalog); + let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(); + let build = build_identity(); + let mut validate = || Ok(()); + + let mut after_execution = initialized_memory_database().await; + let error = apply_governed_migrations( + &mut after_execution, + &catalog, + &wrong_target_catalog, + applied_at, + &build, + &[], + &mut validate, + ) + .await + .expect_err("target schema mismatch must roll back"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert_eq!( + read_state_schema_version(&mut after_execution) + .await + .unwrap(), + 1 + ); + assert!( + read_migration_history(&mut after_execution) + .await + .unwrap() + .is_empty() + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'alpha'", + ) + .fetch_one(&mut after_execution) + .await + .unwrap(), + 0 + ); + + let mut before_execution = initialized_memory_database().await; + sqlx::query("CREATE TABLE unexpected (value INTEGER)") + .execute(&mut before_execution) + .await + .unwrap(); + let expected_target = alpha_schema_catalog(&catalog); + let error = apply_governed_migrations( + &mut before_execution, + &catalog, + &expected_target, + applied_at, + &build, + &[], + &mut validate, + ) + .await + .expect_err("current schema mismatch must fail before execution"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert_eq!( + read_state_schema_version(&mut before_execution) + .await + .unwrap(), + 1 + ); + assert!( + read_migration_history(&mut before_execution) + .await + .unwrap() + .is_empty() + ); + assert_eq!( + sqlx::query_scalar::<_, i64>( + "SELECT COUNT(*) FROM sqlite_schema WHERE type = 'table' AND name = 'alpha'", + ) + .fetch_one(&mut before_execution) + .await + .unwrap(), + 0 + ); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] async fn transaction_control_cannot_escape_schema_ledger_metadata_atomicity() { const COMMIT_ESCAPE_SQL: &str = "CREATE TABLE sql_leaked (id INTEGER PRIMARY KEY); COMMIT; SELECT no_such_function();"; @@ -2076,10 +2241,12 @@ mod tests { let mut sql_connection = initialized_file_database(&sql_path).await; let sql_catalog = MigrationCatalog::new([sql(2, "attempt_commit_escape", COMMIT_ESCAPE_SQL)]).unwrap(); + let sql_schema_catalog = unchanged_schema_catalog(&sql_catalog); let mut validate = || Ok(()); let sql_error = apply_governed_migrations( &mut sql_connection, &sql_catalog, + &sql_schema_catalog, applied_at, &build, &[], @@ -2099,9 +2266,11 @@ mod tests { REPLACEMENT_ESCAPE_SQL, )]) .unwrap(); + let replacement_schema_catalog = unchanged_schema_catalog(&replacement_catalog); let replacement_error = apply_governed_migrations( &mut replacement_connection, &replacement_catalog, + &replacement_schema_catalog, applied_at, &build, &[], @@ -2133,9 +2302,11 @@ mod tests { rollback_escape_callback, ); let callback_catalog = MigrationCatalog::new([callback_descriptor]).unwrap(); + let callback_schema_catalog = unchanged_schema_catalog(&callback_catalog); let callback_error = apply_governed_migrations( &mut callback_connection, &callback_catalog, + &callback_schema_catalog, applied_at, &build, &[callback], @@ -2154,9 +2325,11 @@ mod tests { let mut policy_connection = initialized_file_database(&policy_path).await; let policy_catalog = MigrationCatalog::new([sql(2, "attempt_policy_escape", POLICY_ESCAPE_SQL)]).unwrap(); + let policy_schema_catalog = unchanged_schema_catalog(&policy_catalog); let policy_error = apply_governed_migrations( &mut policy_connection, &policy_catalog, + &policy_schema_catalog, applied_at, &build, &[], @@ -2211,6 +2384,7 @@ mod tests { ) .unwrap(); let catalog = MigrationCatalog::new([sql_descriptor, callback_descriptor.clone()]).unwrap(); + let schema_catalog = alpha_schema_catalog(&catalog); let pending_binding = MigrationCallbackBinding::new( 3, callback_descriptor.name(), @@ -2232,6 +2406,7 @@ mod tests { let mut application = Box::pin(apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &pending_bindings, @@ -2263,6 +2438,7 @@ mod tests { let recovered = apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &[working_binding], @@ -2288,6 +2464,7 @@ mod tests { let first = sql(2, "create_alpha", SQL_TWO); let second = sql(3, "create_beta", "CREATE TABLE beta (id INTEGER);"); let catalog = MigrationCatalog::new([first, second]).unwrap(); + let schema_catalog = alpha_beta_schema_catalog(&catalog, "CREATE TABLE beta (id INTEGER)"); let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(); let build = build_identity(); let mut connection = initialized_memory_database().await; @@ -2297,6 +2474,7 @@ mod tests { let error = apply_governed_migrations_with_observer( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &[], @@ -2315,6 +2493,7 @@ mod tests { let resumed = apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &[], @@ -2347,6 +2526,7 @@ mod tests { ) .unwrap(); let catalog = MigrationCatalog::new([callback_descriptor.clone()]).unwrap(); + let schema_catalog = unchanged_schema_catalog(&catalog); let correct = MigrationCallbackBinding::new( 2, callback_descriptor.name(), @@ -2368,6 +2548,7 @@ mod tests { let error = apply_governed_migrations( &mut connection, &catalog, + &schema_catalog, applied_at, &build, &bindings, @@ -2409,7 +2590,7 @@ mod tests { .execute(&mut connection) .await .unwrap(); - let error = verify_migration_history(&mut connection, &catalog, true) + let error = verify_migration_history(&mut connection, &catalog, &schema_catalog, true) .await .expect_err("mismatched history"); assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); @@ -2421,6 +2602,7 @@ mod tests { let first = sql(2, "create_alpha", SQL_TWO); let second = sql(3, "create_beta", "CREATE TABLE beta (id INTEGER);"); let catalog = MigrationCatalog::new([first.clone(), second.clone()]).unwrap(); + let schema_catalog = unchanged_schema_catalog(&catalog); let mut missing = initialized_memory_database().await; sqlx::query( @@ -2430,7 +2612,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut missing, &catalog, false) + verify_migration_history(&mut missing, &catalog, &schema_catalog, false) .await .expect_err("missing row") .kind(), @@ -2449,7 +2631,7 @@ mod tests { ) .await; assert_eq!( - verify_migration_history(&mut extra, &catalog, false) + verify_migration_history(&mut extra, &catalog, &schema_catalog, false) .await .expect_err("extra row") .kind(), @@ -2483,7 +2665,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut reordered, &catalog, true) + verify_migration_history(&mut reordered, &catalog, &schema_catalog, true) .await .expect_err("reordered names and checksums") .kind(), @@ -2498,7 +2680,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut newer, &catalog, false) + verify_migration_history(&mut newer, &catalog, &schema_catalog, false) .await .expect_err("newer schema") .kind(), @@ -2524,7 +2706,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut corrupt, &catalog, false) + verify_migration_history(&mut corrupt, &catalog, &schema_catalog, false) .await .expect_err("corrupt time or build") .kind(), @@ -2538,6 +2720,7 @@ mod tests { async fn oversized_corrupt_history_is_bounded_before_decode() { let descriptor = sql(2, "create_alpha", SQL_TWO); let catalog = MigrationCatalog::new([descriptor.clone()]).unwrap(); + let schema_catalog = unchanged_schema_catalog(&catalog); let oversized_text = "a".repeat(4 * 1024 * 1024); for (column, update) in [ ( @@ -2592,7 +2775,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut connection, &catalog, true) + verify_migration_history(&mut connection, &catalog, &schema_catalog, true) .await .expect_err("oversized text must fail before decode") .kind(), @@ -2619,7 +2802,7 @@ mod tests { .await .unwrap(); assert_eq!( - verify_migration_history(&mut checksum, &catalog, true) + verify_migration_history(&mut checksum, &catalog, &schema_catalog, true) .await .expect_err("oversized checksum must fail before decode") .kind(), diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs @@ -30,7 +30,7 @@ use sqlx::{ #[cfg(any(target_os = "linux", target_os = "macos"))] use crate::{ - MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, MigrationCatalog, + MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, MigrationCatalog, SchemaCatalog, ServiceDatabaseIdentity, ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind, WriterAuthority, }; @@ -191,11 +191,13 @@ struct PrivateConnectionPool { binding: DirectoryBinding, paths: ServiceSqlitePaths, catalog: MigrationCatalog, + schema_catalog: SchemaCatalog, authority: Option<WriterAuthority>, inspection_guard: Option<ReadOnlyInspectionGuard>, authority_failure: Arc<AtomicBool>, metadata_failure: Arc<AtomicBool>, migration_failure: Arc<AtomicBool>, + integrity_failure: Arc<AtomicBool>, pragma_failure: Arc<AtomicBool>, } @@ -206,17 +208,13 @@ struct PrivateConnectionPool { )] impl PrivateConnectionPool { fn connection_failure_kind(&self) -> ServiceSqliteErrorKind { - if self.authority_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Authority - } else if self.metadata_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Metadata - } else if self.migration_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Migration - } else if self.pragma_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Pragma - } else { - ServiceSqliteErrorKind::Open - } + connection_failure_kind( + self.authority_failure.load(Ordering::Acquire), + self.metadata_failure.load(Ordering::Acquire), + self.migration_failure.load(Ordering::Acquire), + self.integrity_failure.load(Ordering::Acquire), + self.pragma_failure.load(Ordering::Acquire), + ) } async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> { @@ -225,8 +223,13 @@ impl PrivateConnectionPool { self.binding.validate(&self.paths)?; let mut connection = result.map_err(|source| connection_source(self.connection_failure_kind(), source))?; - let history = - crate::migration::verify_migration_history(&mut connection, &self.catalog, true).await; + let history = crate::migration::verify_migration_history( + &mut connection, + &self.catalog, + &self.schema_catalog, + true, + ) + .await; self.binding.validate(&self.paths)?; history?; Ok(connection) @@ -260,6 +263,7 @@ impl PrivateConnectionPool { let result = crate::migration::apply_governed_migrations( &mut connection, &self.catalog, + &self.schema_catalog, applied_at, build, callbacks, @@ -279,6 +283,29 @@ impl PrivateConnectionPool { } #[cfg(any(target_os = "linux", target_os = "macos"))] +const fn connection_failure_kind( + authority: bool, + metadata: bool, + migration: bool, + integrity: bool, + pragma: bool, +) -> ServiceSqliteErrorKind { + if authority { + ServiceSqliteErrorKind::Authority + } else if metadata { + ServiceSqliteErrorKind::Metadata + } else if migration { + ServiceSqliteErrorKind::Migration + } else if integrity { + ServiceSqliteErrorKind::Integrity + } else if pragma { + ServiceSqliteErrorKind::Pragma + } else { + ServiceSqliteErrorKind::Open + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] #[allow( dead_code, reason = "Step 056 keeps pool opening private until the Step 061 host boundary" @@ -287,6 +314,7 @@ async fn open_existing_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, mode: OpenMode, policy: ServiceSqliteConnectionOptions, ) -> Result<PrivateConnectionPool, ServiceSqliteError> { @@ -305,6 +333,7 @@ async fn open_existing_connection_pool( paths, identity, catalog, + schema_catalog, mode, policy, authority, @@ -322,6 +351,7 @@ async fn open_initialized_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, policy: ServiceSqliteConnectionOptions, authority: WriterAuthority, ) -> Result<PrivateConnectionPool, ServiceSqliteError> { @@ -330,6 +360,7 @@ async fn open_initialized_connection_pool( paths, identity, catalog, + schema_catalog, OpenMode::Initialize, policy, Some(authority), @@ -340,6 +371,7 @@ async fn open_initialized_connection_pool( #[cfg(any(target_os = "linux", target_os = "macos"))] #[allow( + clippy::too_many_arguments, dead_code, reason = "Step 056 keeps pool construction private until the Step 061 host boundary" )] @@ -347,6 +379,7 @@ async fn open_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, mode: OpenMode, policy: ServiceSqliteConnectionOptions, authority: Option<WriterAuthority>, @@ -358,6 +391,9 @@ async fn open_connection_pool( if identity.supported_state_schema_version().get() != catalog.current_version() { return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Migration)); } + if !schema_catalog.matches_migrations(catalog) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); + } let binding = match (mode, authority.as_ref(), inspection_guard.as_ref()) { (OpenMode::Initialize | OpenMode::ReadWriteExisting, Some(authority), None) => { authority.validate_for(paths)?; @@ -390,6 +426,7 @@ async fn open_connection_pool( let preflight_history = crate::migration::verify_migration_history( &mut preflight, catalog, + schema_catalog, mode == OpenMode::ReadOnlyInspection, ) .await; @@ -412,19 +449,25 @@ async fn open_connection_pool( let after_metadata = identity.clone(); let before_metadata = identity.clone(); let retained_catalog = catalog.clone(); + let retained_schema_catalog = schema_catalog.clone(); let after_catalog = catalog.clone(); + let after_schema_catalog = schema_catalog.clone(); let before_catalog = catalog.clone(); + let before_schema_catalog = schema_catalog.clone(); let authority_failure = Arc::new(AtomicBool::new(false)); let metadata_failure = Arc::new(AtomicBool::new(false)); let migration_failure = Arc::new(AtomicBool::new(false)); + let integrity_failure = Arc::new(AtomicBool::new(false)); let pragma_failure = Arc::new(AtomicBool::new(false)); let after_authority_failure = Arc::clone(&authority_failure); let after_metadata_failure = Arc::clone(&metadata_failure); let after_migration_failure = Arc::clone(&migration_failure); + let after_integrity_failure = Arc::clone(&integrity_failure); let after_pragma_failure = Arc::clone(&pragma_failure); let before_authority_failure = Arc::clone(&authority_failure); let before_metadata_failure = Arc::clone(&metadata_failure); let before_migration_failure = Arc::clone(&migration_failure); + let before_integrity_failure = Arc::clone(&integrity_failure); let before_pragma_failure = Arc::clone(&pragma_failure); let pool_result = SqlitePoolOptions::new() .min_connections(1) @@ -438,9 +481,11 @@ async fn open_connection_pool( let paths = after_paths.clone(); let metadata = after_metadata.clone(); let catalog = after_catalog.clone(); + let schema_catalog = after_schema_catalog.clone(); let authority_failure = Arc::clone(&after_authority_failure); let metadata_failure = Arc::clone(&after_metadata_failure); let migration_failure = Arc::clone(&after_migration_failure); + let integrity_failure = Arc::clone(&after_integrity_failure); let pragma_failure = Arc::clone(&after_pragma_failure); Box::pin(async move { if binding.validate(&paths).is_err() { @@ -483,6 +528,7 @@ async fn open_connection_pool( let migration_result = crate::migration::verify_migration_history( connection, &catalog, + &schema_catalog, after_mode == OpenMode::ReadOnlyInspection, ) .await; @@ -493,7 +539,14 @@ async fn open_connection_pool( )); } if migration_result.is_err() { - migration_failure.store(true, Ordering::Release); + if migration_result + .as_ref() + .is_err_and(|error| error.kind() == ServiceSqliteErrorKind::Integrity) + { + integrity_failure.store(true, Ordering::Release); + } else { + migration_failure.store(true, Ordering::Release); + } return Err(sqlx::Error::Protocol( "SQLite migration history mismatch".to_owned(), )); @@ -506,9 +559,11 @@ async fn open_connection_pool( let paths = before_paths.clone(); let metadata = before_metadata.clone(); let catalog = before_catalog.clone(); + let schema_catalog = before_schema_catalog.clone(); let authority_failure = Arc::clone(&before_authority_failure); let metadata_failure = Arc::clone(&before_metadata_failure); let migration_failure = Arc::clone(&before_migration_failure); + let integrity_failure = Arc::clone(&before_integrity_failure); let pragma_failure = Arc::clone(&before_pragma_failure); Box::pin(async move { if binding.validate(&paths).is_err() { @@ -549,6 +604,7 @@ async fn open_connection_pool( let migration_result = crate::migration::verify_migration_history( connection, &catalog, + &schema_catalog, before_mode == OpenMode::ReadOnlyInspection, ) .await; @@ -559,7 +615,14 @@ async fn open_connection_pool( )); } if migration_result.is_err() { - migration_failure.store(true, Ordering::Release); + if migration_result + .as_ref() + .is_err_and(|error| error.kind() == ServiceSqliteErrorKind::Integrity) + { + integrity_failure.store(true, Ordering::Release); + } else { + migration_failure.store(true, Ordering::Release); + } return Err(sqlx::Error::Protocol( "SQLite migration history mismatch".to_owned(), )); @@ -571,17 +634,13 @@ async fn open_connection_pool( .await; pool_binding.validate(paths)?; let pool = pool_result.map_err(|source| { - let kind = if authority_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Authority - } else if metadata_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Metadata - } else if migration_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Migration - } else if pragma_failure.load(Ordering::Acquire) { - ServiceSqliteErrorKind::Pragma - } else { - ServiceSqliteErrorKind::Open - }; + let kind = connection_failure_kind( + authority_failure.load(Ordering::Acquire), + metadata_failure.load(Ordering::Acquire), + migration_failure.load(Ordering::Acquire), + integrity_failure.load(Ordering::Acquire), + pragma_failure.load(Ordering::Acquire), + ); connection_source(kind, source) })?; @@ -590,11 +649,13 @@ async fn open_connection_pool( binding: retained_binding, paths: paths.clone(), catalog: retained_catalog, + schema_catalog: retained_schema_catalog, authority, inspection_guard, authority_failure, metadata_failure, migration_failure, + integrity_failure, pragma_failure, }) } @@ -1100,13 +1161,14 @@ impl DirectoryBinding { #[cfg(test)] mod tests { - use std::{num::NonZeroU32, path::PathBuf}; + use std::path::PathBuf; #[cfg(any(target_os = "linux", target_os = "macos"))] use std::{ collections::BTreeMap, convert::Infallible, fs, + num::NonZeroU32, os::unix::fs::{PermissionsExt, symlink}, sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}, time::{Duration, SystemTime}, @@ -1116,8 +1178,10 @@ mod tests { RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, RadrootsPlatform, RuntimeContextBootstrap, RuntimeContextSource, }; + #[cfg(any(target_os = "linux", target_os = "macos"))] use radroots_storage::event::SourceGeneration; + #[cfg(any(target_os = "linux", target_os = "macos"))] use crate::{ServiceDatabaseMetadata, ServiceSqliteApplicationId}; use super::*; @@ -1168,6 +1232,44 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + fn schema_catalog( + migrations: &MigrationCatalog, + versions: Vec<Vec<crate::SchemaObject>>, + ) -> crate::SchemaCatalog { + let versions = versions + .into_iter() + .enumerate() + .map(|(index, objects)| { + let version = u32::try_from(index + 1).expect("schema version"); + let digest = + crate::SchemaVersionCatalog::computed_digest(version, objects.iter().cloned()) + .expect("schema digest"); + crate::SchemaVersionCatalog::new(version, objects, digest).expect("schema version") + }) + .collect::<Vec<_>>(); + crate::SchemaCatalog::new(migrations, versions).expect("schema catalog") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn base_schema_catalog() -> crate::SchemaCatalog { + schema_catalog(&base_catalog(), vec![Vec::new()]) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn single_table_schema_catalog(name: &'static str, sql: &'static str) -> crate::SchemaCatalog { + let table = crate::SchemaObject::new( + crate::SchemaObjectKind::Table, + name, + name, + sql, + crate::SchemaObject::computed_digest(crate::SchemaObjectKind::Table, name, name, sql) + .expect("schema table digest"), + ) + .expect("schema table"); + schema_catalog(&base_catalog(), vec![vec![table]]) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] fn migration_catalog() -> MigrationCatalog { const CREATE: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL);"; const CALLBACK_DEFINITION: &[u8] = b"callback:migration_probe:v1"; @@ -1191,6 +1293,29 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + fn migration_schema_catalog() -> crate::SchemaCatalog { + const SQL: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL)"; + let table = crate::SchemaObject::new( + crate::SchemaObjectKind::Table, + "migration_probe", + "migration_probe", + SQL, + crate::SchemaObject::computed_digest( + crate::SchemaObjectKind::Table, + "migration_probe", + "migration_probe", + SQL, + ) + .expect("migration schema digest"), + ) + .expect("migration schema table"); + schema_catalog( + &migration_catalog(), + vec![Vec::new(), vec![table.clone()], vec![table]], + ) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] fn migration_build() -> MigrationBuildIdentity { MigrationBuildIdentity::new( "0.1.0-alpha", @@ -1277,10 +1402,12 @@ mod tests { fs::create_dir_all(paths.state_database().parent().expect("state directory")) .expect("create state directory"); let metadata = database_metadata(&paths); + let schema_catalog = base_schema_catalog(); let authority = crate::initialize_database( &paths, OpenMode::Initialize, &metadata, + &schema_catalog, |database_path| async move { let options = SqliteConnectOptions::new() .filename(database_path) @@ -1304,10 +1431,16 @@ mod tests { policy: ServiceSqliteConnectionOptions, ) -> (ServiceSqlitePaths, PrivateConnectionPool) { let (paths, identity, authority) = initialized_authority(root, "primary").await; - let pool = - open_initialized_connection_pool(&paths, &identity, &base_catalog(), policy, authority) - .await - .expect("open initialized pool"); + let pool = open_initialized_connection_pool( + &paths, + &identity, + &base_catalog(), + &base_schema_catalog(), + policy, + authority, + ) + .await + .expect("open initialized pool"); (paths, pool) } @@ -1328,9 +1461,17 @@ mod tests { NonZeroU32::new(catalog.current_version()).expect("catalog version"), base_identity.application_id(), ); - let pool = open_initialized_connection_pool(&paths, &identity, catalog, policy, authority) - .await - .expect("open migration pool"); + let schema_catalog = migration_schema_catalog(); + let pool = open_initialized_connection_pool( + &paths, + &identity, + catalog, + &schema_catalog, + policy, + authority, + ) + .await + .expect("open migration pool"); (paths, identity, pool) } @@ -1558,6 +1699,7 @@ mod tests { &paths, &wrong, &base_catalog(), + &base_schema_catalog(), OpenMode::ReadWriteExisting, policy, ) @@ -1569,6 +1711,72 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + #[test] + fn connection_failure_precedence_is_exact() { + assert_eq!( + connection_failure_kind(false, false, false, false, false), + ServiceSqliteErrorKind::Open + ); + assert_eq!( + connection_failure_kind(false, false, false, false, true), + ServiceSqliteErrorKind::Pragma + ); + assert_eq!( + connection_failure_kind(false, false, false, true, true), + ServiceSqliteErrorKind::Integrity + ); + assert_eq!( + connection_failure_kind(false, false, true, true, true), + ServiceSqliteErrorKind::Migration + ); + assert_eq!( + connection_failure_kind(false, true, true, true, true), + ServiceSqliteErrorKind::Metadata + ); + assert_eq!( + connection_failure_kind(true, true, true, true, true), + ServiceSqliteErrorKind::Authority + ); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] + async fn schema_drift_is_rejected_on_checkout_and_fresh_open() { + let directory = tempfile::tempdir().expect("temporary directory"); + let policy = ServiceSqliteConnectionOptions::reviewed(); + let (paths, pool) = initialized_pool(directory.path(), policy).await; + let identity = database_metadata(&paths).identity(); + let mut connection = pool.acquire().await.expect("connection"); + sqlx::query("CREATE TABLE unlisted (value INTEGER)") + .execute(&mut *connection) + .await + .expect("create unlisted object"); + drop(connection); + + let error = pool + .acquire() + .await + .expect_err("checkout must reject schema drift"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + let authority = pool.close().await.expect("writer authority"); + drop(authority); + + let result = open_existing_connection_pool( + &paths, + &identity, + &base_catalog(), + &base_schema_catalog(), + OpenMode::ReadWriteExisting, + policy, + ) + .await; + let Err(error) = result else { + panic!("fresh open must reject schema drift"); + }; + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] #[tokio::test(flavor = "current_thread")] async fn migration_entry_preserves_post_open_ledger_drift_classification() { let directory = tempfile::tempdir().expect("temporary directory"); @@ -1629,6 +1837,7 @@ mod tests { &paths, &identity, &catalog, + &migration_schema_catalog(), OpenMode::ReadOnlyInspection, policy, ) @@ -1642,6 +1851,7 @@ mod tests { &paths, &identity, &catalog, + &migration_schema_catalog(), OpenMode::ReadWriteExisting, policy, ) @@ -1667,6 +1877,7 @@ mod tests { &paths, &identity, &catalog, + &migration_schema_catalog(), OpenMode::ReadOnlyInspection, policy, ) @@ -1783,6 +1994,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &base_schema_catalog(), mode, ServiceSqliteConnectionOptions::reviewed(), ) @@ -1797,6 +2009,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &base_schema_catalog(), OpenMode::Initialize, ServiceSqliteConnectionOptions::reviewed(), ) @@ -1831,6 +2044,7 @@ mod tests { &other, &metadata, &base_catalog(), + &base_schema_catalog(), ServiceSqliteConnectionOptions::reviewed(), authority, ) @@ -1851,6 +2065,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &base_schema_catalog(), ServiceSqliteConnectionOptions::reviewed(), authority, ) @@ -1907,6 +2122,7 @@ mod tests { &symlink_paths, &symlink_metadata, &base_catalog(), + &base_schema_catalog(), OpenMode::ReadWriteExisting, policy, ) @@ -1929,6 +2145,7 @@ mod tests { &hardlink_paths, &hardlink_metadata, &base_catalog(), + &base_schema_catalog(), OpenMode::ReadWriteExisting, policy, ) @@ -1984,11 +2201,16 @@ mod tests { .await .expect("insert fixture"); drop(connection); + let inspection_schema_catalog = single_table_schema_catalog( + "inspection_fixture", + "CREATE TABLE inspection_fixture (value INTEGER NOT NULL)", + ); let contended = open_existing_connection_pool( &paths, &metadata, &base_catalog(), + &inspection_schema_catalog, OpenMode::ReadOnlyInspection, policy, ) @@ -2007,6 +2229,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &inspection_schema_catalog, OpenMode::ReadOnlyInspection, policy, ) @@ -2023,6 +2246,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &inspection_schema_catalog, OpenMode::ReadOnlyInspection, policy, ) @@ -2099,6 +2323,7 @@ mod tests { &paths, &metadata, &base_catalog(), + &base_schema_catalog(), OpenMode::ReadOnlyInspection, ServiceSqliteConnectionOptions::reviewed(), ) diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs @@ -6,6 +6,8 @@ const AUTHORITY_SOURCE: &str = include_str!("../src/authority.rs"); const CONFIG_SOURCE: &str = include_str!("../src/config.rs"); const ERROR_SOURCE: &str = include_str!("../src/error.rs"); const INITIALIZE_SOURCE: &str = include_str!("../src/initialize.rs"); +const INTEGRITY_SOURCE: &str = include_str!("../src/integrity/mod.rs"); +const INTEGRITY_CATALOG_SOURCE: &str = include_str!("../src/integrity/catalog.rs"); const METADATA_SOURCE: &str = include_str!("../src/metadata.rs"); const MIGRATION_SOURCE: &str = include_str!("../src/migration.rs"); const OPEN_SOURCE: &str = include_str!("../src/open.rs"); @@ -48,6 +50,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "config", "error", "initialize", + "integrity", "metadata", "migration", "open", @@ -90,6 +93,12 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "MigrationEvidenceError", "MigrationKind", "MigrationName", + "SchemaCatalog", + "SchemaCatalogContractError", + "SchemaDigest", + "SchemaObject", + "SchemaObjectKind", + "SchemaVersionCatalog", "ServiceSqlitePathError", "ServiceSqlitePaths", "OpenMode", @@ -113,13 +122,6 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub fn for_sql", "pub fn for_callback", ".take(MAX_MIGRATION_COUNT + 1)", - "schema_migrations", - "applied_at_unix_s", - "service_commit TEXT", - "lib_revision TEXT", - "provider_contract_version INTEGER", - "schema_migrations_no_update", - "schema_migrations_no_delete", ".begin_with(\"BEGIN IMMEDIATE\")", "pub(crate) async fn verify_migration_history", "pub(crate) async fn apply_governed_migrations", @@ -142,6 +144,51 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { ); } + for required in [ + "radroots.service_sqlite.schema_object.v1\\0", + "radroots.service_sqlite.schema_snapshot.v1\\0", + "radroots.service_sqlite.schema_catalog.v1\\0", + "MAX_SCHEMA_OBJECT_COUNT", + "MAX_SCHEMA_SQL_UTF8_BYTES", + "MAX_SCHEMA_CATALOG_UTF8_BYTES", + ".take(MAX_SCHEMA_OBJECT_COUNT + 1)", + ".take(MAX_SCHEMA_VERSION_COUNT + 1)", + "CREATE TABLE radroots_service_metadata", + "CREATE TABLE schema_migrations", + "schema_migrations_no_update", + "schema_migrations_no_delete", + "pub fn computed_digest", + "pub fn new<I>", + "pub(crate) async fn verify_schema_catalog", + "FROM main.sqlite_schema", + "LIMIT 4097", + "typeof(sql) = 'text'", + "length(CAST(sql AS BLOB)) BETWEEN 1 AND 1048576", + ] { + assert!( + INTEGRITY_SOURCE.contains(required) || INTEGRITY_CATALOG_SOURCE.contains(required), + "Step 060 integrity source is missing `{required}`" + ); + } + + for forbidden in [ + "pub async fn verify_schema_catalog", + "pub fn verify_schema_catalog", + "pub use sqlx", + "SqlitePool", + "PoolConnection", + "pub fn sql(&self", + "Serialize", + "Deserialize", + "myc_", + "rhi_", + ] { + assert!( + !INTEGRITY_SOURCE.contains(forbidden) && !INTEGRITY_CATALOG_SOURCE.contains(forbidden), + "Step 060 integrity source contains forbidden surface `{forbidden}`" + ); + } + for forbidden in [ "pub fn content(&self", "pub fn callback_definition", @@ -179,13 +226,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { } for required in [ - "radroots_service_metadata", "PRAGMA application_id", - "source_generation BLOB", - "state_schema_version INTEGER", - "created_at_unix_ms INTEGER", - "radroots_service_metadata_guard_update", - "radroots_service_metadata_no_delete", "LIMIT 2", "SourceGeneration", "NonZeroU32", @@ -200,6 +241,25 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { ); } + for required in [ + "radroots_service_metadata", + "source_generation BLOB", + "state_schema_version INTEGER", + "created_at_unix_ms INTEGER", + "radroots_service_metadata_guard_update", + "radroots_service_metadata_no_delete", + "schema_migrations", + "applied_at_unix_s", + "service_commit TEXT", + "lib_revision TEXT", + "provider_contract_version INTEGER", + ] { + assert!( + INTEGRITY_CATALOG_SOURCE.contains(required), + "shared schema authority is missing `{required}`" + ); + } + for forbidden in [ "pub use sqlx", "pub fn write_database_metadata",