lib

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

open.rs (137670B)


      1 //! Instance-bound SQLite paths and declarative open modes.
      2 
      3 use core::fmt;
      4 use std::{
      5     error::Error,
      6     path::{Path, PathBuf},
      7 };
      8 
      9 #[cfg(any(target_os = "linux", target_os = "macos"))]
     10 use std::{
     11     fs::File,
     12     sync::{
     13         Arc, Mutex,
     14         atomic::{AtomicBool, Ordering},
     15     },
     16 };
     17 
     18 #[cfg(any(target_os = "linux", target_os = "macos"))]
     19 use fs2::FileExt;
     20 #[cfg(any(target_os = "linux", target_os = "macos"))]
     21 use futures::future::BoxFuture;
     22 use radroots_runtime_paths::{
     23     InstanceId, RuntimeContext, ServiceId, default_service_instance_artifacts,
     24 };
     25 use serde::Serialize;
     26 #[cfg(any(target_os = "linux", target_os = "macos"))]
     27 use sqlx::{
     28     ConnectOptions, Connection, Sqlite, SqliteConnection, SqlitePool,
     29     pool::PoolConnection,
     30     sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous},
     31 };
     32 
     33 #[cfg(any(target_os = "linux", target_os = "macos"))]
     34 use crate::{
     35     ExistingServiceDatabaseIntent, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity,
     36     MigrationCatalog, SchemaCatalog, ServiceDatabaseIdentity, ServiceDatabaseMetadata,
     37     ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind, WriterAuthority,
     38 };
     39 
     40 #[cfg(any(target_os = "linux", target_os = "macos"))]
     41 const STATEMENT_CACHE_CAPACITY: usize = 100;
     42 #[cfg(any(target_os = "linux", target_os = "macos"))]
     43 const COMMAND_BUFFER_CAPACITY: usize = 50;
     44 #[cfg(any(target_os = "linux", target_os = "macos"))]
     45 const ROW_BUFFER_CAPACITY: usize = 50;
     46 #[cfg(any(target_os = "linux", target_os = "macos"))]
     47 const WAL_FILE_NAME: &str = "state.sqlite-wal";
     48 #[cfg(any(target_os = "linux", target_os = "macos"))]
     49 const SHARED_MEMORY_FILE_NAME: &str = "state.sqlite-shm";
     50 
     51 #[cfg(any(target_os = "linux", target_os = "macos"))]
     52 #[derive(Clone)]
     53 struct PoolConnectionValidation {
     54     binding: DirectoryBinding,
     55     paths: ServiceSqlitePaths,
     56     identity: ServiceDatabaseIdentity,
     57     catalog: MigrationCatalog,
     58     schema_catalog: SchemaCatalog,
     59     mode: OpenMode,
     60     policy: ServiceSqliteConnectionOptions,
     61 }
     62 
     63 #[cfg(any(target_os = "linux", target_os = "macos"))]
     64 #[derive(Clone, Copy)]
     65 enum ServiceDatabaseExpectation<'a> {
     66     Exact(&'a ServiceDatabaseIdentity),
     67     Existing(&'a ExistingServiceDatabaseIntent),
     68 }
     69 
     70 #[cfg(any(target_os = "linux", target_os = "macos"))]
     71 impl ServiceDatabaseExpectation<'_> {
     72     fn matches_paths(self, paths: &ServiceSqlitePaths) -> bool {
     73         match self {
     74             Self::Exact(identity) => identity.matches_paths(paths),
     75             Self::Existing(intent) => intent.matches_paths(paths),
     76         }
     77     }
     78 
     79     fn supported_state_schema_version(self) -> core::num::NonZeroU32 {
     80         match self {
     81             Self::Exact(identity) => identity.supported_state_schema_version(),
     82             Self::Existing(intent) => intent.supported_state_schema_version(),
     83         }
     84     }
     85 
     86     async fn verify_metadata(
     87         self,
     88         connection: &mut SqliteConnection,
     89     ) -> Result<ServiceDatabaseMetadata, ServiceSqliteError> {
     90         match self {
     91             Self::Exact(identity) => {
     92                 crate::metadata::verify_database_metadata(connection, identity).await
     93             }
     94             Self::Existing(intent) => {
     95                 crate::metadata::verify_existing_database_intent(connection, intent).await
     96             }
     97         }
     98     }
     99 
    100     fn exact_identity(self, metadata: &ServiceDatabaseMetadata) -> ServiceDatabaseIdentity {
    101         match self {
    102             Self::Exact(identity) => identity.clone(),
    103             Self::Existing(intent) => intent.identity_for(metadata),
    104         }
    105     }
    106 }
    107 
    108 #[cfg(any(target_os = "linux", target_os = "macos"))]
    109 enum PoolConnectionValidationFailure {
    110     Authority,
    111     Pragma(sqlx::Error),
    112     PolicyMismatch,
    113     Metadata,
    114     Migration,
    115     Integrity,
    116 }
    117 
    118 #[cfg(any(target_os = "linux", target_os = "macos"))]
    119 #[derive(Clone)]
    120 struct PoolConnectionFailureFlags {
    121     authority: Arc<AtomicBool>,
    122     metadata: Arc<AtomicBool>,
    123     migration: Arc<AtomicBool>,
    124     integrity: Arc<AtomicBool>,
    125     pragma: Arc<AtomicBool>,
    126 }
    127 
    128 #[cfg(any(target_os = "linux", target_os = "macos"))]
    129 impl PoolConnectionFailureFlags {
    130     fn record(&self, failure: &PoolConnectionValidationFailure) {
    131         match failure {
    132             PoolConnectionValidationFailure::Authority => &self.authority,
    133             PoolConnectionValidationFailure::Metadata => &self.metadata,
    134             PoolConnectionValidationFailure::Migration => &self.migration,
    135             PoolConnectionValidationFailure::Integrity => &self.integrity,
    136             PoolConnectionValidationFailure::Pragma(_)
    137             | PoolConnectionValidationFailure::PolicyMismatch => &self.pragma,
    138         }
    139         .store(true, Ordering::Release);
    140     }
    141 
    142     fn kind(&self) -> ServiceSqliteErrorKind {
    143         connection_failure_kind(
    144             self.authority.load(Ordering::Acquire),
    145             self.metadata.load(Ordering::Acquire),
    146             self.migration.load(Ordering::Acquire),
    147             self.integrity.load(Ordering::Acquire),
    148             self.pragma.load(Ordering::Acquire),
    149         )
    150     }
    151 }
    152 
    153 #[cfg(any(target_os = "linux", target_os = "macos"))]
    154 impl PoolConnectionValidationFailure {
    155     fn into_sqlx(self) -> sqlx::Error {
    156         match self {
    157             Self::Pragma(source) => source,
    158             Self::Authority => {
    159                 sqlx::Error::Protocol("SQLite connection authority mismatch".to_owned())
    160             }
    161             Self::PolicyMismatch => {
    162                 sqlx::Error::Protocol("SQLite connection policy mismatch".to_owned())
    163             }
    164             Self::Metadata => {
    165                 sqlx::Error::Protocol("SQLite connection metadata mismatch".to_owned())
    166             }
    167             Self::Migration | Self::Integrity => {
    168                 sqlx::Error::Protocol("SQLite migration history mismatch".to_owned())
    169             }
    170         }
    171     }
    172 }
    173 
    174 /// Canonical database and writer-lock paths for one validated service instance.
    175 ///
    176 /// Callers cannot forge paths or rebind the service and instance independently:
    177 ///
    178 /// ```compile_fail
    179 /// use std::path::PathBuf;
    180 /// use radroots_runtime_paths::{InstanceId, ServiceId};
    181 /// use radroots_service_sqlite::ServiceSqlitePaths;
    182 ///
    183 /// let _ = ServiceSqlitePaths {
    184 ///     service: ServiceId::new("example").unwrap(),
    185 ///     instance: InstanceId::new("primary").unwrap(),
    186 ///     state_database: PathBuf::from("/tmp/alternate.sqlite"),
    187 ///     state_lock: PathBuf::from("/tmp/alternate.lock"),
    188 /// };
    189 /// ```
    190 #[derive(Clone, PartialEq, Eq)]
    191 pub struct ServiceSqlitePaths {
    192     service: ServiceId,
    193     instance: InstanceId,
    194     state_database: PathBuf,
    195     state_lock: PathBuf,
    196 }
    197 
    198 impl ServiceSqlitePaths {
    199     /// Derives the fixed SQLite artifacts from one immutable runtime context.
    200     pub fn from_runtime_context(context: &RuntimeContext) -> Result<Self, ServiceSqlitePathError> {
    201         validate_state_directory(context.paths().state())?;
    202         let artifacts = default_service_instance_artifacts(context.paths());
    203         Ok(Self {
    204             service: context.service().clone(),
    205             instance: context.instance().clone(),
    206             state_database: artifacts.state_database().to_path_buf(),
    207             state_lock: artifacts.state_lock().to_path_buf(),
    208         })
    209     }
    210 
    211     /// Returns the validated service identity bound to these paths.
    212     #[must_use]
    213     pub fn service(&self) -> &ServiceId {
    214         &self.service
    215     }
    216 
    217     /// Returns the validated instance identity bound to these paths.
    218     #[must_use]
    219     pub fn instance(&self) -> &InstanceId {
    220         &self.instance
    221     }
    222 
    223     /// Returns the canonical `state.sqlite` path.
    224     #[must_use]
    225     pub fn state_database(&self) -> &Path {
    226         &self.state_database
    227     }
    228 
    229     /// Returns the canonical retained `state.lock` path.
    230     #[must_use]
    231     pub fn state_lock(&self) -> &Path {
    232         &self.state_lock
    233     }
    234 }
    235 
    236 impl fmt::Debug for ServiceSqlitePaths {
    237     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    238         formatter
    239             .debug_struct("ServiceSqlitePaths")
    240             .field("service", &self.service)
    241             .field("instance", &self.instance)
    242             .field("state_database", &"[redacted]")
    243             .field("state_lock", &"[redacted]")
    244             .finish()
    245     }
    246 }
    247 
    248 /// Path-shape failure detected before any filesystem or SQLite operation.
    249 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    250 pub enum ServiceSqlitePathError {
    251     RelativeStateDirectory,
    252     MissingStateDirectoryParent,
    253 }
    254 
    255 impl fmt::Display for ServiceSqlitePathError {
    256     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    257         match self {
    258             Self::RelativeStateDirectory => {
    259                 formatter.write_str("SQLite state directory must be absolute")
    260             }
    261             Self::MissingStateDirectoryParent => {
    262                 formatter.write_str("SQLite state directory must have a parent")
    263             }
    264         }
    265     }
    266 }
    267 
    268 impl Error for ServiceSqlitePathError {}
    269 
    270 fn validate_state_directory(path: &Path) -> Result<(), ServiceSqlitePathError> {
    271     if !path.is_absolute() {
    272         return Err(ServiceSqlitePathError::RelativeStateDirectory);
    273     }
    274     if path.parent().is_none() {
    275         return Err(ServiceSqlitePathError::MissingStateDirectoryParent);
    276     }
    277     Ok(())
    278 }
    279 
    280 /// Declarative behavior for opening one service-owned SQLite database.
    281 #[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
    282 #[serde(rename_all = "snake_case")]
    283 pub enum OpenMode {
    284     Initialize,
    285     ReadWriteExisting,
    286     ReadOnlyInspection,
    287 }
    288 
    289 impl OpenMode {
    290     /// Returns whether this mode permits creating missing state.
    291     #[must_use]
    292     pub const fn may_create(self) -> bool {
    293         matches!(self, Self::Initialize)
    294     }
    295 
    296     /// Returns whether state must already exist before opening.
    297     #[must_use]
    298     pub const fn requires_existing(self) -> bool {
    299         !matches!(self, Self::Initialize)
    300     }
    301 
    302     /// Returns whether exclusive writer authority is required.
    303     #[must_use]
    304     pub const fn requires_writer_authority(self) -> bool {
    305         !matches!(self, Self::ReadOnlyInspection)
    306     }
    307 }
    308 
    309 #[cfg(any(target_os = "linux", target_os = "macos"))]
    310 pub(crate) struct PrivateConnectionPool {
    311     pool: SqlitePool,
    312     binding: DirectoryBinding,
    313     paths: ServiceSqlitePaths,
    314     identity: ServiceDatabaseIdentity,
    315     catalog: MigrationCatalog,
    316     schema_catalog: SchemaCatalog,
    317     mode: OpenMode,
    318     policy: ServiceSqliteConnectionOptions,
    319     resources: Arc<Mutex<PrivateConnectionResources>>,
    320     close_driver: tokio::sync::Mutex<PrivateCloseDriver>,
    321     #[cfg(test)]
    322     close_phase: std::sync::atomic::AtomicU8,
    323     authority_failure: Arc<AtomicBool>,
    324     metadata_failure: Arc<AtomicBool>,
    325     migration_failure: Arc<AtomicBool>,
    326     integrity_failure: Arc<AtomicBool>,
    327     pragma_failure: Arc<AtomicBool>,
    328 }
    329 
    330 #[cfg(any(target_os = "linux", target_os = "macos"))]
    331 struct PrivateConnectionResources {
    332     authority: Option<WriterAuthority>,
    333     inspection_guard: Option<ReadOnlyInspectionGuard>,
    334 }
    335 
    336 #[cfg(any(target_os = "linux", target_os = "macos"))]
    337 #[derive(Clone)]
    338 pub(crate) struct BackupSourceValidator {
    339     binding: DirectoryBinding,
    340     paths: ServiceSqlitePaths,
    341     resources: Arc<Mutex<PrivateConnectionResources>>,
    342 }
    343 
    344 #[cfg(any(target_os = "linux", target_os = "macos"))]
    345 impl BackupSourceValidator {
    346     pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> {
    347         let resources = self.resources.lock().map_err(|_| {
    348             connection_error(
    349                 ServiceSqliteErrorKind::Authority,
    350                 ConnectionFailureKind::AuthorityMismatch,
    351             )
    352         })?;
    353         resources
    354             .authority
    355             .as_ref()
    356             .ok_or_else(|| {
    357                 connection_error(
    358                     ServiceSqliteErrorKind::Authority,
    359                     ConnectionFailureKind::AuthorityMismatch,
    360                 )
    361             })?
    362             .validate_for(&self.paths)?;
    363         self.binding.validate(&self.paths)
    364     }
    365 }
    366 
    367 #[cfg(any(target_os = "linux", target_os = "macos"))]
    368 enum PrivateCloseDriver {
    369     Pending,
    370     Connecting(BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>>),
    371     Connected(SqliteConnection),
    372     Closing {
    373         future: BoxFuture<'static, Result<(), ServiceSqliteError>>,
    374         authority_error: Option<ServiceSqliteError>,
    375         checkpoint_error: Option<ServiceSqliteError>,
    376         connection_close_error: Option<ServiceSqliteError>,
    377     },
    378     Complete(Option<ServiceSqliteErrorKind>),
    379 }
    380 
    381 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    382 pub(crate) const TEST_CLOSE_PHASE_CHECKPOINT: u8 = 4;
    383 
    384 #[cfg(any(target_os = "linux", target_os = "macos"))]
    385 impl PrivateConnectionPool {
    386     fn connection_failure_kind(&self) -> ServiceSqliteErrorKind {
    387         connection_failure_kind(
    388             self.authority_failure.load(Ordering::Acquire),
    389             self.metadata_failure.load(Ordering::Acquire),
    390             self.migration_failure.load(Ordering::Acquire),
    391             self.integrity_failure.load(Ordering::Acquire),
    392             self.pragma_failure.load(Ordering::Acquire),
    393         )
    394     }
    395 
    396     pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> {
    397         let resources = self.resources.lock().map_err(|_| {
    398             connection_error(
    399                 ServiceSqliteErrorKind::Authority,
    400                 ConnectionFailureKind::AuthorityMismatch,
    401             )
    402         })?;
    403         match self.mode {
    404             OpenMode::Initialize | OpenMode::ReadWriteExisting => resources
    405                 .authority
    406                 .as_ref()
    407                 .ok_or_else(|| {
    408                     connection_error(
    409                         ServiceSqliteErrorKind::Authority,
    410                         ConnectionFailureKind::AuthorityMismatch,
    411                     )
    412                 })?
    413                 .validate_for(&self.paths)?,
    414             OpenMode::ReadOnlyInspection => resources
    415                 .inspection_guard
    416                 .as_ref()
    417                 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?
    418                 .validate_for(&self.paths)?,
    419         }
    420         self.binding.validate(&self.paths)
    421     }
    422 
    423     pub(crate) const fn mode(&self) -> OpenMode {
    424         self.mode
    425     }
    426 
    427     pub(crate) fn identity(&self) -> &ServiceDatabaseIdentity {
    428         &self.identity
    429     }
    430 
    431     pub(crate) async fn database_metadata(
    432         &self,
    433     ) -> Result<ServiceDatabaseMetadata, ServiceSqliteError> {
    434         self.validate()?;
    435         let mut connection = self.acquire().await?;
    436         let result =
    437             crate::metadata::verify_database_metadata(&mut connection, &self.identity).await;
    438         self.validate()?;
    439         result
    440     }
    441 
    442     pub(crate) fn backup_source_validator(&self) -> BackupSourceValidator {
    443         BackupSourceValidator {
    444             binding: self.binding.clone(),
    445             paths: self.paths.clone(),
    446             resources: Arc::clone(&self.resources),
    447         }
    448     }
    449 
    450     pub(crate) fn catalog(&self) -> &MigrationCatalog {
    451         &self.catalog
    452     }
    453 
    454     pub(crate) fn schema_catalog(&self) -> &SchemaCatalog {
    455         &self.schema_catalog
    456     }
    457 
    458     #[cfg(test)]
    459     pub(crate) fn close_phase(&self) -> u8 {
    460         self.close_phase.load(Ordering::Acquire)
    461     }
    462 
    463     pub(crate) async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> {
    464         self.validate()?;
    465         let result = self.pool.acquire().await;
    466         self.validate()?;
    467         let mut connection =
    468             result.map_err(|source| connection_source(self.connection_failure_kind(), source))?;
    469         let history = crate::migration::verify_migration_history(
    470             &mut connection,
    471             &self.catalog,
    472             &self.schema_catalog,
    473             true,
    474         )
    475         .await;
    476         self.validate()?;
    477         history?;
    478         Ok(connection)
    479     }
    480 
    481     pub(crate) async fn apply_migrations(
    482         &self,
    483         applied_at: MigrationAppliedAtUnixSeconds,
    484         build: &MigrationBuildIdentity,
    485         callbacks: &[crate::migration::MigrationCallbackBinding],
    486     ) -> Result<crate::migration::MigrationApplicationOutcome, ServiceSqliteError> {
    487         self.validate_writer_authority()?;
    488         let acquired = self.pool.acquire().await;
    489         self.validate_writer_authority()?;
    490         let mut connection =
    491             acquired.map_err(|source| connection_source(self.connection_failure_kind(), source))?;
    492         // Migration execution installs connection-local fail-closed guards. Always
    493         // discard this one-time connection so cancellation cannot return a guarded
    494         // or callback-altered handle to the pool.
    495         connection.close_on_drop();
    496         let mut validate_authority = || self.validate_writer_authority();
    497         let result = crate::migration::apply_governed_migrations(
    498             &mut connection,
    499             &self.catalog,
    500             &self.schema_catalog,
    501             applied_at,
    502             build,
    503             callbacks,
    504             &mut validate_authority,
    505         )
    506         .await;
    507         self.validate_writer_authority()?;
    508         result
    509     }
    510 
    511     fn validate_writer_authority(&self) -> Result<(), ServiceSqliteError> {
    512         let resources = self.resources.lock().map_err(|_| {
    513             connection_error(
    514                 ServiceSqliteErrorKind::Authority,
    515                 ConnectionFailureKind::AuthorityMismatch,
    516             )
    517         })?;
    518         resources
    519             .authority
    520             .as_ref()
    521             .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Authority))?
    522             .validate_for(&self.paths)?;
    523         self.binding.validate(&self.paths)
    524     }
    525 
    526     pub(crate) async fn close(self) -> Option<WriterAuthority> {
    527         self.pool.close().await;
    528         let mut resources = match self.resources.lock() {
    529             Ok(resources) => resources,
    530             Err(poisoned) => poisoned.into_inner(),
    531         };
    532         resources.inspection_guard.take();
    533         resources.authority.take()
    534     }
    535 
    536     /// The outer error means authority release was not proven and close must retry.
    537     /// The inner result is terminal and may be cached by the host.
    538     pub(crate) async fn close_explicit(
    539         &self,
    540         failpoints: &crate::failpoint::DurabilityFailpoints,
    541     ) -> Result<Result<(), ServiceSqliteError>, ServiceSqliteError> {
    542         let mut authority_error = self.validate().err();
    543         let before_drain = failpoints
    544             .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeDrain)
    545             .map_err(|source| {
    546                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    547             });
    548         authority_error = authority_error.or_else(|| self.validate().err());
    549         if authority_error.is_none() {
    550             before_drain?;
    551         }
    552         self.pool.close().await;
    553         let after_drain = failpoints
    554             .hit(crate::failpoint::DurabilityFailpoint::CloseAfterDrain)
    555             .map_err(|source| {
    556                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    557             });
    558         authority_error = authority_error.or_else(|| self.validate().err());
    559         let terminal = match authority_error {
    560             Some(authority) => Err(authority),
    561             None => {
    562                 after_drain?;
    563                 if self.mode.requires_writer_authority() {
    564                     self.drive_writable_close(failpoints).await
    565                 } else {
    566                     self.validate()
    567                 }
    568             }
    569         };
    570         let release_error = self.release_resources(failpoints)?;
    571         Ok(match (terminal, release_error) {
    572             (Err(error), _) if error.kind() == ServiceSqliteErrorKind::Authority => Err(error),
    573             (_, Some(release_error)) => Err(release_error),
    574             (terminal, None) => terminal,
    575         })
    576     }
    577 
    578     async fn drive_writable_close(
    579         &self,
    580         failpoints: &crate::failpoint::DurabilityFailpoints,
    581     ) -> Result<(), ServiceSqliteError> {
    582         let mut driver = self.close_driver.lock().await;
    583         loop {
    584             match &mut *driver {
    585                 PrivateCloseDriver::Pending => {
    586                     if let Err(error) = self.validate_writer_authority() {
    587                         #[cfg(test)]
    588                         self.close_phase.store(6, Ordering::Release);
    589                         *driver = PrivateCloseDriver::Complete(Some(error.kind()));
    590                         return Err(error);
    591                     }
    592                     let options = sqlite_connect_options(&self.paths, self.mode, self.policy);
    593                     let connect: BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>> =
    594                         Box::pin(async move {
    595                             SqliteConnection::connect_with(&options)
    596                                 .await
    597                                 .map_err(|source| {
    598                                     connection_source(ServiceSqliteErrorKind::Open, source)
    599                                 })
    600                         });
    601                     #[cfg(test)]
    602                     self.close_phase.store(1, Ordering::Release);
    603                     *driver = PrivateCloseDriver::Connecting(connect);
    604                 }
    605                 PrivateCloseDriver::Connecting(connect) => {
    606                     let connected = connect.as_mut().await;
    607                     match connected {
    608                         Ok(connection) => {
    609                             #[cfg(test)]
    610                             self.close_phase.store(2, Ordering::Release);
    611                             *driver = PrivateCloseDriver::Connected(connection);
    612                         }
    613                         Err(error) => {
    614                             let authority_error = self.validate_writer_authority().err();
    615                             let error = authority_error.unwrap_or(error);
    616                             #[cfg(test)]
    617                             self.close_phase.store(6, Ordering::Release);
    618                             *driver = PrivateCloseDriver::Complete(Some(error.kind()));
    619                             return Err(error);
    620                         }
    621                     }
    622                 }
    623                 PrivateCloseDriver::Connected(connection) => {
    624                     let mut authority_error = self.validate_writer_authority().err();
    625                     let mut checkpoint_error = None;
    626                     if authority_error.is_none() {
    627                         #[cfg(test)]
    628                         self.close_phase.store(3, Ordering::Release);
    629                         checkpoint_error =
    630                             verify_connection_policy(connection, self.mode, self.policy)
    631                                 .await
    632                                 .map_err(|source| {
    633                                     connection_source(ServiceSqliteErrorKind::Pragma, source)
    634                                 })
    635                                 .err();
    636                         authority_error =
    637                             authority_error.or_else(|| self.validate_writer_authority().err());
    638                     }
    639                     if checkpoint_error.is_none() && authority_error.is_none() {
    640                         #[cfg(test)]
    641                         self.close_phase
    642                             .store(TEST_CLOSE_PHASE_CHECKPOINT, Ordering::Release);
    643                         checkpoint_error = failpoints
    644                             .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeCheckpoint)
    645                             .map_err(|source| {
    646                                 ServiceSqliteError::with_source(
    647                                     ServiceSqliteErrorKind::Pragma,
    648                                     source,
    649                                 )
    650                             })
    651                             .err();
    652                         authority_error =
    653                             authority_error.or_else(|| self.validate_writer_authority().err());
    654                     }
    655                     if checkpoint_error.is_none() && authority_error.is_none() {
    656                         checkpoint_error =
    657                             sqlx::query_as::<_, (i64, i64, i64)>("PRAGMA wal_checkpoint(TRUNCATE)")
    658                                 .fetch_one(&mut *connection)
    659                                 .await
    660                                 .map_err(|source| {
    661                                     connection_source(ServiceSqliteErrorKind::Pragma, source)
    662                                 })
    663                                 .and_then(|(busy, _log_frames, _checkpointed_frames)| {
    664                                     if busy == 0 {
    665                                         Ok(())
    666                                     } else {
    667                                         Err(connection_error(
    668                                             ServiceSqliteErrorKind::Pragma,
    669                                             ConnectionFailureKind::CheckpointBusy,
    670                                         ))
    671                                     }
    672                                 })
    673                                 .err();
    674                         if checkpoint_error.is_none() {
    675                             checkpoint_error = failpoints
    676                                 .hit(crate::failpoint::DurabilityFailpoint::CloseAfterCheckpoint)
    677                                 .map_err(|source| {
    678                                     ServiceSqliteError::with_source(
    679                                         ServiceSqliteErrorKind::Pragma,
    680                                         source,
    681                                     )
    682                                 })
    683                                 .err();
    684                         }
    685                         authority_error =
    686                             authority_error.or_else(|| self.validate_writer_authority().err());
    687                     }
    688 
    689                     let connection_close_error = failpoints
    690                         .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeConnectionClose)
    691                         .map_err(|source| {
    692                             ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    693                         })
    694                         .err();
    695 
    696                     let connected = core::mem::replace(&mut *driver, PrivateCloseDriver::Pending);
    697                     let PrivateCloseDriver::Connected(connection) = connected else {
    698                         unreachable!("close driver retains its connected phase")
    699                     };
    700                     let close: BoxFuture<'static, Result<(), ServiceSqliteError>> =
    701                         Box::pin(async move {
    702                             connection.close().await.map_err(|source| {
    703                                 connection_source(ServiceSqliteErrorKind::Open, source)
    704                             })
    705                         });
    706                     #[cfg(test)]
    707                     self.close_phase.store(5, Ordering::Release);
    708                     *driver = PrivateCloseDriver::Closing {
    709                         future: close,
    710                         authority_error,
    711                         checkpoint_error,
    712                         connection_close_error,
    713                     };
    714                 }
    715                 PrivateCloseDriver::Closing {
    716                     future,
    717                     authority_error,
    718                     checkpoint_error,
    719                     connection_close_error,
    720                 } => {
    721                     let close_error = future.as_mut().await.err();
    722                     let injected_after_close = failpoints
    723                         .hit(crate::failpoint::DurabilityFailpoint::CloseAfterConnectionClose)
    724                         .map_err(|source| {
    725                             ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
    726                         })
    727                         .err();
    728                     let authority_error = authority_error
    729                         .take()
    730                         .or_else(|| self.validate_writer_authority().err());
    731                     let error = authority_error
    732                         .or_else(|| checkpoint_error.take())
    733                         .or_else(|| connection_close_error.take())
    734                         .or(close_error)
    735                         .or(injected_after_close);
    736                     #[cfg(test)]
    737                     self.close_phase.store(6, Ordering::Release);
    738                     *driver =
    739                         PrivateCloseDriver::Complete(error.as_ref().map(ServiceSqliteError::kind));
    740                     return error.map_or(Ok(()), Err);
    741                 }
    742                 PrivateCloseDriver::Complete(kind) => {
    743                     return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind)));
    744                 }
    745             }
    746         }
    747     }
    748 
    749     fn release_resources(
    750         &self,
    751         failpoints: &crate::failpoint::DurabilityFailpoints,
    752     ) -> Result<Option<ServiceSqliteError>, ServiceSqliteError> {
    753         failpoints
    754             .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeAuthorityRelease)
    755             .map_err(|source| {
    756                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source)
    757             })?;
    758         let mut resources = self.resources.lock().map_err(|_| {
    759             connection_error(
    760                 ServiceSqliteErrorKind::Authority,
    761                 ConnectionFailureKind::AuthorityMismatch,
    762             )
    763         })?;
    764         match self.mode {
    765             OpenMode::Initialize | OpenMode::ReadWriteExisting => {
    766                 let authority = resources.authority.as_mut().ok_or_else(|| {
    767                     connection_error(
    768                         ServiceSqliteErrorKind::Authority,
    769                         ConnectionFailureKind::AuthorityMismatch,
    770                     )
    771                 })?;
    772                 authority.release()?;
    773                 resources.authority.take();
    774             }
    775             OpenMode::ReadOnlyInspection => {
    776                 let inspection = resources.inspection_guard.as_mut().ok_or_else(|| {
    777                     inspection_error(ConnectionFailureKind::InspectionUnavailable)
    778                 })?;
    779                 inspection.release()?;
    780                 resources.inspection_guard.take();
    781             }
    782         }
    783         Ok(failpoints
    784             .hit(crate::failpoint::DurabilityFailpoint::CloseAfterAuthorityRelease)
    785             .map_err(|source| {
    786                 ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source)
    787             })
    788             .err())
    789     }
    790 }
    791 
    792 #[cfg(any(target_os = "linux", target_os = "macos"))]
    793 const fn connection_failure_kind(
    794     authority: bool,
    795     metadata: bool,
    796     migration: bool,
    797     integrity: bool,
    798     pragma: bool,
    799 ) -> ServiceSqliteErrorKind {
    800     if authority {
    801         ServiceSqliteErrorKind::Authority
    802     } else if metadata {
    803         ServiceSqliteErrorKind::Metadata
    804     } else if migration {
    805         ServiceSqliteErrorKind::Migration
    806     } else if integrity {
    807         ServiceSqliteErrorKind::Integrity
    808     } else if pragma {
    809         ServiceSqliteErrorKind::Pragma
    810     } else {
    811         ServiceSqliteErrorKind::Open
    812     }
    813 }
    814 
    815 #[cfg(any(target_os = "linux", target_os = "macos"))]
    816 pub(crate) async fn open_existing_connection_pool(
    817     paths: &ServiceSqlitePaths,
    818     identity: &ServiceDatabaseIdentity,
    819     catalog: &MigrationCatalog,
    820     schema_catalog: &SchemaCatalog,
    821     mode: OpenMode,
    822     policy: ServiceSqliteConnectionOptions,
    823 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    824     open_existing_connection_pool_for(
    825         paths,
    826         ServiceDatabaseExpectation::Exact(identity),
    827         catalog,
    828         schema_catalog,
    829         mode,
    830         policy,
    831     )
    832     .await
    833 }
    834 
    835 #[cfg(any(target_os = "linux", target_os = "macos"))]
    836 pub(crate) async fn open_existing_connection_pool_with_intent(
    837     paths: &ServiceSqlitePaths,
    838     intent: &ExistingServiceDatabaseIntent,
    839     catalog: &MigrationCatalog,
    840     schema_catalog: &SchemaCatalog,
    841     mode: OpenMode,
    842     policy: ServiceSqliteConnectionOptions,
    843 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    844     open_existing_connection_pool_for(
    845         paths,
    846         ServiceDatabaseExpectation::Existing(intent),
    847         catalog,
    848         schema_catalog,
    849         mode,
    850         policy,
    851     )
    852     .await
    853 }
    854 
    855 #[cfg(any(target_os = "linux", target_os = "macos"))]
    856 pub(crate) async fn open_existing_connection_pool_with_intent_and_authority(
    857     paths: &ServiceSqlitePaths,
    858     intent: &ExistingServiceDatabaseIntent,
    859     catalog: &MigrationCatalog,
    860     schema_catalog: &SchemaCatalog,
    861     policy: ServiceSqliteConnectionOptions,
    862     authority: WriterAuthority,
    863 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    864     authority.validate_for(paths)?;
    865     open_connection_pool(
    866         paths,
    867         ServiceDatabaseExpectation::Existing(intent),
    868         catalog,
    869         schema_catalog,
    870         OpenMode::ReadWriteExisting,
    871         policy,
    872         Some(authority),
    873         None,
    874     )
    875     .await
    876 }
    877 
    878 #[cfg(any(target_os = "linux", target_os = "macos"))]
    879 async fn open_existing_connection_pool_for(
    880     paths: &ServiceSqlitePaths,
    881     expectation: ServiceDatabaseExpectation<'_>,
    882     catalog: &MigrationCatalog,
    883     schema_catalog: &SchemaCatalog,
    884     mode: OpenMode,
    885     policy: ServiceSqliteConnectionOptions,
    886 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    887     if mode == OpenMode::Initialize {
    888         return Err(connection_error(
    889             ServiceSqliteErrorKind::Open,
    890             ConnectionFailureKind::UnsupportedMode,
    891         ));
    892     }
    893     let (authority, inspection_guard) = match mode {
    894         OpenMode::ReadWriteExisting => (WriterAuthority::acquire(paths, mode)?, None),
    895         OpenMode::ReadOnlyInspection => (None, Some(ReadOnlyInspectionGuard::acquire(paths)?)),
    896         OpenMode::Initialize => unreachable!("initialize mode returned above"),
    897     };
    898     open_connection_pool(
    899         paths,
    900         expectation,
    901         catalog,
    902         schema_catalog,
    903         mode,
    904         policy,
    905         authority,
    906         inspection_guard,
    907     )
    908     .await
    909 }
    910 
    911 #[cfg(any(target_os = "linux", target_os = "macos"))]
    912 pub(crate) async fn open_initialized_connection_pool(
    913     paths: &ServiceSqlitePaths,
    914     identity: &ServiceDatabaseIdentity,
    915     catalog: &MigrationCatalog,
    916     schema_catalog: &SchemaCatalog,
    917     policy: ServiceSqliteConnectionOptions,
    918     authority: WriterAuthority,
    919 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    920     authority.validate_for(paths)?;
    921     open_connection_pool(
    922         paths,
    923         ServiceDatabaseExpectation::Exact(identity),
    924         catalog,
    925         schema_catalog,
    926         OpenMode::Initialize,
    927         policy,
    928         Some(authority),
    929         None,
    930     )
    931     .await
    932 }
    933 
    934 #[cfg(any(target_os = "linux", target_os = "macos"))]
    935 impl PoolConnectionValidation {
    936     async fn validate(
    937         &self,
    938         connection: &mut SqliteConnection,
    939     ) -> Result<(), PoolConnectionValidationFailure> {
    940         self.validate_authority()?;
    941         let policy_result = connection_policy_matches(connection, self.mode, self.policy).await;
    942         self.validate_authority()?;
    943         if !policy_result.map_err(PoolConnectionValidationFailure::Pragma)? {
    944             return Err(PoolConnectionValidationFailure::PolicyMismatch);
    945         }
    946         let metadata_result =
    947             crate::metadata::verify_database_metadata(connection, &self.identity).await;
    948         self.validate_authority()?;
    949         metadata_result.map_err(|_| PoolConnectionValidationFailure::Metadata)?;
    950         let migration_result = crate::migration::verify_migration_history(
    951             connection,
    952             &self.catalog,
    953             &self.schema_catalog,
    954             self.mode == OpenMode::ReadOnlyInspection,
    955         )
    956         .await;
    957         self.validate_authority()?;
    958         migration_result.map_err(|error| {
    959             if error.kind() == ServiceSqliteErrorKind::Integrity {
    960                 PoolConnectionValidationFailure::Integrity
    961             } else {
    962                 PoolConnectionValidationFailure::Migration
    963             }
    964         })?;
    965         Ok(())
    966     }
    967 
    968     fn validate_authority(&self) -> Result<(), PoolConnectionValidationFailure> {
    969         self.binding
    970             .validate(&self.paths)
    971             .map_err(|_| PoolConnectionValidationFailure::Authority)
    972     }
    973 }
    974 
    975 #[cfg(any(target_os = "linux", target_os = "macos"))]
    976 #[allow(clippy::too_many_arguments)]
    977 async fn open_connection_pool(
    978     paths: &ServiceSqlitePaths,
    979     expectation: ServiceDatabaseExpectation<'_>,
    980     catalog: &MigrationCatalog,
    981     schema_catalog: &SchemaCatalog,
    982     mode: OpenMode,
    983     policy: ServiceSqliteConnectionOptions,
    984     authority: Option<WriterAuthority>,
    985     inspection_guard: Option<ReadOnlyInspectionGuard>,
    986 ) -> Result<PrivateConnectionPool, ServiceSqliteError> {
    987     if !expectation.matches_paths(paths) {
    988         return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
    989     }
    990     if expectation.supported_state_schema_version().get() != catalog.current_version() {
    991         return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Migration));
    992     }
    993     if !schema_catalog.matches_migrations(catalog) {
    994         return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
    995     }
    996     let recovery_guard = match (mode, authority.as_ref(), inspection_guard.as_ref()) {
    997         (OpenMode::Initialize, Some(authority), None) => {
    998             authority.validate_for(paths)?;
    999             let result = crate::restore::refuse_unresolved_recovery(authority.directory());
   1000             authority.validate_for(paths)?;
   1001             result
   1002         }
   1003         (OpenMode::ReadWriteExisting, Some(authority), None) => match expectation {
   1004             ServiceDatabaseExpectation::Exact(identity) => {
   1005                 crate::restore::recover_for_open(paths, identity, authority)
   1006             }
   1007             ServiceDatabaseExpectation::Existing(intent) => {
   1008                 crate::restore::recover_for_open_with_intent(paths, intent, authority)
   1009             }
   1010         },
   1011         (OpenMode::ReadOnlyInspection, None, Some(inspection_guard)) => {
   1012             inspection_guard.validate_for(paths)?;
   1013             let result = crate::restore::refuse_unresolved_recovery(&inspection_guard.directory);
   1014             inspection_guard.validate_for(paths)?;
   1015             result
   1016         }
   1017         _ => Err(connection_error(
   1018             ServiceSqliteErrorKind::Authority,
   1019             ConnectionFailureKind::AuthorityMismatch,
   1020         )),
   1021     };
   1022     recovery_guard?;
   1023     let binding = match (mode, authority.as_ref(), inspection_guard.as_ref()) {
   1024         (OpenMode::Initialize | OpenMode::ReadWriteExisting, Some(authority), None) => {
   1025             authority.validate_for(paths)?;
   1026             DirectoryBinding::capture(authority.directory(), paths)?
   1027         }
   1028         (OpenMode::ReadOnlyInspection, None, Some(inspection_guard)) => {
   1029             DirectoryBinding::capture(&inspection_guard.directory, paths)?
   1030         }
   1031         _ => {
   1032             return Err(connection_error(
   1033                 ServiceSqliteErrorKind::Authority,
   1034                 ConnectionFailureKind::AuthorityMismatch,
   1035             ));
   1036         }
   1037     };
   1038 
   1039     let connect_options = sqlite_connect_options(paths, mode, policy);
   1040     binding.validate(paths)?;
   1041     let preflight_result = SqliteConnection::connect_with(&connect_options).await;
   1042     binding.validate(paths)?;
   1043     let mut preflight = preflight_result
   1044         .map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source))?;
   1045     let preflight_policy = verify_connection_policy(&mut preflight, mode, policy).await;
   1046     binding.validate(paths)?;
   1047     preflight_policy.map_err(|source| connection_source(ServiceSqliteErrorKind::Pragma, source))?;
   1048     let preflight_metadata = expectation.verify_metadata(&mut preflight).await;
   1049     binding.validate(paths)?;
   1050     let preflight_metadata = preflight_metadata?;
   1051     let identity = expectation.exact_identity(&preflight_metadata);
   1052     let preflight_history = crate::migration::verify_migration_history(
   1053         &mut preflight,
   1054         catalog,
   1055         schema_catalog,
   1056         mode == OpenMode::ReadOnlyInspection,
   1057     )
   1058     .await;
   1059     binding.validate(paths)?;
   1060     preflight_history?;
   1061     let preflight_close = preflight.close().await;
   1062     binding.validate(paths)?;
   1063     preflight_close.map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source))?;
   1064 
   1065     let retained_binding = binding.clone();
   1066     let pool_binding = binding;
   1067     let retained_catalog = catalog.clone();
   1068     let retained_schema_catalog = schema_catalog.clone();
   1069     let retained_identity = identity.clone();
   1070     let authority_failure = Arc::new(AtomicBool::new(false));
   1071     let metadata_failure = Arc::new(AtomicBool::new(false));
   1072     let migration_failure = Arc::new(AtomicBool::new(false));
   1073     let integrity_failure = Arc::new(AtomicBool::new(false));
   1074     let pragma_failure = Arc::new(AtomicBool::new(false));
   1075     let validation = PoolConnectionValidation {
   1076         binding: retained_binding.clone(),
   1077         paths: paths.clone(),
   1078         identity: identity.clone(),
   1079         catalog: catalog.clone(),
   1080         schema_catalog: schema_catalog.clone(),
   1081         mode,
   1082         policy,
   1083     };
   1084     let after_validation = validation.clone();
   1085     let before_validation = validation;
   1086     let flags = PoolConnectionFailureFlags {
   1087         authority: Arc::clone(&authority_failure),
   1088         metadata: Arc::clone(&metadata_failure),
   1089         migration: Arc::clone(&migration_failure),
   1090         integrity: Arc::clone(&integrity_failure),
   1091         pragma: Arc::clone(&pragma_failure),
   1092     };
   1093     let after_flags = flags.clone();
   1094     let before_flags = flags.clone();
   1095     let pool_result = SqlitePoolOptions::new()
   1096         .min_connections(1)
   1097         .max_connections(policy.max_connections())
   1098         .acquire_timeout(policy.busy_timeout())
   1099         .idle_timeout(None)
   1100         .max_lifetime(None)
   1101         .test_before_acquire(true)
   1102         .after_connect(move |connection, _metadata| {
   1103             let validation = after_validation.clone();
   1104             let flags = after_flags.clone();
   1105             Box::pin(async move {
   1106                 validation.validate(connection).await.map_err(|failure| {
   1107                     flags.record(&failure);
   1108                     failure.into_sqlx()
   1109                 })
   1110             })
   1111         })
   1112         .before_acquire(move |connection, _metadata| {
   1113             let validation = before_validation.clone();
   1114             let flags = before_flags.clone();
   1115             Box::pin(async move {
   1116                 match validation.validate(connection).await {
   1117                     Ok(()) => Ok(true),
   1118                     Err(PoolConnectionValidationFailure::PolicyMismatch) => {
   1119                         flags.record(&PoolConnectionValidationFailure::PolicyMismatch);
   1120                         Ok(false)
   1121                     }
   1122                     Err(failure) => {
   1123                         flags.record(&failure);
   1124                         Err(failure.into_sqlx())
   1125                     }
   1126                 }
   1127             })
   1128         })
   1129         .connect_with(connect_options)
   1130         .await;
   1131     pool_binding.validate(paths)?;
   1132     let pool = pool_result.map_err(|source| connection_source(flags.kind(), source))?;
   1133 
   1134     Ok(PrivateConnectionPool {
   1135         pool,
   1136         binding: retained_binding,
   1137         paths: paths.clone(),
   1138         identity: retained_identity,
   1139         catalog: retained_catalog,
   1140         schema_catalog: retained_schema_catalog,
   1141         mode,
   1142         policy,
   1143         resources: Arc::new(Mutex::new(PrivateConnectionResources {
   1144             authority,
   1145             inspection_guard,
   1146         })),
   1147         close_driver: tokio::sync::Mutex::new(PrivateCloseDriver::Pending),
   1148         #[cfg(test)]
   1149         close_phase: std::sync::atomic::AtomicU8::new(0),
   1150         authority_failure,
   1151         metadata_failure,
   1152         migration_failure,
   1153         integrity_failure,
   1154         pragma_failure,
   1155     })
   1156 }
   1157 
   1158 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1159 #[allow(
   1160     dead_code,
   1161     reason = "Step 056 keeps SQLx options private until the Step 061 host boundary"
   1162 )]
   1163 fn sqlite_connect_options(
   1164     paths: &ServiceSqlitePaths,
   1165     mode: OpenMode,
   1166     policy: ServiceSqliteConnectionOptions,
   1167 ) -> SqliteConnectOptions {
   1168     let mut options = SqliteConnectOptions::new()
   1169         .filename(paths.state_database())
   1170         .read_only(mode == OpenMode::ReadOnlyInspection)
   1171         .create_if_missing(false)
   1172         .foreign_keys(true)
   1173         .busy_timeout(policy.busy_timeout())
   1174         .synchronous(SqliteSynchronous::Full)
   1175         .pragma("trusted_schema", "OFF")
   1176         .pragma(
   1177             "query_only",
   1178             if mode == OpenMode::ReadOnlyInspection {
   1179                 "ON"
   1180             } else {
   1181                 "OFF"
   1182             },
   1183         )
   1184         .statement_cache_capacity(STATEMENT_CACHE_CAPACITY)
   1185         .command_buffer_size(COMMAND_BUFFER_CAPACITY)
   1186         .row_buffer_size(ROW_BUFFER_CAPACITY)
   1187         .disable_statement_logging();
   1188     if mode != OpenMode::ReadOnlyInspection {
   1189         options = options.journal_mode(SqliteJournalMode::Wal);
   1190     } else {
   1191         options = options.immutable(true);
   1192     }
   1193     options
   1194 }
   1195 
   1196 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1197 #[allow(
   1198     dead_code,
   1199     reason = "Step 056 keeps pragma verification private until the Step 061 host boundary"
   1200 )]
   1201 async fn verify_connection_policy(
   1202     connection: &mut SqliteConnection,
   1203     mode: OpenMode,
   1204     policy: ServiceSqliteConnectionOptions,
   1205 ) -> Result<(), sqlx::Error> {
   1206     if connection_policy_matches(connection, mode, policy).await? {
   1207         Ok(())
   1208     } else {
   1209         Err(sqlx::Error::Protocol(
   1210             "SQLite connection policy mismatch".to_owned(),
   1211         ))
   1212     }
   1213 }
   1214 
   1215 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1216 #[allow(
   1217     dead_code,
   1218     reason = "Step 056 keeps pragma verification private until the Step 061 host boundary"
   1219 )]
   1220 async fn connection_policy_matches(
   1221     connection: &mut SqliteConnection,
   1222     mode: OpenMode,
   1223     policy: ServiceSqliteConnectionOptions,
   1224 ) -> Result<bool, sqlx::Error> {
   1225     let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode")
   1226         .fetch_one(&mut *connection)
   1227         .await?;
   1228     let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous")
   1229         .fetch_one(&mut *connection)
   1230         .await?;
   1231     let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys")
   1232         .fetch_one(&mut *connection)
   1233         .await?;
   1234     let trusted_schema = sqlx::query_scalar::<_, i64>("PRAGMA trusted_schema")
   1235         .fetch_one(&mut *connection)
   1236         .await?;
   1237     let busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout")
   1238         .fetch_one(&mut *connection)
   1239         .await?;
   1240     let query_only = sqlx::query_scalar::<_, i64>("PRAGMA query_only")
   1241         .fetch_one(&mut *connection)
   1242         .await?;
   1243     Ok(connection_policy_values_match(
   1244         ConnectionPolicyValues {
   1245             journal_mode: &journal_mode,
   1246             synchronous,
   1247             foreign_keys,
   1248             trusted_schema,
   1249             busy_timeout,
   1250             query_only,
   1251         },
   1252         mode,
   1253         policy,
   1254     ))
   1255 }
   1256 
   1257 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1258 #[derive(Clone, Copy)]
   1259 struct ConnectionPolicyValues<'a> {
   1260     journal_mode: &'a str,
   1261     synchronous: i64,
   1262     foreign_keys: i64,
   1263     trusted_schema: i64,
   1264     busy_timeout: i64,
   1265     query_only: i64,
   1266 }
   1267 
   1268 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1269 fn connection_policy_values_match(
   1270     values: ConnectionPolicyValues<'_>,
   1271     mode: OpenMode,
   1272     policy: ServiceSqliteConnectionOptions,
   1273 ) -> bool {
   1274     // SQLite reports `delete` for immutable handles; the inspection guard
   1275     // independently verifies WAL read/write header bytes before this opens.
   1276     let journal_mode_matches = if mode == OpenMode::ReadOnlyInspection {
   1277         values.journal_mode.eq_ignore_ascii_case("delete")
   1278     } else {
   1279         values.journal_mode.eq_ignore_ascii_case("wal")
   1280     };
   1281     crate::all_constraints([
   1282         journal_mode_matches,
   1283         values.synchronous == 2,
   1284         values.foreign_keys == 1,
   1285         values.trusted_schema == 0,
   1286         values.busy_timeout == policy.busy_timeout_milliseconds(),
   1287         values.query_only == i64::from(mode == OpenMode::ReadOnlyInspection),
   1288     ])
   1289 }
   1290 
   1291 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1292 #[allow(
   1293     dead_code,
   1294     reason = "Step 056 keeps connection failures private until the Step 061 host boundary"
   1295 )]
   1296 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
   1297 enum ConnectionFailureKind {
   1298     UnsupportedMode,
   1299     AuthorityMismatch,
   1300     InspectionUnavailable,
   1301     InspectionContended,
   1302     CheckpointBusy,
   1303 }
   1304 
   1305 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1306 impl fmt::Display for ConnectionFailureKind {
   1307     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1308         formatter.write_str(match self {
   1309             Self::UnsupportedMode => "SQLite initialize mode requires reserved state",
   1310             Self::AuthorityMismatch => "SQLite writer authority is missing or mismatched",
   1311             Self::InspectionUnavailable => "SQLite inspection authority is unavailable",
   1312             Self::InspectionContended => "SQLite inspection requires an offline writer",
   1313             Self::CheckpointBusy => "SQLite close checkpoint could not drain active readers",
   1314         })
   1315     }
   1316 }
   1317 
   1318 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1319 impl Error for ConnectionFailureKind {}
   1320 
   1321 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1322 #[allow(
   1323     dead_code,
   1324     reason = "Step 056 keeps connection failures private until the Step 061 host boundary"
   1325 )]
   1326 fn connection_error(
   1327     kind: ServiceSqliteErrorKind,
   1328     cause: ConnectionFailureKind,
   1329 ) -> ServiceSqliteError {
   1330     ServiceSqliteError::with_source(kind, cause)
   1331 }
   1332 
   1333 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1334 fn require_connection_condition(
   1335     condition: bool,
   1336     kind: ServiceSqliteErrorKind,
   1337     cause: ConnectionFailureKind,
   1338 ) -> Result<(), ServiceSqliteError> {
   1339     condition
   1340         .then_some(())
   1341         .ok_or_else(|| connection_error(kind, cause))
   1342 }
   1343 
   1344 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1345 #[allow(
   1346     dead_code,
   1347     reason = "Step 056 keeps dependency causes private until the Step 061 host boundary"
   1348 )]
   1349 fn connection_source(kind: ServiceSqliteErrorKind, cause: sqlx::Error) -> ServiceSqliteError {
   1350     ServiceSqliteError::with_source(kind, cause)
   1351 }
   1352 
   1353 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1354 struct ReadOnlyInspectionGuard {
   1355     lock: Option<File>,
   1356     lock_device: u64,
   1357     lock_inode: u64,
   1358     directory: File,
   1359     directory_device: u64,
   1360     directory_inode: u64,
   1361     _database: File,
   1362 }
   1363 
   1364 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1365 impl ReadOnlyInspectionGuard {
   1366     fn acquire(paths: &ServiceSqlitePaths) -> Result<Self, ServiceSqliteError> {
   1367         use rustix::{
   1368             fs::{AtFlags, FileType, Mode, OFlags, fstat, open, openat, statat},
   1369             process::geteuid,
   1370         };
   1371 
   1372         let directory = open(
   1373             paths
   1374                 .state_lock()
   1375                 .parent()
   1376                 .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?,
   1377             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1378             Mode::empty(),
   1379         )
   1380         .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1381         let directory_status = fstat(&directory)
   1382             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1383         require_connection_condition(
   1384             crate::native_metadata::secure_directory(
   1385                 FileType::from_raw_mode(directory_status.st_mode).is_dir(),
   1386                 directory_status.st_uid,
   1387                 geteuid().as_raw(),
   1388                 crate::native_metadata::mode(directory_status.st_mode),
   1389             ),
   1390             ServiceSqliteErrorKind::Authority,
   1391             ConnectionFailureKind::InspectionUnavailable,
   1392         )?;
   1393         let directory = File::from(directory);
   1394         let lock = openat(
   1395             &directory,
   1396             radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME,
   1397             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
   1398             Mode::empty(),
   1399         )
   1400         .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1401         let lock_status = fstat(&lock)
   1402             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1403         require_connection_condition(
   1404             crate::native_metadata::exact_regular_file(
   1405                 FileType::from_raw_mode(lock_status.st_mode).is_file(),
   1406                 crate::native_metadata::link_count(lock_status.st_nlink),
   1407                 lock_status.st_uid,
   1408                 geteuid().as_raw(),
   1409                 crate::native_metadata::mode(lock_status.st_mode),
   1410             ),
   1411             ServiceSqliteErrorKind::Authority,
   1412             ConnectionFailureKind::InspectionUnavailable,
   1413         )?;
   1414         let lock = File::from(lock);
   1415         FileExt::try_lock_shared(&lock).map_err(|error| {
   1416             if error.kind() == std::io::ErrorKind::WouldBlock {
   1417                 inspection_error(ConnectionFailureKind::InspectionContended)
   1418             } else {
   1419                 inspection_error(ConnectionFailureKind::InspectionUnavailable)
   1420             }
   1421         })?;
   1422         crate::restore::refuse_unresolved_recovery(&directory)?;
   1423         for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] {
   1424             match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) {
   1425                 Err(error) if error == rustix::io::Errno::NOENT => {}
   1426                 Ok(_) | Err(_) => {
   1427                     return Err(inspection_error(
   1428                         ConnectionFailureKind::InspectionUnavailable,
   1429                     ));
   1430                 }
   1431             }
   1432         }
   1433         let database = openat(
   1434             &directory,
   1435             radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME,
   1436             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1437             Mode::empty(),
   1438         )
   1439         .map_err(|_| {
   1440             connection_error(
   1441                 ServiceSqliteErrorKind::Open,
   1442                 ConnectionFailureKind::InspectionUnavailable,
   1443             )
   1444         })?;
   1445         let database_status = fstat(&database).map_err(|_| {
   1446             connection_error(
   1447                 ServiceSqliteErrorKind::Open,
   1448                 ConnectionFailureKind::InspectionUnavailable,
   1449             )
   1450         })?;
   1451         require_connection_condition(
   1452             crate::native_metadata::exact_regular_file(
   1453                 FileType::from_raw_mode(database_status.st_mode).is_file(),
   1454                 crate::native_metadata::link_count(database_status.st_nlink),
   1455                 database_status.st_uid,
   1456                 geteuid().as_raw(),
   1457                 crate::native_metadata::mode(database_status.st_mode),
   1458             ),
   1459             ServiceSqliteErrorKind::Open,
   1460             ConnectionFailureKind::InspectionUnavailable,
   1461         )?;
   1462         let database = File::from(database);
   1463         let mut sqlite_header = [0_u8; 20];
   1464         std::os::unix::fs::FileExt::read_exact_at(&database, &mut sqlite_header, 0).map_err(
   1465             |_| {
   1466                 connection_error(
   1467                     ServiceSqliteErrorKind::Open,
   1468                     ConnectionFailureKind::InspectionUnavailable,
   1469                 )
   1470             },
   1471         )?;
   1472         require_connection_condition(
   1473             crate::native_metadata::sqlite_wal_header(&sqlite_header),
   1474             ServiceSqliteErrorKind::Pragma,
   1475             ConnectionFailureKind::InspectionUnavailable,
   1476         )?;
   1477         Ok(Self {
   1478             lock: Some(lock),
   1479             lock_device: crate::native_metadata::device(lock_status.st_dev)
   1480                 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?,
   1481             lock_inode: lock_status.st_ino,
   1482             directory,
   1483             directory_device: crate::native_metadata::device(directory_status.st_dev)
   1484                 .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?,
   1485             directory_inode: directory_status.st_ino,
   1486             _database: database,
   1487         })
   1488     }
   1489 
   1490     fn validate_for(&self, paths: &ServiceSqlitePaths) -> Result<(), ServiceSqliteError> {
   1491         use rustix::{
   1492             fs::{AtFlags, FileType, Mode, OFlags, fstat, open, openat, statat},
   1493             process::geteuid,
   1494         };
   1495 
   1496         let directory_path = paths
   1497             .state_lock()
   1498             .parent()
   1499             .filter(|parent| Some(*parent) == paths.state_database().parent())
   1500             .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1501         let directory = open(
   1502             directory_path,
   1503             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1504             Mode::empty(),
   1505         )
   1506         .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1507         let directory_status = fstat(&directory)
   1508             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1509         let held_directory_status = fstat(&self.directory)
   1510             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1511         let directory_device = crate::native_metadata::device(directory_status.st_dev)
   1512             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1513         let held_directory_device = crate::native_metadata::device(held_directory_status.st_dev)
   1514             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1515         require_connection_condition(
   1516             crate::all_constraints([
   1517                 crate::native_metadata::secure_directory(
   1518                     FileType::from_raw_mode(directory_status.st_mode).is_dir(),
   1519                     directory_status.st_uid,
   1520                     geteuid().as_raw(),
   1521                     crate::native_metadata::mode(directory_status.st_mode),
   1522                 ),
   1523                 crate::native_metadata::secure_directory(
   1524                     FileType::from_raw_mode(held_directory_status.st_mode).is_dir(),
   1525                     held_directory_status.st_uid,
   1526                     geteuid().as_raw(),
   1527                     crate::native_metadata::mode(held_directory_status.st_mode),
   1528                 ),
   1529                 crate::native_metadata::identity_pair_matches(
   1530                     held_directory_device,
   1531                     held_directory_status.st_ino,
   1532                     directory_device,
   1533                     directory_status.st_ino,
   1534                     self.directory_device,
   1535                     self.directory_inode,
   1536                 ),
   1537             ]),
   1538             ServiceSqliteErrorKind::Authority,
   1539             ConnectionFailureKind::InspectionUnavailable,
   1540         )?;
   1541 
   1542         let lock = openat(
   1543             &directory,
   1544             radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME,
   1545             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
   1546             Mode::empty(),
   1547         )
   1548         .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1549         let lock_status = fstat(&lock)
   1550             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1551         let held_lock = self
   1552             .lock
   1553             .as_ref()
   1554             .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1555         let held_lock_status = fstat(held_lock)
   1556             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1557         let lock_device = crate::native_metadata::device(lock_status.st_dev)
   1558             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1559         let held_lock_device = crate::native_metadata::device(held_lock_status.st_dev)
   1560             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1561         require_connection_condition(
   1562             crate::all_constraints([
   1563                 crate::native_metadata::exact_regular_file(
   1564                     FileType::from_raw_mode(lock_status.st_mode).is_file(),
   1565                     crate::native_metadata::link_count(lock_status.st_nlink),
   1566                     lock_status.st_uid,
   1567                     geteuid().as_raw(),
   1568                     crate::native_metadata::mode(lock_status.st_mode),
   1569                 ),
   1570                 crate::native_metadata::exact_regular_file(
   1571                     FileType::from_raw_mode(held_lock_status.st_mode).is_file(),
   1572                     crate::native_metadata::link_count(held_lock_status.st_nlink),
   1573                     held_lock_status.st_uid,
   1574                     geteuid().as_raw(),
   1575                     crate::native_metadata::mode(held_lock_status.st_mode),
   1576                 ),
   1577                 crate::native_metadata::identity_pair_matches(
   1578                     held_lock_device,
   1579                     held_lock_status.st_ino,
   1580                     lock_device,
   1581                     lock_status.st_ino,
   1582                     self.lock_device,
   1583                     self.lock_inode,
   1584                 ),
   1585             ]),
   1586             ServiceSqliteErrorKind::Authority,
   1587             ConnectionFailureKind::InspectionUnavailable,
   1588         )?;
   1589         for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] {
   1590             match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) {
   1591                 Err(error) if error == rustix::io::Errno::NOENT => {}
   1592                 Ok(_) | Err(_) => {
   1593                     return Err(inspection_error(
   1594                         ConnectionFailureKind::InspectionUnavailable,
   1595                     ));
   1596                 }
   1597             }
   1598         }
   1599         Ok(())
   1600     }
   1601 
   1602     fn release(&mut self) -> Result<(), ServiceSqliteError> {
   1603         let Some(lock) = self.lock.as_ref() else {
   1604             return Ok(());
   1605         };
   1606         FileExt::unlock(lock)
   1607             .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?;
   1608         self.lock.take();
   1609         Ok(())
   1610     }
   1611 }
   1612 
   1613 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1614 impl Drop for ReadOnlyInspectionGuard {
   1615     fn drop(&mut self) {
   1616         let _ = self.release();
   1617     }
   1618 }
   1619 
   1620 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1621 fn inspection_error(cause: ConnectionFailureKind) -> ServiceSqliteError {
   1622     connection_error(ServiceSqliteErrorKind::Authority, cause)
   1623 }
   1624 
   1625 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1626 #[derive(Clone)]
   1627 struct DirectoryBinding {
   1628     database_path: PathBuf,
   1629     directory: Arc<File>,
   1630     directory_device: u64,
   1631     directory_inode: u64,
   1632     database: Arc<File>,
   1633     database_device: u64,
   1634     database_inode: u64,
   1635 }
   1636 
   1637 #[cfg(any(target_os = "linux", target_os = "macos"))]
   1638 impl DirectoryBinding {
   1639     fn capture(directory: &File, paths: &ServiceSqlitePaths) -> Result<Self, ServiceSqliteError> {
   1640         use rustix::{
   1641             fs::{FileType, Mode, OFlags, fstat, openat},
   1642             process::geteuid,
   1643         };
   1644 
   1645         let directory_status = fstat(directory).map_err(|_| {
   1646             connection_error(
   1647                 ServiceSqliteErrorKind::Authority,
   1648                 ConnectionFailureKind::AuthorityMismatch,
   1649             )
   1650         })?;
   1651         let database = openat(
   1652             directory,
   1653             radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME,
   1654             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1655             Mode::empty(),
   1656         )
   1657         .map_err(|error| {
   1658             connection_error(
   1659                 if error == rustix::io::Errno::NOENT {
   1660                     ServiceSqliteErrorKind::Open
   1661                 } else {
   1662                     ServiceSqliteErrorKind::Authority
   1663                 },
   1664                 ConnectionFailureKind::AuthorityMismatch,
   1665             )
   1666         })?;
   1667         let database_status = fstat(&database).map_err(|_| {
   1668             connection_error(
   1669                 ServiceSqliteErrorKind::Authority,
   1670                 ConnectionFailureKind::AuthorityMismatch,
   1671             )
   1672         })?;
   1673         require_connection_condition(
   1674             crate::native_metadata::exact_regular_file(
   1675                 FileType::from_raw_mode(database_status.st_mode).is_file(),
   1676                 crate::native_metadata::link_count(database_status.st_nlink),
   1677                 database_status.st_uid,
   1678                 geteuid().as_raw(),
   1679                 crate::native_metadata::mode(database_status.st_mode),
   1680             ),
   1681             ServiceSqliteErrorKind::Authority,
   1682             ConnectionFailureKind::AuthorityMismatch,
   1683         )?;
   1684         Ok(Self {
   1685             database_path: paths.state_database().to_path_buf(),
   1686             directory: Arc::new(directory.try_clone().map_err(|_| {
   1687                 connection_error(
   1688                     ServiceSqliteErrorKind::Authority,
   1689                     ConnectionFailureKind::AuthorityMismatch,
   1690                 )
   1691             })?),
   1692             directory_device: crate::native_metadata::device(directory_status.st_dev).map_err(
   1693                 |_| {
   1694                     connection_error(
   1695                         ServiceSqliteErrorKind::Authority,
   1696                         ConnectionFailureKind::AuthorityMismatch,
   1697                     )
   1698                 },
   1699             )?,
   1700             directory_inode: directory_status.st_ino,
   1701             database: Arc::new(File::from(database)),
   1702             database_device: crate::native_metadata::device(database_status.st_dev).map_err(
   1703                 |_| {
   1704                     connection_error(
   1705                         ServiceSqliteErrorKind::Authority,
   1706                         ConnectionFailureKind::AuthorityMismatch,
   1707                     )
   1708                 },
   1709             )?,
   1710             database_inode: database_status.st_ino,
   1711         })
   1712     }
   1713 
   1714     fn validate(&self, paths: &ServiceSqlitePaths) -> Result<(), ServiceSqliteError> {
   1715         use rustix::{
   1716             fs::{FileType, Mode, OFlags, fstat, open, openat},
   1717             process::geteuid,
   1718         };
   1719 
   1720         require_connection_condition(
   1721             self.database_path == paths.state_database(),
   1722             ServiceSqliteErrorKind::Authority,
   1723             ConnectionFailureKind::AuthorityMismatch,
   1724         )?;
   1725         let directory = open(
   1726             paths.state_database().parent().ok_or_else(|| {
   1727                 connection_error(
   1728                     ServiceSqliteErrorKind::Authority,
   1729                     ConnectionFailureKind::AuthorityMismatch,
   1730                 )
   1731             })?,
   1732             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1733             Mode::empty(),
   1734         )
   1735         .map_err(|_| {
   1736             connection_error(
   1737                 ServiceSqliteErrorKind::Authority,
   1738                 ConnectionFailureKind::AuthorityMismatch,
   1739             )
   1740         })?;
   1741         let held_directory_status = fstat(&*self.directory).map_err(|_| {
   1742             connection_error(
   1743                 ServiceSqliteErrorKind::Authority,
   1744                 ConnectionFailureKind::AuthorityMismatch,
   1745             )
   1746         })?;
   1747         let directory_status = fstat(&directory).map_err(|_| {
   1748             connection_error(
   1749                 ServiceSqliteErrorKind::Authority,
   1750                 ConnectionFailureKind::AuthorityMismatch,
   1751             )
   1752         })?;
   1753         let directory_device =
   1754             crate::native_metadata::device(directory_status.st_dev).map_err(|_| {
   1755                 connection_error(
   1756                     ServiceSqliteErrorKind::Authority,
   1757                     ConnectionFailureKind::AuthorityMismatch,
   1758                 )
   1759             })?;
   1760         let held_directory_device = crate::native_metadata::device(held_directory_status.st_dev)
   1761             .map_err(|_| {
   1762                 connection_error(
   1763                     ServiceSqliteErrorKind::Authority,
   1764                     ConnectionFailureKind::AuthorityMismatch,
   1765                 )
   1766             })?;
   1767         require_connection_condition(
   1768             crate::all_constraints([
   1769                 crate::native_metadata::secure_directory(
   1770                     FileType::from_raw_mode(directory_status.st_mode).is_dir(),
   1771                     directory_status.st_uid,
   1772                     geteuid().as_raw(),
   1773                     crate::native_metadata::mode(directory_status.st_mode),
   1774                 ),
   1775                 crate::native_metadata::secure_directory(
   1776                     FileType::from_raw_mode(held_directory_status.st_mode).is_dir(),
   1777                     held_directory_status.st_uid,
   1778                     geteuid().as_raw(),
   1779                     crate::native_metadata::mode(held_directory_status.st_mode),
   1780                 ),
   1781                 crate::native_metadata::identity_pair_matches(
   1782                     held_directory_device,
   1783                     held_directory_status.st_ino,
   1784                     directory_device,
   1785                     directory_status.st_ino,
   1786                     self.directory_device,
   1787                     self.directory_inode,
   1788                 ),
   1789             ]),
   1790             ServiceSqliteErrorKind::Authority,
   1791             ConnectionFailureKind::AuthorityMismatch,
   1792         )?;
   1793 
   1794         let database = openat(
   1795             &directory,
   1796             radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME,
   1797             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1798             Mode::empty(),
   1799         )
   1800         .map_err(|_| {
   1801             connection_error(
   1802                 ServiceSqliteErrorKind::Authority,
   1803                 ConnectionFailureKind::AuthorityMismatch,
   1804             )
   1805         })?;
   1806         let held_database_status = fstat(&*self.database).map_err(|_| {
   1807             connection_error(
   1808                 ServiceSqliteErrorKind::Authority,
   1809                 ConnectionFailureKind::AuthorityMismatch,
   1810             )
   1811         })?;
   1812         let database_status = fstat(&database).map_err(|_| {
   1813             connection_error(
   1814                 ServiceSqliteErrorKind::Authority,
   1815                 ConnectionFailureKind::AuthorityMismatch,
   1816             )
   1817         })?;
   1818         let database_device =
   1819             crate::native_metadata::device(database_status.st_dev).map_err(|_| {
   1820                 connection_error(
   1821                     ServiceSqliteErrorKind::Authority,
   1822                     ConnectionFailureKind::AuthorityMismatch,
   1823                 )
   1824             })?;
   1825         let held_database_device = crate::native_metadata::device(held_database_status.st_dev)
   1826             .map_err(|_| {
   1827                 connection_error(
   1828                     ServiceSqliteErrorKind::Authority,
   1829                     ConnectionFailureKind::AuthorityMismatch,
   1830                 )
   1831             })?;
   1832         require_connection_condition(
   1833             crate::all_constraints([
   1834                 crate::native_metadata::exact_regular_file(
   1835                     FileType::from_raw_mode(database_status.st_mode).is_file(),
   1836                     crate::native_metadata::link_count(database_status.st_nlink),
   1837                     database_status.st_uid,
   1838                     geteuid().as_raw(),
   1839                     crate::native_metadata::mode(database_status.st_mode),
   1840                 ),
   1841                 crate::native_metadata::exact_regular_file(
   1842                     FileType::from_raw_mode(held_database_status.st_mode).is_file(),
   1843                     crate::native_metadata::link_count(held_database_status.st_nlink),
   1844                     held_database_status.st_uid,
   1845                     geteuid().as_raw(),
   1846                     crate::native_metadata::mode(held_database_status.st_mode),
   1847                 ),
   1848                 crate::native_metadata::identity_pair_matches(
   1849                     held_database_device,
   1850                     held_database_status.st_ino,
   1851                     database_device,
   1852                     database_status.st_ino,
   1853                     self.database_device,
   1854                     self.database_inode,
   1855                 ),
   1856             ]),
   1857             ServiceSqliteErrorKind::Authority,
   1858             ConnectionFailureKind::AuthorityMismatch,
   1859         )?;
   1860         Ok(())
   1861     }
   1862 }
   1863 
   1864 #[cfg(test)]
   1865 mod tests {
   1866     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1867     use std::path::PathBuf;
   1868 
   1869     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1870     use std::{
   1871         collections::BTreeMap,
   1872         convert::Infallible,
   1873         fs,
   1874         num::NonZeroU32,
   1875         os::unix::fs::{MetadataExt, PermissionsExt, symlink},
   1876         sync::{
   1877             Arc,
   1878             atomic::{AtomicUsize, Ordering as AtomicOrdering},
   1879         },
   1880         time::{Duration, SystemTime},
   1881     };
   1882 
   1883     use radroots_runtime_paths::{
   1884         RadrootsHostEnvironment, RadrootsPathProfile, RadrootsPathResolver, RadrootsPlatform,
   1885         RuntimeContextBootstrap, RuntimeContextSource,
   1886     };
   1887     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1888     use radroots_storage::event::SourceGeneration;
   1889     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1890     use sha2::{Digest, Sha256};
   1891 
   1892     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1893     use crate::{ServiceDatabaseMetadata, ServiceSqliteApplicationId};
   1894     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1895     use tokio::sync::Notify;
   1896 
   1897     use super::*;
   1898 
   1899     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1900     #[test]
   1901     fn connection_failure_inventory_is_complete_and_source_free() {
   1902         for (kind, message) in [
   1903             (
   1904                 ConnectionFailureKind::UnsupportedMode,
   1905                 "SQLite initialize mode requires reserved state",
   1906             ),
   1907             (
   1908                 ConnectionFailureKind::AuthorityMismatch,
   1909                 "SQLite writer authority is missing or mismatched",
   1910             ),
   1911             (
   1912                 ConnectionFailureKind::InspectionUnavailable,
   1913                 "SQLite inspection authority is unavailable",
   1914             ),
   1915             (
   1916                 ConnectionFailureKind::InspectionContended,
   1917                 "SQLite inspection requires an offline writer",
   1918             ),
   1919             (
   1920                 ConnectionFailureKind::CheckpointBusy,
   1921                 "SQLite close checkpoint could not drain active readers",
   1922             ),
   1923         ] {
   1924             assert_eq!(kind.to_string(), message);
   1925             assert!(kind.source().is_none());
   1926             assert!(format!("{kind:?}").contains(&format!("{kind:?}")));
   1927             let error = connection_error(ServiceSqliteErrorKind::Open, kind);
   1928             assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
   1929             assert!(error.source().is_some());
   1930             assert!(require_connection_condition(true, ServiceSqliteErrorKind::Open, kind).is_ok());
   1931             let rejected =
   1932                 require_connection_condition(false, ServiceSqliteErrorKind::Authority, kind)
   1933                     .expect_err("false condition");
   1934             assert_eq!(rejected.kind(), ServiceSqliteErrorKind::Authority);
   1935             assert!(rejected.source().is_some());
   1936         }
   1937     }
   1938 
   1939     #[cfg(any(target_os = "linux", target_os = "macos"))]
   1940     #[test]
   1941     fn connection_policy_value_matrix_rejects_each_independent_drift() {
   1942         let policy = ServiceSqliteConnectionOptions::reviewed();
   1943         let writable = ConnectionPolicyValues {
   1944             journal_mode: "wal",
   1945             synchronous: 2,
   1946             foreign_keys: 1,
   1947             trusted_schema: 0,
   1948             busy_timeout: policy.busy_timeout_milliseconds(),
   1949             query_only: 0,
   1950         };
   1951         assert!(connection_policy_values_match(
   1952             writable,
   1953             OpenMode::Initialize,
   1954             policy,
   1955         ));
   1956         for values in [
   1957             ConnectionPolicyValues {
   1958                 journal_mode: "delete",
   1959                 ..writable
   1960             },
   1961             ConnectionPolicyValues {
   1962                 synchronous: 1,
   1963                 ..writable
   1964             },
   1965             ConnectionPolicyValues {
   1966                 foreign_keys: 0,
   1967                 ..writable
   1968             },
   1969             ConnectionPolicyValues {
   1970                 trusted_schema: 1,
   1971                 ..writable
   1972             },
   1973             ConnectionPolicyValues {
   1974                 busy_timeout: 1,
   1975                 ..writable
   1976             },
   1977             ConnectionPolicyValues {
   1978                 query_only: 1,
   1979                 ..writable
   1980             },
   1981         ] {
   1982             assert!(!connection_policy_values_match(
   1983                 values,
   1984                 OpenMode::Initialize,
   1985                 policy,
   1986             ));
   1987         }
   1988 
   1989         assert!(connection_policy_values_match(
   1990             ConnectionPolicyValues {
   1991                 journal_mode: "DELETE",
   1992                 query_only: 1,
   1993                 ..writable
   1994             },
   1995             OpenMode::ReadOnlyInspection,
   1996             policy,
   1997         ));
   1998         assert!(!connection_policy_values_match(
   1999             ConnectionPolicyValues {
   2000                 journal_mode: "wal",
   2001                 query_only: 1,
   2002                 ..writable
   2003             },
   2004             OpenMode::ReadOnlyInspection,
   2005             policy,
   2006         ));
   2007     }
   2008 
   2009     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2010     #[derive(Debug, PartialEq, Eq)]
   2011     struct FileSnapshot {
   2012         bytes: Vec<u8>,
   2013         length: u64,
   2014         modified: SystemTime,
   2015         mode: u32,
   2016     }
   2017 
   2018     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2019     fn directory_snapshot(directory: &Path) -> BTreeMap<String, FileSnapshot> {
   2020         fs::read_dir(directory)
   2021             .expect("read state directory")
   2022             .map(|entry| {
   2023                 let entry = entry.expect("state entry");
   2024                 let name = entry.file_name().into_string().expect("UTF-8 state entry");
   2025                 let metadata = entry.metadata().expect("state metadata");
   2026                 let snapshot = FileSnapshot {
   2027                     bytes: fs::read(entry.path()).expect("state bytes"),
   2028                     length: metadata.len(),
   2029                     modified: metadata.modified().expect("modified time"),
   2030                     mode: metadata.permissions().mode() & 0o777,
   2031                 };
   2032                 (name, snapshot)
   2033             })
   2034             .collect()
   2035     }
   2036 
   2037     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2038     fn database_metadata(paths: &ServiceSqlitePaths) -> ServiceDatabaseMetadata {
   2039         ServiceDatabaseMetadata::new(
   2040             paths,
   2041             SourceGeneration::new([7; 32]).expect("source generation"),
   2042             NonZeroU32::new(1).expect("schema version"),
   2043             1_700_000_000_000,
   2044             ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   2045         )
   2046         .expect("database metadata")
   2047     }
   2048 
   2049     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2050     fn restore_expectation(path: &Path) -> crate::restore::RestoreArtifactExpectation {
   2051         let metadata = fs::metadata(path).expect("restore artifact metadata");
   2052         crate::restore::RestoreArtifactExpectation::new(
   2053             metadata.dev(),
   2054             metadata.ino(),
   2055             metadata.len(),
   2056             Sha256::digest(fs::read(path).expect("restore artifact bytes")).into(),
   2057         )
   2058         .expect("restore artifact expectation")
   2059     }
   2060 
   2061     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2062     fn base_catalog() -> MigrationCatalog {
   2063         MigrationCatalog::new([]).expect("empty v1 catalog")
   2064     }
   2065 
   2066     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2067     fn schema_catalog(
   2068         migrations: &MigrationCatalog,
   2069         versions: Vec<Vec<crate::SchemaObject>>,
   2070     ) -> crate::SchemaCatalog {
   2071         let versions = versions
   2072             .into_iter()
   2073             .enumerate()
   2074             .map(|(index, objects)| {
   2075                 let version = u32::try_from(index + 1).expect("schema version");
   2076                 let digest =
   2077                     crate::SchemaVersionCatalog::computed_digest(version, objects.iter().cloned())
   2078                         .expect("schema digest");
   2079                 crate::SchemaVersionCatalog::new(version, objects, digest).expect("schema version")
   2080             })
   2081             .collect::<Vec<_>>();
   2082         crate::SchemaCatalog::new(migrations, versions).expect("schema catalog")
   2083     }
   2084 
   2085     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2086     fn base_schema_catalog() -> crate::SchemaCatalog {
   2087         schema_catalog(&base_catalog(), vec![Vec::new()])
   2088     }
   2089 
   2090     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2091     fn single_table_schema_catalog(name: &'static str, sql: &'static str) -> crate::SchemaCatalog {
   2092         let table = crate::SchemaObject::new(
   2093             crate::SchemaObjectKind::Table,
   2094             name,
   2095             name,
   2096             sql,
   2097             crate::SchemaObject::computed_digest(crate::SchemaObjectKind::Table, name, name, sql)
   2098                 .expect("schema table digest"),
   2099         )
   2100         .expect("schema table");
   2101         schema_catalog(&base_catalog(), vec![vec![table]])
   2102     }
   2103 
   2104     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2105     fn migration_catalog() -> MigrationCatalog {
   2106         const CREATE: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL);";
   2107         const CALLBACK_DEFINITION: &[u8] = b"callback:migration_probe:v1";
   2108         MigrationCatalog::new([
   2109             crate::MigrationDescriptor::sql(
   2110                 2,
   2111                 "create_migration_probe",
   2112                 CREATE,
   2113                 crate::MigrationChecksum::for_sql(CREATE),
   2114             )
   2115             .expect("SQL migration"),
   2116             crate::MigrationDescriptor::callback(
   2117                 3,
   2118                 "populate_migration_probe",
   2119                 CALLBACK_DEFINITION,
   2120                 crate::MigrationChecksum::for_callback(CALLBACK_DEFINITION),
   2121             )
   2122             .expect("callback migration"),
   2123         ])
   2124         .expect("migration catalog")
   2125     }
   2126 
   2127     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2128     fn migration_schema_catalog() -> crate::SchemaCatalog {
   2129         const SQL: &str = "CREATE TABLE migration_probe (value INTEGER NOT NULL)";
   2130         let table = crate::SchemaObject::new(
   2131             crate::SchemaObjectKind::Table,
   2132             "migration_probe",
   2133             "migration_probe",
   2134             SQL,
   2135             crate::SchemaObject::computed_digest(
   2136                 crate::SchemaObjectKind::Table,
   2137                 "migration_probe",
   2138                 "migration_probe",
   2139                 SQL,
   2140             )
   2141             .expect("migration schema digest"),
   2142         )
   2143         .expect("migration schema table");
   2144         schema_catalog(
   2145             &migration_catalog(),
   2146             vec![Vec::new(), vec![table.clone()], vec![table]],
   2147         )
   2148     }
   2149 
   2150     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2151     fn migration_build() -> MigrationBuildIdentity {
   2152         MigrationBuildIdentity::new(
   2153             "0.1.0-alpha",
   2154             "0123456789abcdef0123456789abcdef01234567",
   2155             "89abcdef0123456789abcdef0123456789abcdef",
   2156             "1.97.1",
   2157             "x86_64-unknown-linux-gnu",
   2158             "service-host",
   2159             1,
   2160             2,
   2161             3,
   2162             4,
   2163             5,
   2164         )
   2165         .expect("migration build")
   2166     }
   2167 
   2168     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2169     fn migration_callback_binding() -> crate::migration::MigrationCallbackBinding {
   2170         let catalog = migration_catalog();
   2171         let descriptor = &catalog.descriptors()[1];
   2172         crate::migration::MigrationCallbackBinding::new(
   2173             descriptor.target_version(),
   2174             descriptor.name(),
   2175             descriptor.checksum(),
   2176             migration_callback,
   2177         )
   2178     }
   2179 
   2180     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2181     fn migration_callback<'a>(
   2182         executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>,
   2183     ) -> crate::migration::MigrationCallbackFuture<'a> {
   2184         Box::pin(async move {
   2185             executor
   2186                 .execute("INSERT INTO migration_probe (value) VALUES (41)")
   2187                 .await
   2188         })
   2189     }
   2190 
   2191     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2192     fn counted_migration_callback<'a>(
   2193         executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>,
   2194     ) -> crate::migration::MigrationCallbackFuture<'a> {
   2195         Box::pin(async move {
   2196             CONCURRENT_CALLBACK_COUNT.fetch_add(1, AtomicOrdering::SeqCst);
   2197             executor
   2198                 .execute("INSERT INTO migration_probe (value) VALUES (41)")
   2199                 .await
   2200         })
   2201     }
   2202 
   2203     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2204     static CONCURRENT_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0);
   2205 
   2206     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2207     fn yielding_migration_callback<'a>(
   2208         executor: &'a mut crate::migration::MigrationTransactionExecutor<'_>,
   2209     ) -> crate::migration::MigrationCallbackFuture<'a> {
   2210         Box::pin(async move {
   2211             AUTHORITY_CALLBACK_COUNT.fetch_add(1, AtomicOrdering::SeqCst);
   2212             tokio::task::yield_now().await;
   2213             executor
   2214                 .execute("INSERT INTO migration_probe (value) VALUES (41)")
   2215                 .await
   2216         })
   2217     }
   2218 
   2219     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2220     static AUTHORITY_CALLBACK_COUNT: AtomicUsize = AtomicUsize::new(0);
   2221 
   2222     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2223     async fn initialized_authority(
   2224         root: &Path,
   2225         instance: &str,
   2226     ) -> (ServiceSqlitePaths, ServiceDatabaseIdentity, WriterAuthority) {
   2227         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2228             RadrootsPathProfile::RepoLocal,
   2229             Some(root.to_path_buf()),
   2230             "myc",
   2231             instance,
   2232         ))
   2233         .expect("SQLite paths");
   2234         fs::create_dir_all(paths.state_database().parent().expect("state directory"))
   2235             .expect("create state directory");
   2236         let metadata = database_metadata(&paths);
   2237         let schema_catalog = base_schema_catalog();
   2238         let authority = crate::initialize_database(
   2239             &paths,
   2240             OpenMode::Initialize,
   2241             &metadata,
   2242             &schema_catalog,
   2243             |_| Box::pin(async move { Ok::<_, Infallible>(()) }),
   2244         )
   2245         .await
   2246         .expect("initialize database");
   2247         let identity = metadata.identity();
   2248         (paths, identity, authority)
   2249     }
   2250 
   2251     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2252     async fn initialized_pool(
   2253         root: &Path,
   2254         policy: ServiceSqliteConnectionOptions,
   2255     ) -> (ServiceSqlitePaths, PrivateConnectionPool) {
   2256         let (paths, identity, authority) = initialized_authority(root, "primary").await;
   2257         let pool = open_initialized_connection_pool(
   2258             &paths,
   2259             &identity,
   2260             &base_catalog(),
   2261             &base_schema_catalog(),
   2262             policy,
   2263             authority,
   2264         )
   2265         .await
   2266         .expect("open initialized pool");
   2267         (paths, pool)
   2268     }
   2269 
   2270     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2271     async fn initialized_migration_pool(
   2272         root: &Path,
   2273         policy: ServiceSqliteConnectionOptions,
   2274         catalog: &MigrationCatalog,
   2275     ) -> (
   2276         ServiceSqlitePaths,
   2277         ServiceDatabaseIdentity,
   2278         PrivateConnectionPool,
   2279     ) {
   2280         let (paths, base_identity, authority) = initialized_authority(root, "migrations").await;
   2281         let identity = ServiceDatabaseIdentity::new(
   2282             &paths,
   2283             base_identity.source_generation(),
   2284             NonZeroU32::new(catalog.current_version()).expect("catalog version"),
   2285             base_identity.application_id(),
   2286         );
   2287         let schema_catalog = migration_schema_catalog();
   2288         let pool = open_initialized_connection_pool(
   2289             &paths,
   2290             &identity,
   2291             catalog,
   2292             &schema_catalog,
   2293             policy,
   2294             authority,
   2295         )
   2296         .await
   2297         .expect("open migration pool");
   2298         (paths, identity, pool)
   2299     }
   2300 
   2301     fn runtime_context(
   2302         profile: RadrootsPathProfile,
   2303         repo_local_root: Option<PathBuf>,
   2304         service: &str,
   2305         instance: &str,
   2306     ) -> RuntimeContext {
   2307         let profile_source = if matches!(profile, RadrootsPathProfile::RepoLocal) {
   2308             RuntimeContextSource::BootstrapCli
   2309         } else {
   2310             RuntimeContextSource::SafeDefault
   2311         };
   2312         RuntimeContext::resolve(
   2313             &RadrootsPathResolver::new(RadrootsPlatform::Linux, RadrootsHostEnvironment::default()),
   2314             RuntimeContextBootstrap::new(
   2315                 profile,
   2316                 repo_local_root,
   2317                 profile_source,
   2318                 RuntimeContextSource::BootstrapCli,
   2319             )
   2320             .expect("valid bootstrap"),
   2321             ServiceId::new(service).expect("valid service"),
   2322             InstanceId::new(instance).expect("valid instance"),
   2323         )
   2324         .expect("valid runtime context")
   2325     }
   2326 
   2327     #[test]
   2328     fn paths_bind_exact_service_host_and_repo_local_artifacts() {
   2329         let myc = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2330             RadrootsPathProfile::ServiceHost,
   2331             None,
   2332             "myc",
   2333             "primary",
   2334         ))
   2335         .expect("Myc paths");
   2336         assert_eq!(myc.service().as_str(), "myc");
   2337         assert_eq!(myc.instance().as_str(), "primary");
   2338         assert_eq!(
   2339             myc.state_database(),
   2340             Path::new("/var/lib/radroots/services/myc/primary/state.sqlite")
   2341         );
   2342         assert_eq!(
   2343             myc.state_lock(),
   2344             Path::new("/var/lib/radroots/services/myc/primary/state.lock")
   2345         );
   2346 
   2347         let rhi = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2348             RadrootsPathProfile::RepoLocal,
   2349             Some(PathBuf::from("/repo/.local/radroots")),
   2350             "rhi",
   2351             "north-01",
   2352         ))
   2353         .expect("RHI paths");
   2354         assert_eq!(rhi.service().as_str(), "rhi");
   2355         assert_eq!(rhi.instance().as_str(), "north-01");
   2356         assert_eq!(
   2357             rhi.state_database(),
   2358             Path::new("/repo/.local/radroots/data/services/rhi/north-01/state.sqlite")
   2359         );
   2360         assert_eq!(
   2361             rhi.state_lock(),
   2362             Path::new("/repo/.local/radroots/data/services/rhi/north-01/state.lock")
   2363         );
   2364 
   2365         let second = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2366             RadrootsPathProfile::RepoLocal,
   2367             Some(PathBuf::from("/repo/.local/radroots")),
   2368             "rhi",
   2369             "south-02",
   2370         ))
   2371         .expect("second RHI paths");
   2372         assert_ne!(rhi, second);
   2373         assert_ne!(rhi.state_database(), second.state_database());
   2374         assert_ne!(rhi.state_lock(), second.state_lock());
   2375     }
   2376 
   2377     #[test]
   2378     fn path_shape_failures_are_typed_path_free_and_debug_is_redacted() {
   2379         assert_eq!(
   2380             validate_state_directory(Path::new("relative/state")),
   2381             Err(ServiceSqlitePathError::RelativeStateDirectory)
   2382         );
   2383         assert_eq!(
   2384             validate_state_directory(Path::new("/")),
   2385             Err(ServiceSqlitePathError::MissingStateDirectoryParent)
   2386         );
   2387 
   2388         let error = ServiceSqlitePathError::RelativeStateDirectory;
   2389         assert_eq!(error.to_string(), "SQLite state directory must be absolute");
   2390         assert_eq!(format!("{error:?}"), "RelativeStateDirectory");
   2391 
   2392         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2393             RadrootsPathProfile::RepoLocal,
   2394             Some(PathBuf::from("/sensitive/project-root")),
   2395             "myc",
   2396             "private-instance",
   2397         ))
   2398         .expect("redacted paths");
   2399         let debug = format!("{paths:?}");
   2400         assert!(debug.contains("service: ServiceId(\"myc\")"));
   2401         assert!(debug.contains("instance: InstanceId(\"private-instance\")"));
   2402         assert!(debug.contains("state_database: \"[redacted]\""));
   2403         assert!(debug.contains("state_lock: \"[redacted]\""));
   2404         assert!(!debug.contains("sensitive"));
   2405         assert!(!debug.contains("project-root"));
   2406         assert!(!debug.contains("state.sqlite"));
   2407         assert!(!debug.contains("state.lock"));
   2408     }
   2409 
   2410     #[test]
   2411     fn open_mode_wire_inventory_and_semantics_are_exact() {
   2412         let inventory = [
   2413             (OpenMode::Initialize, "initialize", true, false, true),
   2414             (
   2415                 OpenMode::ReadWriteExisting,
   2416                 "read_write_existing",
   2417                 false,
   2418                 true,
   2419                 true,
   2420             ),
   2421             (
   2422                 OpenMode::ReadOnlyInspection,
   2423                 "read_only_inspection",
   2424                 false,
   2425                 true,
   2426                 false,
   2427             ),
   2428         ];
   2429         for (mode, wire, may_create, requires_existing, requires_writer) in inventory {
   2430             assert_eq!(
   2431                 serde_json::to_string(&mode).unwrap(),
   2432                 format!(r#""{wire}""#)
   2433             );
   2434             assert_eq!(mode.may_create(), may_create);
   2435             assert_eq!(mode.requires_existing(), requires_existing);
   2436             assert_eq!(mode.requires_writer_authority(), requires_writer);
   2437         }
   2438     }
   2439 
   2440     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2441     #[tokio::test(flavor = "current_thread")]
   2442     async fn every_connection_uses_the_exact_reviewed_pragma_policy() {
   2443         let directory = tempfile::tempdir().expect("temporary directory");
   2444         let policy = ServiceSqliteConnectionOptions::reviewed();
   2445         let (_paths, pool) = initialized_pool(directory.path(), policy).await;
   2446 
   2447         let mut connections = Vec::with_capacity(8);
   2448         for _ in 0..8 {
   2449             connections.push(pool.acquire().await.expect("pooled connection"));
   2450         }
   2451         assert_eq!(pool.pool.size(), 8);
   2452         for connection in &mut connections {
   2453             assert!(
   2454                 connection_policy_matches(connection, OpenMode::Initialize, policy)
   2455                     .await
   2456                     .unwrap()
   2457             );
   2458         }
   2459         drop(connections);
   2460         assert!(pool.close().await.expect("writer authority").is_held());
   2461     }
   2462 
   2463     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2464     #[tokio::test(flavor = "current_thread")]
   2465     async fn reused_connection_drift_is_rejected_before_checkout() {
   2466         let directory = tempfile::tempdir().expect("temporary directory");
   2467         let policy = ServiceSqliteConnectionOptions::reviewed();
   2468         let (_paths, pool) = initialized_pool(directory.path(), policy).await;
   2469 
   2470         let mut connection = pool.acquire().await.expect("pooled connection");
   2471         sqlx::query("PRAGMA foreign_keys = OFF")
   2472             .execute(&mut *connection)
   2473             .await
   2474             .expect("drift pragma");
   2475         assert_eq!(
   2476             sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys")
   2477                 .fetch_one(&mut *connection)
   2478                 .await
   2479                 .unwrap(),
   2480             0
   2481         );
   2482         drop(connection);
   2483 
   2484         let mut replacement = pool.acquire().await.expect("replacement connection");
   2485         assert!(
   2486             connection_policy_matches(&mut replacement, OpenMode::Initialize, policy)
   2487                 .await
   2488                 .unwrap()
   2489         );
   2490         drop(replacement);
   2491         let _authority = pool.close().await;
   2492     }
   2493 
   2494     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2495     #[tokio::test(flavor = "current_thread")]
   2496     async fn metadata_mismatch_fails_open_and_checkout_before_use() {
   2497         let directory = tempfile::tempdir().expect("temporary directory");
   2498         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap();
   2499         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   2500 
   2501         let mut connection = pool.acquire().await.expect("pooled connection");
   2502         sqlx::query("PRAGMA application_id = 1380209490")
   2503             .execute(&mut *connection)
   2504             .await
   2505             .expect("drift application ID");
   2506         drop(connection);
   2507         let error = pool
   2508             .acquire()
   2509             .await
   2510             .expect_err("metadata drift must prevent checkout");
   2511         assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata);
   2512         let authority = pool.close().await.expect("writer authority retained");
   2513         drop(authority);
   2514 
   2515         let wrong = ServiceDatabaseIdentity::new(
   2516             &paths,
   2517             SourceGeneration::new([8; 32]).expect("wrong generation"),
   2518             NonZeroU32::new(1).expect("schema version"),
   2519             ServiceSqliteApplicationId::new(0x5244_5351).expect("application ID"),
   2520         );
   2521         let result = open_existing_connection_pool(
   2522             &paths,
   2523             &wrong,
   2524             &base_catalog(),
   2525             &base_schema_catalog(),
   2526             OpenMode::ReadWriteExisting,
   2527             policy,
   2528         )
   2529         .await;
   2530         let Err(error) = result else {
   2531             panic!("wrong generation must fail open");
   2532         };
   2533         assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata);
   2534     }
   2535 
   2536     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2537     #[tokio::test(flavor = "current_thread")]
   2538     async fn writable_open_recovers_while_read_only_and_initialize_paths_do_not_mutate() {
   2539         let directory = tempfile::tempdir().expect("temporary directory");
   2540         let policy = ServiceSqliteConnectionOptions::reviewed();
   2541         let (paths, identity, mut authority) =
   2542             initialized_authority(directory.path(), "restore-recovery").await;
   2543         authority
   2544             .release()
   2545             .expect("release initialization authority");
   2546 
   2547         let staged = paths
   2548             .state_database()
   2549             .with_file_name(crate::restore::STAGED_FILE_NAME);
   2550         fs::copy(paths.state_database(), &staged).expect("copy exact restore stage");
   2551         fs::set_permissions(&staged, fs::Permissions::from_mode(0o600))
   2552             .expect("restrict staged database");
   2553         let authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   2554             .expect("writer authority")
   2555             .expect("writable mode retains authority");
   2556         let marker = crate::restore::RestoreRecoveryMarker::prepared(
   2557             &database_metadata(&paths),
   2558             crate::BackupManifestSha256::from_bytes([23; 32]),
   2559             restore_expectation(paths.state_database()),
   2560             restore_expectation(&staged),
   2561         )
   2562         .expect("prepared restore marker");
   2563         crate::restore::RestoreMarkerBinding::create(&paths, &authority, &marker)
   2564             .expect("persist prepared marker");
   2565         drop(authority);
   2566 
   2567         let state_directory = paths.state_database().parent().expect("state directory");
   2568         let before = directory_snapshot(state_directory);
   2569         let read_only = open_existing_connection_pool(
   2570             &paths,
   2571             &identity,
   2572             &base_catalog(),
   2573             &base_schema_catalog(),
   2574             OpenMode::ReadOnlyInspection,
   2575             policy,
   2576         )
   2577         .await;
   2578         let Err(error) = read_only else {
   2579             panic!("read-only open must not recover");
   2580         };
   2581         assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery);
   2582         assert_eq!(directory_snapshot(state_directory), before);
   2583 
   2584         let initialize = crate::initialize_database(
   2585             &paths,
   2586             OpenMode::Initialize,
   2587             &database_metadata(&paths),
   2588             &base_schema_catalog(),
   2589             |_| Box::pin(async { Ok::<_, Infallible>(()) }),
   2590         )
   2591         .await
   2592         .expect_err("initialize must not recover existing evidence");
   2593         assert_eq!(initialize.kind(), ServiceSqliteErrorKind::Recovery);
   2594         assert_eq!(directory_snapshot(state_directory), before);
   2595 
   2596         let initialized_authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   2597             .expect("writer authority")
   2598             .expect("writable mode retains authority");
   2599         let initialized = open_initialized_connection_pool(
   2600             &paths,
   2601             &identity,
   2602             &base_catalog(),
   2603             &base_schema_catalog(),
   2604             policy,
   2605             initialized_authority,
   2606         )
   2607         .await;
   2608         let Err(error) = initialized else {
   2609             panic!("initialized open must not recover");
   2610         };
   2611         assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery);
   2612         assert_eq!(directory_snapshot(state_directory), before);
   2613 
   2614         let writable = open_existing_connection_pool(
   2615             &paths,
   2616             &identity,
   2617             &base_catalog(),
   2618             &base_schema_catalog(),
   2619             OpenMode::ReadWriteExisting,
   2620             policy,
   2621         )
   2622         .await
   2623         .expect("writable open recovers before SQLite");
   2624         assert!(!staged.exists());
   2625         assert!(
   2626             !paths
   2627                 .state_database()
   2628                 .with_file_name(crate::restore::MARKER_FILE_NAME)
   2629                 .exists()
   2630         );
   2631         drop(writable.close().await);
   2632     }
   2633 
   2634     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2635     #[test]
   2636     fn connection_failure_precedence_is_exact() {
   2637         assert_eq!(
   2638             connection_failure_kind(false, false, false, false, false),
   2639             ServiceSqliteErrorKind::Open
   2640         );
   2641         assert_eq!(
   2642             connection_failure_kind(false, false, false, false, true),
   2643             ServiceSqliteErrorKind::Pragma
   2644         );
   2645         assert_eq!(
   2646             connection_failure_kind(false, false, false, true, true),
   2647             ServiceSqliteErrorKind::Integrity
   2648         );
   2649         assert_eq!(
   2650             connection_failure_kind(false, false, true, true, true),
   2651             ServiceSqliteErrorKind::Migration
   2652         );
   2653         assert_eq!(
   2654             connection_failure_kind(false, true, true, true, true),
   2655             ServiceSqliteErrorKind::Metadata
   2656         );
   2657         assert_eq!(
   2658             connection_failure_kind(true, true, true, true, true),
   2659             ServiceSqliteErrorKind::Authority
   2660         );
   2661 
   2662         for (failure, expected) in [
   2663             (
   2664                 PoolConnectionValidationFailure::Authority,
   2665                 ServiceSqliteErrorKind::Authority,
   2666             ),
   2667             (
   2668                 PoolConnectionValidationFailure::Metadata,
   2669                 ServiceSqliteErrorKind::Metadata,
   2670             ),
   2671             (
   2672                 PoolConnectionValidationFailure::Migration,
   2673                 ServiceSqliteErrorKind::Migration,
   2674             ),
   2675             (
   2676                 PoolConnectionValidationFailure::Integrity,
   2677                 ServiceSqliteErrorKind::Integrity,
   2678             ),
   2679             (
   2680                 PoolConnectionValidationFailure::PolicyMismatch,
   2681                 ServiceSqliteErrorKind::Pragma,
   2682             ),
   2683             (
   2684                 PoolConnectionValidationFailure::Pragma(sqlx::Error::Protocol(
   2685                     "test pragma query failure".to_owned(),
   2686                 )),
   2687                 ServiceSqliteErrorKind::Pragma,
   2688             ),
   2689         ] {
   2690             let flags = PoolConnectionFailureFlags {
   2691                 authority: Arc::new(AtomicBool::new(false)),
   2692                 metadata: Arc::new(AtomicBool::new(false)),
   2693                 migration: Arc::new(AtomicBool::new(false)),
   2694                 integrity: Arc::new(AtomicBool::new(false)),
   2695                 pragma: Arc::new(AtomicBool::new(false)),
   2696             };
   2697             flags.record(&failure);
   2698             assert_eq!(flags.kind(), expected);
   2699             let _ = failure.into_sqlx();
   2700         }
   2701     }
   2702 
   2703     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2704     #[tokio::test(flavor = "current_thread")]
   2705     async fn schema_drift_is_rejected_on_checkout_and_fresh_open() {
   2706         let directory = tempfile::tempdir().expect("temporary directory");
   2707         let policy = ServiceSqliteConnectionOptions::reviewed();
   2708         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   2709         let identity = database_metadata(&paths).identity();
   2710         let mut connection = pool.acquire().await.expect("connection");
   2711         sqlx::query("CREATE TABLE unlisted (value INTEGER)")
   2712             .execute(&mut *connection)
   2713             .await
   2714             .expect("create unlisted object");
   2715         drop(connection);
   2716 
   2717         let error = pool
   2718             .acquire()
   2719             .await
   2720             .expect_err("checkout must reject schema drift");
   2721         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   2722         let authority = pool.close().await.expect("writer authority");
   2723         drop(authority);
   2724 
   2725         let result = open_existing_connection_pool(
   2726             &paths,
   2727             &identity,
   2728             &base_catalog(),
   2729             &base_schema_catalog(),
   2730             OpenMode::ReadWriteExisting,
   2731             policy,
   2732         )
   2733         .await;
   2734         let Err(error) = result else {
   2735             panic!("fresh open must reject schema drift");
   2736         };
   2737         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   2738     }
   2739 
   2740     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2741     #[tokio::test(flavor = "current_thread")]
   2742     async fn migration_entry_preserves_post_open_ledger_drift_classification() {
   2743         let directory = tempfile::tempdir().expect("temporary directory");
   2744         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap();
   2745         let (_paths, pool) = initialized_pool(directory.path(), policy).await;
   2746         let build = migration_build();
   2747         let mut connection = pool.acquire().await.expect("pooled connection");
   2748         sqlx::query(
   2749             "INSERT INTO schema_migrations (
   2750                 version, name, checksum, applied_at_unix_s,
   2751                 service_version, service_commit, lib_revision, rust_version, target,
   2752                 feature_profile, config_contract_version, state_contract_version,
   2753                 admin_contract_version, status_contract_version, provider_contract_version
   2754              ) VALUES (2, 'unexpected_row', zeroblob(32), 0, ?, ?, ?, ?, ?, ?, 1, 2, 3, 4, 5)",
   2755         )
   2756         .bind(build.service_version())
   2757         .bind(build.service_commit())
   2758         .bind(build.lib_revision())
   2759         .bind(build.rust_version())
   2760         .bind(build.target())
   2761         .bind(build.feature_profile())
   2762         .execute(&mut *connection)
   2763         .await
   2764         .expect("inject post-open ledger drift");
   2765         drop(connection);
   2766 
   2767         let error = pool
   2768             .apply_migrations(
   2769                 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(),
   2770                 &build,
   2771                 &[],
   2772             )
   2773             .await
   2774             .expect_err("migration entry must reject drift before execution");
   2775         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2776         let authority = pool.close().await.expect("writer authority retained");
   2777         drop(authority);
   2778     }
   2779 
   2780     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2781     #[tokio::test(flavor = "current_thread")]
   2782     async fn pending_history_is_not_exposed_and_read_only_requires_current_state() {
   2783         let directory = tempfile::tempdir().expect("temporary directory");
   2784         let catalog = migration_catalog();
   2785         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap();
   2786         let (paths, identity, pending) =
   2787             initialized_migration_pool(directory.path(), policy, &catalog).await;
   2788 
   2789         let error = pending
   2790             .acquire()
   2791             .await
   2792             .expect_err("pending schema must not escape the private pool");
   2793         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2794         let authority = pending.close().await.expect("writer authority");
   2795         drop(authority);
   2796 
   2797         let read_only = open_existing_connection_pool(
   2798             &paths,
   2799             &identity,
   2800             &catalog,
   2801             &migration_schema_catalog(),
   2802             OpenMode::ReadOnlyInspection,
   2803             policy,
   2804         )
   2805         .await;
   2806         let Err(error) = read_only else {
   2807             panic!("read-only pending state must fail closed");
   2808         };
   2809         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   2810 
   2811         let writable = open_existing_connection_pool(
   2812             &paths,
   2813             &identity,
   2814             &catalog,
   2815             &migration_schema_catalog(),
   2816             OpenMode::ReadWriteExisting,
   2817             policy,
   2818         )
   2819         .await
   2820         .expect("reopen writable prefix");
   2821         let outcome = writable
   2822             .apply_migrations(
   2823                 MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(),
   2824                 &migration_build(),
   2825                 &[migration_callback_binding()],
   2826             )
   2827             .await
   2828             .expect("apply pending migrations");
   2829         assert_eq!(outcome.initial_version(), 1);
   2830         assert_eq!(outcome.final_version(), 3);
   2831         assert_eq!(outcome.applied_count(), 2);
   2832         let connection = writable.acquire().await.expect("current checkout");
   2833         drop(connection);
   2834         let authority = writable.close().await.expect("writer authority");
   2835         drop(authority);
   2836 
   2837         let current = open_existing_connection_pool(
   2838             &paths,
   2839             &identity,
   2840             &catalog,
   2841             &migration_schema_catalog(),
   2842             OpenMode::ReadOnlyInspection,
   2843             policy,
   2844         )
   2845         .await
   2846         .expect("current read-only state");
   2847         assert!(current.close().await.is_none());
   2848     }
   2849 
   2850     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2851     #[tokio::test(flavor = "current_thread")]
   2852     async fn concurrent_migration_attempts_serialize_and_execute_callback_once() {
   2853         CONCURRENT_CALLBACK_COUNT.store(0, AtomicOrdering::SeqCst);
   2854         let directory = tempfile::tempdir().expect("temporary directory");
   2855         let catalog = migration_catalog();
   2856         let policy = ServiceSqliteConnectionOptions::new(Duration::from_secs(2), 2).unwrap();
   2857         let (_paths, _identity, pool) =
   2858             initialized_migration_pool(directory.path(), policy, &catalog).await;
   2859         let applied_at = MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap();
   2860         let build = migration_build();
   2861         let descriptor = &catalog.descriptors()[1];
   2862         let callbacks = [crate::migration::MigrationCallbackBinding::new(
   2863             descriptor.target_version(),
   2864             descriptor.name(),
   2865             descriptor.checksum(),
   2866             counted_migration_callback,
   2867         )];
   2868 
   2869         let (first, second) = tokio::join!(
   2870             pool.apply_migrations(applied_at, &build, &callbacks),
   2871             pool.apply_migrations(applied_at, &build, &callbacks),
   2872         );
   2873         let first = first.expect("first migration attempt");
   2874         let second = second.expect("second migration attempt");
   2875         assert_eq!(
   2876             first.applied_count() + second.applied_count(),
   2877             2,
   2878             "exact catalog is committed only once"
   2879         );
   2880         assert_eq!(CONCURRENT_CALLBACK_COUNT.load(AtomicOrdering::SeqCst), 1);
   2881         let mut connection = pool.acquire().await.expect("current connection");
   2882         assert_eq!(
   2883             sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM migration_probe")
   2884                 .fetch_one(&mut *connection)
   2885                 .await
   2886                 .unwrap(),
   2887             1
   2888         );
   2889         drop(connection);
   2890         let _authority = pool.close().await;
   2891     }
   2892 
   2893     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2894     #[tokio::test(flavor = "current_thread")]
   2895     async fn authority_replacement_wins_over_in_flight_migration_failure() {
   2896         AUTHORITY_CALLBACK_COUNT.store(0, AtomicOrdering::SeqCst);
   2897         let directory = tempfile::tempdir().expect("temporary directory");
   2898         let catalog = migration_catalog();
   2899         let policy = ServiceSqliteConnectionOptions::new(Duration::from_secs(2), 1).unwrap();
   2900         let (paths, _identity, pool) =
   2901             initialized_migration_pool(directory.path(), policy, &catalog).await;
   2902         let descriptor = &catalog.descriptors()[1];
   2903         let callbacks = [crate::migration::MigrationCallbackBinding::new(
   2904             descriptor.target_version(),
   2905             descriptor.name(),
   2906             descriptor.checksum(),
   2907             yielding_migration_callback,
   2908         )];
   2909         let state_directory = paths.state_database().parent().unwrap().to_path_buf();
   2910         let displaced = directory.path().join("displaced-migration-state");
   2911         let build = migration_build();
   2912 
   2913         let application = pool.apply_migrations(
   2914             MigrationAppliedAtUnixSeconds::new(1_800_000_000).unwrap(),
   2915             &build,
   2916             &callbacks,
   2917         );
   2918         let replace = async {
   2919             while AUTHORITY_CALLBACK_COUNT.load(AtomicOrdering::SeqCst) == 0 {
   2920                 tokio::task::yield_now().await;
   2921             }
   2922             fs::rename(&state_directory, &displaced).expect("displace state directory");
   2923             fs::create_dir_all(&state_directory).expect("replace state directory");
   2924             fs::copy(displaced.join("state.sqlite"), paths.state_database())
   2925                 .expect("copy replacement database");
   2926             fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o600))
   2927                 .expect("secure replacement database");
   2928         };
   2929         let (result, ()) = tokio::join!(application, replace);
   2930         assert_eq!(
   2931             result.expect_err("authority replacement must fail").kind(),
   2932             ServiceSqliteErrorKind::Authority
   2933         );
   2934         let authority = pool.close().await.expect("retained writer authority");
   2935         assert!(authority.is_held());
   2936     }
   2937 
   2938     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2939     #[tokio::test(flavor = "current_thread")]
   2940     async fn existing_modes_never_create_missing_state_and_enforce_authority() {
   2941         let directory = tempfile::tempdir().expect("temporary directory");
   2942         let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2943             RadrootsPathProfile::RepoLocal,
   2944             Some(directory.path().to_path_buf()),
   2945             "rhi",
   2946             "default",
   2947         ))
   2948         .expect("SQLite paths");
   2949         fs::create_dir_all(paths.state_database().parent().expect("state directory"))
   2950             .expect("create state directory");
   2951         let metadata = database_metadata(&paths).identity();
   2952 
   2953         for mode in [OpenMode::ReadWriteExisting, OpenMode::ReadOnlyInspection] {
   2954             let result = open_existing_connection_pool(
   2955                 &paths,
   2956                 &metadata,
   2957                 &base_catalog(),
   2958                 &base_schema_catalog(),
   2959                 mode,
   2960                 ServiceSqliteConnectionOptions::reviewed(),
   2961             )
   2962             .await;
   2963             let Err(error) = result else {
   2964                 panic!("missing database must fail");
   2965             };
   2966             assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
   2967             assert!(!paths.state_database().exists());
   2968         }
   2969         let result = open_existing_connection_pool(
   2970             &paths,
   2971             &metadata,
   2972             &base_catalog(),
   2973             &base_schema_catalog(),
   2974             OpenMode::Initialize,
   2975             ServiceSqliteConnectionOptions::reviewed(),
   2976         )
   2977         .await;
   2978         let Err(error) = result else {
   2979             panic!("initialize needs reserved state");
   2980         };
   2981         assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
   2982     }
   2983 
   2984     #[cfg(any(target_os = "linux", target_os = "macos"))]
   2985     #[tokio::test(flavor = "current_thread")]
   2986     async fn initialized_pool_rejects_mismatched_paths_and_rebound_directory() {
   2987         let directory = tempfile::tempdir().expect("temporary directory");
   2988         let (_paths, metadata, authority) =
   2989             initialized_authority(directory.path(), "primary").await;
   2990         let other = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   2991             RadrootsPathProfile::RepoLocal,
   2992             Some(directory.path().to_path_buf()),
   2993             "myc",
   2994             "secondary",
   2995         ))
   2996         .expect("other SQLite paths");
   2997         fs::create_dir_all(
   2998             other
   2999                 .state_database()
   3000                 .parent()
   3001                 .expect("other state directory"),
   3002         )
   3003         .expect("create other state directory");
   3004         let result = open_initialized_connection_pool(
   3005             &other,
   3006             &metadata,
   3007             &base_catalog(),
   3008             &base_schema_catalog(),
   3009             ServiceSqliteConnectionOptions::reviewed(),
   3010             authority,
   3011         )
   3012         .await;
   3013         let Err(error) = result else {
   3014             panic!("mismatched paths must fail");
   3015         };
   3016         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3017 
   3018         let (paths, metadata, authority) = initialized_authority(directory.path(), "rebound").await;
   3019         let state_directory = paths.state_database().parent().expect("state directory");
   3020         let displaced = directory.path().join("displaced-state");
   3021         fs::rename(state_directory, &displaced).expect("displace state directory");
   3022         fs::create_dir_all(state_directory).expect("replace state directory");
   3023         fs::copy(displaced.join("state.sqlite"), paths.state_database())
   3024             .expect("copy replacement database");
   3025         let result = open_initialized_connection_pool(
   3026             &paths,
   3027             &metadata,
   3028             &base_catalog(),
   3029             &base_schema_catalog(),
   3030             ServiceSqliteConnectionOptions::reviewed(),
   3031             authority,
   3032         )
   3033         .await;
   3034         let Err(error) = result else {
   3035             panic!("rebound state directory must fail");
   3036         };
   3037         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3038     }
   3039 
   3040     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3041     #[tokio::test(flavor = "current_thread")]
   3042     async fn pool_preflight_rejects_identity_version_and_schema_catalog_drift() {
   3043         let directory = tempfile::tempdir().expect("temporary directory");
   3044 
   3045         let (paths, identity, authority) =
   3046             initialized_authority(directory.path(), "identity-drift").await;
   3047         let other_paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(
   3048             RadrootsPathProfile::RepoLocal,
   3049             Some(directory.path().to_path_buf()),
   3050             "myc",
   3051             "other-identity",
   3052         ))
   3053         .expect("other paths");
   3054         let wrong_identity = ServiceDatabaseIdentity::new(
   3055             &other_paths,
   3056             identity.source_generation(),
   3057             identity.supported_state_schema_version(),
   3058             identity.application_id(),
   3059         );
   3060         let Err(error) = open_connection_pool(
   3061             &paths,
   3062             ServiceDatabaseExpectation::Exact(&wrong_identity),
   3063             &base_catalog(),
   3064             &base_schema_catalog(),
   3065             OpenMode::Initialize,
   3066             ServiceSqliteConnectionOptions::reviewed(),
   3067             Some(authority),
   3068             None,
   3069         )
   3070         .await
   3071         else {
   3072             panic!("identity drift must fail");
   3073         };
   3074         assert_eq!(error.kind(), ServiceSqliteErrorKind::Metadata);
   3075 
   3076         let (paths, identity, authority) =
   3077             initialized_authority(directory.path(), "version-drift").await;
   3078         let newer_identity = ServiceDatabaseIdentity::new(
   3079             &paths,
   3080             identity.source_generation(),
   3081             NonZeroU32::new(2).expect("newer schema"),
   3082             identity.application_id(),
   3083         );
   3084         let Err(error) = open_connection_pool(
   3085             &paths,
   3086             ServiceDatabaseExpectation::Exact(&newer_identity),
   3087             &base_catalog(),
   3088             &base_schema_catalog(),
   3089             OpenMode::Initialize,
   3090             ServiceSqliteConnectionOptions::reviewed(),
   3091             Some(authority),
   3092             None,
   3093         )
   3094         .await
   3095         else {
   3096             panic!("version drift must fail");
   3097         };
   3098         assert_eq!(error.kind(), ServiceSqliteErrorKind::Migration);
   3099 
   3100         let (paths, identity, authority) =
   3101             initialized_authority(directory.path(), "schema-drift-preflight").await;
   3102         let Err(error) = open_connection_pool(
   3103             &paths,
   3104             ServiceDatabaseExpectation::Exact(&identity),
   3105             &base_catalog(),
   3106             &migration_schema_catalog(),
   3107             OpenMode::Initialize,
   3108             ServiceSqliteConnectionOptions::reviewed(),
   3109             Some(authority),
   3110             None,
   3111         )
   3112         .await
   3113         else {
   3114             panic!("schema catalog drift must fail");
   3115         };
   3116         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   3117     }
   3118 
   3119     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3120     #[tokio::test(flavor = "current_thread")]
   3121     async fn pool_checkout_rejects_state_directory_replacement_before_growth() {
   3122         let directory = tempfile::tempdir().expect("temporary directory");
   3123         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap();
   3124         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   3125         assert_eq!(pool.pool.size(), 1);
   3126 
   3127         let state_directory = paths.state_database().parent().expect("state directory");
   3128         let displaced = directory.path().join("displaced-live-state");
   3129         fs::rename(state_directory, &displaced).expect("displace live state directory");
   3130         fs::create_dir_all(state_directory).expect("replace live state directory");
   3131         fs::copy(displaced.join("state.sqlite"), paths.state_database())
   3132             .expect("copy replacement database");
   3133 
   3134         let error = pool
   3135             .acquire()
   3136             .await
   3137             .expect_err("rebound directory must prevent checkout and pool growth");
   3138         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3139         assert_eq!(pool.pool.size(), 1);
   3140         let authority = pool.close().await.expect("writer authority retained");
   3141         assert!(authority.is_held());
   3142     }
   3143 
   3144     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3145     #[tokio::test(flavor = "current_thread")]
   3146     async fn writable_existing_rejects_database_symlink_and_hardlink() {
   3147         let directory = tempfile::tempdir().expect("temporary directory");
   3148         let policy = ServiceSqliteConnectionOptions::reviewed();
   3149 
   3150         let (symlink_paths, symlink_metadata, symlink_authority) =
   3151             initialized_authority(directory.path(), "symlink-database").await;
   3152         drop(symlink_authority);
   3153         let symlink_backing = symlink_paths
   3154             .state_database()
   3155             .parent()
   3156             .expect("state directory")
   3157             .join("backing.sqlite");
   3158         fs::rename(symlink_paths.state_database(), &symlink_backing)
   3159             .expect("displace symlink database");
   3160         symlink(&symlink_backing, symlink_paths.state_database()).expect("database symlink");
   3161         let symlink_result = open_existing_connection_pool(
   3162             &symlink_paths,
   3163             &symlink_metadata,
   3164             &base_catalog(),
   3165             &base_schema_catalog(),
   3166             OpenMode::ReadWriteExisting,
   3167             policy,
   3168         )
   3169         .await;
   3170         let Err(symlink_error) = symlink_result else {
   3171             panic!("database symlink must fail");
   3172         };
   3173         assert_eq!(symlink_error.kind(), ServiceSqliteErrorKind::Authority);
   3174 
   3175         let (hardlink_paths, hardlink_metadata, hardlink_authority) =
   3176             initialized_authority(directory.path(), "hardlink-database").await;
   3177         drop(hardlink_authority);
   3178         let hardlink_alias = hardlink_paths
   3179             .state_database()
   3180             .parent()
   3181             .expect("state directory")
   3182             .join("alias.sqlite");
   3183         fs::hard_link(hardlink_paths.state_database(), hardlink_alias).expect("database hard link");
   3184         let hardlink_result = open_existing_connection_pool(
   3185             &hardlink_paths,
   3186             &hardlink_metadata,
   3187             &base_catalog(),
   3188             &base_schema_catalog(),
   3189             OpenMode::ReadWriteExisting,
   3190             policy,
   3191         )
   3192         .await;
   3193         let Err(hardlink_error) = hardlink_result else {
   3194             panic!("database hard link must fail");
   3195         };
   3196         assert_eq!(hardlink_error.kind(), ServiceSqliteErrorKind::Authority);
   3197     }
   3198 
   3199     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3200     #[tokio::test(flavor = "current_thread")]
   3201     async fn lazy_pool_growth_rejects_same_directory_database_replacement() {
   3202         let directory = tempfile::tempdir().expect("temporary directory");
   3203         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 2).unwrap();
   3204         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   3205         let held = pool.acquire().await.expect("hold initial connection");
   3206 
   3207         let displaced = paths
   3208             .state_database()
   3209             .parent()
   3210             .expect("state directory")
   3211             .join("displaced.sqlite");
   3212         fs::rename(paths.state_database(), &displaced).expect("displace live database");
   3213         fs::copy(&displaced, paths.state_database()).expect("copy replacement database");
   3214         fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o600))
   3215             .expect("secure replacement database");
   3216 
   3217         let error = pool
   3218             .acquire()
   3219             .await
   3220             .expect_err("database replacement must prevent lazy pool growth");
   3221         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3222         drop(held);
   3223         let authority = pool.close().await.expect("writer authority retained");
   3224         assert!(authority.is_held());
   3225     }
   3226 
   3227     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3228     #[tokio::test(flavor = "current_thread")]
   3229     async fn read_only_inspection_is_offline_query_only_and_side_effect_free() {
   3230         let directory = tempfile::tempdir().expect("temporary directory");
   3231         let policy = ServiceSqliteConnectionOptions::reviewed();
   3232         let (paths, writable) = initialized_pool(directory.path(), policy).await;
   3233         let metadata = database_metadata(&paths).identity();
   3234         let mut connection = writable.acquire().await.expect("writable connection");
   3235         sqlx::query("CREATE TABLE inspection_fixture (value INTEGER NOT NULL)")
   3236             .execute(&mut *connection)
   3237             .await
   3238             .expect("create fixture");
   3239         sqlx::query("INSERT INTO inspection_fixture (value) VALUES (41)")
   3240             .execute(&mut *connection)
   3241             .await
   3242             .expect("insert fixture");
   3243         drop(connection);
   3244         let inspection_schema_catalog = single_table_schema_catalog(
   3245             "inspection_fixture",
   3246             "CREATE TABLE inspection_fixture (value INTEGER NOT NULL)",
   3247         );
   3248 
   3249         let contended = open_existing_connection_pool(
   3250             &paths,
   3251             &metadata,
   3252             &base_catalog(),
   3253             &inspection_schema_catalog,
   3254             OpenMode::ReadOnlyInspection,
   3255             policy,
   3256         )
   3257         .await;
   3258         let Err(error) = contended else {
   3259             panic!("inspection must reject an active writer");
   3260         };
   3261         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3262 
   3263         let authority = writable.close().await.expect("writer authority");
   3264         drop(authority);
   3265         let state_directory = paths.state_database().parent().expect("state directory");
   3266         let stale_wal = state_directory.join(WAL_FILE_NAME);
   3267         fs::write(&stale_wal, b"stale-wal-evidence").expect("write stale WAL evidence");
   3268         let stale = open_existing_connection_pool(
   3269             &paths,
   3270             &metadata,
   3271             &base_catalog(),
   3272             &inspection_schema_catalog,
   3273             OpenMode::ReadOnlyInspection,
   3274             policy,
   3275         )
   3276         .await;
   3277         let Err(error) = stale else {
   3278             panic!("inspection must reject stale WAL state");
   3279         };
   3280         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3281         assert_eq!(fs::read(&stale_wal).unwrap(), b"stale-wal-evidence");
   3282         fs::remove_file(stale_wal).expect("remove test WAL evidence");
   3283         let before = directory_snapshot(state_directory);
   3284 
   3285         let read_only = open_existing_connection_pool(
   3286             &paths,
   3287             &metadata,
   3288             &base_catalog(),
   3289             &inspection_schema_catalog,
   3290             OpenMode::ReadOnlyInspection,
   3291             policy,
   3292         )
   3293         .await
   3294         .expect("offline read-only inspection");
   3295         let mut connection = read_only.acquire().await.expect("inspection connection");
   3296         assert_eq!(
   3297             sqlx::query_scalar::<_, i64>("SELECT value FROM inspection_fixture")
   3298                 .fetch_one(&mut *connection)
   3299                 .await
   3300                 .expect("read fixture"),
   3301             41
   3302         );
   3303         assert!(
   3304             sqlx::query("INSERT INTO inspection_fixture (value) VALUES (42)")
   3305                 .execute(&mut *connection)
   3306                 .await
   3307                 .is_err()
   3308         );
   3309         assert!(
   3310             connection_policy_matches(&mut connection, OpenMode::ReadOnlyInspection, policy)
   3311                 .await
   3312                 .unwrap()
   3313         );
   3314         drop(connection);
   3315         assert!(read_only.close().await.is_none());
   3316 
   3317         let after = directory_snapshot(state_directory);
   3318         assert_eq!(
   3319             after.keys().collect::<Vec<_>>(),
   3320             before.keys().collect::<Vec<_>>()
   3321         );
   3322         for (name, before_file) in &before {
   3323             let after_file = after.get(name).expect("same state entry");
   3324             assert!(
   3325                 after_file.bytes == before_file.bytes,
   3326                 "{name} bytes changed"
   3327             );
   3328             assert_eq!(
   3329                 after_file.length, before_file.length,
   3330                 "{name} length changed"
   3331             );
   3332             assert_eq!(
   3333                 after_file.modified, before_file.modified,
   3334                 "{name} mtime changed"
   3335             );
   3336             assert_eq!(after_file.mode, before_file.mode, "{name} mode changed");
   3337         }
   3338         assert!(!after.contains_key("state.sqlite-wal"));
   3339         assert!(!after.contains_key("state.sqlite-shm"));
   3340     }
   3341 
   3342     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3343     async fn offline_inspection_fixture() -> (tempfile::TempDir, ServiceSqlitePaths) {
   3344         let directory = tempfile::tempdir().expect("temporary directory");
   3345         let (paths, writable) =
   3346             initialized_pool(directory.path(), ServiceSqliteConnectionOptions::reviewed()).await;
   3347         let authority = writable.close().await.expect("writer authority");
   3348         drop(authority);
   3349         (directory, paths)
   3350     }
   3351 
   3352     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3353     #[tokio::test(flavor = "current_thread")]
   3354     async fn read_only_inspection_guard_rejects_each_filesystem_admission_drift() {
   3355         let (_directory, paths) = offline_inspection_fixture().await;
   3356         let state_directory = paths.state_database().parent().expect("state directory");
   3357         fs::set_permissions(state_directory, fs::Permissions::from_mode(0o775))
   3358             .expect("make state directory writable by group");
   3359         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3360 
   3361         let (_directory, paths) = offline_inspection_fixture().await;
   3362         fs::set_permissions(paths.state_lock(), fs::Permissions::from_mode(0o644))
   3363             .expect("weaken state lock mode");
   3364         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3365 
   3366         let (_directory, paths) = offline_inspection_fixture().await;
   3367         let lock_alias = paths
   3368             .state_lock()
   3369             .parent()
   3370             .expect("state directory")
   3371             .join("state-lock-alias");
   3372         fs::hard_link(paths.state_lock(), lock_alias).expect("hard-link state lock");
   3373         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3374 
   3375         let (_directory, paths) = offline_inspection_fixture().await;
   3376         let authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3377             .expect("writer authority");
   3378         let Err(error) = ReadOnlyInspectionGuard::acquire(&paths) else {
   3379             panic!("held writer lock must prevent inspection");
   3380         };
   3381         assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
   3382         drop(authority);
   3383 
   3384         let (_directory, paths) = offline_inspection_fixture().await;
   3385         fs::set_permissions(paths.state_database(), fs::Permissions::from_mode(0o644))
   3386             .expect("weaken database mode");
   3387         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3388 
   3389         let (_directory, paths) = offline_inspection_fixture().await;
   3390         let database_alias = paths
   3391             .state_database()
   3392             .parent()
   3393             .expect("state directory")
   3394             .join("state-database-alias");
   3395         fs::hard_link(paths.state_database(), database_alias).expect("hard-link database");
   3396         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3397 
   3398         let (_directory, paths) = offline_inspection_fixture().await;
   3399         let database = fs::OpenOptions::new()
   3400             .write(true)
   3401             .open(paths.state_database())
   3402             .expect("open database header");
   3403         std::os::unix::fs::FileExt::write_all_at(&database, &[1], 18)
   3404             .expect("corrupt SQLite write version");
   3405         assert!(ReadOnlyInspectionGuard::acquire(&paths).is_err());
   3406     }
   3407 
   3408     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3409     #[tokio::test(flavor = "current_thread")]
   3410     async fn read_only_inspection_guard_revalidates_live_directory_and_sidecars() {
   3411         let (_directory, paths) = offline_inspection_fixture().await;
   3412         let state_directory = paths.state_database().parent().expect("state directory");
   3413         let original_mode = fs::metadata(state_directory)
   3414             .expect("state directory metadata")
   3415             .permissions()
   3416             .mode();
   3417         let guard = ReadOnlyInspectionGuard::acquire(&paths).expect("inspection guard");
   3418         fs::set_permissions(state_directory, fs::Permissions::from_mode(0o775))
   3419             .expect("make live directory unsafe");
   3420         assert!(guard.validate_for(&paths).is_err());
   3421         fs::set_permissions(
   3422             state_directory,
   3423             fs::Permissions::from_mode(original_mode & 0o777),
   3424         )
   3425         .expect("restore directory mode");
   3426         fs::write(state_directory.join(WAL_FILE_NAME), b"stale")
   3427             .expect("create stale WAL evidence");
   3428         assert!(guard.validate_for(&paths).is_err());
   3429     }
   3430 
   3431     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3432     #[tokio::test(flavor = "current_thread")]
   3433     async fn pool_saturation_recovers_and_explicit_close_finishes() {
   3434         let directory = tempfile::tempdir().expect("temporary directory");
   3435         let policy = ServiceSqliteConnectionOptions::new(Duration::from_millis(500), 1).unwrap();
   3436         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   3437         let metadata = database_metadata(&paths).identity();
   3438         let observer = pool.pool.clone();
   3439         let held = pool.acquire().await.expect("only connection");
   3440         let saturated = pool.pool.try_acquire();
   3441         assert!(saturated.is_none());
   3442         drop(held);
   3443         let recovered = pool.acquire().await.expect("recovered connection");
   3444         drop(recovered);
   3445 
   3446         let authority = pool.close().await.expect("writer authority retained");
   3447         assert!(observer.is_closed());
   3448         assert!(authority.is_held());
   3449         drop(authority);
   3450 
   3451         let read_only = open_existing_connection_pool(
   3452             &paths,
   3453             &metadata,
   3454             &base_catalog(),
   3455             &base_schema_catalog(),
   3456             OpenMode::ReadOnlyInspection,
   3457             ServiceSqliteConnectionOptions::reviewed(),
   3458         )
   3459         .await
   3460         .expect("read-only inspection");
   3461         assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection);
   3462         assert!(read_only.close().await.is_none());
   3463     }
   3464 
   3465     #[cfg(any(target_os = "linux", target_os = "macos"))]
   3466     #[tokio::test]
   3467     async fn cancellation_during_explicit_connection_close_retains_the_close_driver() {
   3468         let directory = tempfile::tempdir().expect("temporary directory");
   3469         let policy = ServiceSqliteConnectionOptions::reviewed();
   3470         let (paths, pool) = initialized_pool(directory.path(), policy).await;
   3471         let connection = SqliteConnection::connect_with(&sqlite_connect_options(
   3472             &paths,
   3473             OpenMode::Initialize,
   3474             policy,
   3475         ))
   3476         .await
   3477         .expect("open private checkpoint connection");
   3478         let entered = Arc::new(Notify::new());
   3479         let release = Arc::new(Notify::new());
   3480         let close: BoxFuture<'static, Result<(), ServiceSqliteError>> = Box::pin({
   3481             let entered = Arc::clone(&entered);
   3482             let release = Arc::clone(&release);
   3483             async move {
   3484                 entered.notify_one();
   3485                 release.notified().await;
   3486                 connection
   3487                     .close()
   3488                     .await
   3489                     .map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source))
   3490             }
   3491         });
   3492         {
   3493             let mut driver = pool.close_driver.lock().await;
   3494             *driver = PrivateCloseDriver::Closing {
   3495                 future: close,
   3496                 authority_error: None,
   3497                 checkpoint_error: None,
   3498                 connection_close_error: None,
   3499             };
   3500         }
   3501         let pool = Arc::new(pool);
   3502         let close_task = tokio::spawn({
   3503             let pool = Arc::clone(&pool);
   3504             async move {
   3505                 pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default())
   3506                     .await
   3507             }
   3508         });
   3509         entered.notified().await;
   3510         close_task.abort();
   3511         assert!(
   3512             close_task
   3513                 .await
   3514                 .expect_err("explicit close task is cancelled")
   3515                 .is_cancelled()
   3516         );
   3517         let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting);
   3518         assert!(matches!(
   3519             retained,
   3520             Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority
   3521         ));
   3522 
   3523         release.notify_one();
   3524         pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default())
   3525             .await
   3526             .expect("authority release is proven")
   3527             .expect("retained connection closes explicitly");
   3528         let mut reacquired = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
   3529             .expect("authority reacquisition")
   3530             .expect("writer mode yields authority");
   3531         reacquired.release().expect("release reacquired authority");
   3532     }
   3533 }