lib

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

commit 7bfb2347ce1ac26a9880403b2be158e1b77c3ccf
parent 90dbe592dee3fcab7dfafc9c285491dafc94d87f
Author: triesap <tyson@radroots.org>
Date:   Tue, 11 Aug 2026 16:33:20 +0000

service-sqlite: encapsulate database connections

- seal pool and connection authority behind typed host transactions
- enforce rollback, cancellation, policy, and schema revalidation
- bind live lock identities and reject attached-database escapes
- verify native, workspace, and cross-target service SQLite gates

Diffstat:
MAGENTS.md | 9+++++++++
MCargo.lock | 1+
Mcrates/service_sqlite/Cargo.toml | 1+
Mcrates/service_sqlite/README.md | 21+++++++++++++++++++--
Mcrates/service_sqlite/src/authority.rs | 138++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-------------------
Acrates/service_sqlite/src/connection.rs | 1289+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/service_sqlite/src/lib.rs | 13++++++++++---
Mcrates/service_sqlite/src/migration.rs | 346+++++++++++++++++++++++++++++++++++++++++++++----------------------------------
Mcrates/service_sqlite/src/open.rs | 186+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Acrates/service_sqlite/src/transaction_control.rs | 109+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++
Mcrates/service_sqlite/tests/package_boundary.rs | 138+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
11 files changed, 2014 insertions(+), 237 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -157,6 +157,15 @@ Before editing code: the owning store, keep live mutable state daemon-owned, and route live-state mutations from local tools through the typed, permissioned Unix-socket local-admin boundary. +- `radroots_service_sqlite::ServiceSqliteHost` is the sole public owner of the + private SQLx pool. Service code may execute typed SQLx queries only through + the sealed `&mut ServiceSqliteTransaction` executor passed to + `ServiceSqliteHost::transaction`; do not expose or reconstruct raw pools, + pooled connections, SQLx transactions, commit/rollback handles, or inner + accessors. Do not attach or detach secondary SQLite databases through the + transaction executor. Writable host construction must finish governed + migrations before returning, while read-only inspection must require current + migration and schema state. - Runtime-management flows consume a sealed `RuntimeContext` for every service instance. They must not reconstruct service paths from raw identifiers, ambient selectors, or manager-owned roots, and registries must not persist diff --git a/Cargo.lock b/Cargo.lock @@ -3910,6 +3910,7 @@ name = "radroots_service_sqlite" version = "0.1.0-alpha" dependencies = [ "fs2", + "futures", "radroots_runtime_paths", "radroots_storage", "rustix 1.1.4", diff --git a/crates/service_sqlite/Cargo.toml b/crates/service_sqlite/Cargo.toml @@ -13,6 +13,7 @@ readme = "README.md" [dependencies] fs2 = { workspace = true } +futures = { workspace = true } radroots_runtime_paths = { workspace = true } radroots_storage = { workspace = true } rustix = { workspace = true } diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md @@ -6,10 +6,27 @@ writer authority, instance locking, versioned schema mechanics, bounded transactions, immutable service-instance database identity, integrity checks, backup, restore, and passive storage status. +`ServiceSqliteHost` is the only public connection host. Its SQLx pool and raw +connections are sealed inside the crate. Services run typed SQLx queries through +the borrowed `ServiceSqliteTransaction` executor supplied by +`ServiceSqliteHost::transaction`; transaction begin, commit, rollback, policy +validation, attached-database exclusion, and cancellation quarantine remain +runner-owned. Writable host opening finishes every pending governed migration +before returning, and read-only inspection opens only current migration and +schema state. + +Cancelling a host transaction before the runner enables outer commit +quarantines its connection and leaves no authoritative transaction effect. A +service-operation error is returned only after rollback is confirmed; an +unconfirmed rollback is reported as `RollbackFailed`. Cancelling once outer +commit begins yields no result and must be treated as an unknown commit outcome. +Both that case and `CommitOutcomeUnknown` require rereading authoritative state +before an idempotent retry. + The crate owns mechanics only. Service-specific tables, SQL, repositories, backup content policy, identity material, process lifecycle, and readiness -policy remain with the consuming service. Its connection pool stays private; -the crate does not provide callers with raw database authority. +policy remain with the consuming service. The crate does not provide callers +with raw database authority. Publication is disabled. The package is not part of the public Radroots crate release closure. diff --git a/crates/service_sqlite/src/authority.rs b/crates/service_sqlite/src/authority.rs @@ -1,7 +1,10 @@ //! Lifetime authority for the sole writable service database owner. use core::fmt; -use std::{error::Error, fs::File, path::PathBuf}; +use std::{error::Error, fs::File}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use std::path::PathBuf; use fs2::FileExt; @@ -20,13 +23,18 @@ use crate::{OpenMode, ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqliteP /// ``` pub struct WriterAuthority { file: Option<File>, - #[allow( - dead_code, - reason = "Step 056 binds the authority for the private Step 061 pool host" - )] + #[cfg(any(target_os = "linux", target_os = "macos"))] database_path: PathBuf, #[cfg(any(target_os = "linux", target_os = "macos"))] directory: File, + #[cfg(any(target_os = "linux", target_os = "macos"))] + directory_device: u64, + #[cfg(any(target_os = "linux", target_os = "macos"))] + directory_inode: u64, + #[cfg(any(target_os = "linux", target_os = "macos"))] + lock_device: u64, + #[cfg(any(target_os = "linux", target_os = "macos"))] + lock_inode: u64, } impl WriterAuthority { @@ -55,10 +63,7 @@ impl WriterAuthority { &self.directory } - #[allow( - dead_code, - reason = "Step 056 binds the authority for the private Step 061 pool host" - )] + #[cfg(any(target_os = "linux", target_os = "macos"))] pub(crate) fn validate_for( &self, paths: &ServiceSqlitePaths, @@ -68,7 +73,7 @@ impl WriterAuthority { } #[cfg(any(target_os = "linux", target_os = "macos"))] - validate_directory_binding(&self.directory, paths).map_err(authority_error)?; + validate_authority_binding(self, paths).map_err(authority_error)?; Ok(()) } @@ -123,10 +128,7 @@ enum WriterAuthorityCause { LockWrongOwner, #[cfg(any(target_os = "linux", target_os = "macos"))] Contended, - #[allow( - dead_code, - reason = "Step 056 binds the authority for the private Step 061 pool host" - )] + #[cfg(any(target_os = "linux", target_os = "macos"))] Mismatched, UnlockFailed, } @@ -156,6 +158,7 @@ impl fmt::Display for WriterAuthorityCause { Self::LockWrongOwner => "SQLite writer lock has the wrong owner", #[cfg(any(target_os = "linux", target_os = "macos"))] Self::Contended => "another SQLite writer is active", + #[cfg(any(target_os = "linux", target_os = "macos"))] Self::Mismatched => "SQLite writer authority does not match this database", Self::UnlockFailed => "SQLite writer authority could not be released", }) @@ -211,14 +214,30 @@ fn acquire_supported(paths: &ServiceSqlitePaths) -> Result<WriterAuthority, Writ fchmod(&descriptor, Mode::RUSR | Mode::WUSR) .map_err(|_| WriterAuthorityCause::LockUnavailable)?; + let lock_status = fstat(&descriptor).map_err(|_| WriterAuthorityCause::LockUnavailable)?; + if u32::from(lock_status.st_mode) & 0o777 != 0o600 { + return Err(WriterAuthorityCause::LockUnavailable); + } + let directory_device = u64::try_from(directory_status.st_dev) + .map_err(|_| WriterAuthorityCause::StateDirectoryUnavailable)?; + let lock_device = + u64::try_from(lock_status.st_dev).map_err(|_| WriterAuthorityCause::LockUnavailable)?; let file = File::from(descriptor); let directory = File::from(directory); match FileExt::try_lock_exclusive(&file) { - Ok(()) => Ok(WriterAuthority { - file: Some(file), - database_path: paths.state_database().to_path_buf(), - directory, - }), + Ok(()) => { + let authority = WriterAuthority { + file: Some(file), + database_path: paths.state_database().to_path_buf(), + directory, + directory_device, + directory_inode: directory_status.st_ino, + lock_device, + lock_inode: lock_status.st_ino, + }; + validate_authority_binding(&authority, paths)?; + Ok(authority) + } Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => { Err(WriterAuthorityCause::Contended) } @@ -270,28 +289,83 @@ fn validate_lock( } #[cfg(any(target_os = "linux", target_os = "macos"))] -fn validate_directory_binding( - held_directory: &File, +fn validate_authority_binding( + authority: &WriterAuthority, paths: &ServiceSqlitePaths, ) -> Result<(), WriterAuthorityCause> { - use rustix::fs::{Mode, OFlags, fstat, open}; + use rustix::{ + fs::{FileType, Mode, OFlags, fstat, open, openat}, + process::geteuid, + }; + let directory_path = paths + .state_database() + .parent() + .filter(|parent| Some(*parent) == paths.state_lock().parent()) + .ok_or(WriterAuthorityCause::Mismatched)?; let current_directory = open( - paths - .state_database() - .parent() - .ok_or(WriterAuthorityCause::Mismatched)?, + directory_path, OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, Mode::empty(), ) .map_err(|_| WriterAuthorityCause::Mismatched)?; - let held = fstat(held_directory).map_err(|_| WriterAuthorityCause::Mismatched)?; - let current = fstat(&current_directory).map_err(|_| WriterAuthorityCause::Mismatched)?; - if held.st_dev == current.st_dev && held.st_ino == current.st_ino { - Ok(()) - } else { - Err(WriterAuthorityCause::Mismatched) + let held_directory = + fstat(&authority.directory).map_err(|_| WriterAuthorityCause::Mismatched)?; + let current_directory_status = + fstat(&current_directory).map_err(|_| WriterAuthorityCause::Mismatched)?; + let held_directory_device = + u64::try_from(held_directory.st_dev).map_err(|_| WriterAuthorityCause::Mismatched)?; + let current_directory_device = u64::try_from(current_directory_status.st_dev) + .map_err(|_| WriterAuthorityCause::Mismatched)?; + if !FileType::from_raw_mode(held_directory.st_mode).is_dir() + || held_directory.st_uid != geteuid().as_raw() + || u32::from(held_directory.st_mode) & 0o022 != 0 + || !FileType::from_raw_mode(current_directory_status.st_mode).is_dir() + || current_directory_status.st_uid != geteuid().as_raw() + || u32::from(current_directory_status.st_mode) & 0o022 != 0 + || held_directory_device != authority.directory_device + || held_directory.st_ino != authority.directory_inode + || current_directory_device != authority.directory_device + || current_directory_status.st_ino != authority.directory_inode + { + return Err(WriterAuthorityCause::Mismatched); + } + + let current_lock = openat( + &current_directory, + radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME, + OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK, + Mode::empty(), + ) + .map_err(|_| WriterAuthorityCause::Mismatched)?; + let held_lock = fstat( + authority + .file + .as_ref() + .ok_or(WriterAuthorityCause::Mismatched)?, + ) + .map_err(|_| WriterAuthorityCause::Mismatched)?; + let current_lock_status = fstat(&current_lock).map_err(|_| WriterAuthorityCause::Mismatched)?; + let held_lock_device = + u64::try_from(held_lock.st_dev).map_err(|_| WriterAuthorityCause::Mismatched)?; + let current_lock_device = + u64::try_from(current_lock_status.st_dev).map_err(|_| WriterAuthorityCause::Mismatched)?; + if !FileType::from_raw_mode(held_lock.st_mode).is_file() + || u64::from(held_lock.st_nlink) != 1 + || held_lock.st_uid != geteuid().as_raw() + || u32::from(held_lock.st_mode) & 0o777 != 0o600 + || !FileType::from_raw_mode(current_lock_status.st_mode).is_file() + || u64::from(current_lock_status.st_nlink) != 1 + || current_lock_status.st_uid != geteuid().as_raw() + || u32::from(current_lock_status.st_mode) & 0o777 != 0o600 + || held_lock_device != authority.lock_device + || held_lock.st_ino != authority.lock_inode + || current_lock_device != authority.lock_device + || current_lock_status.st_ino != authority.lock_inode + { + return Err(WriterAuthorityCause::Mismatched); } + Ok(()) } #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] diff --git a/crates/service_sqlite/src/connection.rs b/crates/service_sqlite/src/connection.rs @@ -0,0 +1,1289 @@ +//! Narrow service-store host and transaction execution boundary. + +use core::fmt; +use std::{ + error::Error, + future::Future, + pin::Pin, + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, +}; + +use futures::{future::BoxFuture, stream::BoxStream}; +use sqlx::{ + Either, Execute, Executor, SqlStr, Sqlite, SqliteConnection, + sqlite::{SqliteQueryResult, SqliteRow, SqliteStatement, SqliteTypeInfo}, +}; + +use crate::{ + MigrationApplicationOutcome, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, + MigrationCallbackBinding, MigrationCatalog, OpenMode, SchemaCatalog, ServiceDatabaseIdentity, + ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqlitePaths, + WriterAuthority, +}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use sqlx::{Connection, pool::PoolConnection}; + +/// One service-owned SQLite host whose raw pool remains inaccessible. +/// +/// The host intentionally has no raw-pool accessor: +/// +/// ```compile_fail +/// use radroots_service_sqlite::ServiceSqliteHost; +/// +/// fn leak_pool(host: &ServiceSqliteHost) { +/// let _ = host.pool(); +/// } +/// ``` +pub struct ServiceSqliteHost { + mode: OpenMode, + #[cfg(any(target_os = "linux", target_os = "macos"))] + pool: crate::open::PrivateConnectionPool, +} + +impl ServiceSqliteHost { + /// Opens existing writable state and finishes every pending governed migration. + #[allow(clippy::too_many_arguments)] + pub async fn open_read_write_existing( + paths: &ServiceSqlitePaths, + identity: &ServiceDatabaseIdentity, + migrations: &MigrationCatalog, + schema: &SchemaCatalog, + options: ServiceSqliteConnectionOptions, + applied_at: MigrationAppliedAtUnixSeconds, + build: &MigrationBuildIdentity, + callbacks: &[MigrationCallbackBinding], + ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + let pool = crate::open::open_existing_connection_pool( + paths, + identity, + migrations, + schema, + OpenMode::ReadWriteExisting, + options, + ) + .await?; + match pool.apply_migrations(applied_at, build, callbacks).await { + Ok(outcome) => Ok(( + Self { + mode: OpenMode::ReadWriteExisting, + pool, + }, + outcome, + )), + Err(error) => { + drop(pool.close().await); + Err(error) + } + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = ( + paths, identity, migrations, schema, options, applied_at, build, callbacks, + ); + Err(unsupported_host()) + } + } + + /// Opens state created under a retained initialization writer authority. + #[allow(clippy::too_many_arguments)] + pub async fn open_initialized( + paths: &ServiceSqlitePaths, + identity: &ServiceDatabaseIdentity, + migrations: &MigrationCatalog, + schema: &SchemaCatalog, + options: ServiceSqliteConnectionOptions, + authority: WriterAuthority, + applied_at: MigrationAppliedAtUnixSeconds, + build: &MigrationBuildIdentity, + callbacks: &[MigrationCallbackBinding], + ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + let pool = crate::open::open_initialized_connection_pool( + paths, identity, migrations, schema, options, authority, + ) + .await?; + match pool.apply_migrations(applied_at, build, callbacks).await { + Ok(outcome) => Ok(( + Self { + mode: OpenMode::Initialize, + pool, + }, + outcome, + )), + Err(error) => { + drop(pool.close().await); + Err(error) + } + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = ( + paths, identity, migrations, schema, options, authority, applied_at, build, + callbacks, + ); + Err(unsupported_host()) + } + } + + /// Opens an immutable, current-schema inspection host without writer authority. + pub async fn open_read_only_inspection( + paths: &ServiceSqlitePaths, + identity: &ServiceDatabaseIdentity, + migrations: &MigrationCatalog, + schema: &SchemaCatalog, + options: ServiceSqliteConnectionOptions, + ) -> Result<Self, ServiceSqliteError> { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + let pool = crate::open::open_existing_connection_pool( + paths, + identity, + migrations, + schema, + OpenMode::ReadOnlyInspection, + options, + ) + .await?; + Ok(Self { + mode: OpenMode::ReadOnlyInspection, + pool, + }) + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = (paths, identity, migrations, schema, options); + Err(unsupported_host()) + } + } + + /// Returns the fixed mode selected when the host was opened. + #[must_use] + pub const fn mode(&self) -> OpenMode { + self.mode + } + + /// Executes one runner-owned transaction without exposing its connection or pool. + /// + /// Dropping this future before the runner enables its outer commit quarantines + /// the connection and leaves no authoritative transaction effect. An operation + /// error is returned as `OperationRolledBack` only after rollback is confirmed; + /// an unconfirmed rollback is `RollbackFailed`. Once outer commit begins, + /// cancelling the future yields no result and must be treated as an unknown + /// commit outcome. Callers receiving `CommitOutcomeUnknown`, or cancelling after + /// commit begins, must reread authoritative state before any idempotent retry. + pub async fn transaction<T, E, F>( + &self, + operation: F, + ) -> Result<T, ServiceSqliteTransactionError<E>> + where + T: Send + 'static, + E: Send + 'static, + F: for<'a> FnOnce( + &'a mut ServiceSqliteTransaction<'_>, + ) -> ServiceSqliteTransactionFuture<'a, T, E> + + Send, + { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + self.transaction_supported(operation).await + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + drop(operation); + Err(ServiceSqliteTransactionError::not_committed( + unsupported_host(), + )) + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + async fn transaction_supported<T, E, F>( + &self, + operation: F, + ) -> Result<T, ServiceSqliteTransactionError<E>> + where + T: Send + 'static, + E: Send + 'static, + F: for<'a> FnOnce( + &'a mut ServiceSqliteTransaction<'_>, + ) -> ServiceSqliteTransactionFuture<'a, T, E> + + Send, + { + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::not_committed)?; + let connection = self + .pool + .acquire() + .await + .map_err(ServiceSqliteTransactionError::not_committed)?; + let mut connection = QuarantinedConnection::new(connection); + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::not_committed)?; + let initial_policy = crate::migration::read_connection_policy(&mut connection) + .await + .map_err(ServiceSqliteTransactionError::not_committed)?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::not_committed)?; + let gate = crate::transaction_control::TransactionControlGate::install(&mut connection) + .await + .map_err(|source| { + ServiceSqliteTransactionError::not_committed(sqlite_source(source)) + })?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::not_committed)?; + let mut transaction = match match self.pool.mode() { + OpenMode::Initialize | OpenMode::ReadWriteExisting => { + connection.begin_with("BEGIN IMMEDIATE").await + } + OpenMode::ReadOnlyInspection => connection.begin().await, + } { + Ok(transaction) => transaction, + Err(source) => { + return Err(ServiceSqliteTransactionError::not_committed(sqlite_source( + source, + ))); + } + }; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::not_committed)?; + + let operation_result = { + let database_control_rejected = Arc::new(AtomicBool::new(false)); + let mut executor = ServiceSqliteTransaction { + connection: &mut transaction, + database_control_rejected: Arc::clone(&database_control_rejected), + }; + (operation(&mut executor).await, database_control_rejected) + }; + let (operation_result, database_control_rejected) = operation_result; + if let Err(error) = self.pool.validate() { + let operation_error = operation_result.err(); + let permit = gate.permit_runner_rollback(); + let rollback = transaction.rollback().await.map_err(sqlite_source); + drop(permit); + let rollback_was_confirmed = + gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); + let remove = gate.remove(&mut connection).await.map_err(sqlite_source); + let rollback_error = rollback + .err() + .filter(|_| !rollback_was_confirmed) + .or_else(|| remove.err()); + return Err(match rollback_error { + Some(rollback_error) => { + ServiceSqliteTransactionError::rollback_failed(operation_error, rollback_error) + } + None => ServiceSqliteTransactionError::not_committed_with_operation( + operation_error, + error, + ), + }); + } + let value = match operation_result { + Ok(value) => value, + Err(operation_error) => { + let permit = gate.permit_runner_rollback(); + let rollback = transaction.rollback().await.map_err(sqlite_source); + drop(permit); + let rollback_was_confirmed = + gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); + let remove = gate.remove(&mut connection).await.map_err(sqlite_source); + let authority = self.pool.validate(); + if let Some(error) = rollback + .err() + .filter(|_| !rollback_was_confirmed) + .or_else(|| remove.err()) + .or_else(|| authority.err()) + { + return Err(ServiceSqliteTransactionError::rollback_failed( + Some(operation_error), + error, + )); + } + return Err(ServiceSqliteTransactionError::operation_rolled_back( + operation_error, + )); + } + }; + + let precommit = self + .verify_before_commit( + &mut transaction, + &gate, + &initial_policy, + &database_control_rejected, + ) + .await; + if let Err(error) = precommit { + let permit = gate.permit_runner_rollback(); + let rollback = transaction.rollback().await.map_err(sqlite_source); + drop(permit); + let rollback_was_confirmed = + gate.rejected_commit_rolled_back() && !connection.is_in_transaction(); + let remove = gate.remove(&mut connection).await.map_err(sqlite_source); + let authority = self.pool.validate(); + if let Some(rollback_error) = rollback + .err() + .filter(|_| !rollback_was_confirmed) + .or_else(|| remove.err()) + .or_else(|| authority.err()) + { + return Err(ServiceSqliteTransactionError::rollback_failed( + None, + rollback_error, + )); + } + return Err(ServiceSqliteTransactionError::not_committed(error)); + } + + let permit = gate.permit_outer_commit(); + let commit = transaction.commit().await.map_err(sqlite_source); + drop(permit); + let commit = match commit { + Ok(()) => Ok(()), + Err(error) => Err(ServiceSqliteTransactionError::commit_outcome_unknown(error)), + }; + let remove = gate.remove(&mut connection).await.map_err(sqlite_source); + commit?; + remove.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + let final_policy = crate::migration::read_connection_policy(&mut connection) + .await + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + if final_policy != initial_policy { + return Err(ServiceSqliteTransactionError::commit_outcome_unknown( + ServiceSqliteError::new(ServiceSqliteErrorKind::Pragma), + )); + } + crate::metadata::verify_database_metadata(&mut connection, self.pool.identity()) + .await + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + crate::migration::verify_migration_history( + &mut connection, + self.pool.catalog(), + self.pool.schema_catalog(), + true, + ) + .await + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + self.pool + .validate() + .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?; + connection.trust(); + Ok(value) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + async fn verify_before_commit( + &self, + connection: &mut SqliteConnection, + gate: &crate::transaction_control::TransactionControlGate, + initial_policy: &crate::migration::MigrationConnectionPolicy, + database_control_rejected: &AtomicBool, + ) -> Result<(), ServiceSqliteError> { + self.pool.validate()?; + if gate.control_violation_observed() || database_control_rejected.load(Ordering::Acquire) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open)); + } + crate::migration::assert_governed_transaction(connection).await?; + self.pool.validate()?; + if &crate::migration::read_connection_policy(connection).await? != initial_policy { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Pragma)); + } + self.pool.validate()?; + crate::metadata::verify_database_metadata(connection, self.pool.identity()).await?; + self.pool.validate()?; + crate::migration::verify_migration_history_snapshot( + connection, + self.pool.catalog(), + self.pool.schema_catalog(), + true, + ) + .await?; + self.pool.validate()?; + crate::migration::assert_governed_transaction(connection).await?; + if gate.control_violation_observed() { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open)); + } + Ok(()) + } + + #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] + async fn close_pool(&self) { + self.pool.close_pool().await; + } +} + +impl fmt::Debug for ServiceSqliteHost { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ServiceSqliteHost") + .field("mode", &self.mode) + .field("pool", &"[redacted]") + .finish() + } +} + +/// A sealed transaction executor that never exposes its raw SQLite connection. +/// +/// Service repositories may use ordinary typed SQLx queries through the +/// borrowed executor: +/// +/// ``` +/// use radroots_service_sqlite::ServiceSqliteTransaction; +/// +/// async fn row_count( +/// transaction: &mut ServiceSqliteTransaction<'_>, +/// ) -> Result<i64, sqlx::Error> { +/// sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM service_items") +/// .fetch_one(transaction) +/// .await +/// } +/// ``` +/// +/// Transaction control remains runner-owned: +/// +/// ```compile_fail +/// use radroots_service_sqlite::ServiceSqliteTransaction; +/// +/// async fn bypass(transaction: ServiceSqliteTransaction<'_>) { +/// transaction.commit().await.unwrap(); +/// } +/// ``` +pub struct ServiceSqliteTransaction<'connection> { + connection: &'connection mut SqliteConnection, + database_control_rejected: Arc<AtomicBool>, +} + +struct RestrictedExecute<Q> { + query: Q, + database_control_rejected: Arc<AtomicBool>, +} + +impl<'query, Q> Execute<'query, Sqlite> for RestrictedExecute<Q> +where + Q: Execute<'query, Sqlite>, +{ + fn sql(self) -> SqlStr { + restricted_sql(self.query.sql(), &self.database_control_rejected) + } + + fn statement(&self) -> Option<&SqliteStatement> { + None + } + + fn take_arguments( + &mut self, + ) -> Result<Option<<Sqlite as sqlx::Database>::Arguments>, sqlx::error::BoxDynError> { + self.query.take_arguments() + } + + fn persistent(&self) -> bool { + self.query.persistent() + } +} + +fn restricted_sql(sql: SqlStr, database_control_rejected: &AtomicBool) -> SqlStr { + if contains_database_control(sql.as_str()) { + database_control_rejected.store(true, Ordering::Release); + SqlStr::from_static("RADROOTS_FORBIDDEN_DATABASE_CONTROL") + } else { + sql + } +} + +pub(crate) fn contains_database_control(sql: &str) -> bool { + sql.as_bytes() + .split(|byte| !byte.is_ascii_alphanumeric() && *byte != b'_') + .any(|token| token.eq_ignore_ascii_case(b"attach") || token.eq_ignore_ascii_case(b"detach")) +} + +impl fmt::Debug for ServiceSqliteTransaction<'_> { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("ServiceSqliteTransaction([redacted])") + } +} + +impl<'executor, 'connection> Executor<'executor> + for &'executor mut ServiceSqliteTransaction<'connection> +where + 'connection: 'executor, +{ + type Database = Sqlite; + + fn fetch_many<'e, 'q: 'e, Q>( + self, + query: Q, + ) -> BoxStream<'e, Result<Either<SqliteQueryResult, SqliteRow>, sqlx::Error>> + where + 'executor: 'e, + Q: 'q + Execute<'q, Self::Database>, + { + (&mut *self.connection).fetch_many(RestrictedExecute { + query, + database_control_rejected: Arc::clone(&self.database_control_rejected), + }) + } + + fn fetch_optional<'e, 'q: 'e, Q>( + self, + query: Q, + ) -> BoxFuture<'e, Result<Option<SqliteRow>, sqlx::Error>> + where + 'executor: 'e, + Q: 'q + Execute<'q, Self::Database>, + { + (&mut *self.connection).fetch_optional(RestrictedExecute { + query, + database_control_rejected: Arc::clone(&self.database_control_rejected), + }) + } + + fn prepare_with<'e>( + self, + sql: SqlStr, + parameters: &'e [SqliteTypeInfo], + ) -> BoxFuture<'e, Result<SqliteStatement, sqlx::Error>> + where + 'executor: 'e, + { + (&mut *self.connection).prepare_with( + restricted_sql(sql, &self.database_control_rejected), + parameters, + ) + } +} + +/// Boxed callback future tied to the borrowed transaction executor. +pub type ServiceSqliteTransactionFuture<'a, T, E> = + Pin<Box<dyn Future<Output = Result<T, E>> + Send + 'a>>; + +/// Stable transaction completion phases without expanding SQLite error kinds. +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +pub enum ServiceSqliteTransactionErrorKind { + NotCommitted, + OperationRolledBack, + RollbackFailed, + CommitOutcomeUnknown, +} + +/// Transaction failure retaining trusted details behind redacted diagnostics. +pub struct ServiceSqliteTransactionError<E> { + kind: ServiceSqliteTransactionErrorKind, + operation_error: Option<E>, + sqlite_error: Option<ServiceSqliteError>, +} + +impl<E> ServiceSqliteTransactionError<E> { + fn not_committed(error: ServiceSqliteError) -> Self { + Self::not_committed_with_operation(None, error) + } + + fn not_committed_with_operation( + operation_error: Option<E>, + sqlite_error: ServiceSqliteError, + ) -> Self { + Self { + kind: ServiceSqliteTransactionErrorKind::NotCommitted, + operation_error, + sqlite_error: Some(sqlite_error), + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn operation_rolled_back(error: E) -> Self { + Self { + kind: ServiceSqliteTransactionErrorKind::OperationRolledBack, + operation_error: Some(error), + sqlite_error: None, + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn rollback_failed(operation_error: Option<E>, sqlite_error: ServiceSqliteError) -> Self { + Self { + kind: ServiceSqliteTransactionErrorKind::RollbackFailed, + operation_error, + sqlite_error: Some(sqlite_error), + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn commit_outcome_unknown(error: ServiceSqliteError) -> Self { + Self { + kind: ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown, + operation_error: None, + sqlite_error: Some(error), + } + } + + #[must_use] + pub const fn kind(&self) -> ServiceSqliteTransactionErrorKind { + self.kind + } + + #[must_use] + pub const fn operation_error(&self) -> Option<&E> { + self.operation_error.as_ref() + } + + #[must_use] + pub const fn sqlite_error(&self) -> Option<&ServiceSqliteError> { + self.sqlite_error.as_ref() + } +} + +impl<E> fmt::Debug for ServiceSqliteTransactionError<E> { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter + .debug_struct("ServiceSqliteTransactionError") + .field("kind", &self.kind) + .field( + "operation_error", + &self.operation_error.as_ref().map(|_| "[redacted]"), + ) + .field( + "sqlite_error", + &self.sqlite_error.as_ref().map(|_| "[redacted]"), + ) + .finish() + } +} + +impl<E> fmt::Display for ServiceSqliteTransactionError<E> { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str(match self.kind { + ServiceSqliteTransactionErrorKind::NotCommitted => { + "SQLite transaction did not reach commit" + } + ServiceSqliteTransactionErrorKind::OperationRolledBack => { + "SQLite transaction operation was rolled back" + } + ServiceSqliteTransactionErrorKind::RollbackFailed => { + "SQLite transaction rollback could not be confirmed" + } + ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown => { + "SQLite transaction commit outcome is unknown" + } + }) + } +} + +impl<E: Send + Sync + 'static> Error for ServiceSqliteTransactionError<E> {} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +struct QuarantinedConnection { + connection: Option<PoolConnection<Sqlite>>, + trusted: bool, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl QuarantinedConnection { + fn new(connection: PoolConnection<Sqlite>) -> Self { + Self { + connection: Some(connection), + trusted: false, + } + } + + fn trust(&mut self) { + self.trusted = true; + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl core::ops::Deref for QuarantinedConnection { + type Target = SqliteConnection; + + fn deref(&self) -> &Self::Target { + self.connection.as_deref().expect("connection is retained") + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl core::ops::DerefMut for QuarantinedConnection { + fn deref_mut(&mut self) -> &mut Self::Target { + self.connection + .as_deref_mut() + .expect("connection is retained") + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl Drop for QuarantinedConnection { + fn drop(&mut self) { + if !self.trusted + && let Some(connection) = self.connection.as_mut() + { + connection.close_on_drop(); + } + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +fn sqlite_source(source: sqlx::Error) -> ServiceSqliteError { + ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source) +} + +#[cfg(not(any(target_os = "linux", target_os = "macos")))] +fn unsupported_host() -> ServiceSqliteError { + ServiceSqliteError::new(ServiceSqliteErrorKind::Open) +} + +#[cfg(test)] +mod tests { + #[cfg(any(target_os = "linux", target_os = "macos"))] + use std::{ + convert::Infallible, + fs, + num::NonZeroU32, + os::unix::fs::PermissionsExt, + sync::{ + Arc, + atomic::{AtomicBool, Ordering}, + }, + }; + + #[cfg(any(target_os = "linux", target_os = "macos"))] + use radroots_runtime_paths::{ + InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, + RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId, + }; + #[cfg(any(target_os = "linux", target_os = "macos"))] + use radroots_storage::event::SourceGeneration; + #[cfg(any(target_os = "linux", target_os = "macos"))] + use sqlx::{Connection, sqlite::SqliteConnectOptions}; + + use super::*; + + #[cfg(any(target_os = "linux", target_os = "macos"))] + const HOST_TABLE_SQL: &str = "CREATE TABLE host_probe (value INTEGER NOT NULL)"; + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn runtime_context(root: &std::path::Path) -> RuntimeContext { + RuntimeContext::resolve( + &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()), + RuntimeContextBootstrap::new( + RadrootsPathProfile::RepoLocal, + Some(root.to_path_buf()), + RuntimeContextSource::BootstrapCli, + RuntimeContextSource::BootstrapCli, + ) + .expect("valid bootstrap"), + ServiceId::new("myc").expect("service ID"), + InstanceId::new("host-boundary").expect("instance ID"), + ) + .expect("runtime context") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn migration_catalog() -> MigrationCatalog { + MigrationCatalog::new([]).expect("empty v1 migration catalog") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn schema_catalog(migrations: &MigrationCatalog) -> SchemaCatalog { + let table = crate::SchemaObject::new( + crate::SchemaObjectKind::Table, + "host_probe", + "host_probe", + HOST_TABLE_SQL, + crate::SchemaObject::computed_digest( + crate::SchemaObjectKind::Table, + "host_probe", + "host_probe", + HOST_TABLE_SQL, + ) + .expect("table digest"), + ) + .expect("table descriptor"); + let version_digest = crate::SchemaVersionCatalog::computed_digest(1, [table.clone()]) + .expect("version digest"); + let version = + crate::SchemaVersionCatalog::new(1, [table], version_digest).expect("schema version"); + SchemaCatalog::new(migrations, [version]).expect("schema catalog") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn build_identity() -> MigrationBuildIdentity { + MigrationBuildIdentity::new( + "0.1.0-alpha", + "0123456789abcdef0123456789abcdef01234567", + "89abcdef0123456789abcdef0123456789abcdef", + "1.97.1", + "x86_64-unknown-linux-gnu", + "service-host", + 1, + 2, + 3, + 4, + 5, + ) + .expect("build identity") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + async fn initialized_host() -> ( + tempfile::TempDir, + ServiceSqlitePaths, + ServiceDatabaseIdentity, + MigrationCatalog, + SchemaCatalog, + ServiceSqliteHost, + ) { + let root = tempfile::tempdir().expect("temporary root"); + let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path())) + .expect("SQLite paths"); + fs::create_dir_all(paths.state_database().parent().expect("state directory")) + .expect("create state directory"); + let metadata = crate::ServiceDatabaseMetadata::new( + &paths, + SourceGeneration::new([9; 32]).expect("source generation"), + NonZeroU32::new(1).expect("schema version"), + 1_700_000_000_000, + crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"), + ) + .expect("database metadata"); + let migrations = migration_catalog(); + let schema = schema_catalog(&migrations); + let authority = crate::initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema, + |database_path| async move { + let options = SqliteConnectOptions::new() + .filename(database_path) + .create_if_missing(false); + let mut connection = SqliteConnection::connect_with(&options) + .await + .expect("open reserved database"); + sqlx::query(HOST_TABLE_SQL) + .execute(&mut connection) + .await + .expect("create host table"); + connection.close().await.expect("close reserved database"); + Ok::<_, Infallible>(()) + }, + ) + .await + .expect("initialize database"); + let identity = metadata.identity(); + let (host, outcome) = ServiceSqliteHost::open_initialized( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + authority, + MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"), + &build_identity(), + &[], + ) + .await + .expect("open initialized host"); + assert_eq!(outcome.initial_version(), 1); + assert_eq!(outcome.final_version(), 1); + assert_eq!(outcome.applied_count(), 0); + (root, paths, identity, migrations, schema, host) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + async fn row_count(host: &ServiceSqliteHost) -> i64 { + host.transaction(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *transaction) + .await + }) + }) + .await + .expect("count rows") + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn typed_execution_commits_and_operation_failure_rolls_back() { + let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; + assert_eq!(host.mode(), OpenMode::Initialize); + + let inserted = host + .transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (?)") + .bind(41_i64) + .execute(&mut *transaction) + .await?; + sqlx::query_scalar::<_, i64>("SELECT value FROM host_probe") + .fetch_one(&mut *transaction) + .await + }) + }) + .await + .expect("commit typed operation"); + assert_eq!(inserted, 41); + + let error = host + .transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (99)") + .execute(&mut *transaction) + .await + .map_err(|_| "query-failure-secret")?; + Err::<(), _>("operation-secret") + }) + }) + .await + .expect_err("operation rejection must roll back"); + assert_eq!( + error.kind(), + ServiceSqliteTransactionErrorKind::OperationRolledBack + ); + assert_eq!(error.operation_error(), Some(&"operation-secret")); + assert!(error.sqlite_error().is_none()); + assert!(!format!("{error:?}").contains("operation-secret")); + assert!(!error.to_string().contains("operation-secret")); + assert_eq!(row_count(&host).await, 1); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn cancellation_quarantines_connection_and_pool_recovers() { + let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; + let host = Arc::new(host); + let entered = Arc::new(AtomicBool::new(false)); + let task = tokio::spawn({ + let host = Arc::clone(&host); + let entered = Arc::clone(&entered); + async move { + host.transaction::<(), Infallible, _>(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (77)") + .execute(&mut *transaction) + .await + .expect("tentative insert"); + entered.store(true, Ordering::Release); + std::future::pending::<()>().await; + Ok(()) + }) + }) + .await + } + }); + while !entered.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + task.abort(); + assert!( + task.await + .expect_err("task must be cancelled") + .is_cancelled() + ); + assert_eq!(row_count(&host).await, 0); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn transaction_control_policy_and_attachment_escapes_fail_closed() { + for (statement, expected_kind) in [ + ( + "INSERT INTO host_probe (value) VALUES (0); COMMIT", + ServiceSqliteTransactionErrorKind::RollbackFailed, + ), + ( + "ROLLBACK; BEGIN DEFERRED; INSERT INTO host_probe (value) VALUES (1)", + ServiceSqliteTransactionErrorKind::NotCommitted, + ), + ( + "PRAGMA trusted_schema=ON; INSERT INTO host_probe (value) VALUES (2)", + ServiceSqliteTransactionErrorKind::NotCommitted, + ), + ( + "ATTACH DATABASE ':memory:' AS extra; INSERT INTO host_probe (value) VALUES (3)", + ServiceSqliteTransactionErrorKind::NotCommitted, + ), + ] { + let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; + let error = host + .transaction(|transaction| { + Box::pin(async move { + let _ = sqlx::raw_sql(statement).execute(&mut *transaction).await; + Ok::<_, Infallible>(()) + }) + }) + .await + .expect_err("escape attempt must not commit"); + assert_eq!(error.kind(), expected_kind); + assert!(error.sqlite_error().is_some()); + assert_eq!(row_count(&host).await, 0); + } + + let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; + let error = host + .transaction(|transaction| { + Box::pin(async move { + let _ = sqlx::raw_sql("INSERT INTO host_probe (value) VALUES (4); COMMIT") + .execute(&mut *transaction) + .await; + let _ = + sqlx::raw_sql("BEGIN DEFERRED; INSERT INTO host_probe (value) VALUES (5)") + .execute(&mut *transaction) + .await; + Ok::<_, Infallible>(()) + }) + }) + .await + .expect_err("replacement transaction after denied COMMIT must not escape"); + assert_eq!( + error.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + assert_eq!(row_count(&host).await, 0); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn attach_detach_is_rejected_before_it_can_create_external_state() { + let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; + let external_database = root.path().join("forbidden-attachment.sqlite"); + let statement = format!( + "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", + external_database.display() + ); + let error = host + .transaction(|transaction| { + Box::pin(async move { + sqlx::raw_sql(sqlx::AssertSqlSafe(statement)) + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect_err("ATTACH and DETACH must be rejected before SQLite compilation"); + assert_eq!( + error.kind(), + ServiceSqliteTransactionErrorKind::OperationRolledBack + ); + assert!(!external_database.exists()); + assert_eq!(row_count(&host).await, 0); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn writer_lock_replacement_and_insecure_directory_revoke_live_host() { + let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; + let retired_lock = paths + .state_lock() + .parent() + .expect("state directory") + .join("retired-state.lock"); + fs::rename(paths.state_lock(), &retired_lock).expect("replace canonical lock name"); + let replacement_authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("replacement authority acquisition") + .expect("new writer authority"); + let replaced = host + .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) + .await + .expect_err("old writer authority must reject the replacement lock"); + assert_eq!( + replaced.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + assert_eq!( + replaced.sqlite_error().map(ServiceSqliteError::kind), + Some(ServiceSqliteErrorKind::Authority) + ); + drop(replacement_authority); + drop(host); + + let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; + let directory = paths.state_database().parent().expect("state directory"); + let original_mode = fs::metadata(directory) + .expect("state directory metadata") + .permissions() + .mode() + & 0o777; + fs::set_permissions(directory, fs::Permissions::from_mode(0o770)) + .expect("make directory insecure"); + let insecure = host + .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) + .await + .expect_err("insecure live directory must revoke authority"); + assert_eq!( + insecure.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + assert_eq!( + insecure.sqlite_error().map(ServiceSqliteError::kind), + Some(ServiceSqliteErrorKind::Authority) + ); + fs::set_permissions(directory, fs::Permissions::from_mode(original_mode)) + .expect("restore directory mode"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn inspection_lock_replacement_and_new_writer_revoke_live_inspection() { + let (_root, paths, identity, migrations, schema, host) = initialized_host().await; + host.close_pool().await; + drop(host); + let inspection = ServiceSqliteHost::open_read_only_inspection( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + ) + .await + .expect("open inspection host"); + let retired_lock = paths + .state_lock() + .parent() + .expect("state directory") + .join("inspection-state.lock"); + fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock name"); + let (writer, _outcome) = ServiceSqliteHost::open_read_write_existing( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), + &build_identity(), + &[], + ) + .await + .expect("open replacement writer"); + writer + .transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (88)") + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect("replacement writer commits"); + let stale = inspection + .transaction(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *transaction) + .await + }) + }) + .await + .expect_err("stale inspection authority must refuse work"); + assert_eq!( + stale.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + assert_eq!( + stale.sqlite_error().map(ServiceSqliteError::kind), + Some(ServiceSqliteErrorKind::Authority) + ); + } + + #[test] + fn database_control_token_screen_is_closed_and_case_insensitive() { + for forbidden in [ + "ATTACH DATABASE 'x' AS extra", + "detach database extra", + "SELECT 1; /* ignored */ AtTaCh ':memory:' AS x", + "SELECT 'attach'", + ] { + assert!(contains_database_control(forbidden)); + } + for allowed in [ + "SELECT attachment FROM items", + "SELECT detached FROM items", + "SELECT COUNT(*) FROM host_probe", + ] { + assert!(!contains_database_control(allowed)); + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn read_only_writes_roll_back_and_internal_pool_close_refuses_new_work() { + let (root, paths, identity, migrations, schema, host) = initialized_host().await; + host.close_pool().await; + let closed = host + .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) + .await + .expect_err("closed pool must refuse work"); + assert_eq!( + closed.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + drop(host); + + let inspection = ServiceSqliteHost::open_read_only_inspection( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + ) + .await + .expect("open read-only host"); + let error = inspection + .transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (5)") + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect_err("read-only write must fail"); + assert_eq!( + error.kind(), + ServiceSqliteTransactionErrorKind::OperationRolledBack + ); + assert_eq!(row_count(&inspection).await, 0); + drop(inspection); + drop(root); + } + + #[test] + fn host_and_transaction_errors_are_redacted_and_source_free() { + let error = ServiceSqliteTransactionError::rollback_failed( + Some("operation-secret"), + ServiceSqliteError::new(ServiceSqliteErrorKind::Open), + ); + let debug = format!("{error:?}"); + assert!(!debug.contains("operation-secret")); + assert!(!debug.contains("state.sqlite")); + assert!(error.source().is_none()); + assert_eq!(error.operation_error(), Some(&"operation-secret")); + assert_eq!( + error.sqlite_error().map(ServiceSqliteError::kind), + Some(ServiceSqliteErrorKind::Open) + ); + } +} diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs @@ -4,6 +4,7 @@ mod authority; mod config; +mod connection; mod error; mod initialize; mod integrity; @@ -11,9 +12,14 @@ mod metadata; mod migration; mod open; mod status; +mod transaction_control; pub use authority::WriterAuthority; pub use config::{ServiceSqliteConnectionOptions, ServiceSqliteConnectionOptionsError}; +pub use connection::{ + ServiceSqliteHost, ServiceSqliteTransaction, ServiceSqliteTransactionError, + ServiceSqliteTransactionErrorKind, ServiceSqliteTransactionFuture, +}; pub use error::{ SafeServiceSqliteError, ServiceSqliteError, ServiceSqliteErrorCode, ServiceSqliteErrorKind, }; @@ -27,9 +33,10 @@ pub use metadata::{ ServiceSqliteMetadataValueError, }; pub use migration::{ - MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, MigrationCatalog, MigrationChecksum, - MigrationContractError, MigrationDescriptor, MigrationEvidenceError, MigrationKind, - MigrationName, + MigrationApplicationOutcome, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity, + MigrationCallback, MigrationCallbackBinding, MigrationCallbackFuture, MigrationCatalog, + MigrationChecksum, MigrationContractError, MigrationDescriptor, MigrationEvidenceError, + MigrationKind, MigrationName, MigrationTransactionExecutor, }; pub use open::{OpenMode, ServiceSqlitePathError, ServiceSqlitePaths}; pub use status::{StorageHealth, StorageIntegrity, StorageStatus}; diff --git a/crates/service_sqlite/src/migration.rs b/crates/service_sqlite/src/migration.rs @@ -1,26 +1,20 @@ //! Deterministic migration identity, ledger, and governed execution mechanics. use core::fmt; -use std::{collections::BTreeSet, error::Error}; +use std::{collections::BTreeSet, error::Error, future::Future, pin::Pin}; #[cfg(any(target_os = "linux", target_os = "macos"))] -use std::{ - collections::BTreeMap, - future::Future, - pin::Pin, - sync::{ - Arc, - atomic::{AtomicBool, Ordering}, - }, -}; +use std::collections::BTreeMap; use sha2::{Digest, Sha256}; +use sqlx::SqliteConnection; #[cfg(any(target_os = "linux", target_os = "macos"))] -use sqlx::{Connection, Row, SqliteConnection}; +use sqlx::{Connection, Row}; +use crate::ServiceSqliteError; #[cfg(any(target_os = "linux", target_os = "macos"))] -use crate::{SchemaCatalog, ServiceSqliteError, ServiceSqliteErrorKind}; +use crate::{SchemaCatalog, 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"; @@ -479,19 +473,13 @@ impl fmt::Display for MigrationEvidenceError { impl Error for MigrationEvidenceError {} -#[cfg(any(target_os = "linux", target_os = "macos"))] #[derive(Clone, Copy, Debug, PartialEq, Eq)] -pub(crate) struct MigrationApplicationOutcome { +pub struct MigrationApplicationOutcome { initial_version: u32, final_version: u32, applied_count: u32, } -#[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 059 keeps application outcomes private until the Step 061 host boundary" -)] impl MigrationApplicationOutcome { /// Returns the schema version observed under the migration transaction. #[must_use] @@ -634,128 +622,48 @@ fn catalog_digest(descriptors: &[MigrationDescriptor]) -> MigrationChecksum { MigrationChecksum(hasher.finalize().into()) } -#[cfg(any(target_os = "linux", target_os = "macos"))] -pub(crate) type MigrationCallbackFuture<'a> = +pub type MigrationCallbackFuture<'a> = Pin<Box<dyn Future<Output = Result<(), ServiceSqliteError>> + Send + 'a>>; -#[cfg(any(target_os = "linux", target_os = "macos"))] -type MigrationCallback = +pub type MigrationCallback = for<'a> fn(&'a mut MigrationTransactionExecutor<'_>) -> MigrationCallbackFuture<'a>; -#[cfg(any(target_os = "linux", target_os = "macos"))] -pub(crate) struct MigrationTransactionExecutor<'a> { +pub struct MigrationTransactionExecutor<'a> { connection: &'a mut SqliteConnection, + database_control_rejected: bool, } -#[cfg(any(target_os = "linux", target_os = "macos"))] impl MigrationTransactionExecutor<'_> { - pub(crate) async fn execute(&mut self, sql: &'static str) -> Result<(), ServiceSqliteError> { - assert_governed_transaction(self.connection).await?; - let execution = sqlx::raw_sql(sql).execute(&mut *self.connection).await; - let transaction = assert_governed_transaction(self.connection).await; - transaction?; - execution - .map(|_| ()) - .map_err(|source| migration_source(MigrationFailureKind::Execution, source)) - } -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -struct MigrationCommitGate { - allow_commit: Arc<AtomicBool>, - allow_runner_rollback: Arc<AtomicBool>, - rollback_observed: Arc<AtomicBool>, -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -impl MigrationCommitGate { - async fn install(connection: &mut SqliteConnection) -> Result<Self, ServiceSqliteError> { - let allow_commit = Arc::new(AtomicBool::new(false)); - let allow_runner_rollback = Arc::new(AtomicBool::new(false)); - let rollback_observed = Arc::new(AtomicBool::new(false)); - let hook_permission = Arc::clone(&allow_commit); - let rollback_permission = Arc::clone(&allow_runner_rollback); - let rollback_epoch = Arc::clone(&rollback_observed); - let mut handle = connection - .lock_handle() - .await - .map_err(|source| migration_source(MigrationFailureKind::Execution, source))?; - handle.set_commit_hook(move || hook_permission.load(Ordering::Acquire)); - handle.set_rollback_hook(move || { - if !rollback_permission.load(Ordering::Acquire) { - rollback_epoch.store(true, Ordering::Release); - } - }); - drop(handle); - Ok(Self { - allow_commit, - allow_runner_rollback, - rollback_observed, - }) - } - - fn permit_outer_commit(&self) -> MigrationCommitPermit { - self.allow_commit.store(true, Ordering::Release); - MigrationCommitPermit { - allow_commit: Arc::clone(&self.allow_commit), + pub async fn execute(&mut self, sql: &'static str) -> Result<(), ServiceSqliteError> { + if crate::connection::contains_database_control(sql) { + self.database_control_rejected = true; + return Err(ServiceSqliteError::new( + crate::ServiceSqliteErrorKind::Migration, + )); } - } - - fn permit_runner_rollback(&self) -> MigrationRollbackPermit { - self.allow_runner_rollback.store(true, Ordering::Release); - MigrationRollbackPermit { - allow_runner_rollback: Arc::clone(&self.allow_runner_rollback), + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + assert_governed_transaction(self.connection).await?; + let execution = sqlx::raw_sql(sql).execute(&mut *self.connection).await; + let transaction = assert_governed_transaction(self.connection).await; + transaction?; + execution + .map(|_| ()) + .map_err(|source| migration_source(MigrationFailureKind::Execution, source)) } - } - - fn reject_observed_rollback(&self) -> Result<(), ServiceSqliteError> { - if self.rollback_observed.load(Ordering::Acquire) { - Err(migration_error(MigrationFailureKind::Execution)) - } else { - Ok(()) + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = (&mut self.connection, sql); + Err(ServiceSqliteError::new( + crate::ServiceSqliteErrorKind::Migration, + )) } } - - async fn remove(self, connection: &mut SqliteConnection) -> Result<(), ServiceSqliteError> { - self.allow_commit.store(false, Ordering::Release); - self.allow_runner_rollback.store(false, Ordering::Release); - let mut handle = connection - .lock_handle() - .await - .map_err(|source| migration_source(MigrationFailureKind::Commit, source))?; - handle.remove_commit_hook(); - handle.remove_rollback_hook(); - Ok(()) - } -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -struct MigrationCommitPermit { - allow_commit: Arc<AtomicBool>, -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -impl Drop for MigrationCommitPermit { - fn drop(&mut self) { - self.allow_commit.store(false, Ordering::Release); - } -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -struct MigrationRollbackPermit { - allow_runner_rollback: Arc<AtomicBool>, -} - -#[cfg(any(target_os = "linux", target_os = "macos"))] -impl Drop for MigrationRollbackPermit { - fn drop(&mut self) { - self.allow_runner_rollback.store(false, Ordering::Release); - } } #[cfg(any(target_os = "linux", target_os = "macos"))] #[derive(PartialEq, Eq)] -struct MigrationConnectionPolicy { +pub(crate) struct MigrationConnectionPolicy { application_id: i64, journal_mode: String, synchronous: i64, @@ -765,32 +673,38 @@ struct MigrationConnectionPolicy { query_only: i64, } -#[cfg(any(target_os = "linux", target_os = "macos"))] #[derive(Clone, Copy)] -pub(crate) struct MigrationCallbackBinding { +pub struct MigrationCallbackBinding { + #[cfg(any(target_os = "linux", target_os = "macos"))] target_version: u32, + #[cfg(any(target_os = "linux", target_os = "macos"))] name: MigrationName, + #[cfg(any(target_os = "linux", target_os = "macos"))] checksum: MigrationChecksum, + #[cfg(any(target_os = "linux", target_os = "macos"))] callback: MigrationCallback, } -#[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 059 keeps callback bindings private until the Step 061 host boundary" -)] impl MigrationCallbackBinding { - pub(crate) const fn new( + pub const fn new( target_version: u32, name: MigrationName, checksum: MigrationChecksum, callback: MigrationCallback, ) -> Self { - Self { - target_version, - name, - checksum, - callback, + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + Self { + target_version, + name, + checksum, + callback, + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + let _ = (target_version, name, checksum, callback); + Self {} } } } @@ -922,7 +836,7 @@ pub(crate) async fn verify_migration_history( } #[cfg(any(target_os = "linux", target_os = "macos"))] -async fn verify_migration_history_snapshot( +pub(crate) async fn verify_migration_history_snapshot( connection: &mut SqliteConnection, catalog: &MigrationCatalog, schema_catalog: &SchemaCatalog, @@ -994,7 +908,9 @@ where loop { validate_authority()?; - let gate_result = MigrationCommitGate::install(connection).await; + let gate_result = crate::transaction_control::TransactionControlGate::install(connection) + .await + .map_err(|source| migration_source(MigrationFailureKind::Execution, source)); validate_authority()?; let commit_gate = gate_result?; let transaction_result = connection.begin_with("BEGIN IMMEDIATE").await; @@ -1014,7 +930,10 @@ where validate_authority()?; rollback_result .map_err(|source| migration_source(MigrationFailureKind::Commit, source))?; - let remove_result = commit_gate.remove(connection).await; + let remove_result = commit_gate + .remove(connection) + .await + .map_err(|source| migration_source(MigrationFailureKind::Commit, source)); validate_authority()?; remove_result?; break; @@ -1029,7 +948,9 @@ where let execution_result = execute_descriptor(&mut transaction, descriptor, &callbacks).await; validate_authority()?; execution_result?; - commit_gate.reject_observed_rollback()?; + if commit_gate.control_violation_observed() { + return Err(migration_error(MigrationFailureKind::Execution)); + } let transaction_result = assert_governed_transaction(&mut transaction).await; validate_authority()?; transaction_result?; @@ -1063,7 +984,9 @@ where let transaction_result = assert_governed_transaction(&mut transaction).await; validate_authority()?; transaction_result?; - commit_gate.reject_observed_rollback()?; + if commit_gate.control_violation_observed() { + return Err(migration_error(MigrationFailureKind::Execution)); + } let policy_result = read_connection_policy(&mut transaction).await; validate_authority()?; if policy_result? != initial_policy { @@ -1072,13 +995,18 @@ where let transaction_result = assert_governed_transaction(&mut transaction).await; validate_authority()?; transaction_result?; - commit_gate.reject_observed_rollback()?; + if commit_gate.control_violation_observed() { + return Err(migration_error(MigrationFailureKind::Execution)); + } let permit = commit_gate.permit_outer_commit(); let commit_result = transaction.commit().await; drop(permit); validate_authority()?; commit_result.map_err(|source| migration_source(MigrationFailureKind::Commit, source))?; - let remove_result = commit_gate.remove(connection).await; + let remove_result = commit_gate + .remove(connection) + .await + .map_err(|source| migration_source(MigrationFailureKind::Commit, source)); validate_authority()?; remove_result?; let observed = after_commit(); @@ -1137,7 +1065,10 @@ async fn execute_descriptor( descriptor: &MigrationDescriptor, callbacks: &BTreeMap<u32, MigrationCallback>, ) -> Result<(), ServiceSqliteError> { - let mut executor = MigrationTransactionExecutor { connection }; + let mut executor = MigrationTransactionExecutor { + connection, + database_control_rejected: false, + }; match descriptor.kind() { MigrationKind::Sql => { let sql = core::str::from_utf8(descriptor.content) @@ -1154,11 +1085,14 @@ async fn execute_descriptor( assert_governed_transaction(executor.connection).await?; } } + if executor.database_control_rejected { + return Err(migration_error(MigrationFailureKind::Execution)); + } Ok(()) } #[cfg(any(target_os = "linux", target_os = "macos"))] -async fn assert_governed_transaction( +pub(crate) async fn assert_governed_transaction( connection: &mut SqliteConnection, ) -> Result<(), ServiceSqliteError> { sqlx::raw_sql( @@ -1172,7 +1106,7 @@ async fn assert_governed_transaction( } #[cfg(any(target_os = "linux", target_os = "macos"))] -async fn read_connection_policy( +pub(crate) async fn read_connection_policy( connection: &mut SqliteConnection, ) -> Result<MigrationConnectionPolicy, ServiceSqliteError> { let text = |source| migration_source(MigrationFailureKind::Execution, source); @@ -1504,7 +1438,10 @@ mod tests { use std::{ num::NonZeroU32, path::{Path, PathBuf}, - sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}, + sync::{ + Mutex, + atomic::{AtomicUsize, Ordering as AtomicOrdering}, + }, }; #[cfg(any(target_os = "linux", target_os = "macos"))] @@ -1712,9 +1649,26 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + fn ignored_attachment_callback<'a>( + executor: &'a mut MigrationTransactionExecutor<'_>, + ) -> MigrationCallbackFuture<'a> { + let sql = IGNORED_ATTACHMENT_SQL + .lock() + .expect("attachment SQL mutex") + .expect("attachment SQL is installed"); + Box::pin(async move { + let _ = executor.execute(sql).await; + Ok(()) + }) + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] static PENDING_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0); #[cfg(any(target_os = "linux", target_os = "macos"))] + static IGNORED_ATTACHMENT_SQL: Mutex<Option<&'static str>> = Mutex::new(None); + + #[cfg(any(target_os = "linux", target_os = "macos"))] async fn replace_with_permissive_ledger(connection: &mut SqliteConnection) { sqlx::raw_sql( "DROP TRIGGER schema_migrations_no_update; @@ -2343,6 +2297,100 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] + async fn migration_executor_rejects_transient_attachment_before_file_creation() { + let directory = tempfile::tempdir().unwrap(); + let database_path = directory.path().join("main.sqlite"); + let external_path = directory.path().join("external.sqlite"); + let migration_sql = Box::leak( + format!( + "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", + external_path.display() + ) + .into_boxed_str(), + ); + let mut connection = initialized_file_database(&database_path).await; + let catalog = + MigrationCatalog::new([sql(2, "reject_transient_attachment", migration_sql)]).unwrap(); + let schema_catalog = unchanged_schema_catalog(&catalog); + let mut validate = || Ok(()); + let error = apply_governed_migrations( + &mut connection, + &catalog, + &schema_catalog, + MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(), + &build_identity(), + &[], + &mut validate, + ) + .await + .expect_err("migration ATTACH/DETACH must fail before SQLite compilation"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration); + assert!(!external_path.exists()); + assert_eq!(read_state_schema_version(&mut connection).await.unwrap(), 1); + assert!( + read_migration_history(&mut connection) + .await + .unwrap() + .is_empty() + ); + + const CALLBACK_DEFINITION: &[u8] = b"callback:reject_transient_attachment:v1"; + let callback_database_path = directory.path().join("callback-main.sqlite"); + let callback_external_path = directory.path().join("callback-external.sqlite"); + let callback_sql = Box::leak( + format!( + "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", + callback_external_path.display() + ) + .into_boxed_str(), + ); + *IGNORED_ATTACHMENT_SQL.lock().expect("attachment SQL mutex") = Some(callback_sql); + let mut callback_connection = initialized_file_database(&callback_database_path).await; + let callback_descriptor = MigrationDescriptor::callback( + 2, + "reject_callback_attachment", + CALLBACK_DEFINITION, + MigrationChecksum::for_callback(CALLBACK_DEFINITION), + ) + .unwrap(); + let callback = MigrationCallbackBinding::new( + callback_descriptor.target_version(), + callback_descriptor.name(), + callback_descriptor.checksum(), + ignored_attachment_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, + MigrationAppliedAtUnixSeconds::new(1_800_000_001).unwrap(), + &build_identity(), + &[callback], + &mut validate, + ) + .await + .expect_err("ignored callback ATTACH/DETACH must still fail the migration"); + *IGNORED_ATTACHMENT_SQL.lock().expect("attachment SQL mutex") = None; + assert_eq!(callback_error.kind(), ServiceSqliteErrorKind::Migration); + assert!(!callback_external_path.exists()); + assert_eq!( + read_state_schema_version(&mut callback_connection) + .await + .unwrap(), + 1 + ); + assert!( + read_migration_history(&mut callback_connection) + .await + .unwrap() + .is_empty() + ); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] async fn assert_fresh_connection_has_no_migration_effect(path: &Path, table: &str) { let mut connection = SqliteConnection::connect_with( &SqliteConnectOptions::new() diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs @@ -182,16 +182,14 @@ impl OpenMode { } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps the pool private until the Step 061 host boundary" -)] -struct PrivateConnectionPool { +pub(crate) struct PrivateConnectionPool { pool: SqlitePool, binding: DirectoryBinding, paths: ServiceSqlitePaths, + identity: ServiceDatabaseIdentity, catalog: MigrationCatalog, schema_catalog: SchemaCatalog, + mode: OpenMode, authority: Option<WriterAuthority>, inspection_guard: Option<ReadOnlyInspectionGuard>, authority_failure: Arc<AtomicBool>, @@ -202,10 +200,6 @@ struct PrivateConnectionPool { } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps pool lifecycle private until the Step 061 host boundary" -)] impl PrivateConnectionPool { fn connection_failure_kind(&self) -> ServiceSqliteErrorKind { connection_failure_kind( @@ -217,10 +211,36 @@ impl PrivateConnectionPool { ) } - async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> { - self.binding.validate(&self.paths)?; + pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> { + if let Some(authority) = self.authority.as_ref() { + authority.validate_for(&self.paths)?; + } + if let Some(inspection_guard) = self.inspection_guard.as_ref() { + inspection_guard.validate_for(&self.paths)?; + } + self.binding.validate(&self.paths) + } + + pub(crate) const fn mode(&self) -> OpenMode { + self.mode + } + + pub(crate) fn identity(&self) -> &ServiceDatabaseIdentity { + &self.identity + } + + pub(crate) fn catalog(&self) -> &MigrationCatalog { + &self.catalog + } + + pub(crate) fn schema_catalog(&self) -> &SchemaCatalog { + &self.schema_catalog + } + + pub(crate) async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> { + self.validate()?; let result = self.pool.acquire().await; - self.binding.validate(&self.paths)?; + self.validate()?; let mut connection = result.map_err(|source| connection_source(self.connection_failure_kind(), source))?; let history = crate::migration::verify_migration_history( @@ -230,12 +250,12 @@ impl PrivateConnectionPool { true, ) .await; - self.binding.validate(&self.paths)?; + self.validate()?; history?; Ok(connection) } - async fn apply_migrations( + pub(crate) async fn apply_migrations( &self, applied_at: MigrationAppliedAtUnixSeconds, build: &MigrationBuildIdentity, @@ -275,11 +295,16 @@ impl PrivateConnectionPool { result } - async fn close(mut self) -> Option<WriterAuthority> { + pub(crate) async fn close(mut self) -> Option<WriterAuthority> { self.pool.close().await; self.inspection_guard.take(); self.authority.take() } + + #[cfg(test)] + pub(crate) async fn close_pool(&self) { + self.pool.close().await; + } } #[cfg(any(target_os = "linux", target_os = "macos"))] @@ -306,11 +331,7 @@ const fn connection_failure_kind( } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps pool opening private until the Step 061 host boundary" -)] -async fn open_existing_connection_pool( +pub(crate) async fn open_existing_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, catalog: &MigrationCatalog, @@ -343,11 +364,7 @@ async fn open_existing_connection_pool( } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps pool opening private until the Step 061 host boundary" -)] -async fn open_initialized_connection_pool( +pub(crate) async fn open_initialized_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, catalog: &MigrationCatalog, @@ -370,11 +387,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" -)] +#[allow(clippy::too_many_arguments)] async fn open_connection_pool( paths: &ServiceSqlitePaths, identity: &ServiceDatabaseIdentity, @@ -450,6 +463,7 @@ async fn open_connection_pool( let before_metadata = identity.clone(); let retained_catalog = catalog.clone(); let retained_schema_catalog = schema_catalog.clone(); + let retained_identity = identity.clone(); let after_catalog = catalog.clone(); let after_schema_catalog = schema_catalog.clone(); let before_catalog = catalog.clone(); @@ -648,8 +662,10 @@ async fn open_connection_pool( pool, binding: retained_binding, paths: paths.clone(), + identity: retained_identity, catalog: retained_catalog, schema_catalog: retained_schema_catalog, + mode, authority, inspection_guard, authority_failure, @@ -810,21 +826,17 @@ fn connection_source(kind: ServiceSqliteErrorKind, cause: sqlx::Error) -> Servic } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps inspection authority private until the Step 061 host boundary" -)] struct ReadOnlyInspectionGuard { lock: File, + lock_device: u64, + lock_inode: u64, directory: File, + directory_device: u64, + directory_inode: u64, _database: File, } #[cfg(any(target_os = "linux", target_os = "macos"))] -#[allow( - dead_code, - reason = "Step 056 keeps inspection authority private until the Step 061 host boundary" -)] impl ReadOnlyInspectionGuard { fn acquire(paths: &ServiceSqlitePaths) -> Result<Self, ServiceSqliteError> { use rustix::{ @@ -936,10 +948,102 @@ impl ReadOnlyInspectionGuard { } Ok(Self { lock, + lock_device: u64::try_from(lock_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, + lock_inode: lock_status.st_ino, directory: File::from(directory), + directory_device: u64::try_from(directory_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, + directory_inode: directory_status.st_ino, _database: database, }) } + + fn validate_for(&self, paths: &ServiceSqlitePaths) -> Result<(), ServiceSqliteError> { + use rustix::{ + fs::{AtFlags, FileType, Mode, OFlags, fstat, open, openat, statat}, + process::geteuid, + }; + + let directory_path = paths + .state_lock() + .parent() + .filter(|parent| Some(*parent) == paths.state_database().parent()) + .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let directory = open( + directory_path, + OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC, + Mode::empty(), + ) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let directory_status = fstat(&directory) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let held_directory_status = fstat(&self.directory) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let directory_device = u64::try_from(directory_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let held_directory_device = u64::try_from(held_directory_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + if !FileType::from_raw_mode(directory_status.st_mode).is_dir() + || directory_status.st_uid != geteuid().as_raw() + || u32::from(directory_status.st_mode) & 0o022 != 0 + || !FileType::from_raw_mode(held_directory_status.st_mode).is_dir() + || held_directory_status.st_uid != geteuid().as_raw() + || u32::from(held_directory_status.st_mode) & 0o022 != 0 + || directory_device != self.directory_device + || directory_status.st_ino != self.directory_inode + || held_directory_device != self.directory_device + || held_directory_status.st_ino != self.directory_inode + { + return Err(inspection_error( + ConnectionFailureKind::InspectionUnavailable, + )); + } + + let lock = openat( + &directory, + radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME, + OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK, + Mode::empty(), + ) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let lock_status = fstat(&lock) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let held_lock_status = fstat(&self.lock) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let lock_device = u64::try_from(lock_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let held_lock_device = u64::try_from(held_lock_status.st_dev) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + if !FileType::from_raw_mode(lock_status.st_mode).is_file() + || u64::from(lock_status.st_nlink) != 1 + || lock_status.st_uid != geteuid().as_raw() + || u32::from(lock_status.st_mode) & 0o777 != 0o600 + || !FileType::from_raw_mode(held_lock_status.st_mode).is_file() + || u64::from(held_lock_status.st_nlink) != 1 + || held_lock_status.st_uid != geteuid().as_raw() + || u32::from(held_lock_status.st_mode) & 0o777 != 0o600 + || lock_device != self.lock_device + || lock_status.st_ino != self.lock_inode + || held_lock_device != self.lock_device + || held_lock_status.st_ino != self.lock_inode + { + return Err(inspection_error( + ConnectionFailureKind::InspectionUnavailable, + )); + } + for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] { + match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) { + Err(error) if error == rustix::io::Errno::NOENT => {} + Ok(_) | Err(_) => { + return Err(inspection_error( + ConnectionFailureKind::InspectionUnavailable, + )); + } + } + } + Ok(()) + } } #[cfg(any(target_os = "linux", target_os = "macos"))] @@ -1094,6 +1198,12 @@ impl DirectoryBinding { || directory_status.st_ino != self.directory_inode || held_directory_device != self.directory_device || held_directory_status.st_ino != self.directory_inode + || !FileType::from_raw_mode(directory_status.st_mode).is_dir() + || directory_status.st_uid != geteuid().as_raw() + || u32::from(directory_status.st_mode) & 0o022 != 0 + || !FileType::from_raw_mode(held_directory_status.st_mode).is_dir() + || held_directory_status.st_uid != geteuid().as_raw() + || u32::from(held_directory_status.st_mode) & 0o022 != 0 { return Err(connection_error( ServiceSqliteErrorKind::Authority, diff --git a/crates/service_sqlite/src/transaction_control.rs b/crates/service_sqlite/src/transaction_control.rs @@ -0,0 +1,109 @@ +//! Private deny-by-default SQLite transaction-control fencing. + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use std::sync::{ + Arc, + atomic::{AtomicBool, Ordering}, +}; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +use sqlx::SqliteConnection; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) struct TransactionControlGate { + allow_commit: Arc<AtomicBool>, + allow_runner_rollback: Arc<AtomicBool>, + rejected_commit: Arc<AtomicBool>, + rollback_observed: Arc<AtomicBool>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl TransactionControlGate { + pub(crate) async fn install(connection: &mut SqliteConnection) -> Result<Self, sqlx::Error> { + let allow_commit = Arc::new(AtomicBool::new(false)); + let allow_runner_rollback = Arc::new(AtomicBool::new(false)); + let rejected_commit = Arc::new(AtomicBool::new(false)); + let rollback_observed = Arc::new(AtomicBool::new(false)); + let hook_permission = Arc::clone(&allow_commit); + let rejected_commit_epoch = Arc::clone(&rejected_commit); + let rollback_permission = Arc::clone(&allow_runner_rollback); + let rollback_epoch = Arc::clone(&rollback_observed); + let mut handle = connection.lock_handle().await?; + handle.set_commit_hook(move || { + let permitted = hook_permission.load(Ordering::Acquire); + if !permitted { + rejected_commit_epoch.store(true, Ordering::Release); + } + permitted + }); + handle.set_rollback_hook(move || { + if !rollback_permission.load(Ordering::Acquire) { + rollback_epoch.store(true, Ordering::Release); + } + }); + drop(handle); + Ok(Self { + allow_commit, + allow_runner_rollback, + rejected_commit, + rollback_observed, + }) + } + + pub(crate) fn permit_outer_commit(&self) -> TransactionCommitPermit { + self.allow_commit.store(true, Ordering::Release); + TransactionCommitPermit { + allow_commit: Arc::clone(&self.allow_commit), + } + } + + pub(crate) fn permit_runner_rollback(&self) -> TransactionRollbackPermit { + self.allow_runner_rollback.store(true, Ordering::Release); + TransactionRollbackPermit { + allow_runner_rollback: Arc::clone(&self.allow_runner_rollback), + } + } + + pub(crate) fn control_violation_observed(&self) -> bool { + self.rejected_commit.load(Ordering::Acquire) + || self.rollback_observed.load(Ordering::Acquire) + } + + pub(crate) fn rejected_commit_rolled_back(&self) -> bool { + self.rejected_commit.load(Ordering::Acquire) + && self.rollback_observed.load(Ordering::Acquire) + } + + pub(crate) async fn remove(self, connection: &mut SqliteConnection) -> Result<(), sqlx::Error> { + self.allow_commit.store(false, Ordering::Release); + self.allow_runner_rollback.store(false, Ordering::Release); + let mut handle = connection.lock_handle().await?; + handle.remove_commit_hook(); + handle.remove_rollback_hook(); + Ok(()) + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) struct TransactionCommitPermit { + allow_commit: Arc<AtomicBool>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl Drop for TransactionCommitPermit { + fn drop(&mut self) { + self.allow_commit.store(false, Ordering::Release); + } +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) struct TransactionRollbackPermit { + allow_runner_rollback: Arc<AtomicBool>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +impl Drop for TransactionRollbackPermit { + fn drop(&mut self) { + self.allow_runner_rollback.store(false, Ordering::Release); + } +} diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs @@ -1,9 +1,11 @@ use std::collections::BTreeSet; const MANIFEST: &str = include_str!("../Cargo.toml"); +const README: &str = include_str!("../README.md"); const ROOT: &str = include_str!("../src/lib.rs"); const AUTHORITY_SOURCE: &str = include_str!("../src/authority.rs"); const CONFIG_SOURCE: &str = include_str!("../src/config.rs"); +const CONNECTION_SOURCE: &str = include_str!("../src/connection.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"); @@ -12,6 +14,7 @@ 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"); const STATUS_SOURCE: &str = include_str!("../src/status.rs"); +const TRANSACTION_CONTROL_SOURCE: &str = include_str!("../src/transaction_control.rs"); #[test] fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { @@ -31,6 +34,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { dependency_keys(MANIFEST, "[dependencies]"), BTreeSet::from([ "fs2", + "futures", "radroots_runtime_paths", "radroots_storage", "rustix", @@ -48,16 +52,40 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { BTreeSet::from([ "authority", "config", + "connection", "error", "initialize", "integrity", "metadata", "migration", "open", - "status" + "status", + "transaction_control" ]) ); assert!(public_modules(ROOT).is_empty()); + for required in [ + "`ServiceSqliteHost` is the only public connection host", + "borrowed `ServiceSqliteTransaction` executor", + "transaction begin, commit, rollback, policy", + "attached-database exclusion", + "Writable host opening finishes every pending governed migration", + "read-only inspection opens only current migration and", + "with raw database authority", + "before the runner enables outer commit", + "leaves no authoritative transaction effect", + "only after rollback is confirmed", + "unconfirmed rollback is reported as `RollbackFailed`", + "Cancelling once outer", + "must be treated as an unknown commit outcome", + "require rereading authoritative state", + "before an idempotent retry", + ] { + assert!( + README.contains(required), + "Step 061 README contract is missing `{required}`" + ); + } let authority_production = AUTHORITY_SOURCE .split_once("#[cfg(all(test") .map(|(production, _)| production) @@ -85,7 +113,11 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "ServiceSqliteApplicationId", "ServiceSqliteMetadataValueError", "MigrationAppliedAtUnixSeconds", + "MigrationApplicationOutcome", "MigrationBuildIdentity", + "MigrationCallback", + "MigrationCallbackBinding", + "MigrationCallbackFuture", "MigrationCatalog", "MigrationChecksum", "MigrationContractError", @@ -93,6 +125,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "MigrationEvidenceError", "MigrationKind", "MigrationName", + "MigrationTransactionExecutor", "SchemaCatalog", "SchemaCatalogContractError", "SchemaDigest", @@ -102,6 +135,11 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "ServiceSqlitePathError", "ServiceSqlitePaths", "OpenMode", + "ServiceSqliteHost", + "ServiceSqliteTransaction", + "ServiceSqliteTransactionError", + "ServiceSqliteTransactionErrorKind", + "ServiceSqliteTransactionFuture", "StorageHealth", "StorageIntegrity", "StorageStatus", @@ -127,12 +165,9 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub(crate) async fn apply_governed_migrations", "validate_callback_bindings", "advance_schema_version", - "pub(crate) struct MigrationTransactionExecutor", - "set_commit_hook", - "set_rollback_hook", - "permit_outer_commit", - "permit_runner_rollback", - "reject_observed_rollback", + "pub struct MigrationTransactionExecutor", + "contains_database_control", + "database_control_rejected", "SAVEPOINT radroots_migration_transaction_probe", "FROM pragma_database_list", "CASE WHEN typeof(name) = 'text'", @@ -145,6 +180,67 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { } for required in [ + "pub struct ServiceSqliteHost", + "pub struct ServiceSqliteTransaction<'connection>", + "impl<'executor, 'connection> Executor<'executor>", + "pub enum ServiceSqliteTransactionErrorKind", + "pub struct ServiceSqliteTransactionError<E>", + "pub type ServiceSqliteTransactionFuture", + "BEGIN IMMEDIATE", + "OpenMode::ReadOnlyInspection => connection.begin().await", + "verify_before_commit", + "connection.close_on_drop()", + "connection.trust()", + "RestrictedExecute", + "contains_database_control", + "RADROOTS_FORBIDDEN_DATABASE_CONTROL", + "OperationRolledBack", + "RollbackFailed", + "CommitOutcomeUnknown", + "cancelling the future yields no result", + "must reread authoritative state before any idempotent retry", + ] { + assert!( + CONNECTION_SOURCE.contains(required), + "Step 061 connection source is missing `{required}`" + ); + } + + for required in [ + "set_commit_hook", + "set_rollback_hook", + "permit_outer_commit", + "permit_runner_rollback", + "control_violation_observed", + "rejected_commit", + "rejected_commit_rolled_back", + "remove_commit_hook", + "remove_rollback_hook", + ] { + assert!( + TRANSACTION_CONTROL_SOURCE.contains(required), + "Step 061 transaction-control source is missing `{required}`" + ); + } + + for forbidden in [ + "pub fn pool(", + "pub fn connection(", + "pub fn into_inner(", + "Deref for ServiceSqliteTransaction", + "AsRef<SqliteConnection>", + "pub async fn begin(", + "pub async fn commit(", + "pub async fn rollback(", + "pub use sqlx", + ] { + assert!( + !CONNECTION_SOURCE.contains(forbidden) && !ROOT.contains(forbidden), + "Step 061 source exposes forbidden raw authority `{forbidden}`" + ); + } + + for required in [ "radroots.service_sqlite.schema_object.v1\\0", "radroots.service_sqlite.schema_snapshot.v1\\0", "radroots.service_sqlite.schema_catalog.v1\\0", @@ -193,9 +289,6 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub fn content(&self", "pub fn callback_definition", "pub fn migration_sql", - "pub type MigrationCallback", - "pub struct MigrationCallbackBinding", - "pub struct MigrationTransactionExecutor", "pub fn apply_governed_migrations", "pub async fn apply_governed_migrations", "pub use sqlx", @@ -317,9 +410,6 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub struct PrivateConnectionPool", "pub fn open_connection_pool", "pub async fn open_connection_pool", - "pub struct MigrationCallbackBinding", - "pub type MigrationCallback", - "pub struct MigrationApplicationOutcome", "tokio::runtime", "Runtime::new", ] { @@ -379,6 +469,12 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub fn release(&mut self)", "database_path: paths.state_database().to_path_buf()", "pub(crate) fn validate_for", + "validate_authority_binding", + "directory_device", + "directory_inode", + "lock_device", + "lock_inode", + "current_lock_status", ] { assert!( AUTHORITY_SOURCE.contains(required), @@ -386,6 +482,22 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { ); } + for required in [ + "inspection_guard.validate_for(&self.paths)", + "fn validate_for(&self, paths: &ServiceSqlitePaths)", + "held_lock_status", + "lock_device", + "directory_device", + "WAL_FILE_NAME", + "SHARED_MEMORY_FILE_NAME", + "u32::from(directory_status.st_mode) & 0o022", + ] { + assert!( + OPEN_SOURCE.contains(required), + "Step 061 live authority source is missing `{required}`" + ); + } + for forbidden in [ "remove_file", "create_dir",