lib

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

commit 053d0c750bf9cd683c6ea37cefe7e79617ba629f
parent 74b0181f4585932d20e8b60097eed1dfc18d77b5
Author: triesap <tyson@radroots.org>
Date:   Thu, 27 Aug 2026 20:50:47 +0000

service-sqlite: seal initialization transaction

Diffstat:
Mcontracts/api_baselines/radroots_service_sqlite.txt | 12+++++++++++-
Mcrates/service_sqlite/README.md | 27+++++++++++++++++++++++----
Mcrates/service_sqlite/src/connection.rs | 237++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
Mcrates/service_sqlite/src/initialize.rs | 692+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----------------
Mcrates/service_sqlite/src/lib.rs | 4+++-
Mcrates/service_sqlite/src/metadata.rs | 47+++++++++++++++++++++++++++++++----------------
Mcrates/service_sqlite/src/open.rs | 36+++++++++++++++++++++++++-----------
Mcrates/service_sqlite/src/restore/process_tests.rs | 22++++++----------------
Mcrates/service_sqlite/src/restore/stage.rs | 11++---------
Mcrates/service_sqlite/tests/package_boundary.rs | 58++++++++++++++++++++++++++++++++++++++++++++++++++++++++--
10 files changed, 938 insertions(+), 208 deletions(-)

diff --git a/contracts/api_baselines/radroots_service_sqlite.txt b/contracts/api_baselines/radroots_service_sqlite.txt @@ -383,6 +383,7 @@ pub async fn radroots_service_sqlite::ServiceSqliteHost::close(&self) -> core::r pub async fn radroots_service_sqlite::ServiceSqliteHost::inspect_integrity(&self, radroots_service_sqlite::IntegrityCheckedAtUnixMs) -> core::result::Result<radroots_service_sqlite::ServiceSqliteIntegrityReport, radroots_service_sqlite::ServiceSqliteError> pub const fn radroots_service_sqlite::ServiceSqliteHost::mode(&self) -> radroots_service_sqlite::OpenMode pub async fn radroots_service_sqlite::ServiceSqliteHost::open_initialized(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ServiceDatabaseIdentity, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::ServiceSqliteConnectionOptions, radroots_service_sqlite::WriterAuthority, radroots_service_sqlite::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::MigrationBuildIdentity, &[radroots_service_sqlite::MigrationCallbackBinding]) -> core::result::Result<(Self, radroots_service_sqlite::MigrationApplicationOutcome), radroots_service_sqlite::ServiceSqliteError> +pub async fn radroots_service_sqlite::ServiceSqliteHost::open_or_initialize<F, E>(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ServiceDatabaseMetadata, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::ServiceSqliteConnectionOptions, radroots_service_sqlite::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::MigrationBuildIdentity, &[radroots_service_sqlite::MigrationCallbackBinding], F) -> core::result::Result<(Self, radroots_service_sqlite::MigrationApplicationOutcome), radroots_service_sqlite::ServiceSqliteError> where F: for<'a> core::ops::function::FnOnce(&'a mut radroots_service_sqlite::ServiceSqliteInitializer<'_>) -> radroots_service_sqlite::ServiceSqliteInitializerFuture<'a, E>, E: core::error::Error + core::marker::Send + core::marker::Sync + 'static pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_only_inspection(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ServiceDatabaseIdentity, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::ServiceSqliteConnectionOptions) -> core::result::Result<Self, radroots_service_sqlite::ServiceSqliteError> pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_only_inspection_with_intent(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ExistingServiceDatabaseIntent, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::ServiceSqliteConnectionOptions) -> core::result::Result<radroots_service_sqlite::OpenedExistingServiceDatabase, radroots_service_sqlite::ServiceSqliteError> pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_write_existing(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ServiceDatabaseIdentity, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::ServiceSqliteConnectionOptions, radroots_service_sqlite::MigrationAppliedAtUnixSeconds, &radroots_service_sqlite::MigrationBuildIdentity, &[radroots_service_sqlite::MigrationCallbackBinding]) -> core::result::Result<(Self, radroots_service_sqlite::MigrationApplicationOutcome), radroots_service_sqlite::ServiceSqliteError> @@ -390,6 +391,14 @@ pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_write_existin pub async fn radroots_service_sqlite::ServiceSqliteHost::transaction<T, E, F>(&self, F) -> core::result::Result<T, radroots_service_sqlite::ServiceSqliteTransactionError<E>> where T: core::marker::Send + 'static, E: core::marker::Send + 'static, F: for<'a> core::ops::function::FnOnce(&'a mut radroots_service_sqlite::ServiceSqliteTransaction<'_>) -> radroots_service_sqlite::ServiceSqliteTransactionFuture<'a, T, E> + core::marker::Send impl core::fmt::Debug for radroots_service_sqlite::ServiceSqliteHost pub fn radroots_service_sqlite::ServiceSqliteHost::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +pub struct radroots_service_sqlite::ServiceSqliteInitializer<'connection> +impl core::fmt::Debug for radroots_service_sqlite::ServiceSqliteInitializer<'_> +pub fn radroots_service_sqlite::ServiceSqliteInitializer<'_>::fmt(&self, &mut core::fmt::Formatter<'_>) -> core::fmt::Result +impl<'executor, 'connection> sqlx_core::executor::Executor<'executor> for &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection> where 'connection: 'executor +pub type &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection>::Database = sqlx_sqlite::database::Sqlite +pub fn &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection>::fetch_many<'e, 'q: 'e, Q>(self, Q) -> futures_core::stream::BoxStream<'e, core::result::Result<either::Either<sqlx_sqlite::query_result::SqliteQueryResult, sqlx_sqlite::row::SqliteRow>, sqlx_core::error::Error>> where Q: 'q + sqlx_core::executor::Execute<'q, Self::Database>, 'executor: 'e +pub fn &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection>::fetch_optional<'e, 'q: 'e, Q>(self, Q) -> futures_core::future::BoxFuture<'e, core::result::Result<core::option::Option<sqlx_sqlite::row::SqliteRow>, sqlx_core::error::Error>> where Q: 'q + sqlx_core::executor::Execute<'q, Self::Database>, 'executor: 'e +pub fn &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection>::prepare_with<'e>(self, sqlx_core::sql_str::SqlStr, &'e [sqlx_sqlite::type_info::SqliteTypeInfo]) -> futures_core::future::BoxFuture<'e, core::result::Result<sqlx_sqlite::statement::SqliteStatement, sqlx_core::error::Error>> where 'executor: 'e pub struct radroots_service_sqlite::ServiceSqliteIntegrityReport impl radroots_service_sqlite::ServiceSqliteIntegrityReport pub const fn radroots_service_sqlite::ServiceSqliteIntegrityReport::checked_at_unix_ms(&self) -> radroots_service_sqlite::IntegrityCheckedAtUnixMs @@ -464,10 +473,11 @@ pub fn radroots_service_sqlite::StateFilesystemCapacitySource::available_bytes(& impl radroots_service_sqlite::StateFilesystemCapacitySource for radroots_service_sqlite::PlatformStateFilesystemCapacitySource pub fn radroots_service_sqlite::PlatformStateFilesystemCapacitySource::available_bytes(&self, &radroots_service_sqlite::ServiceSqlitePaths) -> core::result::Result<u64, radroots_service_sqlite::StateFilesystemCapacityError> pub async fn radroots_service_sqlite::finalize_staged_restore(radroots_service_sqlite::StagedServiceRestore) -> core::result::Result<(), radroots_service_sqlite::ServiceSqliteError> -pub async fn radroots_service_sqlite::initialize_database<F, Fut, E>(&radroots_service_sqlite::ServiceSqlitePaths, radroots_service_sqlite::OpenMode, &radroots_service_sqlite::ServiceDatabaseMetadata, &radroots_service_sqlite::SchemaCatalog, F) -> core::result::Result<radroots_service_sqlite::WriterAuthority, radroots_service_sqlite::ServiceSqliteError> where F: core::ops::function::FnOnce(std::path::PathBuf) -> Fut, Fut: core::future::future::Future<Output = core::result::Result<(), E>>, E: core::error::Error + core::marker::Send + core::marker::Sync + 'static +pub async fn radroots_service_sqlite::initialize_database<F, E>(&radroots_service_sqlite::ServiceSqlitePaths, radroots_service_sqlite::OpenMode, &radroots_service_sqlite::ServiceDatabaseMetadata, &radroots_service_sqlite::SchemaCatalog, F) -> core::result::Result<radroots_service_sqlite::WriterAuthority, radroots_service_sqlite::ServiceSqliteError> where F: for<'a> core::ops::function::FnOnce(&'a mut radroots_service_sqlite::ServiceSqliteInitializer<'_>) -> radroots_service_sqlite::ServiceSqliteInitializerFuture<'a, E>, E: core::error::Error + core::marker::Send + core::marker::Sync + 'static pub fn radroots_service_sqlite::inspect_state_filesystem_capacity<S: radroots_service_sqlite::StateFilesystemCapacitySource + ?core::marker::Sized>(&radroots_service_sqlite::ServiceSqlitePaths, radroots_service_sqlite::MinimumFreeBytes, &S) -> core::result::Result<radroots_service_sqlite::StateFilesystemCapacity, radroots_service_sqlite::StateFilesystemCapacityError> pub async fn radroots_service_sqlite::stage_verified_restore(&radroots_service_sqlite::ServiceSqlitePaths, &radroots_service_sqlite::ServiceDatabaseIdentity, &radroots_service_sqlite::MigrationCatalog, &radroots_service_sqlite::SchemaCatalog, radroots_service_sqlite::VerifiedServiceBackup) -> core::result::Result<radroots_service_sqlite::StagedServiceRestore, radroots_service_sqlite::ServiceSqliteError> pub fn radroots_service_sqlite::verify_backup_bundle(&[u8], radroots_service_sqlite::BackupManifestSha256, &std::path::Path, &radroots_service_sqlite::ServiceDatabaseIdentity, core::num::nonzero::NonZeroU64) -> core::result::Result<radroots_service_sqlite::VerifiedServiceBackup, radroots_service_sqlite::ServiceSqliteError> pub type radroots_service_sqlite::MigrationCallback = for<'a> fn(&'a mut radroots_service_sqlite::MigrationTransactionExecutor<'_>) -> radroots_service_sqlite::MigrationCallbackFuture<'a> pub type radroots_service_sqlite::MigrationCallbackFuture<'a> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = core::result::Result<(), radroots_service_sqlite::ServiceSqliteError>> + core::marker::Send + 'a)>> +pub type radroots_service_sqlite::ServiceSqliteInitializerFuture<'a, E> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = core::result::Result<(), E>> + core::marker::Send + 'a)>> pub type radroots_service_sqlite::ServiceSqliteTransactionFuture<'a, T, E> = core::pin::Pin<alloc::boxed::Box<(dyn core::future::future::Future<Output = core::result::Result<T, E>> + core::marker::Send + 'a)>> diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md @@ -15,6 +15,25 @@ runner-owned. Writable host opening finishes every pending governed migration before returning, and read-only inspection opens only current migration and schema state. +Create-new state uses the borrowed `ServiceSqliteInitializer` executor rather +than a database path or raw connection. The runner reserves and retains the +exact canonical file, opens SQLite only through that retained descriptor, owns +`BEGIN IMMEDIATE`, and uses a private memory journal so descriptor-bound +initialization creates no path-derived SQLite sidecar. It commits the product +schema, shared `radroots_service_metadata`, empty v1 `schema_migrations` +ledger, and exact schema-catalog verification in one transaction. The +initializer screens the same closed transaction-control and attachment +inventory as host transactions; an ignored rejection still prevents commit. +Callback failure or cancellation cannot return a reusable connection or +publish a partial database. + +Interactive callers may use `ServiceSqliteHost::open_or_initialize`. One held +writer authority and an exclusive create decide whether the sealed initializer +runs or the exact existing database is opened. The existing branch never runs +the initializer. Callers therefore do not use pathname probes, error-text +matching, recursive directory creation, permission repair, or direct SQLx +connections to choose the bootstrap branch. + Existing databases can be admitted without a caller guessing their stored source generation. `ExistingServiceDatabaseIntent` seals the canonical service and instance, supported schema ceiling, and SQLite application ID; @@ -331,7 +350,7 @@ The crate-root exports are frozen in the reviewed [service-SQLite API baseline](../../contracts/api_baselines/radroots_service_sqlite.txt). Raw pools, pooled or direct connections, transaction-control handles, and dependency re-exports are forbidden. The deliberate narrow exception is the -`sqlx::Executor` implementation for a borrowed -`&mut ServiceSqliteTransaction<'_>`: it permits compile-time typed queries -while the crate retains connection ownership and sole begin, commit, rollback, -policy, and cancellation authority. +`sqlx::Executor` implementation for borrowed +`&mut ServiceSqliteInitializer<'_>` and `&mut ServiceSqliteTransaction<'_>` +values. They permit compile-time typed queries. The crate retains connection +ownership and sole begin, commit, rollback, policy, and cancellation authority. diff --git a/crates/service_sqlite/src/connection.rs b/crates/service_sqlite/src/connection.rs @@ -351,6 +351,100 @@ impl ServiceSqliteHost { } } + /// Atomically creates interactive state or opens the exact existing database. + /// + /// The runner holds writer authority while an exclusive create decides the + /// branch. Callers never probe the filesystem or inspect error text. On the + /// create branch, the sealed initializer, shared metadata, empty v1 ledger, + /// and schema-catalog verification commit together before the host opens and + /// applies governed migrations. On exact create collision, the same retained + /// authority is transferred to an existing-only open and the initialization + /// callback is never invoked. + #[allow(clippy::too_many_arguments)] + pub async fn open_or_initialize<F, E>( + paths: &ServiceSqlitePaths, + initialization_metadata: &ServiceDatabaseMetadata, + migrations: &MigrationCatalog, + schema: &SchemaCatalog, + options: ServiceSqliteConnectionOptions, + applied_at: MigrationAppliedAtUnixSeconds, + build: &MigrationBuildIdentity, + callbacks: &[MigrationCallbackBinding], + initialize_schema: F, + ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> + where + F: for<'a> FnOnce( + &'a mut crate::ServiceSqliteInitializer<'_>, + ) -> crate::ServiceSqliteInitializerFuture<'a, E>, + E: Error + Send + Sync + 'static, + { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + if !schema.matches_migrations(migrations) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)); + } + let supported_version = core::num::NonZeroU32::new(migrations.current_version()) + .expect("migration catalogs always have a nonzero current version"); + let initialized = crate::initialize::initialize_or_existing_database( + paths, + initialization_metadata, + schema, + initialize_schema, + ) + .await?; + let (mode, pool) = match initialized { + crate::initialize::InitializeDatabaseOutcome::Initialized(authority) => { + let identity = ServiceDatabaseIdentity::new( + paths, + initialization_metadata.source_generation(), + supported_version, + initialization_metadata.application_id(), + ); + let pool = crate::open::open_initialized_connection_pool( + paths, &identity, migrations, schema, options, authority, + ) + .await?; + (OpenMode::Initialize, pool) + } + crate::initialize::InitializeDatabaseOutcome::Existing(authority) => { + let intent = ExistingServiceDatabaseIntent::new( + paths, + supported_version, + initialization_metadata.application_id(), + ); + let pool = + crate::open::open_existing_connection_pool_with_intent_and_authority( + paths, &intent, migrations, schema, options, authority, + ) + .await?; + (OpenMode::ReadWriteExisting, pool) + } + }; + match pool.apply_migrations(applied_at, build, callbacks).await { + Ok(outcome) => Ok((Self::from_pool(mode, pool), outcome)), + Err(error) => { + drop(pool.close().await); + Err(error) + } + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + drop(( + paths, + initialization_metadata, + migrations, + schema, + options, + applied_at, + build, + callbacks, + initialize_schema, + )); + Err(unsupported_host()) + } + } + /// Opens state created under a retained initialization writer authority. #[allow(clippy::too_many_arguments)] pub async fn open_initialized( @@ -1535,19 +1629,14 @@ mod tests { 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>(()) + |initializer| { + Box::pin(async move { + sqlx::query(HOST_TABLE_SQL) + .execute(initializer) + .await + .expect("create host table"); + Ok::<_, Infallible>(()) + }) }, ) .await @@ -1668,6 +1757,128 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] + async fn open_or_initialize_uses_exclusive_create_without_filesystem_probing() { + 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 = ServiceDatabaseMetadata::new( + &paths, + SourceGeneration::new([31; 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 initialized_callback = Arc::new(AtomicBool::new(false)); + let called = Arc::clone(&initialized_callback); + let (initialized, outcome) = ServiceSqliteHost::open_or_initialize( + &paths, + &metadata, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"), + &build_identity(), + &[], + move |initializer| { + called.store(true, Ordering::Release); + Box::pin(async move { + sqlx::query(HOST_TABLE_SQL).execute(initializer).await?; + Ok::<(), sqlx::Error>(()) + }) + }, + ) + .await + .expect("initialize missing state"); + assert!(initialized_callback.load(Ordering::Acquire)); + assert_eq!(initialized.mode(), OpenMode::Initialize); + assert_eq!(outcome.applied_count(), 0); + initialized.close().await.expect("close initialized host"); + + let existing_callback = Arc::new(AtomicBool::new(false)); + let called = Arc::clone(&existing_callback); + let (existing, outcome) = ServiceSqliteHost::open_or_initialize( + &paths, + &metadata, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"), + &build_identity(), + &[], + move |_| { + called.store(true, Ordering::Release); + Box::pin(async { Ok::<(), Infallible>(()) }) + }, + ) + .await + .expect("open exact existing state"); + assert!(!existing_callback.load(Ordering::Acquire)); + assert_eq!(existing.mode(), OpenMode::ReadWriteExisting); + assert_eq!(outcome.applied_count(), 0); + assert_eq!(row_count(&existing).await, 0); + existing.close().await.expect("close existing host"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test(flavor = "current_thread")] + async fn open_or_initialize_rejects_catalog_mismatch_before_reservation() { + const MIGRATION_SQL: &str = "CREATE TABLE later_probe (value INTEGER NOT NULL)"; + + 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 = ServiceDatabaseMetadata::new( + &paths, + SourceGeneration::new([32; 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 base_migrations = migration_catalog(); + let schema = schema_catalog(&base_migrations); + let different_migrations = MigrationCatalog::new([crate::MigrationDescriptor::sql( + 2, + "create_later_probe", + MIGRATION_SQL, + crate::MigrationChecksum::for_sql(MIGRATION_SQL), + ) + .expect("different migration")]) + .expect("different migration catalog"); + let callback_called = Arc::new(AtomicBool::new(false)); + let called = Arc::clone(&callback_called); + + let error = ServiceSqliteHost::open_or_initialize( + &paths, + &metadata, + &different_migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + MigrationAppliedAtUnixSeconds::new(1_700_000_002).expect("migration time"), + &build_identity(), + &[], + move |_| { + called.store(true, Ordering::Release); + Box::pin(async { Ok::<(), Infallible>(()) }) + }, + ) + .await + .expect_err("catalog mismatch must fail before reservation"); + + assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity); + assert!(!callback_called.load(Ordering::Acquire)); + assert!(!paths.state_database().exists()); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] #[tokio::test] async fn integrity_inspection_is_explicit_safe_and_available_in_every_host_mode() { let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await; diff --git a/crates/service_sqlite/src/initialize.rs b/crates/service_sqlite/src/initialize.rs @@ -1,22 +1,172 @@ //! Create-new initialization for one service-owned SQLite database. use core::{fmt, future::Future}; -use std::{error::Error, path::PathBuf}; +use std::{ + error::Error, + 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::{ OpenMode, SchemaCatalog, ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqlitePaths, WriterAuthority, }; +/// A sealed initialization executor that never exposes its SQLite connection. +/// +/// Service repositories may create their product schema with ordinary typed +/// SQLx queries through a mutable borrow. The service-SQLite runner exclusively +/// owns transaction begin, commit, rollback, shared metadata, and the migration +/// ledger. +/// +/// ``` +/// use radroots_service_sqlite::ServiceSqliteInitializer; +/// +/// async fn create_product_schema( +/// initializer: &mut ServiceSqliteInitializer<'_>, +/// ) -> Result<(), sqlx::Error> { +/// sqlx::query(concat!("CREATE ", "TABLE product_items (id INTEGER PRIMARY KEY)")) +/// .execute(initializer) +/// .await?; +/// Ok(()) +/// } +/// ``` +/// +/// Transaction control and the raw connection remain inaccessible: +/// +/// ```compile_fail +/// use radroots_service_sqlite::ServiceSqliteInitializer; +/// +/// async fn bypass(initializer: ServiceSqliteInitializer<'_>) { +/// initializer.commit().await.unwrap(); +/// } +/// ``` +pub struct ServiceSqliteInitializer<'connection> { + connection: &'connection mut SqliteConnection, + statement_control_rejected: Arc<AtomicBool>, +} + +struct RestrictedInitializationExecute<Q> { + query: Q, + statement_control_rejected: Arc<AtomicBool>, +} + +impl<'query, Q> Execute<'query, Sqlite> for RestrictedInitializationExecute<Q> +where + Q: Execute<'query, Sqlite>, +{ + fn sql(self) -> SqlStr { + restricted_initialization_sql(self.query.sql(), &self.statement_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_initialization_sql(sql: SqlStr, statement_control_rejected: &AtomicBool) -> SqlStr { + if crate::statement_policy::contains_forbidden_statement_control(sql.as_str()) { + statement_control_rejected.store(true, Ordering::Release); + SqlStr::from_static("RADROOTS_FORBIDDEN_INITIALIZATION_STATEMENT_CONTROL") + } else { + sql + } +} + +impl fmt::Debug for ServiceSqliteInitializer<'_> { + fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result { + formatter.write_str("ServiceSqliteInitializer([redacted])") + } +} + +impl<'executor, 'connection> Executor<'executor> + for &'executor mut ServiceSqliteInitializer<'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(RestrictedInitializationExecute { + query, + statement_control_rejected: Arc::clone(&self.statement_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(RestrictedInitializationExecute { + query, + statement_control_rejected: Arc::clone(&self.statement_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_initialization_sql(sql, &self.statement_control_rejected), + parameters, + ) + } +} + +/// Boxed callback future tied to the borrowed initialization executor. +pub type ServiceSqliteInitializerFuture<'a, E> = + Pin<Box<dyn Future<Output = Result<(), E>> + Send + 'a>>; + +#[cfg(any(target_os = "linux", target_os = "macos"))] +#[derive(Debug)] +pub(crate) enum InitializeDatabaseOutcome { + Initialized(WriterAuthority), + Existing(WriterAuthority), +} + /// Creates and initializes a missing service database while holding sole writer authority. /// -/// The callback receives the already-reserved canonical database path. It must -/// open that file without create or replacement flags, initialize its schema, -/// close every database handle, and only then resolve its future. The supplied -/// metadata must derive from the same paths; it is written and verified after -/// the callback but before the database becomes durable. Cancellation or -/// failure removes only the exact inode reserved by this call. -pub async fn initialize_database<F, Fut, E>( +/// The callback receives only a mutable borrow of the sealed initialization +/// executor. Product schema, shared metadata, the empty v1 migration ledger, +/// and catalog verification occur in one runner-owned transaction. Cancellation +/// or failure quarantines the one-shot connection and removes only the exact +/// inode reserved by this call. +pub async fn initialize_database<F, E>( paths: &ServiceSqlitePaths, mode: OpenMode, metadata: &ServiceDatabaseMetadata, @@ -24,8 +174,9 @@ pub async fn initialize_database<F, Fut, E>( initialize_schema: F, ) -> Result<WriterAuthority, ServiceSqliteError> where - F: FnOnce(PathBuf) -> Fut, - Fut: Future<Output = Result<(), E>>, + F: for<'a> FnOnce( + &'a mut ServiceSqliteInitializer<'_>, + ) -> ServiceSqliteInitializerFuture<'a, E>, E: Error + Send + Sync + 'static, { if mode != OpenMode::Initialize { @@ -54,7 +205,7 @@ where #[cfg(any(target_os = "linux", target_os = "macos"))] { let failpoints = crate::failpoint::DurabilityFailpoints::default(); - initialize_with_ops( + match initialize_with_ops( paths, authority, metadata, @@ -63,7 +214,13 @@ where &SystemInitializationOperations, &failpoints, ) - .await + .await? + { + InitializeDatabaseOutcome::Initialized(authority) => Ok(authority), + InitializeDatabaseOutcome::Existing(_authority) => Err(initialization_error( + InitializationCause::new(InitializationFailureKind::StateAlreadyExists), + )), + } } #[cfg(not(any(target_os = "linux", target_os = "macos")))] @@ -75,6 +232,40 @@ where } } +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn initialize_or_existing_database<F, E>( + paths: &ServiceSqlitePaths, + metadata: &ServiceDatabaseMetadata, + schema_catalog: &SchemaCatalog, + initialize_schema: F, +) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> +where + F: for<'a> FnOnce( + &'a mut ServiceSqliteInitializer<'_>, + ) -> ServiceSqliteInitializerFuture<'a, E>, + E: Error + Send + Sync + 'static, +{ + if !metadata.matches_paths(paths) { + return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata)); + } + let authority = WriterAuthority::acquire(paths, OpenMode::Initialize)?.ok_or_else(|| { + initialization_error(InitializationCause::new( + InitializationFailureKind::UnsupportedMode, + )) + })?; + let failpoints = crate::failpoint::DurabilityFailpoints::default(); + initialize_with_ops( + paths, + authority, + metadata, + schema_catalog, + initialize_schema, + &SystemInitializationOperations, + &failpoints, + ) + .await +} + #[derive(Clone, Copy, Debug, PartialEq, Eq)] enum InitializationFailureKind { UnsupportedMode, @@ -271,7 +462,7 @@ mod failure_tests { #[cfg(any(target_os = "linux", target_os = "macos"))] mod supported { - use std::fs::File; + use std::{fs::File, os::fd::AsRawFd}; use rustix::{ fs::{AtFlags, FileType, Mode, OFlags, fchmod, fstat, lstat, openat, statat, unlinkat}, @@ -368,6 +559,15 @@ mod supported { self.validate_entry() } + fn sqlite_descriptor_path(&self) -> String { + let descriptor = self.database.as_raw_fd(); + #[cfg(target_os = "linux")] + let path = format!("/proc/self/fd/{descriptor}"); + #[cfg(target_os = "macos")] + let path = format!("/dev/fd/{descriptor}"); + path + } + fn validate_entry(&self) -> Result<(), InitializationCause> { let status = statat( self.directory, @@ -553,14 +753,14 @@ mod supported { async fn fail_with_rollback<O: InitializationOperations>( mut pending: PendingDatabase<'_, O>, primary: InitializationCause, - ) -> Result<WriterAuthority, ServiceSqliteError> { + ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> { Err(rollback_failure(primary, pending.rollback())) } async fn fail_metadata_with_rollback<O: InitializationOperations>( mut pending: PendingDatabase<'_, O>, primary: ServiceSqliteError, - ) -> Result<WriterAuthority, ServiceSqliteError> { + ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> { match pending.rollback() { Ok(()) => Err(primary), Err(cleanup) => Err(initialization_error(InitializationCause::with_source( @@ -570,7 +770,7 @@ mod supported { } } - pub(super) async fn initialize_with_ops<F, Fut, E, O>( + pub(super) async fn initialize_with_ops<F, E, O>( paths: &ServiceSqlitePaths, authority: WriterAuthority, metadata: &ServiceDatabaseMetadata, @@ -578,10 +778,11 @@ mod supported { initialize_schema: F, operations: &O, failpoints: &crate::failpoint::DurabilityFailpoints, - ) -> Result<WriterAuthority, ServiceSqliteError> + ) -> Result<InitializeDatabaseOutcome, ServiceSqliteError> where - F: FnOnce(PathBuf) -> Fut, - Fut: Future<Output = Result<(), E>>, + F: for<'a> FnOnce( + &'a mut ServiceSqliteInitializer<'_>, + ) -> ServiceSqliteInitializerFuture<'a, E>, E: Error + Send + Sync + 'static, O: InitializationOperations, { @@ -590,8 +791,19 @@ mod supported { crate::failpoint::DurabilityFailpoint::InitializeBeforeCreate, ) .map_err(initialization_error)?; - let mut pending = PendingDatabase::create(authority.directory(), operations) - .map_err(initialization_error)?; + let initialization_directory = authority.directory().try_clone().map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::CreateUnavailable, + source, + )) + })?; + let mut pending = match PendingDatabase::create(&initialization_directory, operations) { + Ok(pending) => pending, + Err(error) if error.kind == InitializationFailureKind::StateAlreadyExists => { + return Ok(InitializeDatabaseOutcome::Existing(authority)); + } + Err(error) => return Err(initialization_error(error)), + }; if let Err(error) = hit( failpoints, crate::failpoint::DurabilityFailpoint::InitializeAfterCreate, @@ -613,47 +825,25 @@ mod supported { ) { return fail_with_rollback(pending, error).await; } - if let Err(error) = pending.validate_canonical_path(paths.state_database()) { - return fail_with_rollback(pending, error).await; - } - let callback_result = initialize_schema(paths.state_database().to_path_buf()).await; - if let Err(error) = callback_result { - return fail_with_rollback( - pending, - InitializationCause::with_source( - InitializationFailureKind::SchemaInitializationFailed, - error, - ), - ) - .await; - } - if let Err(error) = pending.validate() { - return fail_with_rollback(pending, error).await; + authority.validate_for(paths)?; + let recovery = crate::restore::refuse_unresolved_recovery(authority.directory()); + authority.validate_for(paths)?; + if let Err(error) = recovery { + return fail_metadata_with_rollback(pending, error).await; } if let Err(error) = pending.validate_canonical_path(paths.state_database()) { return fail_with_rollback(pending, error).await; } - let metadata_result = async { - use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions}; - - let options = SqliteConnectOptions::new() - .filename(paths.state_database()) - .create_if_missing(false) - .disable_statement_logging(); - let mut connection = sqlx::SqliteConnection::connect_with(&options) - .await - .map_err(|source| { - ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source) - })?; - let write_result = - 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) - }); - write_result.and(close_result) - } + let metadata_result = initialize_transaction( + paths, + &authority, + &pending, + metadata, + schema_catalog, + initialize_schema, + ) .await; + authority.validate_for(paths)?; if let Err(error) = pending.validate() { return fail_with_rollback(pending, error).await; } @@ -673,7 +863,188 @@ mod supported { return fail_with_rollback(pending, error).await; } drop(pending); - Ok(authority) + Ok(InitializeDatabaseOutcome::Initialized(authority)) + } + + async fn initialize_transaction<F, E, O>( + paths: &ServiceSqlitePaths, + authority: &WriterAuthority, + pending: &PendingDatabase<'_, O>, + metadata: &ServiceDatabaseMetadata, + schema_catalog: &SchemaCatalog, + initialize_schema: F, + ) -> Result<(), ServiceSqliteError> + where + F: for<'a> FnOnce( + &'a mut ServiceSqliteInitializer<'_>, + ) -> ServiceSqliteInitializerFuture<'a, E>, + E: Error + Send + Sync + 'static, + O: InitializationOperations, + { + use sqlx::{ + ConnectOptions, Connection, + sqlite::{SqliteConnectOptions, SqliteJournalMode}, + }; + + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let options = SqliteConnectOptions::new() + .filename(pending.sqlite_descriptor_path()) + .create_if_missing(false) + .foreign_keys(true) + .journal_mode(SqliteJournalMode::Memory) + .pragma("trusted_schema", "OFF") + .disable_statement_logging(); + let connected = SqliteConnection::connect_with(&options).await; + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let mut connection = connected.map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + )) + })?; + let installed = + crate::transaction_control::TransactionControlGate::install(&mut connection).await; + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let gate = installed.map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + )) + })?; + let begun = connection.begin_with("BEGIN IMMEDIATE").await; + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let mut transaction = begun.map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + )) + })?; + + let statement_control_rejected = Arc::new(AtomicBool::new(false)); + let callback = { + let mut initializer = ServiceSqliteInitializer { + connection: &mut transaction, + statement_control_rejected: Arc::clone(&statement_control_rejected), + }; + initialize_schema(&mut initializer).await + }; + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let governed_after_callback = + crate::migration::assert_governed_transaction(&mut transaction) + .await + .is_ok(); + authority.validate_for(paths)?; + + let operation = match callback { + Ok(()) + if !statement_control_rejected.load(Ordering::Acquire) + && !gate.control_violation_observed() + && governed_after_callback => + { + crate::metadata::write_database_metadata_in_transaction( + &mut transaction, + metadata, + schema_catalog, + ) + .await + } + Ok(()) => Err(initialization_error(InitializationCause::new( + InitializationFailureKind::SchemaInitializationFailed, + ))), + Err(source) => Err(initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + ))), + }; + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let governed_before_commit = + crate::migration::assert_governed_transaction(&mut transaction) + .await + .is_ok(); + authority.validate_for(paths)?; + + if let Err(primary) = operation { + let permit = gate.permit_runner_rollback(); + let rollback = transaction.rollback().await; + drop(permit); + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let remove = gate.remove(&mut connection).await; + authority.validate_for(paths)?; + let close = connection.close().await; + authority.validate_for(paths)?; + if let Err(source) = rollback.or(remove).or(close) { + return Err(initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + ))); + } + return Err(primary); + } + + if gate.control_violation_observed() || !governed_before_commit { + let permit = gate.permit_runner_rollback(); + let rollback = transaction.rollback().await; + drop(permit); + let remove = gate.remove(&mut connection).await; + let close = connection.close().await; + authority.validate_for(paths)?; + rollback.or(remove).or(close).map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + )) + })?; + return Err(initialization_error(InitializationCause::new( + InitializationFailureKind::SchemaInitializationFailed, + ))); + } + + let permit = gate.permit_outer_commit(); + let committed = transaction.commit().await; + drop(permit); + authority.validate_for(paths)?; + pending + .validate() + .and_then(|()| pending.validate_canonical_path(paths.state_database())) + .map_err(initialization_error)?; + let removed = gate.remove(&mut connection).await; + authority.validate_for(paths)?; + let closed = connection.close().await; + authority.validate_for(paths)?; + committed.or(removed).or(closed).map_err(|source| { + initialization_error(InitializationCause::with_source( + InitializationFailureKind::SchemaInitializationFailed, + source, + )) + }) } fn hit( @@ -689,14 +1060,14 @@ mod supported { mod tests { use std::{ cell::{Cell, RefCell}, + convert::Infallible, fs, - future::{Future, pending, ready}, + future::{pending, ready}, io, num::NonZeroU32, os::unix::fs::{MetadataExt, PermissionsExt, symlink}, path::Path, - pin::Pin, - task::{Context, Poll, Waker}, + sync::Arc, }; use radroots_runtime_paths::{ @@ -705,6 +1076,7 @@ mod supported { ServiceId, }; use radroots_storage::event::SourceGeneration; + use tokio::sync::Notify; use super::*; @@ -764,13 +1136,6 @@ mod supported { } } - fn poll_once<F: Future>(future: F) -> (Poll<F::Output>, Pin<Box<F>>) { - let mut future = Box::pin(future); - let mut context = Context::from_waker(Waker::noop()); - let result = future.as_mut().poll(&mut context); - (result, future) - } - fn paths(root: &Path, instance: &str) -> ServiceSqlitePaths { let context = RuntimeContext::resolve( &RadrootsPathResolver::new( @@ -845,7 +1210,6 @@ mod supported { let root = tempfile::tempdir().expect("root"); let paths = paths(root.path(), "success"); prepare(&paths); - let expected_path = paths.state_database().to_path_buf(); let metadata = metadata(&paths); let schema_catalog = service_schema_catalog(); let mut authority = initialize_database( @@ -853,20 +1217,17 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - |path| async move { - assert_eq!(path, expected_path); - 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 service_schema (id INTEGER PRIMARY KEY)") - .execute(&mut connection) - .await?; - connection.close().await?; - Ok::<(), sqlx::Error>(()) + |initializer| { + Box::pin(async move { + assert_eq!( + format!("{initializer:?}"), + "ServiceSqliteInitializer([redacted])" + ); + sqlx::query("CREATE TABLE service_schema (id INTEGER PRIMARY KEY)") + .execute(initializer) + .await?; + Ok::<(), sqlx::Error>(()) + }) }, ) .await @@ -910,7 +1271,7 @@ mod supported { &schema_catalog, |_| { called.set(true); - ready(Ok::<(), CallbackFailure>(())) + Box::pin(ready(Ok::<(), CallbackFailure>(()))) }, ) .await @@ -931,7 +1292,7 @@ mod supported { let schema_catalog = base_schema_catalog(); let error = initialize_database(&paths, mode, &metadata, &schema_catalog, |_| { called.set(true); - ready(Ok::<(), CallbackFailure>(())) + Box::pin(ready(Ok::<(), CallbackFailure>(()))) }) .await .expect_err("mode must reject"); @@ -957,7 +1318,7 @@ mod supported { &base_schema_catalog(), |_| { called.set(true); - ready(Ok::<(), CallbackFailure>(())) + Box::pin(ready(Ok::<(), CallbackFailure>(()))) }, ) .await @@ -980,7 +1341,7 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - |_| ready(Err::<(), _>(CallbackFailure)), + |_| Box::pin(ready(Err::<(), _>(CallbackFailure))), ) .await .expect_err("callback failure"); @@ -1008,7 +1369,7 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - |_| ready(Ok::<(), CallbackFailure>(())), + |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), ) .await .expect("retry after cleanup"); @@ -1017,6 +1378,60 @@ mod supported { } #[tokio::test(flavor = "current_thread")] + async fn transaction_control_and_attachments_are_rejected_before_sqlite_compilation() { + for scenario in ["commit", "rollback_begin", "attach_detach"] { + let root = tempfile::tempdir().expect("root"); + let paths = paths(root.path(), scenario); + prepare(&paths); + let external = root.path().join("must-not-exist.sqlite"); + let external_for_callback = external.clone(); + let metadata = metadata(&paths); + let schema_catalog = base_schema_catalog(); + let error = initialize_database( + &paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + move |initializer| { + Box::pin(async move { + let sql = match scenario { + "commit" => "COMMIT".to_owned(), + "rollback_begin" => { + "ROLLBACK; BEGIN DEFERRED; CREATE TABLE escaped(value INTEGER)" + .to_owned() + } + "attach_detach" => format!( + "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra", + external_for_callback.display() + ), + _ => unreachable!(), + }; + let _ = sqlx::query(sqlx::AssertSqlSafe(sql.as_str())) + .execute(initializer) + .await; + Ok::<(), Infallible>(()) + }) + }, + ) + .await + .expect_err("forbidden statement control must fail initialization"); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); + for rendered in [error.to_string(), format!("{error:?}")] { + assert!(!rendered.contains("must-not-exist")); + assert!(!rendered.contains("ATTACH")); + assert!(!rendered.contains("ROLLBACK")); + } + assert!(!paths.state_database().exists()); + assert!(!external.exists()); + assert!( + WriterAuthority::acquire(&paths, OpenMode::Initialize) + .expect("reacquire") + .is_some() + ); + } + } + + #[tokio::test(flavor = "current_thread")] async fn metadata_or_migration_ledger_conflict_cleans_the_exact_reserved_database() { let root = tempfile::tempdir().expect("root"); for (instance, statement, expected_kind) in [ @@ -1040,17 +1455,11 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - move |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(statement).execute(&mut connection).await?; - connection.close().await?; - Ok::<(), sqlx::Error>(()) + move |initializer| { + Box::pin(async move { + sqlx::query(statement).execute(initializer).await?; + Ok::<(), sqlx::Error>(()) + }) }, ) .await @@ -1078,19 +1487,13 @@ mod supported { 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>(()) + |initializer| { + Box::pin(async move { + sqlx::query("CREATE TABLE unexpected (value INTEGER)") + .execute(initializer) + .await?; + Ok::<(), sqlx::Error>(()) + }) }, ) .await @@ -1111,17 +1514,32 @@ mod supported { 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()); + let task_paths = paths.clone(); + let reached = Arc::new(Notify::new()); + let reached_from_callback = Arc::clone(&reached); + let task = tokio::spawn(async move { + initialize_database( + &task_paths, + OpenMode::Initialize, + &metadata, + &schema_catalog, + move |initializer| { + Box::pin(async move { + sqlx::query("CREATE TABLE cancelled_probe (value INTEGER)") + .execute(initializer) + .await?; + reached_from_callback.notify_one(); + pending::<Result<(), sqlx::Error>>().await + }) + }, + ) + .await + }); + reached.notified().await; assert!(paths.state_database().exists()); assert!(WriterAuthority::acquire(&paths, OpenMode::Initialize).is_err()); - drop(future); + task.abort(); + task.await.expect_err("initialization task is cancelled"); assert!(!paths.state_database().exists()); assert!( WriterAuthority::acquire(&paths, OpenMode::Initialize) @@ -1143,10 +1561,12 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - move |path| async move { - fs::remove_file(&path)?; - fs::write(&replacement_path, b"replacement")?; - Ok::<(), io::Error>(()) + move |_| { + Box::pin(async move { + fs::remove_file(&replacement_path)?; + fs::write(&replacement_path, b"replacement")?; + Ok::<(), io::Error>(()) + }) }, ) .await @@ -1172,18 +1592,20 @@ mod supported { OpenMode::Initialize, &metadata, &schema_catalog, - move |_| async move { - fs::rename(&state_directory, &displaced_for_callback)?; - fs::create_dir(&state_directory)?; - fs::write(&replacement_path, b"replacement")?; - fs::set_permissions(&replacement_path, fs::Permissions::from_mode(0o600))?; - Ok::<(), io::Error>(()) + move |_| { + Box::pin(async move { + fs::rename(&state_directory, &displaced_for_callback)?; + fs::create_dir(&state_directory)?; + fs::write(&replacement_path, b"replacement")?; + fs::set_permissions(&replacement_path, fs::Permissions::from_mode(0o600))?; + Ok::<(), io::Error>(()) + }) }, ) .await .expect_err("canonical path replacement must fail"); - assert_eq!(error.kind(), ServiceSqliteErrorKind::Create); + assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority); assert_eq!(fs::read(paths.state_database()).unwrap(), b"replacement"); assert!(!displaced_directory.join("state.sqlite").exists()); } @@ -1223,7 +1645,7 @@ mod supported { authority, &metadata, &schema_catalog, - |_| ready(Err::<(), _>(CallbackFailure)), + |_| Box::pin(ready(Err::<(), _>(CallbackFailure))), &operations, &crate::failpoint::DurabilityFailpoints::default(), ) @@ -1234,7 +1656,7 @@ mod supported { authority, &metadata, &schema_catalog, - |_| ready(Ok::<(), CallbackFailure>(())), + |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), &operations, &crate::failpoint::DurabilityFailpoints::default(), ) @@ -1386,7 +1808,7 @@ mod supported { authority, &metadata, &schema_catalog, - |_| ready(Ok::<(), CallbackFailure>(())), + |_| Box::pin(ready(Ok::<(), CallbackFailure>(()))), &operations, &failpoints, ) diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs @@ -62,7 +62,9 @@ pub use connection::{ pub use error::{ SafeServiceSqliteError, ServiceSqliteError, ServiceSqliteErrorCode, ServiceSqliteErrorKind, }; -pub use initialize::initialize_database; +pub use initialize::{ + ServiceSqliteInitializer, ServiceSqliteInitializerFuture, initialize_database, +}; pub use integrity::{ IntegrityCheckOutcome, IntegrityCheckedAtUnixMs, IntegrityDiagnosticCode, SchemaCatalog, SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind, SchemaVersionCatalog, diff --git a/crates/service_sqlite/src/metadata.rs b/crates/service_sqlite/src/metadata.rs @@ -12,7 +12,10 @@ use crate::ServiceSqlitePaths; use crate::{ServiceSqliteError, ServiceSqliteErrorKind}; #[cfg(any(target_os = "linux", target_os = "macos"))] -use sqlx::{Connection, Row, SqliteConnection}; +use sqlx::{Row, SqliteConnection}; + +#[cfg(all(test, any(target_os = "linux", target_os = "macos")))] +use sqlx::Connection; const MAX_APPLICATION_ID: u32 = i32::MAX as u32; const MAX_CREATED_AT_UNIX_MS: u64 = i64::MAX as u64; @@ -409,12 +412,33 @@ impl fmt::Display for MigrationLedgerInitializationFailure { #[cfg(any(target_os = "linux", target_os = "macos"))] impl Error for MigrationLedgerInitializationFailure {} -#[cfg(any(target_os = "linux", target_os = "macos"))] +#[cfg(all(test, any(target_os = "linux", target_os = "macos")))] pub(crate) async fn write_database_metadata( connection: &mut SqliteConnection, expected: &ServiceDatabaseMetadata, schema_catalog: &crate::SchemaCatalog, ) -> Result<(), ServiceSqliteError> { + let mut transaction = connection + .begin() + .await + .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; + write_database_metadata_in_transaction(&mut transaction, expected, schema_catalog).await?; + transaction + .commit() + .await + .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; + + let actual = read_database_metadata(connection).await?; + require_metadata_condition(actual == *expected, MetadataFailureKind::Mismatch)?; + Ok(()) +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn write_database_metadata_in_transaction( + connection: &mut SqliteConnection, + expected: &ServiceDatabaseMetadata, + schema_catalog: &crate::SchemaCatalog, +) -> Result<(), ServiceSqliteError> { if expected.state_schema_version().get() != 1 { return Err(metadata_error(MetadataFailureKind::Mismatch)); } @@ -422,19 +446,15 @@ pub(crate) async fn write_database_metadata( return Err(metadata_error(MetadataFailureKind::AlreadyPresent)); } - let mut transaction = connection - .begin() - .await - .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; for statement in crate::integrity::catalog::METADATA_SCHEMA_SQL { sqlx::query(statement) - .execute(&mut *transaction) + .execute(&mut *connection) .await .map_err(|_| metadata_error(MetadataFailureKind::AlreadyPresent))?; } for statement in crate::integrity::catalog::MIGRATION_LEDGER_SCHEMA_SQL { sqlx::query(statement) - .execute(&mut *transaction) + .execute(&mut *connection) .await .map_err(|_source| { ServiceSqliteError::with_source( @@ -457,7 +477,7 @@ pub(crate) async fn write_database_metadata( i64::try_from(expected.created_at_unix_ms()) .map_err(|_| metadata_error(MetadataFailureKind::Corrupt))?, ) - .execute(&mut *transaction) + .execute(&mut *connection) .await .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; let set_application_id = format!( @@ -466,20 +486,15 @@ pub(crate) async fn write_database_metadata( ); // The only dynamic token is a validated decimal u31 value. sqlx::query(sqlx::AssertSqlSafe(set_application_id.as_str())) - .execute(&mut *transaction) + .execute(&mut *connection) .await .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; crate::integrity::verify_schema_catalog( - &mut transaction, + &mut *connection, schema_catalog, expected.state_schema_version().get(), ) .await?; - transaction - .commit() - .await - .map_err(|_| metadata_error(MetadataFailureKind::Storage))?; - let actual = read_database_metadata(connection).await?; require_metadata_condition(actual == *expected, MetadataFailureKind::Mismatch)?; Ok(()) diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs @@ -853,6 +853,29 @@ pub(crate) async fn open_existing_connection_pool_with_intent( } #[cfg(any(target_os = "linux", target_os = "macos"))] +pub(crate) async fn open_existing_connection_pool_with_intent_and_authority( + paths: &ServiceSqlitePaths, + intent: &ExistingServiceDatabaseIntent, + catalog: &MigrationCatalog, + schema_catalog: &SchemaCatalog, + policy: ServiceSqliteConnectionOptions, + authority: WriterAuthority, +) -> Result<PrivateConnectionPool, ServiceSqliteError> { + authority.validate_for(paths)?; + open_connection_pool( + paths, + ServiceDatabaseExpectation::Existing(intent), + catalog, + schema_catalog, + OpenMode::ReadWriteExisting, + policy, + Some(authority), + None, + ) + .await +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] async fn open_existing_connection_pool_for( paths: &ServiceSqlitePaths, expectation: ServiceDatabaseExpectation<'_>, @@ -2217,16 +2240,7 @@ mod tests { OpenMode::Initialize, &metadata, &schema_catalog, - |database_path| async move { - let options = SqliteConnectOptions::new() - .filename(database_path) - .create_if_missing(false); - let connection = SqliteConnection::connect_with(&options) - .await - .expect("open reserved database"); - connection.close().await.expect("close reserved database"); - Ok::<_, Infallible>(()) - }, + |_| Box::pin(async move { Ok::<_, Infallible>(()) }), ) .await .expect("initialize database"); @@ -2572,7 +2586,7 @@ mod tests { OpenMode::Initialize, &database_metadata(&paths), &base_schema_catalog(), - |_| async { Ok::<_, Infallible>(()) }, + |_| Box::pin(async { Ok::<_, Infallible>(()) }), ) .await .expect_err("initialize must not recover existing evidence"); diff --git a/crates/service_sqlite/src/restore/process_tests.rs b/crates/service_sqlite/src/restore/process_tests.rs @@ -60,22 +60,12 @@ impl Fixture { .expect("restrict state directory"); let metadata = metadata(&paths); let (migrations, schema) = catalogs(); - let mut authority = initialize_database( - &paths, - OpenMode::Initialize, - &metadata, - &schema, - |path| async move { - let options = SqliteConnectOptions::new() - .filename(path) - .create_if_missing(false) - .disable_statement_logging(); - let connection = sqlx::SqliteConnection::connect_with(&options).await?; - connection.close().await - }, - ) - .await - .expect("initialize process database"); + let mut authority = + initialize_database(&paths, OpenMode::Initialize, &metadata, &schema, |_| { + Box::pin(async move { Ok::<(), sqlx::Error>(()) }) + }) + .await + .expect("initialize process database"); authority .release() .expect("release initialization authority"); diff --git a/crates/service_sqlite/src/restore/stage.rs b/crates/service_sqlite/src/restore/stage.rs @@ -1624,14 +1624,7 @@ mod tests { OpenMode::Initialize, &metadata, &schema, - |path| async move { - let options = SqliteConnectOptions::new() - .filename(path) - .create_if_missing(false) - .disable_statement_logging(); - let connection = SqliteConnection::connect_with(&options).await?; - connection.close().await - }, + |_| Box::pin(async move { Ok::<(), sqlx::Error>(()) }), ) .await .expect("initialize"); @@ -2605,7 +2598,7 @@ mod tests { OpenMode::Initialize, &fixture.metadata, &fixture.schema, - |_| async { Ok::<(), std::io::Error>(()) }, + |_| Box::pin(async { Ok::<(), std::io::Error>(()) }), ) .await .expect_err("unresolved recovery must reject initialization"); diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs @@ -399,6 +399,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { for required in [ "pub struct radroots_service_sqlite::ServiceSqliteHost", + "pub struct radroots_service_sqlite::ServiceSqliteInitializer", "pub struct radroots_service_sqlite::ServiceSqliteTransaction", "pub struct radroots_service_sqlite::ServiceSqlitePaths", "pub struct radroots_service_sqlite::ExistingServiceDatabaseIntent", @@ -413,6 +414,8 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub async fn radroots_service_sqlite::stage_verified_restore", "pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_write_existing_with_intent", "pub async fn radroots_service_sqlite::ServiceSqliteHost::open_read_only_inspection_with_intent", + "pub async fn radroots_service_sqlite::ServiceSqliteHost::open_or_initialize", + "impl<'executor, 'connection> sqlx_core::executor::Executor<'executor> for &'executor mut radroots_service_sqlite::ServiceSqliteInitializer<'connection>", "impl<'executor, 'connection> sqlx_core::executor::Executor<'executor> for &'executor mut radroots_service_sqlite::ServiceSqliteTransaction<'connection>", ] { assert!( @@ -456,8 +459,10 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "Every service-owned table, index, trigger,", "[service-SQLite API baseline](../../contracts/api_baselines/radroots_service_sqlite.txt)", "Raw pools, pooled or direct connections, transaction-control handles, and", - "`sqlx::Executor` implementation for a borrowed", - "while the crate retains connection ownership and sole begin, commit, rollback,", + "`sqlx::Executor` implementation for borrowed", + "`&mut ServiceSqliteInitializer<'_>` and `&mut ServiceSqliteTransaction<'_>`", + "values. They permit compile-time typed queries. The crate retains connection", + "ownership and sole begin, commit, rollback, policy, and cancellation authority.", ] { assert!(README.contains(required), "README is missing `{required}`"); } @@ -511,6 +516,8 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "ServiceSqliteConnectionOptions", "ServiceSqliteConnectionOptionsError", "initialize_database", + "ServiceSqliteInitializer", + "ServiceSqliteInitializerFuture", "ServiceDatabaseIdentity", "ServiceDatabaseMetadata", "ServiceSqliteApplicationId", @@ -1561,6 +1568,18 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "sync_directory", "validate_entry", "unlink_database", + "pub struct ServiceSqliteInitializer<'connection>", + "impl<'executor, 'connection> Executor<'executor>", + "ServiceSqliteInitializerFuture", + ".begin_with(\"BEGIN IMMEDIATE\")", + "write_database_metadata_in_transaction", + "permit_outer_commit", + "permit_runner_rollback", + "contains_forbidden_statement_control", + "RADROOTS_FORBIDDEN_INITIALIZATION_STATEMENT_CONTROL", + "pending.sqlite_descriptor_path()", + "journal_mode(SqliteJournalMode::Memory)", + "connection.close().await", ] { assert!( INITIALIZE_SOURCE.contains(required), @@ -1578,6 +1597,14 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "Command::new", "std::process", "pub fn directory", + "FnOnce(PathBuf)", + "FnOnce(std::path::PathBuf)", + "to_path_buf()).await", + ".filename(paths.state_database())", + "Deref for ServiceSqliteInitializer", + "AsRef<SqliteConnection>", + "pub fn connection(", + "pub fn into_inner(", ] { let production = INITIALIZE_SOURCE .split_once("#[cfg(test)]") @@ -1590,6 +1617,33 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { } for required in [ + "pub async fn open_or_initialize", + "initialize_or_existing_database", + "InitializeDatabaseOutcome::Initialized", + "InitializeDatabaseOutcome::Existing", + "schema.matches_migrations(migrations)", + "open_existing_connection_pool_with_intent_and_authority", + "The existing branch never runs", + "the initializer. Callers therefore do not use pathname probes, error-text", + ] { + assert!( + CONNECTION_SOURCE.contains(required) + || OPEN_SOURCE.contains(required) + || README.contains(required), + "Step 243 atomic bootstrap boundary is missing `{required}`" + ); + } + for forbidden in ["try_exists()", "error.to_string()", "create_dir_all"] { + let connection_production = CONNECTION_SOURCE + .split_once("#[cfg(test)]") + .map_or(CONNECTION_SOURCE, |(production, _)| production); + assert!( + !connection_production.contains(forbidden), + "Step 243 atomic bootstrap uses forbidden branch surface `{forbidden}`" + ); + } + + for required in [ "fs2::FileExt", "rustix", "try_lock_exclusive",