lib

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

connection.rs (160648B)


      1 //! Narrow service-store host and transaction execution boundary.
      2 
      3 use core::fmt;
      4 use std::{
      5     error::Error,
      6     future::Future,
      7     path::Path,
      8     pin::Pin,
      9     sync::{
     10         Arc,
     11         atomic::{AtomicBool, Ordering},
     12     },
     13 };
     14 
     15 use futures::{future::BoxFuture, stream::BoxStream};
     16 use sqlx::{
     17     Either, Execute, Executor, SqlStr, Sqlite, SqliteConnection,
     18     sqlite::{SqliteQueryResult, SqliteRow, SqliteStatement, SqliteTypeInfo},
     19 };
     20 
     21 use crate::{
     22     ExistingServiceDatabaseIntent, MigrationApplicationOutcome, MigrationAppliedAtUnixSeconds,
     23     MigrationBuildIdentity, MigrationCallbackBinding, MigrationCatalog, OpenMode, SchemaCatalog,
     24     ServiceDatabaseIdentity, ServiceDatabaseMetadata, ServiceSqliteConnectionOptions,
     25     ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqliteIntegrityReport, ServiceSqlitePaths,
     26     WriterAuthority,
     27 };
     28 
     29 #[cfg(any(target_os = "linux", target_os = "macos"))]
     30 use sqlx::{Connection, pool::PoolConnection};
     31 
     32 /// One service-owned SQLite host whose raw pool remains inaccessible.
     33 ///
     34 /// The host intentionally has no raw-pool accessor:
     35 ///
     36 /// ```compile_fail
     37 /// use radroots_service_sqlite::ServiceSqliteHost;
     38 ///
     39 /// fn leak_pool(host: &ServiceSqliteHost) {
     40 ///     let _ = host.pool();
     41 /// }
     42 /// ```
     43 pub struct ServiceSqliteHost {
     44     mode: OpenMode,
     45     #[cfg(any(target_os = "linux", target_os = "macos"))]
     46     pool: crate::open::PrivateConnectionPool,
     47     #[cfg(any(target_os = "linux", target_os = "macos"))]
     48     closing: AtomicBool,
     49     #[cfg(any(target_os = "linux", target_os = "macos"))]
     50     close_state: tokio::sync::Mutex<ServiceSqliteHostCloseState>,
     51     #[cfg(any(target_os = "linux", target_os = "macos"))]
     52     backup_active: Arc<AtomicBool>,
     53     #[cfg(any(target_os = "linux", target_os = "macos"))]
     54     integrity_driver: tokio::sync::Mutex<IntegrityInspectionDriver>,
     55     #[cfg(any(target_os = "linux", target_os = "macos"))]
     56     failpoints: crate::failpoint::DurabilityFailpoints,
     57 }
     58 
     59 /// Database opened under retained authority with its verified metadata.
     60 ///
     61 /// This result cannot be assembled independently from a host and metadata:
     62 ///
     63 /// ```compile_fail
     64 /// use radroots_service_sqlite::{OpenedServiceDatabase, ServiceSqliteHost};
     65 ///
     66 /// fn forge(host: ServiceSqliteHost) {
     67 ///     let _ = OpenedServiceDatabase { host };
     68 /// }
     69 /// ```
     70 pub struct OpenedServiceDatabase {
     71     host: ServiceSqliteHost,
     72     metadata: ServiceDatabaseMetadata,
     73 }
     74 
     75 impl OpenedServiceDatabase {
     76     fn new(host: ServiceSqliteHost, metadata: ServiceDatabaseMetadata) -> Self {
     77         Self { host, metadata }
     78     }
     79 
     80     /// Borrows the authority-retaining host.
     81     #[must_use]
     82     pub const fn host(&self) -> &ServiceSqliteHost {
     83         &self.host
     84     }
     85 
     86     /// Borrows the metadata discovered and verified by the retained open.
     87     #[must_use]
     88     pub const fn database_metadata(&self) -> &ServiceDatabaseMetadata {
     89         &self.metadata
     90     }
     91 
     92     /// Consumes the binding into the authority-retaining host and actual metadata.
     93     #[must_use]
     94     pub fn into_parts(self) -> (ServiceSqliteHost, ServiceDatabaseMetadata) {
     95         (self.host, self.metadata)
     96     }
     97 }
     98 
     99 impl fmt::Debug for OpenedServiceDatabase {
    100     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    101         formatter
    102             .debug_struct("OpenedServiceDatabase")
    103             .field("mode", &self.host.mode())
    104             .field("database_metadata", &"[redacted]")
    105             .finish()
    106     }
    107 }
    108 
    109 #[cfg(any(target_os = "linux", target_os = "macos"))]
    110 enum ServiceSqliteHostCloseState {
    111     Pending,
    112     Complete(Option<ServiceSqliteErrorKind>),
    113 }
    114 
    115 #[cfg(any(target_os = "linux", target_os = "macos"))]
    116 enum IntegrityInspectionDriver {
    117     Idle,
    118     Connected(QuarantinedConnection),
    119     Closing(BoxFuture<'static, Result<(), sqlx::Error>>),
    120 }
    121 
    122 #[cfg(any(target_os = "linux", target_os = "macos"))]
    123 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    124 enum IntegrityInspectionDriverFailure {
    125     Invariant,
    126     ConnectionClose,
    127 }
    128 
    129 #[cfg(any(target_os = "linux", target_os = "macos"))]
    130 fn integrity_driver_close_result(
    131     result: Result<(), sqlx::Error>,
    132     injected_failure: bool,
    133 ) -> Result<(), IntegrityInspectionDriverFailure> {
    134     if injected_failure {
    135         Err(IntegrityInspectionDriverFailure::ConnectionClose)
    136     } else {
    137         result.map_err(|_| IntegrityInspectionDriverFailure::ConnectionClose)
    138     }
    139 }
    140 
    141 #[cfg(any(target_os = "linux", target_os = "macos"))]
    142 fn final_connection_policy_matches(
    143     initial: &crate::migration::MigrationConnectionPolicy,
    144     final_policy: &crate::migration::MigrationConnectionPolicy,
    145 ) -> Result<(), ServiceSqliteError> {
    146     (final_policy == initial)
    147         .then_some(())
    148         .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Pragma))
    149 }
    150 
    151 #[cfg(any(target_os = "linux", target_os = "macos"))]
    152 fn unconfirmed_rollback_error(
    153     rollback: Option<ServiceSqliteError>,
    154     rollback_was_confirmed: bool,
    155 ) -> Option<ServiceSqliteError> {
    156     rollback.filter(|_| !rollback_was_confirmed)
    157 }
    158 
    159 #[cfg(any(target_os = "linux", target_os = "macos"))]
    160 fn precondition_rollback_failure(
    161     authority: Option<ServiceSqliteError>,
    162     rollback: Option<ServiceSqliteError>,
    163     rollback_was_confirmed: bool,
    164     hook_removal: Option<ServiceSqliteError>,
    165 ) -> Option<ServiceSqliteError> {
    166     authority
    167         .or_else(|| unconfirmed_rollback_error(rollback, rollback_was_confirmed))
    168         .or(hook_removal)
    169 }
    170 
    171 #[cfg(any(target_os = "linux", target_os = "macos"))]
    172 fn authority_drift_rollback_failure(
    173     rollback: Option<ServiceSqliteError>,
    174     rollback_was_confirmed: bool,
    175     hook_removal: Option<ServiceSqliteError>,
    176 ) -> Option<ServiceSqliteError> {
    177     unconfirmed_rollback_error(rollback, rollback_was_confirmed).or(hook_removal)
    178 }
    179 
    180 #[cfg(any(target_os = "linux", target_os = "macos"))]
    181 fn operation_rollback_failure(
    182     rollback: Option<ServiceSqliteError>,
    183     rollback_was_confirmed: bool,
    184     hook_removal: Option<ServiceSqliteError>,
    185     authority: Option<ServiceSqliteError>,
    186 ) -> Option<ServiceSqliteError> {
    187     unconfirmed_rollback_error(rollback, rollback_was_confirmed)
    188         .or(hook_removal)
    189         .or(authority)
    190 }
    191 
    192 #[cfg(any(target_os = "linux", target_os = "macos"))]
    193 impl IntegrityInspectionDriver {
    194     async fn close_retained(&mut self) -> Result<(), IntegrityInspectionDriverFailure> {
    195         loop {
    196             match self {
    197                 Self::Idle => return Ok(()),
    198                 Self::Connected(_) => {
    199                     let Self::Connected(connection) = core::mem::replace(self, Self::Idle) else {
    200                         return Err(IntegrityInspectionDriverFailure::Invariant);
    201                     };
    202                     *self = Self::Closing(
    203                         connection
    204                             .into_close_future()
    205                             .ok_or(IntegrityInspectionDriverFailure::Invariant)?,
    206                     );
    207                 }
    208                 Self::Closing(close) => {
    209                     #[cfg(test)]
    210                     crate::integrity::integrity_test_seam::pause(
    211                         crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING,
    212                     )
    213                     .await;
    214                     let result = close.await;
    215                     *self = Self::Idle;
    216                     #[cfg(test)]
    217                     let injected_failure =
    218                         crate::integrity::integrity_test_seam::take_connection_close_failure();
    219                     #[cfg(not(test))]
    220                     let injected_failure = false;
    221                     return integrity_driver_close_result(result, injected_failure);
    222                 }
    223             }
    224         }
    225     }
    226 
    227     fn connection_mut(&mut self) -> Result<&mut SqliteConnection, ServiceSqliteError> {
    228         match self {
    229             Self::Connected(connection) => Ok(connection),
    230             Self::Idle | Self::Closing(_) => {
    231                 Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))
    232             }
    233         }
    234     }
    235 
    236     fn return_to_pool(&mut self) -> Result<(), ServiceSqliteError> {
    237         let Self::Connected(mut connection) = core::mem::replace(self, Self::Idle) else {
    238             return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
    239         };
    240         connection.trust();
    241         drop(connection);
    242         Ok(())
    243     }
    244 
    245     #[cfg(test)]
    246     const fn is_idle(&self) -> bool {
    247         matches!(self, Self::Idle)
    248     }
    249 }
    250 
    251 impl ServiceSqliteHost {
    252     /// Opens existing writable state and finishes every pending governed migration.
    253     ///
    254     /// Before opening SQLite, this path holds exclusive writer authority and
    255     /// synchronously reconciles any exact interrupted-restore topology. The
    256     /// recovery sequence has no await point: cancellation cannot split one
    257     /// filesystem step from its authority check. If the surrounding open is
    258     /// cancelled later, a retry re-reads the already durable filesystem state.
    259     #[allow(clippy::too_many_arguments)]
    260     pub async fn open_read_write_existing(
    261         paths: &ServiceSqlitePaths,
    262         identity: &ServiceDatabaseIdentity,
    263         migrations: &MigrationCatalog,
    264         schema: &SchemaCatalog,
    265         options: ServiceSqliteConnectionOptions,
    266         applied_at: MigrationAppliedAtUnixSeconds,
    267         build: &MigrationBuildIdentity,
    268         callbacks: &[MigrationCallbackBinding],
    269     ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> {
    270         #[cfg(any(target_os = "linux", target_os = "macos"))]
    271         {
    272             let pool = crate::open::open_existing_connection_pool(
    273                 paths,
    274                 identity,
    275                 migrations,
    276                 schema,
    277                 OpenMode::ReadWriteExisting,
    278                 options,
    279             )
    280             .await?;
    281             match pool.apply_migrations(applied_at, build, callbacks).await {
    282                 Ok(outcome) => Ok((Self::from_pool(OpenMode::ReadWriteExisting, pool), outcome)),
    283                 Err(error) => {
    284                     drop(pool.close().await);
    285                     Err(error)
    286                 }
    287             }
    288         }
    289         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    290         {
    291             let _ = (
    292                 paths, identity, migrations, schema, options, applied_at, build, callbacks,
    293             );
    294             Err(unsupported_host())
    295         }
    296     }
    297 
    298     /// Opens existing writable state and discovers its stored generation under authority.
    299     ///
    300     /// The intent binds service, instance, application ID, and the supported
    301     /// schema ceiling before filesystem or SQLite admission. The returned
    302     /// metadata is read from the same authority-retaining host after governed
    303     /// migrations finish, so callers never guess a source generation or reopen
    304     /// the database between discovery and use.
    305     #[allow(clippy::too_many_arguments)]
    306     pub async fn open_read_write_existing_with_intent(
    307         paths: &ServiceSqlitePaths,
    308         intent: &ExistingServiceDatabaseIntent,
    309         migrations: &MigrationCatalog,
    310         schema: &SchemaCatalog,
    311         options: ServiceSqliteConnectionOptions,
    312         applied_at: MigrationAppliedAtUnixSeconds,
    313         build: &MigrationBuildIdentity,
    314         callbacks: &[MigrationCallbackBinding],
    315     ) -> Result<(OpenedServiceDatabase, MigrationApplicationOutcome), ServiceSqliteError> {
    316         #[cfg(any(target_os = "linux", target_os = "macos"))]
    317         {
    318             let pool = crate::open::open_existing_connection_pool_with_intent(
    319                 paths,
    320                 intent,
    321                 migrations,
    322                 schema,
    323                 OpenMode::ReadWriteExisting,
    324                 options,
    325             )
    326             .await?;
    327             let outcome = match pool.apply_migrations(applied_at, build, callbacks).await {
    328                 Ok(outcome) => outcome,
    329                 Err(error) => {
    330                     drop(pool.close().await);
    331                     return Err(error);
    332                 }
    333             };
    334             let metadata = match pool.database_metadata().await {
    335                 Ok(metadata) => metadata,
    336                 Err(error) => {
    337                     drop(pool.close().await);
    338                     return Err(error);
    339                 }
    340             };
    341             let host = Self::from_pool(OpenMode::ReadWriteExisting, pool);
    342             Ok((OpenedServiceDatabase::new(host, metadata), outcome))
    343         }
    344         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    345         {
    346             let _ = (
    347                 paths, intent, migrations, schema, options, applied_at, build, callbacks,
    348             );
    349             Err(unsupported_host())
    350         }
    351     }
    352 
    353     /// Atomically creates interactive state or opens the exact existing database.
    354     ///
    355     /// The runner holds writer authority while an exclusive create decides the
    356     /// branch. Callers never probe the filesystem or inspect error text. On the
    357     /// create branch, the sealed initializer, shared metadata, empty v1 ledger,
    358     /// and schema-catalog verification commit together before the host opens and
    359     /// applies governed migrations. On exact create collision, the same retained
    360     /// authority is transferred to an existing-only open and the initialization
    361     /// callback is never invoked. Success binds the retained host to the actual
    362     /// verified metadata selected by that atomic decision.
    363     #[allow(clippy::too_many_arguments)]
    364     pub async fn open_or_initialize<F, E>(
    365         paths: &ServiceSqlitePaths,
    366         initialization_metadata: &ServiceDatabaseMetadata,
    367         migrations: &MigrationCatalog,
    368         schema: &SchemaCatalog,
    369         options: ServiceSqliteConnectionOptions,
    370         applied_at: MigrationAppliedAtUnixSeconds,
    371         build: &MigrationBuildIdentity,
    372         callbacks: &[MigrationCallbackBinding],
    373         initialize_schema: F,
    374     ) -> Result<(OpenedServiceDatabase, MigrationApplicationOutcome), ServiceSqliteError>
    375     where
    376         F: for<'a> FnOnce(
    377             &'a mut crate::ServiceSqliteInitializer<'_>,
    378         ) -> crate::ServiceSqliteInitializerFuture<'a, E>,
    379         E: Error + Send + Sync + 'static,
    380     {
    381         #[cfg(any(target_os = "linux", target_os = "macos"))]
    382         {
    383             if !schema.matches_migrations(migrations) {
    384                 return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
    385             }
    386             let supported_version = core::num::NonZeroU32::new(migrations.current_version())
    387                 .expect("migration catalogs always have a nonzero current version");
    388             let initialized = crate::initialize::initialize_or_existing_database(
    389                 paths,
    390                 initialization_metadata,
    391                 schema,
    392                 initialize_schema,
    393             )
    394             .await?;
    395             let (mode, pool) = match initialized {
    396                 crate::initialize::InitializeDatabaseOutcome::Initialized(authority) => {
    397                     let identity = ServiceDatabaseIdentity::new(
    398                         paths,
    399                         initialization_metadata.source_generation(),
    400                         supported_version,
    401                         initialization_metadata.application_id(),
    402                     );
    403                     let pool = crate::open::open_initialized_connection_pool(
    404                         paths, &identity, migrations, schema, options, authority,
    405                     )
    406                     .await?;
    407                     (OpenMode::Initialize, pool)
    408                 }
    409                 crate::initialize::InitializeDatabaseOutcome::Existing(authority) => {
    410                     let intent = ExistingServiceDatabaseIntent::new(
    411                         paths,
    412                         supported_version,
    413                         initialization_metadata.application_id(),
    414                     );
    415                     let pool =
    416                         crate::open::open_existing_connection_pool_with_intent_and_authority(
    417                             paths, &intent, migrations, schema, options, authority,
    418                         )
    419                         .await?;
    420                     (OpenMode::ReadWriteExisting, pool)
    421                 }
    422             };
    423             let outcome = match pool.apply_migrations(applied_at, build, callbacks).await {
    424                 Ok(outcome) => outcome,
    425                 Err(error) => {
    426                     drop(pool.close().await);
    427                     return Err(error);
    428                 }
    429             };
    430             let metadata = match pool.database_metadata().await {
    431                 Ok(metadata) => metadata,
    432                 Err(error) => {
    433                     drop(pool.close().await);
    434                     return Err(error);
    435                 }
    436             };
    437             let host = Self::from_pool(mode, pool);
    438             Ok((OpenedServiceDatabase::new(host, metadata), outcome))
    439         }
    440         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    441         {
    442             drop((
    443                 paths,
    444                 initialization_metadata,
    445                 migrations,
    446                 schema,
    447                 options,
    448                 applied_at,
    449                 build,
    450                 callbacks,
    451                 initialize_schema,
    452             ));
    453             Err(unsupported_host())
    454         }
    455     }
    456 
    457     /// Opens state created under a retained initialization writer authority.
    458     #[allow(clippy::too_many_arguments)]
    459     pub async fn open_initialized(
    460         paths: &ServiceSqlitePaths,
    461         identity: &ServiceDatabaseIdentity,
    462         migrations: &MigrationCatalog,
    463         schema: &SchemaCatalog,
    464         options: ServiceSqliteConnectionOptions,
    465         authority: WriterAuthority,
    466         applied_at: MigrationAppliedAtUnixSeconds,
    467         build: &MigrationBuildIdentity,
    468         callbacks: &[MigrationCallbackBinding],
    469     ) -> Result<(Self, MigrationApplicationOutcome), ServiceSqliteError> {
    470         #[cfg(any(target_os = "linux", target_os = "macos"))]
    471         {
    472             let pool = crate::open::open_initialized_connection_pool(
    473                 paths, identity, migrations, schema, options, authority,
    474             )
    475             .await?;
    476             match pool.apply_migrations(applied_at, build, callbacks).await {
    477                 Ok(outcome) => Ok((Self::from_pool(OpenMode::Initialize, pool), outcome)),
    478                 Err(error) => {
    479                     drop(pool.close().await);
    480                     Err(error)
    481                 }
    482             }
    483         }
    484         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    485         {
    486             let _ = (
    487                 paths, identity, migrations, schema, options, authority, applied_at, build,
    488                 callbacks,
    489             );
    490             Err(unsupported_host())
    491         }
    492     }
    493 
    494     /// Opens an immutable, current-schema inspection host without writer authority.
    495     pub async fn open_read_only_inspection(
    496         paths: &ServiceSqlitePaths,
    497         identity: &ServiceDatabaseIdentity,
    498         migrations: &MigrationCatalog,
    499         schema: &SchemaCatalog,
    500         options: ServiceSqliteConnectionOptions,
    501     ) -> Result<Self, ServiceSqliteError> {
    502         #[cfg(any(target_os = "linux", target_os = "macos"))]
    503         {
    504             let pool = crate::open::open_existing_connection_pool(
    505                 paths,
    506                 identity,
    507                 migrations,
    508                 schema,
    509                 OpenMode::ReadOnlyInspection,
    510                 options,
    511             )
    512             .await?;
    513             Ok(Self::from_pool(OpenMode::ReadOnlyInspection, pool))
    514         }
    515         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    516         {
    517             let _ = (paths, identity, migrations, schema, options);
    518             Err(unsupported_host())
    519         }
    520     }
    521 
    522     /// Opens existing state for immutable inspection and discovers its metadata.
    523     pub async fn open_read_only_inspection_with_intent(
    524         paths: &ServiceSqlitePaths,
    525         intent: &ExistingServiceDatabaseIntent,
    526         migrations: &MigrationCatalog,
    527         schema: &SchemaCatalog,
    528         options: ServiceSqliteConnectionOptions,
    529     ) -> Result<OpenedServiceDatabase, ServiceSqliteError> {
    530         #[cfg(any(target_os = "linux", target_os = "macos"))]
    531         {
    532             let pool = crate::open::open_existing_connection_pool_with_intent(
    533                 paths,
    534                 intent,
    535                 migrations,
    536                 schema,
    537                 OpenMode::ReadOnlyInspection,
    538                 options,
    539             )
    540             .await?;
    541             let metadata = match pool.database_metadata().await {
    542                 Ok(metadata) => metadata,
    543                 Err(error) => {
    544                     drop(pool.close().await);
    545                     return Err(error);
    546                 }
    547             };
    548             let host = Self::from_pool(OpenMode::ReadOnlyInspection, pool);
    549             Ok(OpenedServiceDatabase::new(host, metadata))
    550         }
    551         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    552         {
    553             let _ = (paths, intent, migrations, schema, options);
    554             Err(unsupported_host())
    555         }
    556     }
    557 
    558     /// Returns the fixed mode selected when the host was opened.
    559     #[must_use]
    560     pub const fn mode(&self) -> OpenMode {
    561         self.mode
    562     }
    563 
    564     /// Closes all connections and explicitly releases retained instance authority.
    565     ///
    566     /// Close rejects new transactions as soon as it starts and waits for already
    567     /// admitted transactions to finish. Writable hosts then perform the fixed
    568     /// governed `TRUNCATE` WAL checkpoint before releasing writer authority;
    569     /// read-only inspection performs no checkpoint or filesystem mutation.
    570     ///
    571     /// Cancelling this future leaves the host permanently non-admitting and retains
    572     /// authority until a later call resumes close. A completed result is cached, so
    573     /// sequential or concurrent later calls return the same stable outer outcome.
    574     pub async fn close(&self) -> Result<(), ServiceSqliteError> {
    575         #[cfg(any(target_os = "linux", target_os = "macos"))]
    576         {
    577             self.closing.store(true, Ordering::Release);
    578             let mut state = self.close_state.lock().await;
    579             if let ServiceSqliteHostCloseState::Complete(kind) = *state {
    580                 return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind)));
    581             }
    582             let (integrity_cleanup, integrity_validation) = {
    583                 let mut driver = self.integrity_driver.lock().await;
    584                 let cleanup = driver
    585                     .close_retained()
    586                     .await
    587                     .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
    588                 let validation = self.pool.validate();
    589                 (cleanup, validation)
    590             };
    591             let close = self.pool.close_explicit(&self.failpoints).await;
    592             match close {
    593                 Err(retryable) => Err(retryable),
    594                 Ok(terminal) => {
    595                     let terminal = integrity_validation.and(terminal).and(integrity_cleanup);
    596                     *state = ServiceSqliteHostCloseState::Complete(
    597                         terminal.as_ref().err().map(ServiceSqliteError::kind),
    598                     );
    599                     terminal
    600                 }
    601             }
    602         }
    603         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    604         {
    605             Err(unsupported_host())
    606         }
    607     }
    608 
    609     /// Captures one point-in-time SQLite backup into a new staging directory.
    610     ///
    611     /// The caller supplies the exact new absolute staging-directory path and an
    612     /// injected creation time. Capture is available only on writable hosts and
    613     /// admits at most one active capture per host. The canonical manifest and a
    614     /// successful result are returned only after the visible staging directory's
    615     /// sole `state.sqlite` member has passed metadata, integrity, digest, and
    616     /// durability checks. The manifest remains in memory and is not written into
    617     /// the staging directory.
    618     ///
    619     /// Dropping this future requests cancellation. The admitted worker retains
    620     /// host authority until it has closed SQLite handles and either completed or
    621     /// cleaned the exact staging artifacts, so `close` drains that work before
    622     /// releasing writer authority.
    623     pub async fn capture_online_backup(
    624         &self,
    625         staging_directory: &Path,
    626         created_at_unix_ms: crate::BackupCreatedAtUnixMs,
    627     ) -> Result<crate::ServiceBackupManifest, ServiceSqliteError> {
    628         #[cfg(any(target_os = "linux", target_os = "macos"))]
    629         {
    630             crate::backup::capture_online_backup(
    631                 &self.pool,
    632                 &self.closing,
    633                 &self.backup_active,
    634                 staging_directory,
    635                 created_at_unix_ms,
    636                 &self.failpoints,
    637             )
    638             .await
    639         }
    640         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    641         {
    642             let _ = (staging_directory, created_at_unix_ms);
    643             Err(unsupported_host())
    644         }
    645     }
    646 
    647     /// Runs one explicit bounded integrity inspection over a single read snapshot.
    648     ///
    649     /// The caller injects the wall-clock completion time and owns any monotonic
    650     /// deadline. The host admits at most one inspection at a time. Dropping this
    651     /// future before it returns publishes no report, persists no status, and
    652     /// leaves the checked-out connection in a host-owned explicit-close driver.
    653     /// Retry or host close finishes that close before another check or authority
    654     /// release; a retry must inject a new time.
    655     /// Completed SQLite and foreign-key failures are returned only as fixed safe
    656     /// diagnostic codes. An inability to execute or decode either check is an
    657     /// `Integrity` error.
    658     pub async fn inspect_integrity(
    659         &self,
    660         checked_at: crate::IntegrityCheckedAtUnixMs,
    661     ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
    662         #[cfg(any(target_os = "linux", target_os = "macos"))]
    663         {
    664             self.inspect_integrity_supported(checked_at).await
    665         }
    666         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    667         {
    668             let _ = checked_at;
    669             Err(unsupported_host())
    670         }
    671     }
    672 
    673     /// Executes one runner-owned transaction without exposing its connection or pool.
    674     ///
    675     /// Dropping this future before the runner enables its outer commit quarantines
    676     /// the connection and leaves no authoritative transaction effect. An operation
    677     /// error is returned as `OperationRolledBack` only after rollback is confirmed;
    678     /// an unconfirmed rollback is `RollbackFailed`. Once outer commit begins,
    679     /// cancelling the future yields no result and must be treated as an unknown
    680     /// commit outcome. Callers receiving `CommitOutcomeUnknown`, or cancelling after
    681     /// commit begins, must reread authoritative state before any idempotent retry.
    682     pub async fn transaction<T, E, F>(
    683         &self,
    684         operation: F,
    685     ) -> Result<T, ServiceSqliteTransactionError<E>>
    686     where
    687         T: Send + 'static,
    688         E: Send + 'static,
    689         F: for<'a> FnOnce(
    690                 &'a mut ServiceSqliteTransaction<'_>,
    691             ) -> ServiceSqliteTransactionFuture<'a, T, E>
    692             + Send,
    693     {
    694         #[cfg(any(target_os = "linux", target_os = "macos"))]
    695         {
    696             self.transaction_supported(operation).await
    697         }
    698         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    699         {
    700             drop(operation);
    701             Err(ServiceSqliteTransactionError::not_committed(
    702                 unsupported_host(),
    703             ))
    704         }
    705     }
    706 
    707     #[cfg(any(target_os = "linux", target_os = "macos"))]
    708     async fn transaction_supported<T, E, F>(
    709         &self,
    710         operation: F,
    711     ) -> Result<T, ServiceSqliteTransactionError<E>>
    712     where
    713         T: Send + 'static,
    714         E: Send + 'static,
    715         F: for<'a> FnOnce(
    716                 &'a mut ServiceSqliteTransaction<'_>,
    717             ) -> ServiceSqliteTransactionFuture<'a, T, E>
    718             + Send,
    719     {
    720         crate::require_condition(
    721             !self.closing.load(Ordering::Acquire),
    722             ServiceSqliteErrorKind::Open,
    723         )
    724         .map_err(ServiceSqliteTransactionError::not_committed)?;
    725         self.pool
    726             .validate()
    727             .map_err(ServiceSqliteTransactionError::not_committed)?;
    728         let connection = self
    729             .pool
    730             .acquire()
    731             .await
    732             .map_err(ServiceSqliteTransactionError::not_committed)?;
    733         let mut connection = QuarantinedConnection::new(connection);
    734         self.pool
    735             .validate()
    736             .map_err(ServiceSqliteTransactionError::not_committed)?;
    737         let initial_policy = crate::migration::read_connection_policy(&mut connection)
    738             .await
    739             .map_err(ServiceSqliteTransactionError::not_committed)?;
    740         self.pool
    741             .validate()
    742             .map_err(ServiceSqliteTransactionError::not_committed)?;
    743         let gate = crate::transaction_control::TransactionControlGate::install(&mut connection)
    744             .await
    745             .map_err(|source| {
    746                 ServiceSqliteTransactionError::not_committed(sqlite_source(source))
    747             })?;
    748         self.pool
    749             .validate()
    750             .map_err(ServiceSqliteTransactionError::not_committed)?;
    751         let before_begin = self
    752             .failpoints
    753             .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeBegin)
    754             .map_err(|source| {
    755                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    756             });
    757         self.pool
    758             .validate()
    759             .map_err(ServiceSqliteTransactionError::not_committed)?;
    760         before_begin.map_err(ServiceSqliteTransactionError::not_committed)?;
    761         let mut transaction = match match self.pool.mode() {
    762             OpenMode::Initialize | OpenMode::ReadWriteExisting => {
    763                 connection.begin_with("BEGIN IMMEDIATE").await
    764             }
    765             OpenMode::ReadOnlyInspection => connection.begin().await,
    766         } {
    767             Ok(transaction) => transaction,
    768             Err(source) => {
    769                 return Err(ServiceSqliteTransactionError::not_committed(sqlite_source(
    770                     source,
    771                 )));
    772             }
    773         };
    774         let injected_after_begin = self
    775             .failpoints
    776             .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterBegin)
    777             .map_err(|source| {
    778                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    779             });
    780         let after_begin = self.pool.validate().and(injected_after_begin);
    781         if let Err(error) = after_begin {
    782             let permit = gate.permit_runner_rollback();
    783             let rollback = transaction.rollback().await.map_err(sqlite_source);
    784             drop(permit);
    785             let rollback_was_confirmed =
    786                 gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
    787             let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
    788             let authority = self.pool.validate();
    789             if let Some(rollback_error) = precondition_rollback_failure(
    790                 authority.err(),
    791                 rollback.err(),
    792                 rollback_was_confirmed,
    793                 remove.err(),
    794             ) {
    795                 return Err(ServiceSqliteTransactionError::rollback_failed(
    796                     None,
    797                     rollback_error,
    798                 ));
    799             }
    800             return Err(ServiceSqliteTransactionError::not_committed(error));
    801         }
    802 
    803         let operation_result = {
    804             let statement_control_rejected = Arc::new(AtomicBool::new(false));
    805             let mut executor = ServiceSqliteTransaction {
    806                 connection: &mut transaction,
    807                 statement_control_rejected: Arc::clone(&statement_control_rejected),
    808             };
    809             (operation(&mut executor).await, statement_control_rejected)
    810         };
    811         let (operation_result, statement_control_rejected) = operation_result;
    812         if let Err(error) = self.pool.validate() {
    813             let operation_error = operation_result.err();
    814             let permit = gate.permit_runner_rollback();
    815             let rollback = transaction.rollback().await.map_err(sqlite_source);
    816             drop(permit);
    817             let rollback_was_confirmed =
    818                 gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
    819             let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
    820             let rollback_error = authority_drift_rollback_failure(
    821                 rollback.err(),
    822                 rollback_was_confirmed,
    823                 remove.err(),
    824             );
    825             return Err(match rollback_error {
    826                 Some(rollback_error) => {
    827                     ServiceSqliteTransactionError::rollback_failed(operation_error, rollback_error)
    828                 }
    829                 None => ServiceSqliteTransactionError::not_committed_with_operation(
    830                     operation_error,
    831                     error,
    832                 ),
    833             });
    834         }
    835         let value = match operation_result {
    836             Ok(value) => value,
    837             Err(operation_error) => {
    838                 let permit = gate.permit_runner_rollback();
    839                 let rollback = transaction.rollback().await.map_err(sqlite_source);
    840                 drop(permit);
    841                 let rollback_was_confirmed =
    842                     gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
    843                 let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
    844                 let authority = self.pool.validate();
    845                 if let Some(error) = operation_rollback_failure(
    846                     rollback.err(),
    847                     rollback_was_confirmed,
    848                     remove.err(),
    849                     authority.err(),
    850                 ) {
    851                     return Err(ServiceSqliteTransactionError::rollback_failed(
    852                         Some(operation_error),
    853                         error,
    854                     ));
    855                 }
    856                 return Err(ServiceSqliteTransactionError::operation_rolled_back(
    857                     operation_error,
    858                 ));
    859             }
    860         };
    861 
    862         let precommit = self
    863             .verify_before_commit(
    864                 &mut transaction,
    865                 &gate,
    866                 &initial_policy,
    867                 &statement_control_rejected,
    868             )
    869             .await;
    870         let precommit = match precommit {
    871             Ok(()) => {
    872                 let injected = self
    873                     .failpoints
    874                     .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeCommit)
    875                     .map_err(|source| {
    876                         ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    877                     });
    878                 self.pool.validate().and(injected)
    879             }
    880             Err(error) => Err(error),
    881         };
    882         if let Err(error) = precommit {
    883             let permit = gate.permit_runner_rollback();
    884             let rollback = transaction.rollback().await.map_err(sqlite_source);
    885             drop(permit);
    886             let rollback_was_confirmed =
    887                 gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
    888             let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
    889             let authority = self.pool.validate();
    890             if let Some(rollback_error) = precondition_rollback_failure(
    891                 authority.err(),
    892                 rollback.err(),
    893                 rollback_was_confirmed,
    894                 remove.err(),
    895             ) {
    896                 return Err(ServiceSqliteTransactionError::rollback_failed(
    897                     None,
    898                     rollback_error,
    899                 ));
    900             }
    901             return Err(ServiceSqliteTransactionError::not_committed(error));
    902         }
    903 
    904         let permit = gate.permit_outer_commit();
    905         let commit = transaction.commit().await.map_err(sqlite_source);
    906         drop(permit);
    907         let injected_after_commit = if commit.is_ok() {
    908             self.failpoints
    909                 .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterCommit)
    910                 .map_err(|source| {
    911                     ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    912                 })
    913         } else {
    914             Ok(())
    915         };
    916         let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
    917         self.pool
    918             .validate()
    919             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    920         commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    921         remove.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    922         injected_after_commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    923         let final_policy = crate::migration::read_connection_policy(&mut connection)
    924             .await
    925             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    926         self.pool
    927             .validate()
    928             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    929         final_connection_policy_matches(&initial_policy, &final_policy)
    930             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    931         crate::metadata::verify_database_metadata(&mut connection, self.pool.identity())
    932             .await
    933             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    934         self.pool
    935             .validate()
    936             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    937         crate::migration::verify_migration_history(
    938             &mut connection,
    939             self.pool.catalog(),
    940             self.pool.schema_catalog(),
    941             true,
    942         )
    943         .await
    944         .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    945         self.pool
    946             .validate()
    947             .map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
    948         connection.trust();
    949         Ok(value)
    950     }
    951 
    952     #[cfg(any(target_os = "linux", target_os = "macos"))]
    953     async fn inspect_integrity_supported(
    954         &self,
    955         checked_at: crate::IntegrityCheckedAtUnixMs,
    956     ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
    957         crate::require_condition(
    958             !self.closing.load(Ordering::Acquire),
    959             ServiceSqliteErrorKind::Open,
    960         )?;
    961         let mut driver = self
    962             .integrity_driver
    963             .try_lock()
    964             .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
    965         let cleanup = driver.close_retained().await;
    966         self.pool.validate()?;
    967         cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
    968         crate::require_condition(
    969             !self.closing.load(Ordering::Acquire),
    970             ServiceSqliteErrorKind::Open,
    971         )?;
    972         self.pool.validate()?;
    973         let connection = self.pool.acquire().await;
    974         self.pool.validate()?;
    975         *driver = IntegrityInspectionDriver::Connected(QuarantinedConnection::new(connection?));
    976         crate::require_condition(
    977             !self.closing.load(Ordering::Acquire),
    978             ServiceSqliteErrorKind::Open,
    979         )?;
    980         let report = crate::integrity::inspect_database_integrity(
    981             driver.connection_mut()?,
    982             checked_at,
    983             || self.pool.validate(),
    984         )
    985         .await;
    986         let validation = self.pool.validate();
    987         match (validation, report) {
    988             (Err(error), _) => {
    989                 let cleanup = driver.close_retained().await;
    990                 let validation = self.pool.validate();
    991                 if error.kind() == ServiceSqliteErrorKind::Authority {
    992                     return Err(error);
    993                 }
    994                 validation?;
    995                 cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
    996                 Err(error)
    997             }
    998             (Ok(()), Err(error)) => {
    999                 let cleanup = driver.close_retained().await;
   1000                 self.pool.validate()?;
   1001                 cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
   1002                 Err(error)
   1003             }
   1004             (Ok(()), Ok(report)) => {
   1005                 driver.return_to_pool()?;
   1006                 Ok(report)
   1007             }
   1008         }
   1009     }
   1010 
   1011     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1012     fn from_pool(mode: OpenMode, pool: crate::open::PrivateConnectionPool) -> Self {
   1013         Self {
   1014             mode,
   1015             pool,
   1016             closing: AtomicBool::new(false),
   1017             close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending),
   1018             backup_active: Arc::new(AtomicBool::new(false)),
   1019             integrity_driver: tokio::sync::Mutex::new(IntegrityInspectionDriver::Idle),
   1020             failpoints: crate::failpoint::DurabilityFailpoints::default(),
   1021         }
   1022     }
   1023 
   1024     #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
   1025     fn arm_durability_failpoint(&self, point: crate::failpoint::DurabilityFailpoint) {
   1026         self.failpoints.arm(point);
   1027     }
   1028 
   1029     fn lifecycle_state(&self) -> &'static str {
   1030         #[cfg(any(target_os = "linux", target_os = "macos"))]
   1031         {
   1032             if self.closing.load(Ordering::Acquire) {
   1033                 "closing_or_closed"
   1034             } else {
   1035                 "open"
   1036             }
   1037         }
   1038         #[cfg(not(any(target_os = "linux", target_os = "macos")))]
   1039         {
   1040             "closing_or_closed"
   1041         }
   1042     }
   1043 
   1044     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1045     async fn verify_before_commit(
   1046         &self,
   1047         connection: &mut SqliteConnection,
   1048         gate: &crate::transaction_control::TransactionControlGate,
   1049         initial_policy: &crate::migration::MigrationConnectionPolicy,
   1050         statement_control_rejected: &AtomicBool,
   1051     ) -> Result<(), ServiceSqliteError> {
   1052         self.pool.validate()?;
   1053         crate::require_condition(
   1054             !gate.control_violation_observed()
   1055                 && !statement_control_rejected.load(Ordering::Acquire),
   1056             ServiceSqliteErrorKind::Open,
   1057         )?;
   1058         crate::migration::assert_governed_transaction(connection).await?;
   1059         self.pool.validate()?;
   1060         crate::require_condition(
   1061             &crate::migration::read_connection_policy(connection).await? == initial_policy,
   1062             ServiceSqliteErrorKind::Pragma,
   1063         )?;
   1064         self.pool.validate()?;
   1065         crate::metadata::verify_database_metadata(connection, self.pool.identity()).await?;
   1066         self.pool.validate()?;
   1067         crate::migration::verify_migration_history_snapshot(
   1068             connection,
   1069             self.pool.catalog(),
   1070             self.pool.schema_catalog(),
   1071             true,
   1072         )
   1073         .await?;
   1074         self.pool.validate()?;
   1075         crate::migration::assert_governed_transaction(connection).await?;
   1076         crate::require_condition(
   1077             !gate.control_violation_observed(),
   1078             ServiceSqliteErrorKind::Open,
   1079         )?;
   1080         Ok(())
   1081     }
   1082 }
   1083 
   1084 impl fmt::Debug for ServiceSqliteHost {
   1085     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1086         formatter
   1087             .debug_struct("ServiceSqliteHost")
   1088             .field("mode", &self.mode)
   1089             .field("pool", &"[redacted]")
   1090             .field("lifecycle", &self.lifecycle_state())
   1091             .finish()
   1092     }
   1093 }
   1094 
   1095 /// A sealed transaction executor that never exposes its raw SQLite connection.
   1096 ///
   1097 /// Service repositories may use ordinary typed SQLx queries through the
   1098 /// borrowed executor:
   1099 ///
   1100 /// ```
   1101 /// use radroots_service_sqlite::ServiceSqliteTransaction;
   1102 ///
   1103 /// async fn row_count(
   1104 ///     transaction: &mut ServiceSqliteTransaction<'_>,
   1105 /// ) -> Result<i64, sqlx::Error> {
   1106 ///     sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM service_items")
   1107 ///         .fetch_one(transaction)
   1108 ///         .await
   1109 /// }
   1110 /// ```
   1111 ///
   1112 /// Transaction control remains runner-owned:
   1113 ///
   1114 /// ```compile_fail
   1115 /// use radroots_service_sqlite::ServiceSqliteTransaction;
   1116 ///
   1117 /// async fn bypass(transaction: ServiceSqliteTransaction<'_>) {
   1118 ///     transaction.commit().await.unwrap();
   1119 /// }
   1120 /// ```
   1121 pub struct ServiceSqliteTransaction<'connection> {
   1122     connection: &'connection mut SqliteConnection,
   1123     statement_control_rejected: Arc<AtomicBool>,
   1124 }
   1125 
   1126 struct RestrictedExecute<Q> {
   1127     query: Q,
   1128     statement_control_rejected: Arc<AtomicBool>,
   1129 }
   1130 
   1131 impl<'query, Q> Execute<'query, Sqlite> for RestrictedExecute<Q>
   1132 where
   1133     Q: Execute<'query, Sqlite>,
   1134 {
   1135     fn sql(self) -> SqlStr {
   1136         restricted_sql(self.query.sql(), &self.statement_control_rejected)
   1137     }
   1138 
   1139     fn statement(&self) -> Option<&SqliteStatement> {
   1140         None
   1141     }
   1142 
   1143     fn take_arguments(
   1144         &mut self,
   1145     ) -> Result<Option<<Sqlite as sqlx::Database>::Arguments>, sqlx::error::BoxDynError> {
   1146         self.query.take_arguments()
   1147     }
   1148 
   1149     fn persistent(&self) -> bool {
   1150         self.query.persistent()
   1151     }
   1152 }
   1153 
   1154 fn restricted_sql(sql: SqlStr, statement_control_rejected: &AtomicBool) -> SqlStr {
   1155     if crate::statement_policy::contains_forbidden_statement_control(sql.as_str()) {
   1156         statement_control_rejected.store(true, Ordering::Release);
   1157         SqlStr::from_static("RADROOTS_FORBIDDEN_STATEMENT_CONTROL")
   1158     } else {
   1159         sql
   1160     }
   1161 }
   1162 
   1163 impl fmt::Debug for ServiceSqliteTransaction<'_> {
   1164     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1165         formatter.write_str("ServiceSqliteTransaction([redacted])")
   1166     }
   1167 }
   1168 
   1169 impl<'executor, 'connection> Executor<'executor>
   1170     for &'executor mut ServiceSqliteTransaction<'connection>
   1171 where
   1172     'connection: 'executor,
   1173 {
   1174     type Database = Sqlite;
   1175 
   1176     fn fetch_many<'e, 'q: 'e, Q>(
   1177         self,
   1178         query: Q,
   1179     ) -> BoxStream<'e, Result<Either<SqliteQueryResult, SqliteRow>, sqlx::Error>>
   1180     where
   1181         'executor: 'e,
   1182         Q: 'q + Execute<'q, Self::Database>,
   1183     {
   1184         (&mut *self.connection).fetch_many(RestrictedExecute {
   1185             query,
   1186             statement_control_rejected: Arc::clone(&self.statement_control_rejected),
   1187         })
   1188     }
   1189 
   1190     fn fetch_optional<'e, 'q: 'e, Q>(
   1191         self,
   1192         query: Q,
   1193     ) -> BoxFuture<'e, Result<Option<SqliteRow>, sqlx::Error>>
   1194     where
   1195         'executor: 'e,
   1196         Q: 'q + Execute<'q, Self::Database>,
   1197     {
   1198         (&mut *self.connection).fetch_optional(RestrictedExecute {
   1199             query,
   1200             statement_control_rejected: Arc::clone(&self.statement_control_rejected),
   1201         })
   1202     }
   1203 
   1204     fn prepare_with<'e>(
   1205         self,
   1206         sql: SqlStr,
   1207         parameters: &'e [SqliteTypeInfo],
   1208     ) -> BoxFuture<'e, Result<SqliteStatement, sqlx::Error>>
   1209     where
   1210         'executor: 'e,
   1211     {
   1212         (&mut *self.connection).prepare_with(
   1213             restricted_sql(sql, &self.statement_control_rejected),
   1214             parameters,
   1215         )
   1216     }
   1217 }
   1218 
   1219 /// Boxed callback future tied to the borrowed transaction executor.
   1220 pub type ServiceSqliteTransactionFuture<'a, T, E> =
   1221     Pin<Box<dyn Future<Output = Result<T, E>> + Send + 'a>>;
   1222 
   1223 /// Stable transaction completion phases without expanding SQLite error kinds.
   1224 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
   1225 pub enum ServiceSqliteTransactionErrorKind {
   1226     NotCommitted,
   1227     OperationRolledBack,
   1228     RollbackFailed,
   1229     CommitOutcomeUnknown,
   1230 }
   1231 
   1232 /// Transaction failure retaining trusted details behind redacted diagnostics.
   1233 pub struct ServiceSqliteTransactionError<E> {
   1234     kind: ServiceSqliteTransactionErrorKind,
   1235     operation_error: Option<E>,
   1236     sqlite_error: Option<ServiceSqliteError>,
   1237 }
   1238 
   1239 impl<E> ServiceSqliteTransactionError<E> {
   1240     fn not_committed(error: ServiceSqliteError) -> Self {
   1241         Self::not_committed_with_operation(None, error)
   1242     }
   1243 
   1244     fn not_committed_with_operation(
   1245         operation_error: Option<E>,
   1246         sqlite_error: ServiceSqliteError,
   1247     ) -> Self {
   1248         Self {
   1249             kind: ServiceSqliteTransactionErrorKind::NotCommitted,
   1250             operation_error,
   1251             sqlite_error: Some(sqlite_error),
   1252         }
   1253     }
   1254 
   1255     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1256     fn operation_rolled_back(error: E) -> Self {
   1257         Self {
   1258             kind: ServiceSqliteTransactionErrorKind::OperationRolledBack,
   1259             operation_error: Some(error),
   1260             sqlite_error: None,
   1261         }
   1262     }
   1263 
   1264     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1265     fn rollback_failed(operation_error: Option<E>, sqlite_error: ServiceSqliteError) -> Self {
   1266         Self {
   1267             kind: ServiceSqliteTransactionErrorKind::RollbackFailed,
   1268             operation_error,
   1269             sqlite_error: Some(sqlite_error),
   1270         }
   1271     }
   1272 
   1273     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1274     fn commit_outcome_unknown(error: ServiceSqliteError) -> Self {
   1275         Self {
   1276             kind: ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown,
   1277             operation_error: None,
   1278             sqlite_error: Some(error),
   1279         }
   1280     }
   1281 
   1282     #[must_use]
   1283     pub const fn kind(&self) -> ServiceSqliteTransactionErrorKind {
   1284         self.kind
   1285     }
   1286 
   1287     #[must_use]
   1288     pub const fn operation_error(&self) -> Option<&E> {
   1289         self.operation_error.as_ref()
   1290     }
   1291 
   1292     #[must_use]
   1293     pub const fn sqlite_error(&self) -> Option<&ServiceSqliteError> {
   1294         self.sqlite_error.as_ref()
   1295     }
   1296 }
   1297 
   1298 impl<E> fmt::Debug for ServiceSqliteTransactionError<E> {
   1299     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1300         formatter
   1301             .debug_struct("ServiceSqliteTransactionError")
   1302             .field("kind", &self.kind)
   1303             .field(
   1304                 "operation_error",
   1305                 &self.operation_error.as_ref().map(|_| "[redacted]"),
   1306             )
   1307             .field(
   1308                 "sqlite_error",
   1309                 &self.sqlite_error.as_ref().map(|_| "[redacted]"),
   1310             )
   1311             .finish()
   1312     }
   1313 }
   1314 
   1315 impl<E> fmt::Display for ServiceSqliteTransactionError<E> {
   1316     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1317         formatter.write_str(match self.kind {
   1318             ServiceSqliteTransactionErrorKind::NotCommitted => {
   1319                 "SQLite transaction did not reach commit"
   1320             }
   1321             ServiceSqliteTransactionErrorKind::OperationRolledBack => {
   1322                 "SQLite transaction operation was rolled back"
   1323             }
   1324             ServiceSqliteTransactionErrorKind::RollbackFailed => {
   1325                 "SQLite transaction rollback could not be confirmed"
   1326             }
   1327             ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown => {
   1328                 "SQLite transaction commit outcome is unknown"
   1329             }
   1330         })
   1331     }
   1332 }
   1333 
   1334 impl<E: Send + Sync + 'static> Error for ServiceSqliteTransactionError<E> {}
   1335 
   1336 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1337 struct QuarantinedConnection {
   1338     connection: Option<PoolConnection<Sqlite>>,
   1339     trusted: bool,
   1340 }
   1341 
   1342 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1343 impl QuarantinedConnection {
   1344     fn new(connection: PoolConnection<Sqlite>) -> Self {
   1345         Self {
   1346             connection: Some(connection),
   1347             trusted: false,
   1348         }
   1349     }
   1350 
   1351     fn trust(&mut self) {
   1352         self.trusted = true;
   1353     }
   1354 
   1355     fn into_close_future(mut self) -> Option<BoxFuture<'static, Result<(), sqlx::Error>>> {
   1356         let connection = self.connection.take()?;
   1357         Some(Box::pin(async move { connection.close().await }))
   1358     }
   1359 }
   1360 
   1361 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1362 impl core::ops::Deref for QuarantinedConnection {
   1363     type Target = SqliteConnection;
   1364 
   1365     fn deref(&self) -> &Self::Target {
   1366         self.connection.as_deref().expect("connection is retained")
   1367     }
   1368 }
   1369 
   1370 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1371 impl core::ops::DerefMut for QuarantinedConnection {
   1372     fn deref_mut(&mut self) -> &mut Self::Target {
   1373         self.connection
   1374             .as_deref_mut()
   1375             .expect("connection is retained")
   1376     }
   1377 }
   1378 
   1379 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1380 impl Drop for QuarantinedConnection {
   1381     fn drop(&mut self) {
   1382         if !self.trusted
   1383             && let Some(connection) = self.connection.as_mut()
   1384         {
   1385             connection.close_on_drop();
   1386         }
   1387     }
   1388 }
   1389 
   1390 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1391 fn sqlite_source(source: sqlx::Error) -> ServiceSqliteError {
   1392     ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
   1393 }
   1394 
   1395 #[cfg(not(any(target_os = "linux", target_os = "macos")))]
   1396 fn unsupported_host() -> ServiceSqliteError {
   1397     ServiceSqliteError::new(ServiceSqliteErrorKind::Open)
   1398 }
   1399 
   1400 #[cfg(test)]
   1401 mod tests {
   1402     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1403     use std::{
   1404         collections::BTreeMap,
   1405         convert::Infallible,
   1406         fs,
   1407         num::NonZeroU32,
   1408         os::unix::fs::PermissionsExt,
   1409         path::Path,
   1410         sync::{
   1411             Arc,
   1412             atomic::{AtomicBool, Ordering},
   1413         },
   1414         time::{Duration, SystemTime},
   1415     };
   1416 
   1417     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1418     use radroots_runtime_paths::{
   1419         InstanceId, RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver,
   1420         RadrootsPlatform, RuntimeContext, RuntimeContextBootstrap, RuntimeContextSource, ServiceId,
   1421     };
   1422     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1423     use radroots_storage::event::SourceGeneration;
   1424     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1425     use sqlx::{Connection, sqlite::SqliteConnectOptions};
   1426     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1427     use tokio::sync::Notify;
   1428 
   1429     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1430     static CAPTURE_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
   1431 
   1432     use super::*;
   1433 
   1434     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1435     const HOST_TABLE_SQL: &str = "CREATE TABLE host_probe (value INTEGER NOT NULL)";
   1436 
   1437     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1438     fn runtime_context(root: &std::path::Path) -> RuntimeContext {
   1439         RuntimeContext::resolve(
   1440             &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
   1441             RuntimeContextBootstrap::new(
   1442                 RadrootsPathProfile::RepoLocal,
   1443                 Some(root.to_path_buf()),
   1444                 RuntimeContextSource::BootstrapCli,
   1445                 RuntimeContextSource::BootstrapCli,
   1446             )
   1447             .expect("valid bootstrap"),
   1448             ServiceId::new("myc").expect("service ID"),
   1449             InstanceId::new("host-boundary").expect("instance ID"),
   1450         )
   1451         .expect("runtime context")
   1452     }
   1453 
   1454     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1455     #[test]
   1456     fn rollback_failure_selection_preserves_each_exact_precedence() {
   1457         let error = || ServiceSqliteError::new(ServiceSqliteErrorKind::Open);
   1458         for authority in [false, true] {
   1459             for rollback in [false, true] {
   1460                 for confirmed in [false, true] {
   1461                     for removal in [false, true] {
   1462                         let expected_precondition = if authority {
   1463                             1
   1464                         } else if rollback && !confirmed {
   1465                             2
   1466                         } else if removal {
   1467                             3
   1468                         } else {
   1469                             0
   1470                         };
   1471                         let precondition = precondition_rollback_failure(
   1472                             authority.then(error),
   1473                             rollback.then(error),
   1474                             confirmed,
   1475                             removal.then(error),
   1476                         );
   1477                         assert_eq!(
   1478                             usize::from(precondition.is_some()),
   1479                             usize::from(expected_precondition != 0)
   1480                         );
   1481 
   1482                         let expected_drift = (rollback && !confirmed) || removal;
   1483                         assert_eq!(
   1484                             authority_drift_rollback_failure(
   1485                                 rollback.then(error),
   1486                                 confirmed,
   1487                                 removal.then(error),
   1488                             )
   1489                             .is_some(),
   1490                             expected_drift
   1491                         );
   1492 
   1493                         let expected_operation = (rollback && !confirmed) || removal || authority;
   1494                         assert_eq!(
   1495                             operation_rollback_failure(
   1496                                 rollback.then(error),
   1497                                 confirmed,
   1498                                 removal.then(error),
   1499                                 authority.then(error),
   1500                             )
   1501                             .is_some(),
   1502                             expected_operation
   1503                         );
   1504                     }
   1505                 }
   1506             }
   1507         }
   1508         assert!(unconfirmed_rollback_error(Some(error()), false).is_some());
   1509         assert!(unconfirmed_rollback_error(Some(error()), true).is_none());
   1510         assert!(unconfirmed_rollback_error(None, false).is_none());
   1511 
   1512         assert!(integrity_driver_close_result(Ok(()), false).is_ok());
   1513         assert!(matches!(
   1514             integrity_driver_close_result(Ok(()), true),
   1515             Err(IntegrityInspectionDriverFailure::ConnectionClose)
   1516         ));
   1517         assert!(matches!(
   1518             integrity_driver_close_result(Err(sqlx::Error::Protocol("close".to_owned())), false),
   1519             Err(IntegrityInspectionDriverFailure::ConnectionClose)
   1520         ));
   1521     }
   1522 
   1523     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1524     #[tokio::test(flavor = "current_thread")]
   1525     async fn final_connection_policy_classifier_preserves_pragma_kind() {
   1526         let mut first =
   1527             SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:"))
   1528                 .await
   1529                 .expect("first connection");
   1530         let mut second =
   1531             SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:"))
   1532                 .await
   1533                 .expect("second connection");
   1534         let initial = crate::migration::read_connection_policy(&mut first)
   1535             .await
   1536             .expect("initial policy");
   1537         let same = crate::migration::read_connection_policy(&mut second)
   1538             .await
   1539             .expect("same policy");
   1540         assert!(final_connection_policy_matches(&initial, &same).is_ok());
   1541         sqlx::query("PRAGMA query_only = ON")
   1542             .execute(&mut second)
   1543             .await
   1544             .expect("change policy");
   1545         let changed = crate::migration::read_connection_policy(&mut second)
   1546             .await
   1547             .expect("changed policy");
   1548         let error = final_connection_policy_matches(&initial, &changed).expect_err("policy drift");
   1549         assert_eq!(error.kind(), ServiceSqliteErrorKind::Pragma);
   1550     }
   1551 
   1552     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1553     fn migration_catalog() -> MigrationCatalog {
   1554         MigrationCatalog::new([]).expect("empty v1 migration catalog")
   1555     }
   1556 
   1557     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1558     fn schema_catalog(migrations: &MigrationCatalog) -> SchemaCatalog {
   1559         let table = crate::SchemaObject::new(
   1560             crate::SchemaObjectKind::Table,
   1561             "host_probe",
   1562             "host_probe",
   1563             HOST_TABLE_SQL,
   1564             crate::SchemaObject::computed_digest(
   1565                 crate::SchemaObjectKind::Table,
   1566                 "host_probe",
   1567                 "host_probe",
   1568                 HOST_TABLE_SQL,
   1569             )
   1570             .expect("table digest"),
   1571         )
   1572         .expect("table descriptor");
   1573         let version_digest = crate::SchemaVersionCatalog::computed_digest(1, [table.clone()])
   1574             .expect("version digest");
   1575         let version =
   1576             crate::SchemaVersionCatalog::new(1, [table], version_digest).expect("schema version");
   1577         SchemaCatalog::new(migrations, [version]).expect("schema catalog")
   1578     }
   1579 
   1580     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1581     fn build_identity() -> MigrationBuildIdentity {
   1582         MigrationBuildIdentity::new(
   1583             "0.1.0-alpha",
   1584             "0123456789abcdef0123456789abcdef01234567",
   1585             "89abcdef0123456789abcdef0123456789abcdef",
   1586             "1.97.1",
   1587             "x86_64-unknown-linux-gnu",
   1588             "service-host",
   1589             1,
   1590             2,
   1591             3,
   1592             4,
   1593             5,
   1594         )
   1595         .expect("build identity")
   1596     }
   1597 
   1598     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1599     async fn initialized_host() -> (
   1600         tempfile::TempDir,
   1601         ServiceSqlitePaths,
   1602         ServiceDatabaseIdentity,
   1603         MigrationCatalog,
   1604         SchemaCatalog,
   1605         ServiceSqliteHost,
   1606     ) {
   1607         initialized_host_with_options(ServiceSqliteConnectionOptions::reviewed()).await
   1608     }
   1609 
   1610     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1611     async fn initialized_host_with_options(
   1612         options: ServiceSqliteConnectionOptions,
   1613     ) -> (
   1614         tempfile::TempDir,
   1615         ServiceSqlitePaths,
   1616         ServiceDatabaseIdentity,
   1617         MigrationCatalog,
   1618         SchemaCatalog,
   1619         ServiceSqliteHost,
   1620     ) {
   1621         let root = tempfile::tempdir().expect("temporary root");
   1622         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path()))
   1623             .expect("SQLite paths");
   1624         fs::create_dir_all(paths.state_database().parent().expect("state directory"))
   1625             .expect("create state directory");
   1626         let metadata = crate::ServiceDatabaseMetadata::new(
   1627             &paths,
   1628             SourceGeneration::new([9; 32]).expect("source generation"),
   1629             NonZeroU32::new(1).expect("schema version"),
   1630             1_700_000_000_000,
   1631             crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   1632         )
   1633         .expect("database metadata");
   1634         let migrations = migration_catalog();
   1635         let schema = schema_catalog(&migrations);
   1636         let authority = crate::initialize_database(
   1637             &paths,
   1638             OpenMode::Initialize,
   1639             &metadata,
   1640             &schema,
   1641             |initializer| {
   1642                 Box::pin(async move {
   1643                     sqlx::query(HOST_TABLE_SQL)
   1644                         .execute(initializer)
   1645                         .await
   1646                         .expect("create host table");
   1647                     Ok::<_, Infallible>(())
   1648                 })
   1649             },
   1650         )
   1651         .await
   1652         .expect("initialize database");
   1653         let identity = metadata.identity();
   1654         let (host, outcome) = ServiceSqliteHost::open_initialized(
   1655             &paths,
   1656             &identity,
   1657             &migrations,
   1658             &schema,
   1659             options,
   1660             authority,
   1661             MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"),
   1662             &build_identity(),
   1663             &[],
   1664         )
   1665         .await
   1666         .expect("open initialized host");
   1667         assert_eq!(outcome.initial_version(), 1);
   1668         assert_eq!(outcome.final_version(), 1);
   1669         assert_eq!(outcome.applied_count(), 0);
   1670         (root, paths, identity, migrations, schema, host)
   1671     }
   1672 
   1673     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1674     async fn row_count(host: &ServiceSqliteHost) -> i64 {
   1675         host.transaction(|transaction| {
   1676             Box::pin(async move {
   1677                 sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   1678                     .fetch_one(&mut *transaction)
   1679                     .await
   1680             })
   1681         })
   1682         .await
   1683         .expect("count rows")
   1684     }
   1685 
   1686     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1687     fn integrity_checked_at(value: u64) -> crate::IntegrityCheckedAtUnixMs {
   1688         crate::IntegrityCheckedAtUnixMs::new(value).expect("integrity inspection time")
   1689     }
   1690 
   1691     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1692     #[tokio::test]
   1693     async fn existing_intent_returns_actual_metadata_without_releasing_authority() {
   1694         let (_root, paths, identity, migrations, schema, initialized) = initialized_host().await;
   1695         initialized.close().await.expect("close initialized host");
   1696         let intent = ExistingServiceDatabaseIntent::new(
   1697             &paths,
   1698             identity.supported_state_schema_version(),
   1699             identity.application_id(),
   1700         );
   1701 
   1702         let wrong_application = ExistingServiceDatabaseIntent::new(
   1703             &paths,
   1704             identity.supported_state_schema_version(),
   1705             crate::ServiceSqliteApplicationId::new(7).expect("other application"),
   1706         );
   1707         let error = ServiceSqliteHost::open_read_write_existing_with_intent(
   1708             &paths,
   1709             &wrong_application,
   1710             &migrations,
   1711             &schema,
   1712             ServiceSqliteConnectionOptions::reviewed(),
   1713             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   1714             &build_identity(),
   1715             &[],
   1716         )
   1717         .await
   1718         .expect_err("application mismatch");
   1719         assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata);
   1720 
   1721         let (opened, outcome) = ServiceSqliteHost::open_read_write_existing_with_intent(
   1722             &paths,
   1723             &intent,
   1724             &migrations,
   1725             &schema,
   1726             ServiceSqliteConnectionOptions::reviewed(),
   1727             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   1728             &build_identity(),
   1729             &[],
   1730         )
   1731         .await
   1732         .expect("open writable from intent");
   1733         assert_eq!(outcome.applied_count(), 0);
   1734         assert_eq!(
   1735             opened.database_metadata().source_generation(),
   1736             identity.source_generation()
   1737         );
   1738         assert_eq!(opened.host().mode(), OpenMode::ReadWriteExisting);
   1739         let debug = format!("{opened:?}");
   1740         assert!(debug.contains("OpenedServiceDatabase"));
   1741         assert!(!debug.contains("09090909"));
   1742         assert!(WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting).is_err());
   1743         let (writable, actual) = opened.into_parts();
   1744         assert_eq!(actual.source_generation(), identity.source_generation());
   1745         assert_eq!(row_count(&writable).await, 0);
   1746         writable.close().await.expect("close writable host");
   1747 
   1748         let inspected = ServiceSqliteHost::open_read_only_inspection_with_intent(
   1749             &paths,
   1750             &intent,
   1751             &migrations,
   1752             &schema,
   1753             ServiceSqliteConnectionOptions::reviewed(),
   1754         )
   1755         .await
   1756         .expect("open inspection from intent");
   1757         assert_eq!(inspected.host().mode(), OpenMode::ReadOnlyInspection);
   1758         assert_eq!(
   1759             inspected.database_metadata().source_generation(),
   1760             identity.source_generation()
   1761         );
   1762         assert!(WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting).is_err());
   1763         let (inspection, actual) = inspected.into_parts();
   1764         assert_eq!(actual.source_generation(), identity.source_generation());
   1765         inspection.close().await.expect("close inspection host");
   1766     }
   1767 
   1768     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1769     #[tokio::test(flavor = "current_thread")]
   1770     async fn open_or_initialize_uses_exclusive_create_without_filesystem_probing() {
   1771         let root = tempfile::tempdir().expect("temporary root");
   1772         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path()))
   1773             .expect("SQLite paths");
   1774         fs::create_dir_all(paths.state_database().parent().expect("state directory"))
   1775             .expect("create state directory");
   1776         let metadata = ServiceDatabaseMetadata::new(
   1777             &paths,
   1778             SourceGeneration::new([31; 32]).expect("source generation"),
   1779             NonZeroU32::new(1).expect("schema version"),
   1780             1_700_000_000_000,
   1781             crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   1782         )
   1783         .expect("database metadata");
   1784         let migrations = migration_catalog();
   1785         let schema = schema_catalog(&migrations);
   1786         let initialized_callback = Arc::new(AtomicBool::new(false));
   1787         let called = Arc::clone(&initialized_callback);
   1788         let (initialized, outcome) = ServiceSqliteHost::open_or_initialize(
   1789             &paths,
   1790             &metadata,
   1791             &migrations,
   1792             &schema,
   1793             ServiceSqliteConnectionOptions::reviewed(),
   1794             MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"),
   1795             &build_identity(),
   1796             &[],
   1797             move |initializer| {
   1798                 called.store(true, Ordering::Release);
   1799                 Box::pin(async move {
   1800                     sqlx::query(HOST_TABLE_SQL).execute(initializer).await?;
   1801                     Ok::<(), sqlx::Error>(())
   1802                 })
   1803             },
   1804         )
   1805         .await
   1806         .expect("initialize missing state");
   1807         assert!(initialized_callback.load(Ordering::Acquire));
   1808         assert_eq!(initialized.host().mode(), OpenMode::Initialize);
   1809         assert_eq!(initialized.database_metadata(), &metadata);
   1810         assert_eq!(outcome.applied_count(), 0);
   1811         let (initialized, initialized_metadata) = initialized.into_parts();
   1812         assert_eq!(initialized_metadata, metadata);
   1813         initialized.close().await.expect("close initialized host");
   1814 
   1815         let existing_callback = Arc::new(AtomicBool::new(false));
   1816         let called = Arc::clone(&existing_callback);
   1817         let alternative_metadata = ServiceDatabaseMetadata::new(
   1818             &paths,
   1819             SourceGeneration::new([32; 32]).expect("alternative source generation"),
   1820             NonZeroU32::new(1).expect("schema version"),
   1821             1_700_000_000_001,
   1822             crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   1823         )
   1824         .expect("alternative database metadata");
   1825         let (existing, outcome) = ServiceSqliteHost::open_or_initialize(
   1826             &paths,
   1827             &alternative_metadata,
   1828             &migrations,
   1829             &schema,
   1830             ServiceSqliteConnectionOptions::reviewed(),
   1831             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   1832             &build_identity(),
   1833             &[],
   1834             move |_| {
   1835                 called.store(true, Ordering::Release);
   1836                 Box::pin(async { Ok::<(), Infallible>(()) })
   1837             },
   1838         )
   1839         .await
   1840         .expect("open exact existing state");
   1841         assert!(!existing_callback.load(Ordering::Acquire));
   1842         assert_eq!(existing.host().mode(), OpenMode::ReadWriteExisting);
   1843         assert_eq!(existing.database_metadata(), &metadata);
   1844         assert_eq!(outcome.applied_count(), 0);
   1845         let (existing, existing_metadata) = existing.into_parts();
   1846         assert_eq!(existing_metadata, metadata);
   1847         assert_eq!(row_count(&existing).await, 0);
   1848         existing.close().await.expect("close existing host");
   1849     }
   1850 
   1851     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1852     #[tokio::test(flavor = "current_thread")]
   1853     async fn open_or_initialize_rejects_catalog_mismatch_before_reservation() {
   1854         const MIGRATION_SQL: &str = "CREATE TABLE later_probe (value INTEGER NOT NULL)";
   1855 
   1856         let root = tempfile::tempdir().expect("temporary root");
   1857         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path()))
   1858             .expect("SQLite paths");
   1859         fs::create_dir_all(paths.state_database().parent().expect("state directory"))
   1860             .expect("create state directory");
   1861         let metadata = ServiceDatabaseMetadata::new(
   1862             &paths,
   1863             SourceGeneration::new([32; 32]).expect("source generation"),
   1864             NonZeroU32::new(1).expect("schema version"),
   1865             1_700_000_000_000,
   1866             crate::ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   1867         )
   1868         .expect("database metadata");
   1869         let base_migrations = migration_catalog();
   1870         let schema = schema_catalog(&base_migrations);
   1871         let different_migrations = MigrationCatalog::new([crate::MigrationDescriptor::sql(
   1872             2,
   1873             "create_later_probe",
   1874             MIGRATION_SQL,
   1875             crate::MigrationChecksum::for_sql(MIGRATION_SQL),
   1876         )
   1877         .expect("different migration")])
   1878         .expect("different migration catalog");
   1879         let callback_called = Arc::new(AtomicBool::new(false));
   1880         let called = Arc::clone(&callback_called);
   1881 
   1882         let error = ServiceSqliteHost::open_or_initialize(
   1883             &paths,
   1884             &metadata,
   1885             &different_migrations,
   1886             &schema,
   1887             ServiceSqliteConnectionOptions::reviewed(),
   1888             MigrationAppliedAtUnixSeconds::new(1_700_000_002).expect("migration time"),
   1889             &build_identity(),
   1890             &[],
   1891             move |_| {
   1892                 called.store(true, Ordering::Release);
   1893                 Box::pin(async { Ok::<(), Infallible>(()) })
   1894             },
   1895         )
   1896         .await
   1897         .expect_err("catalog mismatch must fail before reservation");
   1898 
   1899         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   1900         assert!(!callback_called.load(Ordering::Acquire));
   1901         assert!(!paths.state_database().exists());
   1902     }
   1903 
   1904     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1905     #[tokio::test]
   1906     async fn integrity_inspection_is_explicit_safe_and_available_in_every_host_mode() {
   1907         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   1908         crate::integrity::integrity_test_seam::release();
   1909         let (_root, paths, identity, migrations, schema, initialized) = initialized_host().await;
   1910 
   1911         let initialized_report = initialized
   1912             .inspect_integrity(integrity_checked_at(1_700_000_000_500))
   1913             .await
   1914             .expect("inspect initialized host");
   1915         assert_eq!(initialized.mode(), OpenMode::Initialize);
   1916         assert_eq!(
   1917             initialized_report.sqlite(),
   1918             crate::IntegrityCheckOutcome::Verified
   1919         );
   1920         assert_eq!(
   1921             initialized_report.foreign_keys(),
   1922             crate::IntegrityCheckOutcome::Verified
   1923         );
   1924         assert!(initialized_report.diagnostics().is_empty());
   1925         assert_eq!(
   1926             initialized_report.storage_integrity(),
   1927             crate::StorageIntegrity::Verified
   1928         );
   1929         initialized.close().await.expect("close initialized host");
   1930 
   1931         let (writable, outcome) = ServiceSqliteHost::open_read_write_existing(
   1932             &paths,
   1933             &identity,
   1934             &migrations,
   1935             &schema,
   1936             ServiceSqliteConnectionOptions::reviewed(),
   1937             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   1938             &build_identity(),
   1939             &[],
   1940         )
   1941         .await
   1942         .expect("open writable host");
   1943         assert_eq!(outcome.applied_count(), 0);
   1944         let writable_report = writable
   1945             .inspect_integrity(integrity_checked_at(1_700_000_000_501))
   1946             .await
   1947             .expect("inspect writable host");
   1948         assert_eq!(writable.mode(), OpenMode::ReadWriteExisting);
   1949         assert_eq!(
   1950             writable_report.storage_integrity(),
   1951             crate::StorageIntegrity::Verified
   1952         );
   1953         writable.close().await.expect("close writable host");
   1954 
   1955         let read_only = ServiceSqliteHost::open_read_only_inspection(
   1956             &paths,
   1957             &identity,
   1958             &migrations,
   1959             &schema,
   1960             ServiceSqliteConnectionOptions::reviewed(),
   1961         )
   1962         .await
   1963         .expect("open read-only host");
   1964         let read_only_report = read_only
   1965             .inspect_integrity(integrity_checked_at(1_700_000_000_502))
   1966             .await
   1967             .expect("inspect read-only host");
   1968         assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection);
   1969         assert_eq!(
   1970             read_only_report.storage_integrity(),
   1971             crate::StorageIntegrity::Verified
   1972         );
   1973         read_only.close().await.expect("close read-only host");
   1974     }
   1975 
   1976     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1977     #[tokio::test]
   1978     async fn integrity_inspection_is_single_admission_cancel_safe_and_close_drained() {
   1979         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   1980         crate::integrity::integrity_test_seam::release();
   1981         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   1982         let host = Arc::new(host);
   1983 
   1984         crate::integrity::integrity_test_seam::block(
   1985             crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE,
   1986         );
   1987         let first = tokio::spawn({
   1988             let host = Arc::clone(&host);
   1989             async move {
   1990                 host.inspect_integrity(integrity_checked_at(1_700_000_000_510))
   1991                     .await
   1992             }
   1993         });
   1994         while crate::integrity::integrity_test_seam::reached()
   1995             != crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE
   1996         {
   1997             tokio::task::yield_now().await;
   1998         }
   1999         let concurrent = host
   2000             .inspect_integrity(integrity_checked_at(1_700_000_000_511))
   2001             .await
   2002             .expect_err("second integrity inspection is rejected");
   2003         assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Integrity);
   2004         first.abort();
   2005         assert!(
   2006             first
   2007                 .await
   2008                 .expect_err("inspection is cancelled")
   2009                 .is_cancelled()
   2010         );
   2011         crate::integrity::integrity_test_seam::release();
   2012         assert!(
   2013             !host.integrity_driver.lock().await.is_idle(),
   2014             "cancelled SQLx connection remains host-owned until explicit cleanup"
   2015         );
   2016 
   2017         let recovered = host
   2018             .inspect_integrity(integrity_checked_at(1_700_000_000_512))
   2019             .await
   2020             .expect("inspection recovers after cancellation");
   2021         assert_eq!(
   2022             recovered.storage_integrity(),
   2023             crate::StorageIntegrity::Verified
   2024         );
   2025 
   2026         crate::integrity::integrity_test_seam::block(
   2027             crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
   2028         );
   2029         let timed_out = tokio::time::timeout(
   2030             Duration::from_millis(20),
   2031             host.inspect_integrity(integrity_checked_at(1_700_000_000_513)),
   2032         )
   2033         .await;
   2034         assert!(timed_out.is_err(), "caller deadline cancels inspection");
   2035         crate::integrity::integrity_test_seam::release();
   2036         assert!(!host.integrity_driver.lock().await.is_idle());
   2037 
   2038         crate::integrity::integrity_test_seam::block(
   2039             crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK,
   2040         );
   2041         let pre_rollback = tokio::spawn({
   2042             let host = Arc::clone(&host);
   2043             async move {
   2044                 host.inspect_integrity(integrity_checked_at(1_700_000_000_514))
   2045                     .await
   2046             }
   2047         });
   2048         while crate::integrity::integrity_test_seam::reached()
   2049             != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK
   2050         {
   2051             tokio::task::yield_now().await;
   2052         }
   2053         pre_rollback.abort();
   2054         assert!(
   2055             pre_rollback
   2056                 .await
   2057                 .expect_err("pre-rollback inspection is cancelled")
   2058                 .is_cancelled()
   2059         );
   2060         crate::integrity::integrity_test_seam::release();
   2061         assert!(!host.integrity_driver.lock().await.is_idle());
   2062         host.inspect_integrity(integrity_checked_at(1_700_000_000_515))
   2063             .await
   2064             .expect("inspection recovers after pre-rollback cancellation");
   2065 
   2066         crate::integrity::integrity_test_seam::block(
   2067             crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK,
   2068         );
   2069         let admitted = tokio::spawn({
   2070             let host = Arc::clone(&host);
   2071             async move {
   2072                 host.inspect_integrity(integrity_checked_at(1_700_000_000_516))
   2073                     .await
   2074             }
   2075         });
   2076         while crate::integrity::integrity_test_seam::reached()
   2077             != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK
   2078         {
   2079             tokio::task::yield_now().await;
   2080         }
   2081         let close = tokio::spawn({
   2082             let host = Arc::clone(&host);
   2083             async move { host.close().await }
   2084         });
   2085         while !host.closing.load(Ordering::Acquire) {
   2086             tokio::task::yield_now().await;
   2087         }
   2088         let rejected = host
   2089             .inspect_integrity(integrity_checked_at(1_700_000_000_517))
   2090             .await
   2091             .expect_err("closing host rejects inspection");
   2092         assert_eq!(rejected.kind(), ServiceSqliteErrorKind::Open);
   2093         assert!(!close.is_finished());
   2094         crate::integrity::integrity_test_seam::release();
   2095         admitted
   2096             .await
   2097             .expect("admitted inspection task joins")
   2098             .expect("admitted inspection completes");
   2099         close
   2100             .await
   2101             .expect("close task joins")
   2102             .expect("close drains admitted inspection");
   2103         let closed = host
   2104             .inspect_integrity(integrity_checked_at(1_700_000_000_518))
   2105             .await
   2106             .expect_err("closed host rejects inspection");
   2107         assert_eq!(closed.kind(), ServiceSqliteErrorKind::Open);
   2108     }
   2109 
   2110     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2111     #[tokio::test]
   2112     async fn integrity_inspection_preserves_authority_precedence_after_await() {
   2113         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   2114         crate::integrity::integrity_test_seam::release();
   2115         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2116         let host = Arc::new(host);
   2117         crate::integrity::integrity_test_seam::block(
   2118             crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
   2119         );
   2120         let inspection = tokio::spawn({
   2121             let host = Arc::clone(&host);
   2122             async move {
   2123                 host.inspect_integrity(integrity_checked_at(1_700_000_000_520))
   2124                     .await
   2125             }
   2126         });
   2127         while crate::integrity::integrity_test_seam::reached()
   2128             != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS
   2129         {
   2130             tokio::task::yield_now().await;
   2131         }
   2132         let retired_lock = paths
   2133             .state_lock()
   2134             .parent()
   2135             .expect("state directory")
   2136             .join("retired-integrity-state.lock");
   2137         fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
   2138         crate::integrity::integrity_test_seam::release();
   2139         let error = inspection
   2140             .await
   2141             .expect("inspection task joins")
   2142             .expect_err("authority drift rejects inspection");
   2143         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   2144         fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
   2145         host.close().await.expect("close host after restored lock");
   2146     }
   2147 
   2148     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2149     #[tokio::test]
   2150     async fn observed_authority_drift_precedes_transient_restore_and_close_failure() {
   2151         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   2152         crate::integrity::integrity_test_seam::release();
   2153         crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
   2154         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2155         let host = Arc::new(host);
   2156 
   2157         crate::integrity::integrity_test_seam::block(
   2158             crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
   2159         );
   2160         let inspection = tokio::spawn({
   2161             let host = Arc::clone(&host);
   2162             async move {
   2163                 host.inspect_integrity(integrity_checked_at(1_700_000_000_525))
   2164                     .await
   2165             }
   2166         });
   2167         while crate::integrity::integrity_test_seam::reached()
   2168             != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS
   2169         {
   2170             tokio::task::yield_now().await;
   2171         }
   2172 
   2173         let retired_lock = paths
   2174             .state_lock()
   2175             .parent()
   2176             .expect("state directory")
   2177             .join("transient-integrity-state.lock");
   2178         fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
   2179         crate::integrity::integrity_test_seam::block(
   2180             crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING,
   2181         );
   2182         while crate::integrity::integrity_test_seam::reached()
   2183             != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
   2184         {
   2185             tokio::task::yield_now().await;
   2186         }
   2187 
   2188         fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
   2189         crate::integrity::integrity_test_seam::inject_connection_close_failure(true);
   2190         crate::integrity::integrity_test_seam::release();
   2191         let error = inspection
   2192             .await
   2193             .expect("inspection task joins")
   2194             .expect_err("observed authority drift remains terminal");
   2195         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   2196         assert!(host.integrity_driver.lock().await.is_idle());
   2197         host.close().await.expect("close host after restored lock");
   2198         crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
   2199     }
   2200 
   2201     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2202     #[tokio::test]
   2203     async fn cancelled_real_sqlite_work_is_explicitly_closed_before_retry_or_host_close() {
   2204         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   2205         crate::integrity::integrity_test_seam::release();
   2206         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2207         let host = Arc::new(host);
   2208 
   2209         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
   2210         let cancelled = tokio::spawn({
   2211             let host = Arc::clone(&host);
   2212             async move {
   2213                 host.inspect_integrity(integrity_checked_at(1_700_000_000_530))
   2214                     .await
   2215             }
   2216         });
   2217         while crate::integrity::integrity_test_seam::reached()
   2218             != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
   2219         {
   2220             tokio::task::yield_now().await;
   2221         }
   2222         assert!(!cancelled.is_finished(), "SQLite probe remains in flight");
   2223         cancelled.abort();
   2224         assert!(
   2225             cancelled
   2226                 .await
   2227                 .expect_err("real SQLite inspection is cancelled")
   2228                 .is_cancelled()
   2229         );
   2230         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
   2231         assert!(!host.integrity_driver.lock().await.is_idle());
   2232 
   2233         let cancelled_cleanup = tokio::spawn({
   2234             let host = Arc::clone(&host);
   2235             async move {
   2236                 host.inspect_integrity(integrity_checked_at(1_700_000_000_531))
   2237                     .await
   2238             }
   2239         });
   2240         while crate::integrity::integrity_test_seam::reached()
   2241             != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
   2242         {
   2243             tokio::task::yield_now().await;
   2244         }
   2245         cancelled_cleanup.abort();
   2246         assert!(
   2247             cancelled_cleanup
   2248                 .await
   2249                 .expect_err("retained close retry is cancelled")
   2250                 .is_cancelled()
   2251         );
   2252         assert!(!host.integrity_driver.lock().await.is_idle());
   2253 
   2254         let recovered = host
   2255             .inspect_integrity(integrity_checked_at(1_700_000_000_532))
   2256             .await
   2257             .expect("retry explicitly closes prior SQLite worker before inspecting");
   2258         assert_eq!(
   2259             recovered.storage_integrity(),
   2260             crate::StorageIntegrity::Verified
   2261         );
   2262         assert!(host.integrity_driver.lock().await.is_idle());
   2263 
   2264         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
   2265         let cancelled = tokio::spawn({
   2266             let host = Arc::clone(&host);
   2267             async move {
   2268                 host.inspect_integrity(integrity_checked_at(1_700_000_000_533))
   2269                     .await
   2270             }
   2271         });
   2272         while crate::integrity::integrity_test_seam::reached()
   2273             != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
   2274         {
   2275             tokio::task::yield_now().await;
   2276         }
   2277         cancelled.abort();
   2278         assert!(
   2279             cancelled
   2280                 .await
   2281                 .expect_err("second real SQLite inspection is cancelled")
   2282                 .is_cancelled()
   2283         );
   2284         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
   2285         assert!(!host.integrity_driver.lock().await.is_idle());
   2286         host.close()
   2287             .await
   2288             .expect("host close explicitly terminates retained SQLite worker");
   2289         assert!(host.integrity_driver.lock().await.is_idle());
   2290     }
   2291 
   2292     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2293     #[tokio::test]
   2294     async fn retained_integrity_close_releases_authority_and_caches_concurrent_lock_drift() {
   2295         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   2296         crate::integrity::integrity_test_seam::release();
   2297         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2298         let host = Arc::new(host);
   2299 
   2300         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
   2301         let inspection = tokio::spawn({
   2302             let host = Arc::clone(&host);
   2303             async move {
   2304                 host.inspect_integrity(integrity_checked_at(1_700_000_000_540))
   2305                     .await
   2306             }
   2307         });
   2308         while crate::integrity::integrity_test_seam::reached()
   2309             != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
   2310         {
   2311             tokio::task::yield_now().await;
   2312         }
   2313         inspection.abort();
   2314         assert!(
   2315             inspection
   2316                 .await
   2317                 .expect_err("real integrity work is cancelled")
   2318                 .is_cancelled()
   2319         );
   2320         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
   2321 
   2322         let close = tokio::spawn({
   2323             let host = Arc::clone(&host);
   2324             async move { host.close().await }
   2325         });
   2326         while crate::integrity::integrity_test_seam::reached()
   2327             != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
   2328         {
   2329             tokio::task::yield_now().await;
   2330         }
   2331         let retired_lock = paths
   2332             .state_lock()
   2333             .parent()
   2334             .expect("state directory")
   2335             .join("retired-integrity-close-state.lock");
   2336         fs::rename(paths.state_lock(), &retired_lock).expect("retire held writer lock");
   2337         fs::write(paths.state_lock(), b"").expect("create replacement writer lock");
   2338         fs::set_permissions(paths.state_lock(), fs::Permissions::from_mode(0o600))
   2339             .expect("replacement writer lock mode");
   2340 
   2341         let error = close
   2342             .await
   2343             .expect("close task joins")
   2344             .expect_err("authority drift is terminal");
   2345         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   2346         let repeated = host.close().await.expect_err("terminal result is cached");
   2347         assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority);
   2348         assert!(host.integrity_driver.lock().await.is_idle());
   2349 
   2350         fs::remove_file(paths.state_lock()).expect("remove replacement writer lock");
   2351         fs::rename(&retired_lock, paths.state_lock()).expect("restore original writer lock");
   2352         let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   2353             .expect("authority reacquisition")
   2354             .expect("writer authority");
   2355         authority.release().expect("release reacquired authority");
   2356     }
   2357 
   2358     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2359     #[tokio::test]
   2360     async fn retained_integrity_close_failure_is_cached_as_open_and_releases_authority() {
   2361         let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
   2362         crate::integrity::integrity_test_seam::release();
   2363         crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
   2364         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2365         let host = Arc::new(host);
   2366 
   2367         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
   2368         let inspection = tokio::spawn({
   2369             let host = Arc::clone(&host);
   2370             async move {
   2371                 host.inspect_integrity(integrity_checked_at(1_700_000_000_550))
   2372                     .await
   2373             }
   2374         });
   2375         while crate::integrity::integrity_test_seam::reached()
   2376             != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
   2377         {
   2378             tokio::task::yield_now().await;
   2379         }
   2380         inspection.abort();
   2381         assert!(
   2382             inspection
   2383                 .await
   2384                 .expect_err("real integrity work is cancelled")
   2385                 .is_cancelled()
   2386         );
   2387         crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
   2388         assert!(!host.integrity_driver.lock().await.is_idle());
   2389 
   2390         crate::integrity::integrity_test_seam::inject_connection_close_failure(true);
   2391         let error = host
   2392             .close()
   2393             .await
   2394             .expect_err("retained connection close failure is terminal");
   2395         assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
   2396         let repeated = host.close().await.expect_err("terminal result is cached");
   2397         assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Open);
   2398         assert!(host.integrity_driver.lock().await.is_idle());
   2399 
   2400         let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   2401             .expect("authority reacquisition")
   2402             .expect("writer authority");
   2403         authority.release().expect("release reacquired authority");
   2404         crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
   2405     }
   2406 
   2407     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2408     #[tokio::test]
   2409     async fn online_backup_captures_exact_member_manifest_and_preserves_source() {
   2410         use sha2::Digest;
   2411 
   2412         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2413         crate::backup::test_capture_reset();
   2414         let (root, paths, identity, migrations, schema, host) = initialized_host().await;
   2415         host.transaction(|transaction| {
   2416             Box::pin(async move {
   2417                 sqlx::query("INSERT INTO host_probe (value) VALUES (41), (42)")
   2418                     .execute(&mut *transaction)
   2419                     .await
   2420                     .map(|_| ())
   2421             })
   2422         })
   2423         .await
   2424         .expect("seed live WAL state");
   2425 
   2426         let output = root.path().join("backup-output");
   2427         fs::create_dir(&output).expect("backup parent");
   2428         fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2429             .expect("backup parent mode");
   2430         let collision = output.join("collision");
   2431         fs::create_dir(&collision).expect("preexisting collision");
   2432         fs::write(collision.join("foreign"), b"preserve").expect("foreign collision member");
   2433         let collision_error = host
   2434             .capture_online_backup(
   2435                 &collision,
   2436                 crate::BackupCreatedAtUnixMs::new(1_700_000_000_100).expect("backup creation time"),
   2437             )
   2438             .await
   2439             .expect_err("existing destination is rejected");
   2440         assert_eq!(collision_error.kind(), ServiceSqliteErrorKind::Backup);
   2441         assert_eq!(
   2442             fs::read(collision.join("foreign")).expect("preserved collision"),
   2443             b"preserve"
   2444         );
   2445 
   2446         let source_database_before = fs::read(paths.state_database()).expect("source bytes");
   2447         let source_inventory_before =
   2448             fs::read_dir(paths.state_database().parent().expect("source directory"))
   2449                 .expect("source inventory")
   2450                 .map(|entry| entry.expect("source entry").file_name())
   2451                 .collect::<std::collections::BTreeSet<_>>();
   2452         let stage = output.join("successful");
   2453         let created_at =
   2454             crate::BackupCreatedAtUnixMs::new(1_700_000_000_101).expect("backup creation time");
   2455         let manifest = host
   2456             .capture_online_backup(&stage, created_at)
   2457             .await
   2458             .expect("online backup");
   2459 
   2460         let verified = crate::verify_backup_bundle(
   2461             manifest.canonical_bytes(),
   2462             manifest.digest(),
   2463             &stage,
   2464             &identity,
   2465             std::num::NonZeroU64::new(manifest.members()[0].byte_length())
   2466                 .expect("positive captured member length"),
   2467         )
   2468         .expect("independently verify captured bundle");
   2469         assert_eq!(verified.manifest(), &manifest);
   2470         assert_eq!(verified.database_metadata().service(), identity.service());
   2471         assert_eq!(verified.database_metadata().instance(), identity.instance());
   2472 
   2473         assert_eq!(manifest.service(), paths.service());
   2474         assert_eq!(manifest.instance(), paths.instance());
   2475         assert_eq!(manifest.source_generation(), identity.source_generation());
   2476         assert_eq!(
   2477             manifest.state_schema_version(),
   2478             identity.supported_state_schema_version()
   2479         );
   2480         assert_eq!(manifest.created_at_unix_ms(), created_at);
   2481         assert_eq!(manifest.members().len(), 1);
   2482         assert!(!manifest.protected_material_included());
   2483         assert_eq!(manifest.integrity().sqlite(), "ok");
   2484         assert_eq!(manifest.integrity().foreign_keys(), "ok");
   2485 
   2486         let state = stage.join("state.sqlite");
   2487         let bytes = fs::read(&state).expect("captured state bytes");
   2488         let digest: [u8; 32] = sha2::Sha256::digest(&bytes).into();
   2489         assert_eq!(manifest.members()[0].byte_length(), bytes.len() as u64);
   2490         assert_eq!(manifest.members()[0].sha256().as_bytes(), &digest);
   2491         assert_eq!(
   2492             fs::metadata(&stage)
   2493                 .expect("stage metadata")
   2494                 .permissions()
   2495                 .mode()
   2496                 & 0o777,
   2497             0o700
   2498         );
   2499         assert_eq!(
   2500             fs::metadata(&state)
   2501                 .expect("state metadata")
   2502                 .permissions()
   2503                 .mode()
   2504                 & 0o777,
   2505             0o600
   2506         );
   2507         assert_eq!(
   2508             fs::read_dir(&stage)
   2509                 .expect("stage inventory")
   2510                 .map(|entry| entry.expect("stage entry").file_name())
   2511                 .collect::<Vec<_>>(),
   2512             vec![std::ffi::OsString::from("state.sqlite")]
   2513         );
   2514 
   2515         let mut backup = SqliteConnection::connect_with(
   2516             &SqliteConnectOptions::new().filename(&state).read_only(true),
   2517         )
   2518         .await
   2519         .expect("open captured database");
   2520         assert_eq!(
   2521             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   2522                 .fetch_one(&mut backup)
   2523                 .await
   2524                 .expect("captured row count"),
   2525             2
   2526         );
   2527         backup.close().await.expect("close captured database");
   2528 
   2529         assert_eq!(
   2530             fs::read(paths.state_database()).expect("source bytes after capture"),
   2531             source_database_before
   2532         );
   2533         assert_eq!(
   2534             fs::read_dir(paths.state_database().parent().expect("source directory"))
   2535                 .expect("source inventory after")
   2536                 .map(|entry| entry.expect("source entry").file_name())
   2537                 .collect::<std::collections::BTreeSet<_>>(),
   2538             source_inventory_before
   2539         );
   2540 
   2541         host.close().await.expect("close writable host");
   2542         let inspection = ServiceSqliteHost::open_read_only_inspection(
   2543             &paths,
   2544             &identity,
   2545             &migrations,
   2546             &schema,
   2547             ServiceSqliteConnectionOptions::reviewed(),
   2548         )
   2549         .await
   2550         .expect("open inspection");
   2551         let forbidden = output.join("read-only-forbidden");
   2552         let error = inspection
   2553             .capture_online_backup(&forbidden, created_at)
   2554             .await
   2555             .expect_err("read-only capture is unavailable");
   2556         assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
   2557         assert!(!forbidden.exists());
   2558         inspection.close().await.expect("close inspection");
   2559     }
   2560 
   2561     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2562     #[tokio::test]
   2563     async fn cancelled_capture_cleans_up_close_drains_and_next_capture_recovers() {
   2564         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2565         crate::backup::test_capture_reset();
   2566         let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2567         let host = Arc::new(host);
   2568         let output = root.path().join("cancel-output");
   2569         fs::create_dir(&output).expect("backup parent");
   2570         fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2571             .expect("backup parent mode");
   2572         let cancelled_stage = output.join("cancelled");
   2573         crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED);
   2574         let capture = tokio::spawn({
   2575             let host = Arc::clone(&host);
   2576             let stage = cancelled_stage.clone();
   2577             async move {
   2578                 host.capture_online_backup(
   2579                     &stage,
   2580                     crate::BackupCreatedAtUnixMs::new(1_700_000_000_200)
   2581                         .expect("backup creation time"),
   2582                 )
   2583                 .await
   2584             }
   2585         });
   2586         while crate::backup::test_capture_phase()
   2587             != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED
   2588         {
   2589             tokio::task::yield_now().await;
   2590         }
   2591 
   2592         let concurrent_stage = output.join("concurrent");
   2593         let concurrent = host
   2594             .capture_online_backup(
   2595                 &concurrent_stage,
   2596                 crate::BackupCreatedAtUnixMs::new(1_700_000_000_201).expect("backup creation time"),
   2597             )
   2598             .await
   2599             .expect_err("second capture is rejected");
   2600         assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Backup);
   2601         assert!(!concurrent_stage.exists());
   2602 
   2603         capture.abort();
   2604         assert!(
   2605             capture
   2606                 .await
   2607                 .expect_err("capture task is cancelled")
   2608                 .is_cancelled()
   2609         );
   2610         while host.backup_active.load(Ordering::Acquire) {
   2611             tokio::task::yield_now().await;
   2612         }
   2613         crate::backup::test_capture_reset();
   2614         assert!(!cancelled_stage.exists());
   2615 
   2616         let recovered_stage = output.join("recovered");
   2617         host.capture_online_backup(
   2618             &recovered_stage,
   2619             crate::BackupCreatedAtUnixMs::new(1_700_000_000_202).expect("backup creation time"),
   2620         )
   2621         .await
   2622         .expect("capture recovers after cancellation");
   2623         assert!(recovered_stage.join("state.sqlite").exists());
   2624 
   2625         crate::backup::test_capture_reset();
   2626         crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED);
   2627         let close_drained_stage = output.join("close-drained");
   2628         let capture = tokio::spawn({
   2629             let host = Arc::clone(&host);
   2630             let stage = close_drained_stage.clone();
   2631             async move {
   2632                 host.capture_online_backup(
   2633                     &stage,
   2634                     crate::BackupCreatedAtUnixMs::new(1_700_000_000_203)
   2635                         .expect("backup creation time"),
   2636                 )
   2637                 .await
   2638             }
   2639         });
   2640         while crate::backup::test_capture_phase()
   2641             != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED
   2642         {
   2643             tokio::task::yield_now().await;
   2644         }
   2645         capture.abort();
   2646         assert!(
   2647             capture
   2648                 .await
   2649                 .expect_err("capture task is cancelled")
   2650                 .is_cancelled()
   2651         );
   2652         host.close()
   2653             .await
   2654             .expect("close drains every admitted capture");
   2655         crate::backup::test_capture_reset();
   2656         assert!(!close_drained_stage.exists());
   2657         assert!(!host.backup_active.load(Ordering::Acquire));
   2658     }
   2659 
   2660     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2661     #[tokio::test]
   2662     async fn capture_cancellation_is_cleanup_safe_at_every_governed_phase() {
   2663         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2664 
   2665         for (index, phase) in [
   2666             crate::backup::TEST_CAPTURE_PHASE_BEFORE_CREATE,
   2667             crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED,
   2668             crate::backup::TEST_CAPTURE_PHASE_POST_COPY,
   2669             crate::backup::TEST_CAPTURE_PHASE_PRE_FINAL_SYNC,
   2670         ]
   2671         .into_iter()
   2672         .enumerate()
   2673         {
   2674             crate::backup::test_capture_reset();
   2675             let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2676             let host = Arc::new(host);
   2677             if phase == crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED {
   2678                 host.transaction(|transaction| {
   2679                     Box::pin(async move {
   2680                         sqlx::raw_sql(
   2681                             "WITH RECURSIVE sequence(value) AS (
   2682                                  VALUES(1)
   2683                                  UNION ALL
   2684                                  SELECT value + 1 FROM sequence WHERE value < 100000
   2685                              )
   2686                              INSERT INTO host_probe(value) SELECT 0 FROM sequence",
   2687                         )
   2688                         .execute(&mut *transaction)
   2689                         .await
   2690                         .map(|_| ())
   2691                     })
   2692                 })
   2693                 .await
   2694                 .expect("seed multi-batch cancellation fixture");
   2695             }
   2696 
   2697             let output = root.path().join(format!("phase-cancel-output-{index}"));
   2698             fs::create_dir(&output).expect("backup parent");
   2699             fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2700                 .expect("backup parent mode");
   2701             let stage = output.join("cancelled");
   2702             crate::backup::test_capture_block_phase(phase);
   2703             let capture = tokio::spawn({
   2704                 let host = Arc::clone(&host);
   2705                 let stage = stage.clone();
   2706                 async move {
   2707                     host.capture_online_backup(
   2708                         &stage,
   2709                         crate::BackupCreatedAtUnixMs::new(1_700_000_000_220 + index as u64)
   2710                             .expect("backup creation time"),
   2711                     )
   2712                     .await
   2713                 }
   2714             });
   2715             while crate::backup::test_capture_phase() != phase {
   2716                 tokio::task::yield_now().await;
   2717             }
   2718 
   2719             capture.abort();
   2720             assert!(
   2721                 capture
   2722                     .await
   2723                     .expect_err("capture task is cancelled")
   2724                     .is_cancelled()
   2725             );
   2726             host.close()
   2727                 .await
   2728                 .expect("close drains phase-cancelled capture cleanup");
   2729             crate::backup::test_capture_reset();
   2730 
   2731             assert!(!stage.exists(), "phase {phase} must leave no staging tree");
   2732             assert!(!host.backup_active.load(Ordering::Acquire));
   2733             let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   2734                 .expect("writer authority can be reacquired")
   2735                 .expect("writable mode returns authority");
   2736             authority.release().expect("release reacquired authority");
   2737         }
   2738     }
   2739 
   2740     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2741     #[tokio::test]
   2742     async fn online_backup_remains_consistent_with_a_concurrent_wal_writer() {
   2743         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2744         crate::backup::test_capture_reset();
   2745         let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2746         let host = Arc::new(host);
   2747         host.transaction(|transaction| {
   2748             Box::pin(async move {
   2749                 sqlx::raw_sql(
   2750                     "WITH RECURSIVE sequence(value) AS (
   2751                          VALUES(1) UNION ALL SELECT value + 1 FROM sequence WHERE value < 100000
   2752                      )
   2753                      INSERT INTO host_probe(value) SELECT 0 FROM sequence",
   2754                 )
   2755                 .execute(&mut *transaction)
   2756                 .await
   2757                 .map(|_| ())
   2758             })
   2759         })
   2760         .await
   2761         .expect("seed a multi-batch database");
   2762 
   2763         let output = root.path().join("concurrent-output");
   2764         fs::create_dir(&output).expect("backup parent");
   2765         fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2766             .expect("backup parent mode");
   2767         let stage = output.join("consistent");
   2768         crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED);
   2769         let capture = tokio::spawn({
   2770             let host = Arc::clone(&host);
   2771             let stage = stage.clone();
   2772             async move {
   2773                 host.capture_online_backup(
   2774                     &stage,
   2775                     crate::BackupCreatedAtUnixMs::new(1_700_000_000_300)
   2776                         .expect("backup creation time"),
   2777                 )
   2778                 .await
   2779             }
   2780         });
   2781         while crate::backup::test_capture_phase()
   2782             != crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED
   2783         {
   2784             tokio::task::yield_now().await;
   2785         }
   2786         host.transaction(|transaction| {
   2787             Box::pin(async move {
   2788                 sqlx::query("UPDATE host_probe SET value = 1 WHERE rowid IN (1, 2)")
   2789                     .execute(&mut *transaction)
   2790                     .await
   2791                     .map(|_| ())
   2792             })
   2793         })
   2794         .await
   2795         .expect("commit concurrent WAL transaction");
   2796         crate::backup::test_capture_block_phase(0);
   2797         capture
   2798             .await
   2799             .expect("capture joins")
   2800             .expect("capture remains consistent");
   2801         crate::backup::test_capture_reset();
   2802 
   2803         let mut backup = SqliteConnection::connect_with(
   2804             &SqliteConnectOptions::new()
   2805                 .filename(stage.join("state.sqlite"))
   2806                 .read_only(true),
   2807         )
   2808         .await
   2809         .expect("open backup");
   2810         let updated =
   2811             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe WHERE value = 1")
   2812                 .fetch_one(&mut backup)
   2813                 .await
   2814                 .expect("count transaction projection");
   2815         assert!(
   2816             updated == 0 || updated == 2,
   2817             "backup must not tear a transaction"
   2818         );
   2819         assert_eq!(
   2820             sqlx::query_scalar::<_, String>("PRAGMA integrity_check(1)")
   2821                 .fetch_one(&mut backup)
   2822                 .await
   2823                 .expect("backup integrity"),
   2824             "ok"
   2825         );
   2826         backup.close().await.expect("close backup");
   2827         host.close().await.expect("close writer host");
   2828     }
   2829 
   2830     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2831     #[tokio::test]
   2832     async fn backup_storage_full_sync_failures_cleanup_and_leave_host_recoverable() {
   2833         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2834         crate::backup::test_capture_reset();
   2835         let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2836         let output = root.path().join("sync-failure-output");
   2837         fs::create_dir(&output).expect("backup parent");
   2838         fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2839             .expect("backup parent mode");
   2840 
   2841         for (index, failure) in [
   2842             crate::backup::TestCaptureSyncFailure::State,
   2843             crate::backup::TestCaptureSyncFailure::Staging,
   2844             crate::backup::TestCaptureSyncFailure::FinalParent,
   2845         ]
   2846         .into_iter()
   2847         .enumerate()
   2848         {
   2849             let stage = output.join(format!("failure-{index}"));
   2850             let error = crate::backup::test_capture_online_backup_with_sync_failure(
   2851                 &host.pool,
   2852                 &host.closing,
   2853                 &host.backup_active,
   2854                 &stage,
   2855                 crate::BackupCreatedAtUnixMs::new(1_700_000_000_400 + index as u64)
   2856                     .expect("backup creation time"),
   2857                 failure,
   2858             )
   2859             .await
   2860             .expect_err("injected synchronization failure must reject capture");
   2861             assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup);
   2862             let storage = error
   2863                 .source()
   2864                 .and_then(Error::source)
   2865                 .and_then(|source| source.downcast_ref::<std::io::Error>())
   2866                 .expect("storage-full cause");
   2867             assert_eq!(storage.kind(), std::io::ErrorKind::StorageFull);
   2868             assert!(!stage.exists());
   2869             assert!(!host.backup_active.load(Ordering::Acquire));
   2870             assert_eq!(row_count(&host).await, 0);
   2871         }
   2872 
   2873         let recovered = output.join("recovered");
   2874         host.capture_online_backup(
   2875             &recovered,
   2876             crate::BackupCreatedAtUnixMs::new(1_700_000_000_410).expect("backup creation time"),
   2877         )
   2878         .await
   2879         .expect("host remains usable after sync failures");
   2880         assert!(recovered.join("state.sqlite").exists());
   2881         host.close().await.expect("close host");
   2882     }
   2883 
   2884     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2885     #[tokio::test]
   2886     async fn backup_await_boundaries_preserve_authority_precedence() {
   2887         let _serial = CAPTURE_TEST_LOCK.lock().await;
   2888         crate::backup::test_capture_reset();
   2889         let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   2890         let host = Arc::new(host);
   2891         let output = root.path().join("precedence-output");
   2892         fs::create_dir(&output).expect("backup parent");
   2893         fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
   2894             .expect("backup parent mode");
   2895 
   2896         crate::backup::test_capture_inject_metadata_failure(true);
   2897         crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED);
   2898         let metadata_stage = output.join("metadata");
   2899         let capture = tokio::spawn({
   2900             let host = Arc::clone(&host);
   2901             let stage = metadata_stage.clone();
   2902             async move {
   2903                 host.capture_online_backup(
   2904                     &stage,
   2905                     crate::BackupCreatedAtUnixMs::new(1_700_000_000_420)
   2906                         .expect("backup creation time"),
   2907                 )
   2908                 .await
   2909             }
   2910         });
   2911         while crate::backup::test_capture_phase()
   2912             != crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED
   2913         {
   2914             tokio::task::yield_now().await;
   2915         }
   2916         let retired_lock = paths
   2917             .state_lock()
   2918             .parent()
   2919             .expect("state directory")
   2920             .join("retired-backup-precedence.lock");
   2921         fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
   2922         crate::backup::test_capture_block_phase(0);
   2923         let metadata_error = capture
   2924             .await
   2925             .expect("metadata capture joins")
   2926             .expect_err("authority overrides metadata failure");
   2927         assert_eq!(metadata_error.kind(), ServiceSqliteErrorKind::Authority);
   2928         fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
   2929         crate::backup::test_capture_reset();
   2930         assert!(!metadata_stage.exists());
   2931         assert_eq!(row_count(&host).await, 0);
   2932 
   2933         crate::backup::test_capture_panic_worker(true);
   2934         crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED);
   2935         let join_stage = output.join("join");
   2936         let capture = tokio::spawn({
   2937             let host = Arc::clone(&host);
   2938             let stage = join_stage.clone();
   2939             async move {
   2940                 host.capture_online_backup(
   2941                     &stage,
   2942                     crate::BackupCreatedAtUnixMs::new(1_700_000_000_421)
   2943                         .expect("backup creation time"),
   2944                 )
   2945                 .await
   2946             }
   2947         });
   2948         while crate::backup::test_capture_phase() != crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED
   2949         {
   2950             tokio::task::yield_now().await;
   2951         }
   2952         fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock again");
   2953         crate::backup::test_capture_block_phase(0);
   2954         let join_error = capture
   2955             .await
   2956             .expect("join-failure capture joins")
   2957             .expect_err("authority overrides worker join failure");
   2958         assert_eq!(join_error.kind(), ServiceSqliteErrorKind::Authority);
   2959         fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock again");
   2960         crate::backup::test_capture_reset();
   2961         assert!(!join_stage.exists());
   2962         assert_eq!(row_count(&host).await, 0);
   2963         host.close().await.expect("close host");
   2964     }
   2965 
   2966     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2967     #[derive(Debug, PartialEq, Eq)]
   2968     struct StateFileSnapshot {
   2969         bytes: Vec<u8>,
   2970         length: u64,
   2971         modified: SystemTime,
   2972         mode: u32,
   2973     }
   2974 
   2975     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2976     fn state_directory_snapshot(directory: &Path) -> BTreeMap<String, StateFileSnapshot> {
   2977         fs::read_dir(directory)
   2978             .expect("read state directory")
   2979             .map(|entry| {
   2980                 let entry = entry.expect("state entry");
   2981                 let name = entry.file_name().into_string().expect("UTF-8 state entry");
   2982                 let metadata = entry.metadata().expect("state entry metadata");
   2983                 (
   2984                     name,
   2985                     StateFileSnapshot {
   2986                         bytes: fs::read(entry.path()).expect("state entry bytes"),
   2987                         length: metadata.len(),
   2988                         modified: metadata.modified().expect("state entry mtime"),
   2989                         mode: metadata.permissions().mode() & 0o777,
   2990                     },
   2991                 )
   2992             })
   2993             .collect()
   2994     }
   2995 
   2996     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2997     #[tokio::test]
   2998     async fn close_drains_admitted_work_rejects_new_work_and_is_idempotent() {
   2999         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3000         let host = Arc::new(host);
   3001         let entered = Arc::new(Notify::new());
   3002         let release = Arc::new(Notify::new());
   3003         let transaction = tokio::spawn({
   3004             let host = Arc::clone(&host);
   3005             let entered = Arc::clone(&entered);
   3006             let release = Arc::clone(&release);
   3007             async move {
   3008                 host.transaction::<i64, Infallible, _>(|transaction| {
   3009                     Box::pin(async move {
   3010                         let count = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   3011                             .fetch_one(&mut *transaction)
   3012                             .await
   3013                             .expect("read in admitted transaction");
   3014                         entered.notify_one();
   3015                         release.notified().await;
   3016                         Ok(count)
   3017                     })
   3018                 })
   3019                 .await
   3020             }
   3021         });
   3022         entered.notified().await;
   3023 
   3024         let first_close = tokio::spawn({
   3025             let host = Arc::clone(&host);
   3026             async move { host.close().await }
   3027         });
   3028         while !host.closing.load(Ordering::Acquire) {
   3029             tokio::task::yield_now().await;
   3030         }
   3031         let second_close = tokio::spawn({
   3032             let host = Arc::clone(&host);
   3033             async move { host.close().await }
   3034         });
   3035         let rejected = host
   3036             .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) }))
   3037             .await
   3038             .expect_err("close admission must reject new work");
   3039         assert_eq!(
   3040             rejected.kind(),
   3041             ServiceSqliteTransactionErrorKind::NotCommitted
   3042         );
   3043         assert_eq!(
   3044             rejected.sqlite_error().map(ServiceSqliteError::kind),
   3045             Some(ServiceSqliteErrorKind::Open)
   3046         );
   3047         let contended = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting);
   3048         assert!(matches!(
   3049             contended,
   3050             Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority
   3051         ));
   3052 
   3053         release.notify_one();
   3054         assert_eq!(
   3055             transaction
   3056                 .await
   3057                 .expect("transaction task joins")
   3058                 .expect("admitted transaction"),
   3059             0
   3060         );
   3061         first_close
   3062             .await
   3063             .expect("first close task joins")
   3064             .expect("first close succeeds");
   3065         second_close
   3066             .await
   3067             .expect("second close task joins")
   3068             .expect("concurrent close is idempotent");
   3069         host.close().await.expect("sequential close is idempotent");
   3070 
   3071         let mut next = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3072             .expect("authority can be reacquired")
   3073             .expect("writer mode yields authority");
   3074         next.release().expect("release reacquired authority");
   3075     }
   3076 
   3077     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3078     #[tokio::test]
   3079     async fn cancelled_close_retains_authority_and_retry_finishes() {
   3080         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3081         let host = Arc::new(host);
   3082         let entered = Arc::new(Notify::new());
   3083         let release = Arc::new(Notify::new());
   3084         let transaction = tokio::spawn({
   3085             let host = Arc::clone(&host);
   3086             let entered = Arc::clone(&entered);
   3087             let release = Arc::clone(&release);
   3088             async move {
   3089                 host.transaction::<(), Infallible, _>(|transaction| {
   3090                     Box::pin(async move {
   3091                         sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   3092                             .fetch_one(&mut *transaction)
   3093                             .await
   3094                             .expect("admitted read");
   3095                         entered.notify_one();
   3096                         release.notified().await;
   3097                         Ok(())
   3098                     })
   3099                 })
   3100                 .await
   3101             }
   3102         });
   3103         entered.notified().await;
   3104         let close_task = tokio::spawn({
   3105             let host = Arc::clone(&host);
   3106             async move { host.close().await }
   3107         });
   3108         while !host.closing.load(Ordering::Acquire) {
   3109             tokio::task::yield_now().await;
   3110         }
   3111         close_task.abort();
   3112         assert!(
   3113             close_task
   3114                 .await
   3115                 .expect_err("close task is cancelled")
   3116                 .is_cancelled()
   3117         );
   3118         let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting);
   3119         assert!(matches!(
   3120             retained,
   3121             Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority
   3122         ));
   3123         let rejected = host
   3124             .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) }))
   3125             .await
   3126             .expect_err("cancelled close remains non-admitting");
   3127         assert_eq!(
   3128             rejected.kind(),
   3129             ServiceSqliteTransactionErrorKind::NotCommitted
   3130         );
   3131 
   3132         release.notify_one();
   3133         transaction
   3134             .await
   3135             .expect("transaction task joins")
   3136             .expect("admitted transaction finishes");
   3137         host.close().await.expect("close retry succeeds");
   3138         assert!(
   3139             WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3140                 .expect("authority reacquisition after retry")
   3141                 .is_some()
   3142         );
   3143     }
   3144 
   3145     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3146     #[tokio::test]
   3147     async fn writable_close_checkpoints_and_read_only_close_is_side_effect_free() {
   3148         let (_root, paths, identity, migrations, schema, host) = initialized_host().await;
   3149         host.transaction(|transaction| {
   3150             Box::pin(async move {
   3151                 sqlx::query("INSERT INTO host_probe (value) VALUES (41)")
   3152                     .execute(&mut *transaction)
   3153                     .await
   3154                     .map(|_| ())
   3155             })
   3156         })
   3157         .await
   3158         .expect("write WAL frame");
   3159         host.close().await.expect("writable close checkpoints");
   3160         let state_directory = paths.state_database().parent().expect("state directory");
   3161         assert!(!state_directory.join("state.sqlite-wal").exists());
   3162         assert!(!state_directory.join("state.sqlite-shm").exists());
   3163         let before = state_directory_snapshot(state_directory);
   3164 
   3165         let inspection = ServiceSqliteHost::open_read_only_inspection(
   3166             &paths,
   3167             &identity,
   3168             &migrations,
   3169             &schema,
   3170             ServiceSqliteConnectionOptions::reviewed(),
   3171         )
   3172         .await
   3173         .expect("open read-only inspection");
   3174         assert_eq!(row_count(&inspection).await, 1);
   3175         inspection
   3176             .close()
   3177             .await
   3178             .expect("close read-only inspection");
   3179         assert_eq!(state_directory_snapshot(state_directory), before);
   3180 
   3181         let (reopened, outcome) = ServiceSqliteHost::open_read_write_existing(
   3182             &paths,
   3183             &identity,
   3184             &migrations,
   3185             &schema,
   3186             ServiceSqliteConnectionOptions::reviewed(),
   3187             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   3188             &build_identity(),
   3189             &[],
   3190         )
   3191         .await
   3192         .expect("reopen writer after read-only close");
   3193         assert_eq!(outcome.applied_count(), 0);
   3194         assert_eq!(row_count(&reopened).await, 1);
   3195         reopened.close().await.expect("close reopened writer");
   3196     }
   3197 
   3198     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3199     #[tokio::test]
   3200     async fn read_only_close_revalidates_after_drain_before_releasing_stale_authority() {
   3201         let (_root, paths, identity, migrations, schema, writer) = initialized_host().await;
   3202         writer.close().await.expect("close writer host");
   3203         let inspection = Arc::new(
   3204             ServiceSqliteHost::open_read_only_inspection(
   3205                 &paths,
   3206                 &identity,
   3207                 &migrations,
   3208                 &schema,
   3209                 ServiceSqliteConnectionOptions::reviewed(),
   3210             )
   3211             .await
   3212             .expect("open read-only inspection"),
   3213         );
   3214         let entered = Arc::new(Notify::new());
   3215         let release = Arc::new(Notify::new());
   3216         let transaction = tokio::spawn({
   3217             let inspection = Arc::clone(&inspection);
   3218             let entered = Arc::clone(&entered);
   3219             let release = Arc::clone(&release);
   3220             async move {
   3221                 inspection
   3222                     .transaction::<(), Infallible, _>(|transaction| {
   3223                         Box::pin(async move {
   3224                             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   3225                                 .fetch_one(&mut *transaction)
   3226                                 .await
   3227                                 .expect("read through retained inspection");
   3228                             entered.notify_one();
   3229                             release.notified().await;
   3230                             Ok(())
   3231                         })
   3232                     })
   3233                     .await
   3234             }
   3235         });
   3236         entered.notified().await;
   3237         let close_task = tokio::spawn({
   3238             let inspection = Arc::clone(&inspection);
   3239             async move { inspection.close().await }
   3240         });
   3241         while !inspection.closing.load(Ordering::Acquire) {
   3242             tokio::task::yield_now().await;
   3243         }
   3244 
   3245         let retired_lock = paths
   3246             .state_lock()
   3247             .parent()
   3248             .expect("state directory")
   3249             .join("retired-inspection-close.lock");
   3250         fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock");
   3251         let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3252             .expect("replacement acquisition")
   3253             .expect("replacement authority");
   3254         release.notify_one();
   3255         let transaction_error = transaction
   3256             .await
   3257             .expect("inspection transaction task joins")
   3258             .expect_err("binding drift revokes admitted inspection");
   3259         assert_eq!(
   3260             transaction_error.kind(),
   3261             ServiceSqliteTransactionErrorKind::NotCommitted
   3262         );
   3263         let close_error = close_task
   3264             .await
   3265             .expect("inspection close task joins")
   3266             .expect_err("close must report stale inspection authority");
   3267         assert_eq!(close_error.kind(), ServiceSqliteErrorKind::Authority);
   3268         assert_eq!(
   3269             inspection
   3270                 .close()
   3271                 .await
   3272                 .expect_err("terminal authority result is cached")
   3273                 .kind(),
   3274             ServiceSqliteErrorKind::Authority
   3275         );
   3276         replacement
   3277             .release()
   3278             .expect("release replacement authority");
   3279     }
   3280 
   3281     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3282     #[tokio::test]
   3283     async fn cancelled_checkpoint_resumes_to_terminal_error_and_releases_authority() {
   3284         let options = ServiceSqliteConnectionOptions::new(Duration::from_secs(1), 1)
   3285             .expect("short reviewed limits");
   3286         let (_root, paths, _identity, _migrations, _schema, host) =
   3287             initialized_host_with_options(options).await;
   3288         let host = Arc::new(host);
   3289         host.transaction(|transaction| {
   3290             Box::pin(async move {
   3291                 sqlx::query("INSERT INTO host_probe (value) VALUES (1)")
   3292                     .execute(&mut *transaction)
   3293                     .await
   3294                     .map(|_| ())
   3295             })
   3296         })
   3297         .await
   3298         .expect("seed reader snapshot");
   3299 
   3300         let mut reader = SqliteConnection::connect_with(
   3301             &SqliteConnectOptions::new()
   3302                 .filename(paths.state_database())
   3303                 .read_only(true)
   3304                 .create_if_missing(false),
   3305         )
   3306         .await
   3307         .expect("open external reader");
   3308         let mut reader_transaction = reader.begin().await.expect("begin external read");
   3309         assert_eq!(
   3310             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   3311                 .fetch_one(&mut *reader_transaction)
   3312                 .await
   3313                 .expect("establish reader snapshot"),
   3314             1
   3315         );
   3316         host.transaction(|transaction| {
   3317             Box::pin(async move {
   3318                 sqlx::query("INSERT INTO host_probe (value) VALUES (2)")
   3319                     .execute(&mut *transaction)
   3320                     .await
   3321                     .map(|_| ())
   3322             })
   3323         })
   3324         .await
   3325         .expect("append frame after reader snapshot");
   3326 
   3327         let close_task = tokio::spawn({
   3328             let host = Arc::clone(&host);
   3329             async move { host.close().await }
   3330         });
   3331         while host.pool.close_phase() != crate::open::TEST_CLOSE_PHASE_CHECKPOINT {
   3332             tokio::task::yield_now().await;
   3333         }
   3334         close_task.abort();
   3335         assert!(
   3336             close_task
   3337                 .await
   3338                 .expect_err("checkpoint close task is cancelled")
   3339                 .is_cancelled()
   3340         );
   3341         let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting);
   3342         assert!(matches!(
   3343             retained,
   3344             Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority
   3345         ));
   3346 
   3347         let first = host
   3348             .close()
   3349             .await
   3350             .expect_err("active reader prevents TRUNCATE");
   3351         assert_eq!(first.kind(), ServiceSqliteErrorKind::Pragma);
   3352         assert!(!first.to_string().contains("state.sqlite"));
   3353         let repeated = host.close().await.expect_err("terminal result is cached");
   3354         assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Pragma);
   3355         let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3356             .expect("close releases authority despite checkpoint failure")
   3357             .expect("writer mode yields authority");
   3358         replacement
   3359             .release()
   3360             .expect("release replacement authority");
   3361 
   3362         reader_transaction
   3363             .rollback()
   3364             .await
   3365             .expect("release external snapshot");
   3366         reader.close().await.expect("close external reader");
   3367     }
   3368 
   3369     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3370     #[tokio::test]
   3371     async fn close_authority_drift_is_cached_and_releases_the_stale_lock() {
   3372         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3373         let retired_lock = paths
   3374             .state_lock()
   3375             .parent()
   3376             .expect("state directory")
   3377             .join("retired-close-state.lock");
   3378         fs::rename(paths.state_lock(), retired_lock).expect("replace canonical lock");
   3379         let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3380             .expect("replacement acquisition")
   3381             .expect("replacement authority");
   3382 
   3383         let first = host
   3384             .close()
   3385             .await
   3386             .expect_err("close detects authority drift");
   3387         assert_eq!(first.kind(), ServiceSqliteErrorKind::Authority);
   3388         let repeated = host.close().await.expect_err("authority result is cached");
   3389         assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority);
   3390         assert!(!format!("{host:?}").contains(paths.state_database().to_string_lossy().as_ref()));
   3391         replacement
   3392             .release()
   3393             .expect("release replacement authority");
   3394     }
   3395 
   3396     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3397     #[tokio::test]
   3398     async fn typed_execution_commits_and_operation_failure_rolls_back() {
   3399         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3400         assert_eq!(host.mode(), OpenMode::Initialize);
   3401 
   3402         let inserted = host
   3403             .transaction(|transaction| {
   3404                 Box::pin(async move {
   3405                     sqlx::query("INSERT INTO host_probe (value) VALUES (?)")
   3406                         .bind(41_i64)
   3407                         .execute(&mut *transaction)
   3408                         .await?;
   3409                     sqlx::query_scalar::<_, i64>("SELECT value FROM host_probe")
   3410                         .fetch_one(&mut *transaction)
   3411                         .await
   3412                 })
   3413             })
   3414             .await
   3415             .expect("commit typed operation");
   3416         assert_eq!(inserted, 41);
   3417 
   3418         let error = host
   3419             .transaction(|transaction| {
   3420                 Box::pin(async move {
   3421                     sqlx::query("INSERT INTO host_probe (value) VALUES (99)")
   3422                         .execute(&mut *transaction)
   3423                         .await
   3424                         .map_err(|_| "query-failure-secret")?;
   3425                     Err::<(), _>("operation-secret")
   3426                 })
   3427             })
   3428             .await
   3429             .expect_err("operation rejection must roll back");
   3430         assert_eq!(
   3431             error.kind(),
   3432             ServiceSqliteTransactionErrorKind::OperationRolledBack
   3433         );
   3434         assert_eq!(error.operation_error(), Some(&"operation-secret"));
   3435         assert!(error.sqlite_error().is_none());
   3436         assert!(!format!("{error:?}").contains("operation-secret"));
   3437         assert!(!error.to_string().contains("operation-secret"));
   3438         assert_eq!(row_count(&host).await, 1);
   3439     }
   3440 
   3441     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3442     #[tokio::test]
   3443     async fn transaction_durability_edges_preserve_exact_commit_semantics() {
   3444         use crate::failpoint::DurabilityFailpoint;
   3445 
   3446         for (point, expected_reached) in [
   3447             (
   3448                 DurabilityFailpoint::TransactionBeforeBegin,
   3449                 &[DurabilityFailpoint::TransactionBeforeBegin][..],
   3450             ),
   3451             (
   3452                 DurabilityFailpoint::TransactionAfterBegin,
   3453                 &[
   3454                     DurabilityFailpoint::TransactionBeforeBegin,
   3455                     DurabilityFailpoint::TransactionAfterBegin,
   3456                 ][..],
   3457             ),
   3458             (
   3459                 DurabilityFailpoint::TransactionBeforeCommit,
   3460                 &[
   3461                     DurabilityFailpoint::TransactionBeforeBegin,
   3462                     DurabilityFailpoint::TransactionAfterBegin,
   3463                     DurabilityFailpoint::TransactionBeforeCommit,
   3464                 ][..],
   3465             ),
   3466             (
   3467                 DurabilityFailpoint::TransactionAfterCommit,
   3468                 &[
   3469                     DurabilityFailpoint::TransactionBeforeBegin,
   3470                     DurabilityFailpoint::TransactionAfterBegin,
   3471                     DurabilityFailpoint::TransactionBeforeCommit,
   3472                     DurabilityFailpoint::TransactionAfterCommit,
   3473                 ][..],
   3474             ),
   3475         ] {
   3476             let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3477             host.arm_durability_failpoint(point);
   3478             let error = host
   3479                 .transaction(|transaction| {
   3480                     Box::pin(async move {
   3481                         sqlx::query("INSERT INTO host_probe (value) VALUES (72)")
   3482                             .execute(&mut *transaction)
   3483                             .await
   3484                             .map(|_| ())
   3485                     })
   3486                 })
   3487                 .await
   3488                 .expect_err("injected transaction edge");
   3489             let reached = host.failpoints.reached();
   3490             if point == DurabilityFailpoint::TransactionAfterCommit {
   3491                 assert_eq!(
   3492                     error.kind(),
   3493                     ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown
   3494                 );
   3495                 assert_eq!(row_count(&host).await, 1);
   3496             } else {
   3497                 assert_eq!(
   3498                     error.kind(),
   3499                     ServiceSqliteTransactionErrorKind::NotCommitted
   3500                 );
   3501                 assert_eq!(row_count(&host).await, 0);
   3502             }
   3503             assert!(host.failpoints.fired());
   3504             assert_eq!(reached, expected_reached);
   3505             host.close()
   3506                 .await
   3507                 .expect("close host after transaction edge");
   3508         }
   3509 
   3510         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3511         host.arm_durability_failpoint(DurabilityFailpoint::TransactionBeforeCommit);
   3512         let error = host
   3513             .transaction(|transaction| {
   3514                 Box::pin(async move {
   3515                     let _ = sqlx::raw_sql(
   3516                         "PRAGMA trusted_schema=ON; INSERT INTO host_probe (value) VALUES (73)",
   3517                     )
   3518                     .execute(&mut *transaction)
   3519                     .await;
   3520                     Ok::<_, Infallible>(())
   3521                 })
   3522             })
   3523             .await
   3524             .expect_err("precommit policy rejection must precede the commit-edge hook");
   3525         assert_eq!(
   3526             error.kind(),
   3527             ServiceSqliteTransactionErrorKind::NotCommitted
   3528         );
   3529         assert!(!host.failpoints.fired());
   3530         assert_eq!(
   3531             host.failpoints.reached(),
   3532             [
   3533                 DurabilityFailpoint::TransactionBeforeBegin,
   3534                 DurabilityFailpoint::TransactionAfterBegin,
   3535             ]
   3536         );
   3537         host.failpoints.disarm();
   3538         assert_eq!(row_count(&host).await, 0);
   3539         host.close().await.expect("close policy-drift host");
   3540     }
   3541 
   3542     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3543     #[tokio::test]
   3544     async fn backup_durability_edges_fail_once_clean_exact_stage_and_recover() {
   3545         use crate::failpoint::DurabilityFailpoint;
   3546 
   3547         let _serial = CAPTURE_TEST_LOCK.lock().await;
   3548         for (index, (point, expected_phase)) in [
   3549             (DurabilityFailpoint::BackupBeforeCreate, 1),
   3550             (DurabilityFailpoint::BackupAfterCreate, 1),
   3551             (DurabilityFailpoint::BackupBeforeCopy, 2),
   3552             (DurabilityFailpoint::BackupAfterCopy, 3),
   3553             (DurabilityFailpoint::BackupBeforeFileSync, 4),
   3554             (DurabilityFailpoint::BackupAfterFileSync, 4),
   3555             (DurabilityFailpoint::BackupBeforeDirectorySync, 5),
   3556             (DurabilityFailpoint::BackupAfterDirectorySync, 5),
   3557         ]
   3558         .into_iter()
   3559         .enumerate()
   3560         {
   3561             let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3562             let stage = root.path().join(format!("failpoint-backup-{index}"));
   3563             host.arm_durability_failpoint(point);
   3564             let error = host
   3565                 .capture_online_backup(
   3566                     &stage,
   3567                     crate::BackupCreatedAtUnixMs::new(1_700_000_072_000).expect("capture time"),
   3568                 )
   3569                 .await
   3570                 .expect_err("injected backup edge");
   3571             assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup);
   3572             assert!(host.failpoints.fired());
   3573             assert_eq!(host.failpoints.reached().last(), Some(&point));
   3574             assert_eq!(
   3575                 host.failpoints.observation(point),
   3576                 Some(expected_phase),
   3577                 "backup failpoint must fire in its named lifecycle phase"
   3578             );
   3579             assert!(!stage.exists(), "owned failed stage is cleaned");
   3580             let recovery = root.path().join(format!("recovered-backup-{index}"));
   3581             host.capture_online_backup(
   3582                 &recovery,
   3583                 crate::BackupCreatedAtUnixMs::new(1_700_000_072_001).expect("capture time"),
   3584             )
   3585             .await
   3586             .expect("one-shot failpoint permits retry");
   3587             assert!(recovery.join(crate::BACKUP_STATE_MEMBER_NAME).is_file());
   3588             host.close().await.expect("close backup host");
   3589         }
   3590     }
   3591 
   3592     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3593     #[tokio::test]
   3594     async fn close_durability_edges_are_once_only_retryable_or_terminal() {
   3595         use crate::failpoint::DurabilityFailpoint;
   3596 
   3597         let all_close_edges = [
   3598             DurabilityFailpoint::CloseBeforeDrain,
   3599             DurabilityFailpoint::CloseAfterDrain,
   3600             DurabilityFailpoint::CloseBeforeCheckpoint,
   3601             DurabilityFailpoint::CloseAfterCheckpoint,
   3602             DurabilityFailpoint::CloseBeforeConnectionClose,
   3603             DurabilityFailpoint::CloseAfterConnectionClose,
   3604             DurabilityFailpoint::CloseBeforeAuthorityRelease,
   3605             DurabilityFailpoint::CloseAfterAuthorityRelease,
   3606         ];
   3607         for (point, expected, retryable, expected_phase, expected_reached) in [
   3608             (
   3609                 DurabilityFailpoint::CloseBeforeDrain,
   3610                 ServiceSqliteErrorKind::Open,
   3611                 true,
   3612                 0,
   3613                 &all_close_edges[..1],
   3614             ),
   3615             (
   3616                 DurabilityFailpoint::CloseAfterDrain,
   3617                 ServiceSqliteErrorKind::Open,
   3618                 true,
   3619                 0,
   3620                 &all_close_edges[..2],
   3621             ),
   3622             (
   3623                 DurabilityFailpoint::CloseBeforeCheckpoint,
   3624                 ServiceSqliteErrorKind::Pragma,
   3625                 false,
   3626                 6,
   3627                 &[
   3628                     DurabilityFailpoint::CloseBeforeDrain,
   3629                     DurabilityFailpoint::CloseAfterDrain,
   3630                     DurabilityFailpoint::CloseBeforeCheckpoint,
   3631                     DurabilityFailpoint::CloseBeforeConnectionClose,
   3632                     DurabilityFailpoint::CloseAfterConnectionClose,
   3633                     DurabilityFailpoint::CloseBeforeAuthorityRelease,
   3634                     DurabilityFailpoint::CloseAfterAuthorityRelease,
   3635                 ][..],
   3636             ),
   3637             (
   3638                 DurabilityFailpoint::CloseAfterCheckpoint,
   3639                 ServiceSqliteErrorKind::Pragma,
   3640                 false,
   3641                 6,
   3642                 &all_close_edges[..],
   3643             ),
   3644             (
   3645                 DurabilityFailpoint::CloseBeforeConnectionClose,
   3646                 ServiceSqliteErrorKind::Open,
   3647                 false,
   3648                 6,
   3649                 &all_close_edges[..],
   3650             ),
   3651             (
   3652                 DurabilityFailpoint::CloseAfterConnectionClose,
   3653                 ServiceSqliteErrorKind::Open,
   3654                 false,
   3655                 6,
   3656                 &all_close_edges[..],
   3657             ),
   3658             (
   3659                 DurabilityFailpoint::CloseBeforeAuthorityRelease,
   3660                 ServiceSqliteErrorKind::Authority,
   3661                 true,
   3662                 6,
   3663                 &all_close_edges[..7],
   3664             ),
   3665             (
   3666                 DurabilityFailpoint::CloseAfterAuthorityRelease,
   3667                 ServiceSqliteErrorKind::Authority,
   3668                 false,
   3669                 6,
   3670                 &all_close_edges[..],
   3671             ),
   3672         ] {
   3673             let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3674             host.arm_durability_failpoint(point);
   3675             let error = host.close().await.expect_err("injected close edge");
   3676             assert_eq!(error.kind(), expected);
   3677             assert!(host.failpoints.fired());
   3678             assert_eq!(host.failpoints.reached(), expected_reached);
   3679             assert_eq!(
   3680                 host.pool.close_phase(),
   3681                 expected_phase,
   3682                 "close failpoint must fire in its named driver phase"
   3683             );
   3684             if retryable {
   3685                 host.close().await.expect("one-shot close edge resumes");
   3686             } else {
   3687                 assert_eq!(
   3688                     host.close()
   3689                         .await
   3690                         .expect_err("terminal close result is cached")
   3691                         .kind(),
   3692                     expected
   3693                 );
   3694             }
   3695         }
   3696     }
   3697 
   3698     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3699     #[tokio::test]
   3700     async fn cancellation_quarantines_connection_and_pool_recovers() {
   3701         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3702         let host = Arc::new(host);
   3703         let entered = Arc::new(AtomicBool::new(false));
   3704         let task = tokio::spawn({
   3705             let host = Arc::clone(&host);
   3706             let entered = Arc::clone(&entered);
   3707             async move {
   3708                 host.transaction::<(), Infallible, _>(|transaction| {
   3709                     Box::pin(async move {
   3710                         sqlx::query("INSERT INTO host_probe (value) VALUES (77)")
   3711                             .execute(&mut *transaction)
   3712                             .await
   3713                             .expect("tentative insert");
   3714                         entered.store(true, Ordering::Release);
   3715                         std::future::pending::<()>().await;
   3716                         Ok(())
   3717                     })
   3718                 })
   3719                 .await
   3720             }
   3721         });
   3722         while !entered.load(Ordering::Acquire) {
   3723             tokio::task::yield_now().await;
   3724         }
   3725         task.abort();
   3726         assert!(
   3727             task.await
   3728                 .expect_err("task must be cancelled")
   3729                 .is_cancelled()
   3730         );
   3731         assert_eq!(row_count(&host).await, 0);
   3732     }
   3733 
   3734     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3735     #[tokio::test]
   3736     async fn complete_statement_control_inventory_is_sticky_and_fails_closed() {
   3737         for statement in [
   3738             "INSERT INTO host_probe (value) VALUES (0); /* policy */ PrAgMa\ntrusted_schema=ON",
   3739             "INSERT INTO host_probe (value) VALUES (1); ATTACH DATABASE ':memory:' AS extra",
   3740             "INSERT INTO host_probe (value) VALUES (2); DETACH DATABASE extra",
   3741             "INSERT INTO host_probe (value) VALUES (3); BEGIN DEFERRED",
   3742             "INSERT INTO host_probe (value) VALUES (4); COMMIT",
   3743             "INSERT INTO host_probe (value) VALUES (5); END",
   3744             "INSERT INTO host_probe (value) VALUES (6); ROLLBACK",
   3745             "INSERT INTO host_probe (value) VALUES (7); SAVEPOINT escaped",
   3746             "INSERT INTO host_probe (value) VALUES (8); RELEASE SAVEPOINT escaped",
   3747         ] {
   3748             let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3749             let error = host
   3750                 .transaction(|transaction| {
   3751                     Box::pin(async move {
   3752                         let _ = sqlx::raw_sql(statement).execute(&mut *transaction).await;
   3753                         Ok::<_, Infallible>(())
   3754                     })
   3755                 })
   3756                 .await
   3757                 .expect_err("escape attempt must not commit");
   3758             assert_eq!(
   3759                 error.kind(),
   3760                 ServiceSqliteTransactionErrorKind::NotCommitted
   3761             );
   3762             assert!(error.sqlite_error().is_some());
   3763             assert_eq!(row_count(&host).await, 0);
   3764         }
   3765 
   3766         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3767         let error = host
   3768             .transaction(|transaction| {
   3769                 Box::pin(async move {
   3770                     let _ = sqlx::raw_sql("INSERT INTO host_probe (value) VALUES (4); COMMIT")
   3771                         .execute(&mut *transaction)
   3772                         .await;
   3773                     let _ =
   3774                         sqlx::raw_sql("BEGIN DEFERRED; INSERT INTO host_probe (value) VALUES (5)")
   3775                             .execute(&mut *transaction)
   3776                             .await;
   3777                     Ok::<_, Infallible>(())
   3778                 })
   3779             })
   3780             .await
   3781             .expect_err("replacement transaction after denied COMMIT must not escape");
   3782         assert_eq!(
   3783             error.kind(),
   3784             ServiceSqliteTransactionErrorKind::NotCommitted
   3785         );
   3786         assert_eq!(row_count(&host).await, 0);
   3787     }
   3788 
   3789     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3790     #[tokio::test]
   3791     async fn prepared_query_policy_rejection_is_sticky_and_rolls_back_prior_work() {
   3792         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3793         let error = host
   3794             .transaction(|transaction| {
   3795                 Box::pin(async move {
   3796                     sqlx::query("INSERT INTO host_probe (value) VALUES (11)")
   3797                         .execute(&mut *transaction)
   3798                         .await
   3799                         .expect("ordinary service statement");
   3800                     let _ = (&mut *transaction)
   3801                         .prepare(SqlStr::from_static(
   3802                             "SELECT 1; /* ignored */ PRAGMA trusted_schema=ON",
   3803                         ))
   3804                         .await;
   3805                     Ok::<_, Infallible>(())
   3806                 })
   3807             })
   3808             .await
   3809             .expect_err("ignored prepared-query rejection must block commit");
   3810         assert_eq!(
   3811             error.kind(),
   3812             ServiceSqliteTransactionErrorKind::NotCommitted
   3813         );
   3814         assert_eq!(row_count(&host).await, 0);
   3815     }
   3816 
   3817     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3818     #[tokio::test]
   3819     async fn control_words_in_values_and_case_expressions_remain_available() {
   3820         let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3821         let value = host
   3822             .transaction(|transaction| {
   3823                 Box::pin(async move {
   3824                     let value = sqlx::query_scalar::<_, String>(
   3825                         "SELECT CASE WHEN 1 = 1 THEN 'commit' ELSE 'end' END",
   3826                     )
   3827                     .fetch_one(&mut *transaction)
   3828                     .await?;
   3829                     sqlx::query("INSERT INTO host_probe (value) VALUES (12)")
   3830                         .execute(&mut *transaction)
   3831                         .await?;
   3832                     Ok::<_, sqlx::Error>(value)
   3833                 })
   3834             })
   3835             .await
   3836             .expect("ordinary expression commits");
   3837         assert_eq!(value, "commit");
   3838         assert_eq!(row_count(&host).await, 1);
   3839     }
   3840 
   3841     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3842     #[tokio::test]
   3843     async fn attach_detach_is_rejected_before_it_can_create_external_state() {
   3844         let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3845         let external_database = root.path().join("forbidden-attachment.sqlite");
   3846         let statement = format!(
   3847             "ATTACH DATABASE '{}' AS extra; DETACH DATABASE extra",
   3848             external_database.display()
   3849         );
   3850         let error = host
   3851             .transaction(|transaction| {
   3852                 Box::pin(async move {
   3853                     sqlx::raw_sql(sqlx::AssertSqlSafe(statement))
   3854                         .execute(&mut *transaction)
   3855                         .await
   3856                         .map(|_| ())
   3857                 })
   3858             })
   3859             .await
   3860             .expect_err("ATTACH and DETACH must be rejected before SQLite compilation");
   3861         assert_eq!(
   3862             error.kind(),
   3863             ServiceSqliteTransactionErrorKind::OperationRolledBack
   3864         );
   3865         assert!(!external_database.exists());
   3866         assert_eq!(row_count(&host).await, 0);
   3867     }
   3868 
   3869     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3870     #[tokio::test]
   3871     async fn writer_lock_replacement_and_insecure_directory_revoke_live_host() {
   3872         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3873         let retired_lock = paths
   3874             .state_lock()
   3875             .parent()
   3876             .expect("state directory")
   3877             .join("retired-state.lock");
   3878         fs::rename(paths.state_lock(), &retired_lock).expect("replace canonical lock name");
   3879         let replacement_authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3880             .expect("replacement authority acquisition")
   3881             .expect("new writer authority");
   3882         let replaced = host
   3883             .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) }))
   3884             .await
   3885             .expect_err("old writer authority must reject the replacement lock");
   3886         assert_eq!(
   3887             replaced.kind(),
   3888             ServiceSqliteTransactionErrorKind::NotCommitted
   3889         );
   3890         assert_eq!(
   3891             replaced.sqlite_error().map(ServiceSqliteError::kind),
   3892             Some(ServiceSqliteErrorKind::Authority)
   3893         );
   3894         drop(replacement_authority);
   3895         drop(host);
   3896 
   3897         let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
   3898         let directory = paths.state_database().parent().expect("state directory");
   3899         let original_mode = fs::metadata(directory)
   3900             .expect("state directory metadata")
   3901             .permissions()
   3902             .mode()
   3903             & 0o777;
   3904         fs::set_permissions(directory, fs::Permissions::from_mode(0o770))
   3905             .expect("make directory insecure");
   3906         let insecure = host
   3907             .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) }))
   3908             .await
   3909             .expect_err("insecure live directory must revoke authority");
   3910         assert_eq!(
   3911             insecure.kind(),
   3912             ServiceSqliteTransactionErrorKind::NotCommitted
   3913         );
   3914         assert_eq!(
   3915             insecure.sqlite_error().map(ServiceSqliteError::kind),
   3916             Some(ServiceSqliteErrorKind::Authority)
   3917         );
   3918         fs::set_permissions(directory, fs::Permissions::from_mode(original_mode))
   3919             .expect("restore directory mode");
   3920     }
   3921 
   3922     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3923     #[tokio::test]
   3924     async fn inspection_lock_replacement_and_new_writer_revoke_live_inspection() {
   3925         let (_root, paths, identity, migrations, schema, host) = initialized_host().await;
   3926         host.close().await.expect("close writer host");
   3927         drop(host);
   3928         let inspection = ServiceSqliteHost::open_read_only_inspection(
   3929             &paths,
   3930             &identity,
   3931             &migrations,
   3932             &schema,
   3933             ServiceSqliteConnectionOptions::reviewed(),
   3934         )
   3935         .await
   3936         .expect("open inspection host");
   3937         let retired_lock = paths
   3938             .state_lock()
   3939             .parent()
   3940             .expect("state directory")
   3941             .join("inspection-state.lock");
   3942         fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock name");
   3943         let (writer, _outcome) = ServiceSqliteHost::open_read_write_existing(
   3944             &paths,
   3945             &identity,
   3946             &migrations,
   3947             &schema,
   3948             ServiceSqliteConnectionOptions::reviewed(),
   3949             MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
   3950             &build_identity(),
   3951             &[],
   3952         )
   3953         .await
   3954         .expect("open replacement writer");
   3955         writer
   3956             .transaction(|transaction| {
   3957                 Box::pin(async move {
   3958                     sqlx::query("INSERT INTO host_probe (value) VALUES (88)")
   3959                         .execute(&mut *transaction)
   3960                         .await
   3961                         .map(|_| ())
   3962                 })
   3963             })
   3964             .await
   3965             .expect("replacement writer commits");
   3966         let stale = inspection
   3967             .transaction(|transaction| {
   3968                 Box::pin(async move {
   3969                     sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe")
   3970                         .fetch_one(&mut *transaction)
   3971                         .await
   3972                 })
   3973             })
   3974             .await
   3975             .expect_err("stale inspection authority must refuse work");
   3976         assert_eq!(
   3977             stale.kind(),
   3978             ServiceSqliteTransactionErrorKind::NotCommitted
   3979         );
   3980         assert_eq!(
   3981             stale.sqlite_error().map(ServiceSqliteError::kind),
   3982             Some(ServiceSqliteErrorKind::Authority)
   3983         );
   3984     }
   3985 
   3986     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3987     #[tokio::test]
   3988     async fn read_only_writes_roll_back_and_internal_pool_close_refuses_new_work() {
   3989         let (root, paths, identity, migrations, schema, host) = initialized_host().await;
   3990         host.close().await.expect("close writer host");
   3991         let closed = host
   3992             .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) }))
   3993             .await
   3994             .expect_err("closed pool must refuse work");
   3995         assert_eq!(
   3996             closed.kind(),
   3997             ServiceSqliteTransactionErrorKind::NotCommitted
   3998         );
   3999         drop(host);
   4000 
   4001         let inspection = ServiceSqliteHost::open_read_only_inspection(
   4002             &paths,
   4003             &identity,
   4004             &migrations,
   4005             &schema,
   4006             ServiceSqliteConnectionOptions::reviewed(),
   4007         )
   4008         .await
   4009         .expect("open read-only host");
   4010         let error = inspection
   4011             .transaction(|transaction| {
   4012                 Box::pin(async move {
   4013                     sqlx::query("INSERT INTO host_probe (value) VALUES (5)")
   4014                         .execute(&mut *transaction)
   4015                         .await
   4016                         .map(|_| ())
   4017                 })
   4018             })
   4019             .await
   4020             .expect_err("read-only write must fail");
   4021         assert_eq!(
   4022             error.kind(),
   4023             ServiceSqliteTransactionErrorKind::OperationRolledBack
   4024         );
   4025         assert_eq!(row_count(&inspection).await, 0);
   4026         drop(inspection);
   4027         drop(root);
   4028     }
   4029 
   4030     #[test]
   4031     fn host_and_transaction_errors_are_redacted_and_source_free() {
   4032         let error = ServiceSqliteTransactionError::not_committed_with_operation(
   4033             Some("operation-secret"),
   4034             ServiceSqliteError::new(ServiceSqliteErrorKind::Open),
   4035         );
   4036         let debug = format!("{error:?}");
   4037         assert!(!debug.contains("operation-secret"));
   4038         assert!(!debug.contains("state.sqlite"));
   4039         assert!(error.source().is_none());
   4040         assert_eq!(error.operation_error(), Some(&"operation-secret"));
   4041         assert_eq!(
   4042             error.sqlite_error().map(ServiceSqliteError::kind),
   4043             Some(ServiceSqliteErrorKind::Open)
   4044         );
   4045     }
   4046 }