lib

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

capture.rs (71292B)


      1 //! Cancellation-aware online capture into a caller-owned staging directory.
      2 
      3 use core::fmt;
      4 use std::{
      5     error::Error,
      6     ffi::{OsStr, OsString},
      7     fs::File,
      8     io::{self, Read, Seek, SeekFrom},
      9     os::unix::ffi::OsStrExt,
     10     path::{Component, Path, PathBuf},
     11     sync::{
     12         Arc,
     13         atomic::{AtomicBool, Ordering},
     14     },
     15     thread,
     16 };
     17 
     18 #[cfg(test)]
     19 use core::sync::atomic::AtomicU8;
     20 
     21 use rustix::{
     22     fs::{AtFlags, FileType, Mode, OFlags, fchmod, fstat, mkdirat, open, openat, statat, unlinkat},
     23     process::geteuid,
     24 };
     25 use sha2::{Digest, Sha256};
     26 use sqlx::{
     27     ConnectOptions, Connection as _, Row, Sqlite, SqliteConnection, pool::PoolConnection,
     28     sqlite::SqliteConnectOptions,
     29 };
     30 
     31 use crate::{
     32     BackupCreatedAtUnixMs, BackupMemberSha256, OpenMode, ServiceBackupManifest,
     33     ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind,
     34     open::{BackupSourceValidator, PrivateConnectionPool},
     35     sqlite_native_backup::{NativeBackup, NativeBackupStep},
     36 };
     37 
     38 const MAX_STAGING_PATH_BYTES: usize = 4_096;
     39 const BACKUP_PAGES_PER_STEP: i32 = 64;
     40 const HASH_BUFFER_BYTES: usize = 16 * 1_024;
     41 const STATE_FILE_NAME: &str = radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME;
     42 const KNOWN_SIDECARS: [&str; 3] = [
     43     "state.sqlite-wal",
     44     "state.sqlite-shm",
     45     "state.sqlite-journal",
     46 ];
     47 
     48 pub(crate) const TEST_CAPTURE_PHASE_BEFORE_CREATE: u8 = 1;
     49 pub(crate) const TEST_CAPTURE_PHASE_STAGING_CREATED: u8 = 2;
     50 pub(crate) const TEST_CAPTURE_PHASE_BACKUP_STEPPED: u8 = 3;
     51 pub(crate) const TEST_CAPTURE_PHASE_POST_COPY: u8 = 4;
     52 pub(crate) const TEST_CAPTURE_PHASE_PRE_FINAL_SYNC: u8 = 5;
     53 pub(crate) const TEST_CAPTURE_PHASE_METADATA_AWAITED: u8 = 10;
     54 pub(crate) const TEST_CAPTURE_PHASE_JOIN_AWAITED: u8 = 11;
     55 #[cfg(test)]
     56 static TEST_CAPTURE_PHASE: AtomicU8 = AtomicU8::new(0);
     57 #[cfg(test)]
     58 static TEST_CAPTURE_BLOCK_PHASE: AtomicU8 = AtomicU8::new(0);
     59 #[cfg(test)]
     60 static TEST_CAPTURE_INJECT_METADATA_FAILURE: AtomicBool = AtomicBool::new(false);
     61 #[cfg(test)]
     62 static TEST_CAPTURE_PANIC_WORKER: AtomicBool = AtomicBool::new(false);
     63 
     64 #[cfg(test)]
     65 pub(crate) fn test_capture_phase() -> u8 {
     66     TEST_CAPTURE_PHASE.load(Ordering::Acquire)
     67 }
     68 
     69 #[cfg(test)]
     70 pub(crate) fn test_capture_block_phase(phase: u8) {
     71     TEST_CAPTURE_BLOCK_PHASE.store(phase, Ordering::Release);
     72 }
     73 
     74 #[cfg(test)]
     75 pub(crate) fn test_capture_reset() {
     76     TEST_CAPTURE_BLOCK_PHASE.store(0, Ordering::Release);
     77     TEST_CAPTURE_PHASE.store(0, Ordering::Release);
     78     TEST_CAPTURE_INJECT_METADATA_FAILURE.store(false, Ordering::Release);
     79     TEST_CAPTURE_PANIC_WORKER.store(false, Ordering::Release);
     80 }
     81 
     82 #[cfg(test)]
     83 pub(crate) fn test_capture_inject_metadata_failure(enabled: bool) {
     84     TEST_CAPTURE_INJECT_METADATA_FAILURE.store(enabled, Ordering::Release);
     85 }
     86 
     87 #[cfg(test)]
     88 pub(crate) fn test_capture_panic_worker(enabled: bool) {
     89     TEST_CAPTURE_PANIC_WORKER.store(enabled, Ordering::Release);
     90 }
     91 
     92 async fn test_async_phase(phase: u8) {
     93     #[cfg(test)]
     94     {
     95         TEST_CAPTURE_PHASE.store(phase, Ordering::Release);
     96         while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase {
     97             tokio::task::yield_now().await;
     98         }
     99     }
    100     #[cfg(not(test))]
    101     let _ = phase;
    102 }
    103 
    104 trait CaptureOperations: Send + Sync {
    105     fn sync_state(&self, state: &File) -> io::Result<()>;
    106     fn sync_staging(&self, staging: &File) -> io::Result<()>;
    107     fn sync_parent(&self, parent: &File) -> io::Result<()>;
    108 }
    109 
    110 struct SystemCaptureOperations;
    111 
    112 impl CaptureOperations for SystemCaptureOperations {
    113     #[cfg_attr(coverage_nightly, coverage(off))]
    114     fn sync_state(&self, state: &File) -> io::Result<()> {
    115         state.sync_all()
    116     }
    117 
    118     #[cfg_attr(coverage_nightly, coverage(off))]
    119     fn sync_staging(&self, staging: &File) -> io::Result<()> {
    120         staging.sync_all()
    121     }
    122 
    123     #[cfg_attr(coverage_nightly, coverage(off))]
    124     fn sync_parent(&self, parent: &File) -> io::Result<()> {
    125         parent.sync_all()
    126     }
    127 }
    128 
    129 #[cfg(test)]
    130 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    131 pub(crate) enum TestCaptureSyncFailure {
    132     State,
    133     Staging,
    134     FinalParent,
    135 }
    136 
    137 #[cfg(test)]
    138 struct FailingCaptureOperations {
    139     failure: TestCaptureSyncFailure,
    140     parent_syncs: core::sync::atomic::AtomicU8,
    141 }
    142 
    143 #[cfg(test)]
    144 impl CaptureOperations for FailingCaptureOperations {
    145     fn sync_state(&self, state: &File) -> io::Result<()> {
    146         if self.failure == TestCaptureSyncFailure::State {
    147             Err(crate::failpoint::storage_full_error())
    148         } else {
    149             state.sync_all()
    150         }
    151     }
    152 
    153     fn sync_staging(&self, staging: &File) -> io::Result<()> {
    154         if self.failure == TestCaptureSyncFailure::Staging {
    155             Err(crate::failpoint::storage_full_error())
    156         } else {
    157             staging.sync_all()
    158         }
    159     }
    160 
    161     fn sync_parent(&self, parent: &File) -> io::Result<()> {
    162         let occurrence = self.parent_syncs.fetch_add(1, Ordering::AcqRel);
    163         if self.failure == TestCaptureSyncFailure::FinalParent && occurrence == 1 {
    164             Err(crate::failpoint::storage_full_error())
    165         } else {
    166             parent.sync_all()
    167         }
    168     }
    169 }
    170 
    171 pub(crate) async fn capture_online_backup(
    172     pool: &PrivateConnectionPool,
    173     closing: &AtomicBool,
    174     active: &Arc<AtomicBool>,
    175     staging_directory: &Path,
    176     created_at_unix_ms: BackupCreatedAtUnixMs,
    177     failpoints: &crate::failpoint::DurabilityFailpoints,
    178 ) -> Result<ServiceBackupManifest, ServiceSqliteError> {
    179     capture_online_backup_with_operations(
    180         pool,
    181         closing,
    182         active,
    183         staging_directory,
    184         created_at_unix_ms,
    185         Arc::new(SystemCaptureOperations),
    186         failpoints,
    187     )
    188     .await
    189 }
    190 
    191 #[cfg(test)]
    192 pub(crate) async fn test_capture_online_backup_with_sync_failure(
    193     pool: &PrivateConnectionPool,
    194     closing: &AtomicBool,
    195     active: &Arc<AtomicBool>,
    196     staging_directory: &Path,
    197     created_at_unix_ms: BackupCreatedAtUnixMs,
    198     failure: TestCaptureSyncFailure,
    199 ) -> Result<ServiceBackupManifest, ServiceSqliteError> {
    200     capture_online_backup_with_operations(
    201         pool,
    202         closing,
    203         active,
    204         staging_directory,
    205         created_at_unix_ms,
    206         Arc::new(FailingCaptureOperations {
    207             failure,
    208             parent_syncs: core::sync::atomic::AtomicU8::new(0),
    209         }),
    210         &crate::failpoint::DurabilityFailpoints::default(),
    211     )
    212     .await
    213 }
    214 
    215 async fn capture_online_backup_with_operations(
    216     pool: &PrivateConnectionPool,
    217     closing: &AtomicBool,
    218     active: &Arc<AtomicBool>,
    219     staging_directory: &Path,
    220     created_at_unix_ms: BackupCreatedAtUnixMs,
    221     operations: Arc<dyn CaptureOperations>,
    222     failpoints: &crate::failpoint::DurabilityFailpoints,
    223 ) -> Result<ServiceBackupManifest, ServiceSqliteError> {
    224     crate::require_condition(
    225         matches!(
    226             pool.mode(),
    227             OpenMode::Initialize | OpenMode::ReadWriteExisting
    228         ) && !closing.load(Ordering::Acquire),
    229         ServiceSqliteErrorKind::Open,
    230     )?;
    231     let staging = StagingPath::new(staging_directory)?;
    232     let permit = CapturePermit::acquire(Arc::clone(active))?;
    233     pool.validate()?;
    234     let mut admission = pool.acquire().await?;
    235     pool.validate()?;
    236     crate::require_condition(
    237         !closing.load(Ordering::Acquire),
    238         ServiceSqliteErrorKind::Open,
    239     )?;
    240     let metadata = crate::metadata::verify_database_metadata(&mut admission, pool.identity()).await;
    241     #[cfg(test)]
    242     let metadata = if TEST_CAPTURE_INJECT_METADATA_FAILURE.load(Ordering::Acquire) {
    243         Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))
    244     } else {
    245         metadata
    246     };
    247     test_async_phase(TEST_CAPTURE_PHASE_METADATA_AWAITED).await;
    248     pool.validate()?;
    249     let metadata = metadata?;
    250     let validator = pool.backup_source_validator();
    251     validator.validate()?;
    252 
    253     let cancellation = Arc::new(AtomicBool::new(false));
    254     let cancellation_guard = CaptureCancellation::new(Arc::clone(&cancellation));
    255     let worker = CaptureWorker {
    256         admission: Some(admission),
    257         _permit: permit,
    258         validator,
    259         metadata,
    260         staging,
    261         created_at_unix_ms,
    262         cancellation,
    263         operations,
    264         failpoints: failpoints.clone(),
    265         runtime: tokio::runtime::Handle::current(),
    266     };
    267     let joined = tokio::task::spawn_blocking(move || worker.run()).await;
    268     test_async_phase(TEST_CAPTURE_PHASE_JOIN_AWAITED).await;
    269     cancellation_guard.complete();
    270     pool.validate()?;
    271     let result = joined.map_err(|source| backup_source(BackupFailureKind::Join, source))?;
    272     result.map(PendingCapture::commit)
    273 }
    274 
    275 struct CapturePermit {
    276     active: Arc<AtomicBool>,
    277 }
    278 
    279 impl CapturePermit {
    280     fn acquire(active: Arc<AtomicBool>) -> Result<Self, ServiceSqliteError> {
    281         active
    282             .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
    283             .map_err(|_| backup_error(BackupFailureKind::AlreadyActive))?;
    284         Ok(Self { active })
    285     }
    286 }
    287 
    288 impl Drop for CapturePermit {
    289     fn drop(&mut self) {
    290         self.active.store(false, Ordering::Release);
    291     }
    292 }
    293 
    294 struct CaptureCancellation {
    295     cancelled: Arc<AtomicBool>,
    296     complete: AtomicBool,
    297 }
    298 
    299 impl CaptureCancellation {
    300     fn new(cancelled: Arc<AtomicBool>) -> Self {
    301         Self {
    302             cancelled,
    303             complete: AtomicBool::new(false),
    304         }
    305     }
    306 
    307     fn complete(&self) {
    308         self.complete.store(true, Ordering::Release);
    309     }
    310 }
    311 
    312 impl Drop for CaptureCancellation {
    313     fn drop(&mut self) {
    314         if !self.complete.load(Ordering::Acquire) {
    315             self.cancelled.store(true, Ordering::Release);
    316         }
    317     }
    318 }
    319 
    320 struct CaptureWorker {
    321     admission: Option<PoolConnection<Sqlite>>,
    322     _permit: CapturePermit,
    323     validator: BackupSourceValidator,
    324     metadata: ServiceDatabaseMetadata,
    325     staging: StagingPath,
    326     created_at_unix_ms: BackupCreatedAtUnixMs,
    327     cancellation: Arc<AtomicBool>,
    328     operations: Arc<dyn CaptureOperations>,
    329     failpoints: crate::failpoint::DurabilityFailpoints,
    330     runtime: tokio::runtime::Handle,
    331 }
    332 
    333 impl CaptureWorker {
    334     fn run(mut self) -> Result<PendingCapture, ServiceSqliteError> {
    335         self.check_cancelled()?;
    336         self.validator.validate()?;
    337         self.test_phase(TEST_CAPTURE_PHASE_BEFORE_CREATE);
    338         self.check_cancelled()?;
    339         self.hit_checked(
    340             None,
    341             crate::failpoint::DurabilityFailpoint::BackupBeforeCreate,
    342             BackupFailureKind::CreateStaging,
    343         )?;
    344         let mut staging = StagingGuard::create(&self.staging, self.operations.as_ref())?;
    345         self.hit_checked(
    346             Some(&staging),
    347             crate::failpoint::DurabilityFailpoint::BackupAfterCreate,
    348             BackupFailureKind::CreateStaging,
    349         )?;
    350         self.test_phase(TEST_CAPTURE_PHASE_STAGING_CREATED);
    351         #[cfg(test)]
    352         if TEST_CAPTURE_PANIC_WORKER.load(Ordering::Acquire) {
    353             panic!("injected backup worker failure");
    354         }
    355         self.check_cancelled()?;
    356         self.validator.validate()?;
    357         staging.validate()?;
    358 
    359         let source = self
    360             .admission
    361             .as_mut()
    362             .ok_or_else(|| backup_error(BackupFailureKind::Capture))?;
    363         self.runtime.block_on(verify_database_inventory(source))?;
    364         self.runtime
    365             .block_on(verify_database_metadata(source, &self.metadata))?;
    366         self.validator.validate()?;
    367         staging.validate()?;
    368 
    369         let mut destination = self.open_sqlx_destination(&staging)?;
    370         staging.record_sidecars();
    371 
    372         self.hit_checked(
    373             Some(&staging),
    374             crate::failpoint::DurabilityFailpoint::BackupBeforeCopy,
    375             BackupFailureKind::Capture,
    376         )?;
    377         let capture_result = self.copy_with_locked_sqlx_handles(&mut destination, &mut staging);
    378         let close_result = self
    379             .runtime
    380             .block_on(destination.close())
    381             .map_err(|source| backup_source(BackupFailureKind::Capture, source));
    382         staging.record_sidecars();
    383         self.validator.validate()?;
    384         staging.validate()?;
    385         capture_result?;
    386         close_result?;
    387         self.hit_checked(
    388             Some(&staging),
    389             crate::failpoint::DurabilityFailpoint::BackupAfterCopy,
    390             BackupFailureKind::Capture,
    391         )?;
    392         staging.record_sidecars();
    393 
    394         self.test_phase(TEST_CAPTURE_PHASE_POST_COPY);
    395         self.check_cancelled()?;
    396         self.validator.validate()?;
    397         staging.validate()?;
    398         let mut destination = self.open_inspection_destination(&staging)?;
    399         self.runtime
    400             .block_on(verify_database_inventory(&mut destination))?;
    401         self.runtime
    402             .block_on(verify_database_metadata(&mut destination, &self.metadata))?;
    403         self.runtime.block_on(verify_integrity(&mut destination))?;
    404         staging.record_sidecars();
    405         self.check_cancelled()?;
    406         self.validator.validate()?;
    407         staging.validate()?;
    408 
    409         self.runtime
    410             .block_on(destination.close())
    411             .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
    412         staging.record_sidecars();
    413         self.validator.validate()?;
    414         staging.validate()?;
    415         staging.validate_inventory()?;
    416         self.check_cancelled()?;
    417 
    418         self.hit_checked(
    419             Some(&staging),
    420             crate::failpoint::DurabilityFailpoint::BackupBeforeFileSync,
    421             BackupFailureKind::SyncState,
    422         )?;
    423         staging.sync_state(self.operations.as_ref())?;
    424         self.hit_checked(
    425             Some(&staging),
    426             crate::failpoint::DurabilityFailpoint::BackupAfterFileSync,
    427             BackupFailureKind::SyncState,
    428         )?;
    429         let (byte_length, digest) = staging.hash_state(&self.cancellation)?;
    430         self.test_phase(TEST_CAPTURE_PHASE_PRE_FINAL_SYNC);
    431         self.check_cancelled()?;
    432         self.hit_checked(
    433             Some(&staging),
    434             crate::failpoint::DurabilityFailpoint::BackupBeforeDirectorySync,
    435             BackupFailureKind::SyncStaging,
    436         )?;
    437         staging.sync_directories(self.operations.as_ref())?;
    438         self.hit_checked(
    439             Some(&staging),
    440             crate::failpoint::DurabilityFailpoint::BackupAfterDirectorySync,
    441             BackupFailureKind::SyncParent,
    442         )?;
    443         self.validator.validate()?;
    444         staging.validate()?;
    445         staging.validate_inventory()?;
    446         self.check_cancelled()?;
    447 
    448         let manifest = ServiceBackupManifest::from_capture(
    449             &self.metadata,
    450             self.created_at_unix_ms,
    451             byte_length,
    452             BackupMemberSha256::from_bytes(digest),
    453         )
    454         .map_err(|source| backup_source(BackupFailureKind::Manifest, source))?;
    455         Ok(PendingCapture {
    456             staging,
    457             _admission: self
    458                 .admission
    459                 .take()
    460                 .ok_or_else(|| backup_error(BackupFailureKind::Capture))?,
    461             _permit: self._permit,
    462             manifest,
    463         })
    464     }
    465 
    466     fn copy_with_locked_sqlx_handles(
    467         &mut self,
    468         destination: &mut SqliteConnection,
    469         staging: &mut StagingGuard,
    470     ) -> Result<(), ServiceSqliteError> {
    471         let mut admission = self
    472             .admission
    473             .take()
    474             .ok_or_else(|| backup_error(BackupFailureKind::Capture))?;
    475         let result = (|| {
    476             let mut source_handle = self
    477                 .runtime
    478                 .block_on(admission.lock_handle())
    479                 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
    480             let mut destination_handle = self
    481                 .runtime
    482                 .block_on(destination.lock_handle())
    483                 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
    484             let mut backup = NativeBackup::start(&mut destination_handle, &mut source_handle)
    485                 .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
    486             loop {
    487                 self.check_cancelled()?;
    488                 self.validator.validate()?;
    489                 staging.validate()?;
    490                 let step = backup
    491                     .step(BACKUP_PAGES_PER_STEP)
    492                     .map_err(|source| backup_source(BackupFailureKind::Capture, source));
    493                 staging.record_sidecars();
    494                 self.test_phase(TEST_CAPTURE_PHASE_BACKUP_STEPPED);
    495                 let step = step?;
    496                 self.validator.validate()?;
    497                 staging.validate()?;
    498                 match step {
    499                     NativeBackupStep::Done => {
    500                         return backup
    501                             .finish()
    502                             .map_err(|source| backup_source(BackupFailureKind::Capture, source));
    503                     }
    504                     NativeBackupStep::More => {}
    505                     NativeBackupStep::Busy | NativeBackupStep::Locked => thread::yield_now(),
    506                 }
    507             }
    508         })();
    509         self.admission = Some(admission);
    510         result
    511     }
    512 
    513     fn open_sqlx_destination(
    514         &self,
    515         staging: &StagingGuard,
    516     ) -> Result<SqliteConnection, ServiceSqliteError> {
    517         staging.validate()?;
    518         let options = SqliteConnectOptions::new()
    519             .filename(staging.state_path())
    520             .create_if_missing(false)
    521             .foreign_keys(false)
    522             .disable_statement_logging();
    523         let result = self
    524             .runtime
    525             .block_on(SqliteConnection::connect_with(&options));
    526         staging.validate()?;
    527         result.map_err(|source| backup_source(BackupFailureKind::Capture, source))
    528     }
    529 
    530     fn open_inspection_destination(
    531         &self,
    532         staging: &StagingGuard,
    533     ) -> Result<SqliteConnection, ServiceSqliteError> {
    534         staging.validate()?;
    535         let options = SqliteConnectOptions::new()
    536             .filename(staging.state_path())
    537             .create_if_missing(false)
    538             .foreign_keys(false)
    539             .disable_statement_logging();
    540         let result = self
    541             .runtime
    542             .block_on(SqliteConnection::connect_with(&options));
    543         staging.validate()?;
    544         result.map_err(|source| backup_source(BackupFailureKind::Capture, source))
    545     }
    546 
    547     fn check_cancelled(&self) -> Result<(), ServiceSqliteError> {
    548         if self.cancellation.load(Ordering::Acquire) {
    549             Err(backup_error(BackupFailureKind::Cancelled))
    550         } else {
    551             Ok(())
    552         }
    553     }
    554 
    555     fn hit_checked(
    556         &self,
    557         staging: Option<&StagingGuard>,
    558         point: crate::failpoint::DurabilityFailpoint,
    559         failure: BackupFailureKind,
    560     ) -> Result<(), ServiceSqliteError> {
    561         #[cfg(test)]
    562         self.failpoints
    563             .observe(point, TEST_CAPTURE_PHASE.load(Ordering::Acquire));
    564         let injected = self
    565             .failpoints
    566             .hit(point)
    567             .map_err(|source| backup_source(failure, source));
    568         self.validator.validate()?;
    569         if let Some(staging) = staging {
    570             staging.validate()?;
    571         }
    572         injected
    573     }
    574 
    575     fn test_phase(&self, phase: u8) {
    576         #[cfg(test)]
    577         {
    578             TEST_CAPTURE_PHASE.store(phase, Ordering::Release);
    579             while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase
    580                 && !self.cancellation.load(Ordering::Acquire)
    581             {
    582                 thread::yield_now();
    583             }
    584         }
    585         #[cfg(not(test))]
    586         let _ = phase;
    587     }
    588 }
    589 
    590 struct PendingCapture {
    591     staging: StagingGuard,
    592     _admission: PoolConnection<Sqlite>,
    593     _permit: CapturePermit,
    594     manifest: ServiceBackupManifest,
    595 }
    596 
    597 impl PendingCapture {
    598     fn commit(mut self) -> ServiceBackupManifest {
    599         self.staging.commit();
    600         self.manifest
    601     }
    602 }
    603 
    604 #[derive(Clone)]
    605 struct StagingPath {
    606     full: PathBuf,
    607     parent: PathBuf,
    608     name: OsString,
    609 }
    610 
    611 impl StagingPath {
    612     fn new(path: &Path) -> Result<Self, ServiceSqliteError> {
    613         if path.as_os_str().as_bytes().len() > MAX_STAGING_PATH_BYTES || !path.is_absolute() {
    614             return Err(backup_error(BackupFailureKind::InvalidStagingPath));
    615         }
    616         if path
    617             .components()
    618             .any(|component| matches!(component, Component::CurDir | Component::ParentDir))
    619         {
    620             return Err(backup_error(BackupFailureKind::InvalidStagingPath));
    621         }
    622         let name = match path.components().next_back() {
    623             Some(Component::Normal(name)) if !name.as_bytes().is_empty() => name.to_os_string(),
    624             _ => return Err(backup_error(BackupFailureKind::InvalidStagingPath)),
    625         };
    626         let parent = path
    627             .parent()
    628             .filter(|parent| parent.is_absolute())
    629             .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingPath))?
    630             .to_path_buf();
    631         Ok(Self {
    632             full: path.to_path_buf(),
    633             parent,
    634             name,
    635         })
    636     }
    637 }
    638 
    639 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    640 struct FileIdentity {
    641     device: u64,
    642     inode: u64,
    643 }
    644 
    645 struct StagingGuard {
    646     path: StagingPath,
    647     parent: File,
    648     parent_identity: FileIdentity,
    649     directory: File,
    650     directory_identity: FileIdentity,
    651     state: File,
    652     state_identity: FileIdentity,
    653     sidecar_identities: [Option<FileIdentity>; KNOWN_SIDECARS.len()],
    654     committed: bool,
    655 }
    656 
    657 impl StagingGuard {
    658     fn create(
    659         path: &StagingPath,
    660         operations: &dyn CaptureOperations,
    661     ) -> Result<Self, ServiceSqliteError> {
    662         let parent = File::from(
    663             open(
    664                 &path.parent,
    665                 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
    666                 Mode::empty(),
    667             )
    668             .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?,
    669         );
    670         let parent_identity = validate_directory_descriptor(&parent, false)?;
    671         mkdirat(&parent, &path.name, Mode::RUSR | Mode::WUSR | Mode::XUSR).map_err(|source| {
    672             backup_source(
    673                 if source == rustix::io::Errno::EXIST {
    674                     BackupFailureKind::StagingCollision
    675                 } else {
    676                     BackupFailureKind::CreateStaging
    677                 },
    678                 source,
    679             )
    680         })?;
    681         let created_directory_identity = match created_directory_identity(&parent, &path.name) {
    682             Ok(identity) => identity,
    683             Err(error) => {
    684                 // Without a proven identity, the current entry must be preserved.
    685                 let _ = parent.sync_all();
    686                 return Err(error);
    687             }
    688         };
    689         let directory = match openat(
    690             &parent,
    691             &path.name,
    692             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
    693             Mode::empty(),
    694         ) {
    695             Ok(directory) => File::from(directory),
    696             Err(source) => {
    697                 cleanup_partial_staging(
    698                     &parent,
    699                     &path.name,
    700                     Some(created_directory_identity),
    701                     None,
    702                     None,
    703                 );
    704                 return Err(backup_source(BackupFailureKind::CreateStaging, source));
    705             }
    706         };
    707         let opened_directory_identity = match descriptor_identity(&directory) {
    708             Ok(identity) => identity,
    709             Err(error) => {
    710                 cleanup_partial_staging(
    711                     &parent,
    712                     &path.name,
    713                     Some(created_directory_identity),
    714                     Some(&directory),
    715                     None,
    716                 );
    717                 return Err(error);
    718             }
    719         };
    720         if let Err(error) = require_backup_condition(
    721             opened_directory_identity == created_directory_identity,
    722             BackupFailureKind::StagingReplaced,
    723         ) {
    724             cleanup_partial_staging(
    725                 &parent,
    726                 &path.name,
    727                 Some(created_directory_identity),
    728                 Some(&directory),
    729                 None,
    730             );
    731             return Err(error);
    732         }
    733         if let Err(source) = fchmod(&directory, Mode::RUSR | Mode::WUSR | Mode::XUSR) {
    734             cleanup_partial_staging(
    735                 &parent,
    736                 &path.name,
    737                 Some(created_directory_identity),
    738                 Some(&directory),
    739                 None,
    740             );
    741             return Err(backup_source(BackupFailureKind::CreateStaging, source));
    742         }
    743         let directory_identity = match validate_directory_descriptor(&directory, true) {
    744             Ok(identity) => identity,
    745             Err(error) => {
    746                 cleanup_partial_staging(
    747                     &parent,
    748                     &path.name,
    749                     Some(created_directory_identity),
    750                     Some(&directory),
    751                     None,
    752                 );
    753                 return Err(error);
    754             }
    755         };
    756         let state = match openat(
    757             &directory,
    758             STATE_FILE_NAME,
    759             OFlags::RDWR | OFlags::CREATE | OFlags::EXCL | OFlags::NOFOLLOW | OFlags::CLOEXEC,
    760             Mode::RUSR | Mode::WUSR,
    761         ) {
    762             Ok(state) => File::from(state),
    763             Err(source) => {
    764                 cleanup_partial_staging(
    765                     &parent,
    766                     &path.name,
    767                     Some(created_directory_identity),
    768                     Some(&directory),
    769                     None,
    770                 );
    771                 return Err(backup_source(BackupFailureKind::CreateState, source));
    772             }
    773         };
    774         let created_state_identity = match descriptor_identity(&state) {
    775             Ok(identity) => identity,
    776             Err(error) => {
    777                 cleanup_partial_staging(
    778                     &parent,
    779                     &path.name,
    780                     Some(created_directory_identity),
    781                     Some(&directory),
    782                     None,
    783                 );
    784                 return Err(error);
    785             }
    786         };
    787         if let Err(source) = fchmod(&state, Mode::RUSR | Mode::WUSR) {
    788             cleanup_partial_staging(
    789                 &parent,
    790                 &path.name,
    791                 Some(created_directory_identity),
    792                 Some(&directory),
    793                 Some(created_state_identity),
    794             );
    795             return Err(backup_source(BackupFailureKind::CreateState, source));
    796         }
    797         let state_identity = match validate_file_descriptor(&state) {
    798             Ok(identity) => identity,
    799             Err(error) => {
    800                 cleanup_partial_staging(
    801                     &parent,
    802                     &path.name,
    803                     Some(created_directory_identity),
    804                     Some(&directory),
    805                     Some(created_state_identity),
    806                 );
    807                 return Err(error);
    808             }
    809         };
    810         let staging = Self {
    811             path: path.clone(),
    812             parent,
    813             parent_identity,
    814             directory,
    815             directory_identity,
    816             state,
    817             state_identity,
    818             sidecar_identities: [None; KNOWN_SIDECARS.len()],
    819             committed: false,
    820         };
    821         staging.validate()?;
    822         operations
    823             .sync_parent(&staging.parent)
    824             .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?;
    825         Ok(staging)
    826     }
    827 
    828     fn state_path(&self) -> PathBuf {
    829         self.path.full.join(STATE_FILE_NAME)
    830     }
    831 
    832     fn validate(&self) -> Result<(), ServiceSqliteError> {
    833         validate_reopened_directory(&self.path.parent, &self.parent, self.parent_identity, false)?;
    834         validate_directory_entry(
    835             &self.parent,
    836             &self.path.name,
    837             &self.directory,
    838             self.directory_identity,
    839         )?;
    840         validate_file_entry(
    841             &self.directory,
    842             OsStr::new(STATE_FILE_NAME),
    843             &self.state,
    844             self.state_identity,
    845         )
    846     }
    847 
    848     fn validate_inventory(&self) -> Result<(), ServiceSqliteError> {
    849         self.validate()?;
    850         let mut entries = std::fs::read_dir(&self.path.full)
    851             .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?;
    852         let first = entries
    853             .next()
    854             .transpose()
    855             .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?
    856             .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingInventory))?;
    857         require_backup_condition(
    858             crate::all_constraints([
    859                 first.file_name() == OsStr::new(STATE_FILE_NAME),
    860                 entries.next().is_none(),
    861             ]),
    862             BackupFailureKind::InvalidStagingInventory,
    863         )?;
    864         Ok(())
    865     }
    866 
    867     fn sync_state(&self, operations: &dyn CaptureOperations) -> Result<(), ServiceSqliteError> {
    868         self.validate()?;
    869         operations
    870             .sync_state(&self.state)
    871             .map_err(|source| backup_source(BackupFailureKind::SyncState, source))?;
    872         self.validate()
    873     }
    874 
    875     fn hash_state(&self, cancellation: &AtomicBool) -> Result<(u64, [u8; 32]), ServiceSqliteError> {
    876         self.validate()?;
    877         let mut state = self
    878             .state
    879             .try_clone()
    880             .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
    881         state
    882             .seek(SeekFrom::Start(0))
    883             .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
    884         let mut hasher = Sha256::new();
    885         let mut length = 0_u64;
    886         let mut buffer = [0_u8; HASH_BUFFER_BYTES];
    887         loop {
    888             require_backup_condition(
    889                 !cancellation.load(Ordering::Acquire),
    890                 BackupFailureKind::Cancelled,
    891             )?;
    892             let count = state
    893                 .read(&mut buffer)
    894                 .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
    895             if count == 0 {
    896                 break;
    897             }
    898             length = length
    899                 .checked_add(
    900                     u64::try_from(count).map_err(|_| backup_error(BackupFailureKind::HashState))?,
    901                 )
    902                 .ok_or_else(|| backup_error(BackupFailureKind::HashState))?;
    903             require_backup_condition(length <= i64::MAX as u64, BackupFailureKind::HashState)?;
    904             hasher.update(&buffer[..count]);
    905         }
    906         require_backup_condition(length != 0, BackupFailureKind::HashState)?;
    907         self.validate()?;
    908         Ok((length, hasher.finalize().into()))
    909     }
    910 
    911     fn sync_directories(
    912         &self,
    913         operations: &dyn CaptureOperations,
    914     ) -> Result<(), ServiceSqliteError> {
    915         self.validate()?;
    916         operations
    917             .sync_staging(&self.directory)
    918             .map_err(|source| backup_source(BackupFailureKind::SyncStaging, source))?;
    919         self.validate()?;
    920         operations
    921             .sync_parent(&self.parent)
    922             .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?;
    923         self.validate()
    924     }
    925 
    926     fn record_sidecars(&mut self) {
    927         for (index, name) in KNOWN_SIDECARS.iter().enumerate() {
    928             if self.sidecar_identities[index].is_none() {
    929                 self.sidecar_identities[index] = safe_sidecar_identity(&self.directory, name);
    930             }
    931         }
    932     }
    933 
    934     fn commit(&mut self) {
    935         self.committed = true;
    936     }
    937 
    938     fn cleanup(&mut self) {
    939         if self.committed {
    940             return;
    941         }
    942         if validate_directory_entry(
    943             &self.parent,
    944             &self.path.name,
    945             &self.directory,
    946             self.directory_identity,
    947         )
    948         .is_err()
    949         {
    950             return;
    951         }
    952         if current_entry_identity(&self.directory, OsStr::new(STATE_FILE_NAME))
    953             == Some(self.state_identity)
    954         {
    955             let _ = unlinkat(&self.directory, STATE_FILE_NAME, AtFlags::empty());
    956         }
    957         for (sidecar, identity) in KNOWN_SIDECARS.iter().zip(self.sidecar_identities) {
    958             if identity.is_some() && safe_sidecar_identity(&self.directory, sidecar) == identity {
    959                 let _ = unlinkat(&self.directory, *sidecar, AtFlags::empty());
    960             }
    961         }
    962         if current_entry_identity(&self.parent, &self.path.name) == Some(self.directory_identity) {
    963             let _ = unlinkat(&self.parent, &self.path.name, AtFlags::REMOVEDIR);
    964         }
    965         let _ = self.parent.sync_all();
    966         self.committed = true;
    967     }
    968 }
    969 
    970 fn cleanup_partial_staging(
    971     parent: &File,
    972     name: &OsStr,
    973     directory_identity: Option<FileIdentity>,
    974     directory: Option<&File>,
    975     state_identity: Option<FileIdentity>,
    976 ) {
    977     if let (Some(directory), Some(state_identity)) = (directory, state_identity)
    978         && current_entry_identity(directory, OsStr::new(STATE_FILE_NAME)) == Some(state_identity)
    979     {
    980         let _ = unlinkat(directory, STATE_FILE_NAME, AtFlags::empty());
    981     }
    982     if directory_identity.is_some() && current_entry_identity(parent, name) == directory_identity {
    983         let _ = unlinkat(parent, name, AtFlags::REMOVEDIR);
    984     }
    985     let _ = parent.sync_all();
    986 }
    987 
    988 impl Drop for StagingGuard {
    989     fn drop(&mut self) {
    990         self.cleanup();
    991     }
    992 }
    993 
    994 fn created_directory_identity(
    995     parent: &File,
    996     name: &OsStr,
    997 ) -> Result<FileIdentity, ServiceSqliteError> {
    998     let status = statat(parent, name, AtFlags::SYMLINK_NOFOLLOW)
    999         .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?;
   1000     if !crate::native_metadata::secure_directory(
   1001         FileType::from_raw_mode(status.st_mode).is_dir(),
   1002         status.st_uid,
   1003         geteuid().as_raw(),
   1004         crate::native_metadata::mode(status.st_mode),
   1005     ) {
   1006         return Err(backup_error(BackupFailureKind::StagingReplaced));
   1007     }
   1008     identity(&status)
   1009 }
   1010 
   1011 fn descriptor_identity(descriptor: &File) -> Result<FileIdentity, ServiceSqliteError> {
   1012     let status = fstat(descriptor)
   1013         .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?;
   1014     identity(&status)
   1015 }
   1016 
   1017 fn validate_directory_descriptor(
   1018     directory: &File,
   1019     exact_owner_mode: bool,
   1020 ) -> Result<FileIdentity, ServiceSqliteError> {
   1021     let status = fstat(directory)
   1022         .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?;
   1023     let mode = crate::native_metadata::mode(status.st_mode) & 0o777;
   1024     let valid = if exact_owner_mode {
   1025         crate::native_metadata::exact_directory(
   1026             FileType::from_raw_mode(status.st_mode).is_dir(),
   1027             status.st_uid,
   1028             geteuid().as_raw(),
   1029             mode,
   1030         )
   1031     } else {
   1032         crate::native_metadata::secure_directory(
   1033             FileType::from_raw_mode(status.st_mode).is_dir(),
   1034             status.st_uid,
   1035             geteuid().as_raw(),
   1036             mode,
   1037         )
   1038     };
   1039     if !valid {
   1040         return Err(backup_error(BackupFailureKind::InvalidStagingParent));
   1041     }
   1042     identity(&status)
   1043 }
   1044 
   1045 fn validate_file_descriptor(file: &File) -> Result<FileIdentity, ServiceSqliteError> {
   1046     let status = fstat(file)
   1047         .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?;
   1048     if !crate::native_metadata::exact_regular_file(
   1049         FileType::from_raw_mode(status.st_mode).is_file(),
   1050         crate::native_metadata::link_count(status.st_nlink),
   1051         status.st_uid,
   1052         geteuid().as_raw(),
   1053         crate::native_metadata::mode(status.st_mode),
   1054     ) {
   1055         return Err(backup_error(BackupFailureKind::InvalidStagingInventory));
   1056     }
   1057     identity(&status)
   1058 }
   1059 
   1060 fn validate_reopened_directory(
   1061     path: &Path,
   1062     held: &File,
   1063     expected: FileIdentity,
   1064     exact_owner_mode: bool,
   1065 ) -> Result<(), ServiceSqliteError> {
   1066     let current = File::from(
   1067         open(
   1068             path,
   1069             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1070             Mode::empty(),
   1071         )
   1072         .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
   1073     );
   1074     require_backup_condition(
   1075         validate_directory_descriptor(&current, exact_owner_mode)? == expected,
   1076         BackupFailureKind::StagingReplaced,
   1077     )?;
   1078     require_backup_condition(
   1079         validate_directory_descriptor(held, exact_owner_mode)? == expected,
   1080         BackupFailureKind::StagingReplaced,
   1081     )?;
   1082     Ok(())
   1083 }
   1084 
   1085 fn validate_directory_entry(
   1086     parent: &File,
   1087     name: &OsStr,
   1088     held: &File,
   1089     expected: FileIdentity,
   1090 ) -> Result<(), ServiceSqliteError> {
   1091     let current = File::from(
   1092         openat(
   1093             parent,
   1094             name,
   1095             OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1096             Mode::empty(),
   1097         )
   1098         .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
   1099     );
   1100     require_backup_condition(
   1101         validate_directory_descriptor(&current, true)? == expected,
   1102         BackupFailureKind::StagingReplaced,
   1103     )?;
   1104     require_backup_condition(
   1105         validate_directory_descriptor(held, true)? == expected,
   1106         BackupFailureKind::StagingReplaced,
   1107     )?;
   1108     Ok(())
   1109 }
   1110 
   1111 fn validate_file_entry(
   1112     directory: &File,
   1113     name: &OsStr,
   1114     held: &File,
   1115     expected: FileIdentity,
   1116 ) -> Result<(), ServiceSqliteError> {
   1117     let current = File::from(
   1118         openat(
   1119             directory,
   1120             name,
   1121             OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1122             Mode::empty(),
   1123         )
   1124         .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
   1125     );
   1126     require_backup_condition(
   1127         validate_file_descriptor(&current)? == expected,
   1128         BackupFailureKind::StagingReplaced,
   1129     )?;
   1130     require_backup_condition(
   1131         validate_file_descriptor(held)? == expected,
   1132         BackupFailureKind::StagingReplaced,
   1133     )?;
   1134     Ok(())
   1135 }
   1136 
   1137 fn current_entry_identity(directory: &File, name: &OsStr) -> Option<FileIdentity> {
   1138     let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?;
   1139     identity(&status).ok()
   1140 }
   1141 
   1142 fn safe_sidecar_identity(directory: &File, name: &str) -> Option<FileIdentity> {
   1143     let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?;
   1144     if !crate::native_metadata::regular_owner_single_link(
   1145         FileType::from_raw_mode(status.st_mode).is_file(),
   1146         crate::native_metadata::link_count(status.st_nlink),
   1147         status.st_uid,
   1148         geteuid().as_raw(),
   1149     ) {
   1150         return None;
   1151     }
   1152     identity(&status).ok()
   1153 }
   1154 
   1155 fn identity(status: &rustix::fs::Stat) -> Result<FileIdentity, ServiceSqliteError> {
   1156     Ok(FileIdentity {
   1157         device: crate::native_metadata::device(status.st_dev)
   1158             .map_err(|_| backup_error(BackupFailureKind::StagingReplaced))?,
   1159         inode: status.st_ino,
   1160     })
   1161 }
   1162 
   1163 async fn verify_database_inventory(
   1164     connection: &mut SqliteConnection,
   1165 ) -> Result<(), ServiceSqliteError> {
   1166     let rows = sqlx::query(
   1167         "SELECT
   1168             seq,
   1169             typeof(name) = 'text' AS name_type_ok,
   1170             length(CAST(name AS BLOB)) AS name_length,
   1171             substr(CAST(name AS BLOB), 1, 5) AS name_prefix
   1172          FROM pragma_database_list
   1173          LIMIT 2",
   1174     )
   1175     .fetch_all(connection)
   1176     .await
   1177     .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
   1178     let first = rows
   1179         .first()
   1180         .ok_or_else(|| backup_error(BackupFailureKind::Capture))?;
   1181     let sequence = first
   1182         .try_get::<i64, _>(0)
   1183         .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
   1184     let name = crate::persisted_value::bounded_utf8(
   1185         first,
   1186         "name_type_ok",
   1187         "name_length",
   1188         "name_prefix",
   1189         1,
   1190         4,
   1191     )
   1192     .ok_or_else(|| backup_error(BackupFailureKind::Capture))?;
   1193     require_backup_condition(
   1194         database_inventory_matches(sequence, name, rows.len() > 1),
   1195         BackupFailureKind::Capture,
   1196     )?;
   1197     Ok(())
   1198 }
   1199 
   1200 fn database_inventory_matches(sequence: i64, name: &str, has_extra: bool) -> bool {
   1201     crate::all_constraints([sequence == 0, name == "main", !has_extra])
   1202 }
   1203 
   1204 async fn verify_database_metadata(
   1205     connection: &mut SqliteConnection,
   1206     expected: &ServiceDatabaseMetadata,
   1207 ) -> Result<(), ServiceSqliteError> {
   1208     let application_id = sqlx::query_scalar::<_, i64>("PRAGMA application_id")
   1209         .fetch_one(&mut *connection)
   1210         .await
   1211         .map_err(metadata_source)?;
   1212     let row_count = sqlx::query_scalar::<_, i64>(
   1213         "SELECT COUNT(*) FROM (SELECT 1 FROM radroots_service_metadata LIMIT 2)",
   1214     )
   1215     .fetch_one(&mut *connection)
   1216     .await
   1217     .map_err(metadata_source)?;
   1218     crate::require_condition(
   1219         crate::all_constraints([
   1220             row_count == 1,
   1221             application_id == i64::from(expected.application_id().get()),
   1222         ]),
   1223         ServiceSqliteErrorKind::Metadata,
   1224     )?;
   1225     let row = sqlx::query(
   1226         "SELECT
   1227                 typeof(service_id) = 'text' AS service_id_type_ok,
   1228                 length(CAST(service_id AS BLOB)) AS service_id_length,
   1229                 substr(CAST(service_id AS BLOB), 1, 129) AS service_id_prefix,
   1230                 typeof(instance_id) = 'text' AS instance_id_type_ok,
   1231                 length(CAST(instance_id AS BLOB)) AS instance_id_length,
   1232                 substr(CAST(instance_id AS BLOB), 1, 129) AS instance_id_prefix,
   1233                 typeof(source_generation) = 'blob' AS source_generation_type_ok,
   1234                 length(source_generation) AS source_generation_length,
   1235                 substr(source_generation, 1, 33) AS source_generation_prefix,
   1236                 CASE WHEN typeof(state_schema_version) = 'integer'
   1237                      THEN state_schema_version END,
   1238                 CASE WHEN typeof(created_at_unix_ms) = 'integer'
   1239                      THEN created_at_unix_ms END
   1240              FROM radroots_service_metadata
   1241              WHERE singleton = 1
   1242              LIMIT 1",
   1243     )
   1244     .fetch_optional(&mut *connection)
   1245     .await
   1246     .map_err(metadata_source)?;
   1247     let Some(row) = row else {
   1248         return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
   1249     };
   1250     let service = crate::persisted_value::bounded_utf8(
   1251         &row,
   1252         "service_id_type_ok",
   1253         "service_id_length",
   1254         "service_id_prefix",
   1255         1,
   1256         crate::persisted_value::MAX_IDENTIFIER_UTF8_BYTES,
   1257     )
   1258     .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?;
   1259     let instance = crate::persisted_value::bounded_utf8(
   1260         &row,
   1261         "instance_id_type_ok",
   1262         "instance_id_length",
   1263         "instance_id_prefix",
   1264         1,
   1265         crate::persisted_value::MAX_IDENTIFIER_UTF8_BYTES,
   1266     )
   1267     .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?;
   1268     let generation = crate::persisted_value::bounded_bytes(
   1269         &row,
   1270         "source_generation_type_ok",
   1271         "source_generation_length",
   1272         "source_generation_prefix",
   1273         32,
   1274         32,
   1275     )
   1276     .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))?;
   1277     let schema = row.try_get::<Option<i64>, _>(9).map_err(metadata_source)?;
   1278     let created_at = row.try_get::<Option<i64>, _>(10).map_err(metadata_source)?;
   1279     let (Some(schema), Some(created_at)) = (schema, created_at) else {
   1280         return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
   1281     };
   1282     crate::require_condition(
   1283         crate::all_constraints([
   1284             service == expected.service().as_str(),
   1285             instance == expected.instance().as_str(),
   1286             generation == expected.source_generation().as_bytes(),
   1287             schema == i64::from(expected.state_schema_version().get()),
   1288             created_at == i64::try_from(expected.created_at_unix_ms()).unwrap_or(-1),
   1289         ]),
   1290         ServiceSqliteErrorKind::Metadata,
   1291     )?;
   1292     Ok(())
   1293 }
   1294 
   1295 async fn verify_integrity(connection: &mut SqliteConnection) -> Result<(), ServiceSqliteError> {
   1296     let rows = sqlx::query(crate::persisted_value::INTEGRITY_CHECK_SQL)
   1297         .fetch_all(&mut *connection)
   1298         .await
   1299         .map_err(integrity_source)?;
   1300     let row = rows
   1301         .first()
   1302         .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
   1303     let value = crate::persisted_value::bounded_integrity_bytes(row);
   1304     crate::require_condition(
   1305         integrity_projection_is_ok(value) && rows.len() == 1,
   1306         ServiceSqliteErrorKind::Integrity,
   1307     )?;
   1308     let violation = sqlx::query_scalar::<_, i64>("SELECT 1 FROM pragma_foreign_key_check LIMIT 1")
   1309         .fetch_optional(connection)
   1310         .await
   1311         .map_err(integrity_source)?;
   1312     crate::require_condition(violation.is_none(), ServiceSqliteErrorKind::Integrity)?;
   1313     Ok(())
   1314 }
   1315 
   1316 fn integrity_projection_is_ok(value: Option<&[u8]>) -> bool {
   1317     value == Some(b"ok")
   1318 }
   1319 
   1320 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
   1321 enum BackupFailureKind {
   1322     InvalidStagingPath,
   1323     InvalidStagingParent,
   1324     StagingCollision,
   1325     CreateStaging,
   1326     CreateState,
   1327     StagingReplaced,
   1328     InvalidStagingInventory,
   1329     AlreadyActive,
   1330     Capture,
   1331     Cancelled,
   1332     HashState,
   1333     SyncState,
   1334     SyncStaging,
   1335     SyncParent,
   1336     Manifest,
   1337     Join,
   1338 }
   1339 
   1340 struct BackupFailure {
   1341     kind: BackupFailureKind,
   1342     source: Option<Box<dyn Error + Send + Sync + 'static>>,
   1343 }
   1344 
   1345 impl fmt::Debug for BackupFailure {
   1346     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1347         formatter
   1348             .debug_struct("BackupFailure")
   1349             .field("kind", &self.kind)
   1350             .field("source", &self.source.as_ref().map(|_| "[redacted]"))
   1351             .finish()
   1352     }
   1353 }
   1354 
   1355 impl fmt::Display for BackupFailure {
   1356     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
   1357         formatter.write_str(match self.kind {
   1358             BackupFailureKind::InvalidStagingPath => "backup staging path is invalid",
   1359             BackupFailureKind::InvalidStagingParent => "backup staging parent is invalid",
   1360             BackupFailureKind::StagingCollision => "backup staging destination already exists",
   1361             BackupFailureKind::CreateStaging => "backup staging directory could not be created",
   1362             BackupFailureKind::CreateState => "backup state member could not be created",
   1363             BackupFailureKind::StagingReplaced => "backup staging identity changed",
   1364             BackupFailureKind::InvalidStagingInventory => "backup staging inventory is invalid",
   1365             BackupFailureKind::AlreadyActive => "another backup capture is active",
   1366             BackupFailureKind::Capture => "online backup capture failed",
   1367             BackupFailureKind::Cancelled => "online backup capture was cancelled",
   1368             BackupFailureKind::HashState => "backup state member could not be hashed",
   1369             BackupFailureKind::SyncState => "backup state member could not be synchronized",
   1370             BackupFailureKind::SyncStaging => "backup staging directory could not be synchronized",
   1371             BackupFailureKind::SyncParent => "backup staging parent could not be synchronized",
   1372             BackupFailureKind::Manifest => "backup manifest could not be constructed",
   1373             BackupFailureKind::Join => "backup worker could not be joined",
   1374         })
   1375     }
   1376 }
   1377 
   1378 impl Error for BackupFailure {
   1379     fn source(&self) -> Option<&(dyn Error + 'static)> {
   1380         self.source
   1381             .as_deref()
   1382             .map(|source| source as &(dyn Error + 'static))
   1383     }
   1384 }
   1385 
   1386 fn backup_error(kind: BackupFailureKind) -> ServiceSqliteError {
   1387     ServiceSqliteError::with_source(
   1388         ServiceSqliteErrorKind::Backup,
   1389         BackupFailure { kind, source: None },
   1390     )
   1391 }
   1392 
   1393 fn require_backup_condition(
   1394     condition: bool,
   1395     kind: BackupFailureKind,
   1396 ) -> Result<(), ServiceSqliteError> {
   1397     if condition {
   1398         Ok(())
   1399     } else {
   1400         Err(backup_error(kind))
   1401     }
   1402 }
   1403 
   1404 fn backup_source(
   1405     kind: BackupFailureKind,
   1406     source: impl Error + Send + Sync + 'static,
   1407 ) -> ServiceSqliteError {
   1408     ServiceSqliteError::with_source(
   1409         ServiceSqliteErrorKind::Backup,
   1410         BackupFailure {
   1411             kind,
   1412             source: Some(Box::new(source)),
   1413         },
   1414     )
   1415 }
   1416 
   1417 fn integrity_source(source: sqlx::Error) -> ServiceSqliteError {
   1418     ServiceSqliteError::with_source(ServiceSqliteErrorKind::Integrity, source)
   1419 }
   1420 
   1421 fn metadata_source(source: sqlx::Error) -> ServiceSqliteError {
   1422     ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source)
   1423 }
   1424 
   1425 #[cfg(test)]
   1426 mod tests {
   1427     use core::num::NonZeroU32;
   1428     use std::os::unix::fs::{PermissionsExt, symlink};
   1429 
   1430     use radroots_runtime_paths::{InstanceId, ServiceId};
   1431     use radroots_storage::event::SourceGeneration;
   1432     use tempfile::tempdir;
   1433 
   1434     use super::*;
   1435 
   1436     #[test]
   1437     fn backup_failure_inventory_is_complete_and_source_aware() {
   1438         let cases = [
   1439             (
   1440                 BackupFailureKind::InvalidStagingPath,
   1441                 "backup staging path is invalid",
   1442             ),
   1443             (
   1444                 BackupFailureKind::InvalidStagingParent,
   1445                 "backup staging parent is invalid",
   1446             ),
   1447             (
   1448                 BackupFailureKind::StagingCollision,
   1449                 "backup staging destination already exists",
   1450             ),
   1451             (
   1452                 BackupFailureKind::CreateStaging,
   1453                 "backup staging directory could not be created",
   1454             ),
   1455             (
   1456                 BackupFailureKind::CreateState,
   1457                 "backup state member could not be created",
   1458             ),
   1459             (
   1460                 BackupFailureKind::StagingReplaced,
   1461                 "backup staging identity changed",
   1462             ),
   1463             (
   1464                 BackupFailureKind::InvalidStagingInventory,
   1465                 "backup staging inventory is invalid",
   1466             ),
   1467             (
   1468                 BackupFailureKind::AlreadyActive,
   1469                 "another backup capture is active",
   1470             ),
   1471             (BackupFailureKind::Capture, "online backup capture failed"),
   1472             (
   1473                 BackupFailureKind::Cancelled,
   1474                 "online backup capture was cancelled",
   1475             ),
   1476             (
   1477                 BackupFailureKind::HashState,
   1478                 "backup state member could not be hashed",
   1479             ),
   1480             (
   1481                 BackupFailureKind::SyncState,
   1482                 "backup state member could not be synchronized",
   1483             ),
   1484             (
   1485                 BackupFailureKind::SyncStaging,
   1486                 "backup staging directory could not be synchronized",
   1487             ),
   1488             (
   1489                 BackupFailureKind::SyncParent,
   1490                 "backup staging parent could not be synchronized",
   1491             ),
   1492             (
   1493                 BackupFailureKind::Manifest,
   1494                 "backup manifest could not be constructed",
   1495             ),
   1496             (BackupFailureKind::Join, "backup worker could not be joined"),
   1497         ];
   1498         for (kind, message) in cases {
   1499             let plain = BackupFailure { kind, source: None };
   1500             assert_eq!(plain.to_string(), message);
   1501             assert!(plain.source().is_none());
   1502             assert!(format!("{plain:?}").contains("source: None"));
   1503 
   1504             let sourced = BackupFailure {
   1505                 kind,
   1506                 source: Some(Box::new(std::io::Error::other("private-cause"))),
   1507             };
   1508             assert_eq!(sourced.to_string(), message);
   1509             assert_eq!(
   1510                 sourced.source().expect("source").to_string(),
   1511                 "private-cause"
   1512             );
   1513             let debug = format!("{sourced:?}");
   1514             assert!(debug.contains("[redacted]"));
   1515             assert!(!debug.contains("private-cause"));
   1516             assert!(require_backup_condition(true, kind).is_ok());
   1517             assert_eq!(
   1518                 require_backup_condition(false, kind)
   1519                     .expect_err("false condition")
   1520                     .kind(),
   1521                 ServiceSqliteErrorKind::Backup
   1522             );
   1523         }
   1524     }
   1525 
   1526     #[test]
   1527     fn staging_path_rejects_relative_parent_and_oversize_inputs() {
   1528         assert!(StagingPath::new(Path::new("relative/stage")).is_err());
   1529         assert!(StagingPath::new(Path::new("/tmp/../stage")).is_err());
   1530         assert!(StagingPath::new(Path::new("/")).is_err());
   1531         let large = format!("/tmp/{}", "x".repeat(MAX_STAGING_PATH_BYTES));
   1532         assert!(StagingPath::new(Path::new(&large)).is_err());
   1533     }
   1534 
   1535     #[test]
   1536     fn staging_guard_creates_exact_modes_and_cleans_uncommitted_state() {
   1537         let root = tempdir().expect("root");
   1538         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1539             .expect("parent mode");
   1540         let path = StagingPath::new(&root.path().join("backup-stage")).expect("path");
   1541         {
   1542             let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1543             staging.validate().expect("validate");
   1544             assert_eq!(
   1545                 std::fs::metadata(&path.full)
   1546                     .expect("directory metadata")
   1547                     .permissions()
   1548                     .mode()
   1549                     & 0o777,
   1550                 0o700
   1551             );
   1552             assert_eq!(
   1553                 std::fs::metadata(staging.state_path())
   1554                     .expect("state metadata")
   1555                     .permissions()
   1556                     .mode()
   1557                     & 0o777,
   1558                 0o600
   1559             );
   1560         }
   1561         assert!(!path.full.exists());
   1562     }
   1563 
   1564     #[test]
   1565     fn staging_guard_success_inventory_hash_sync_commit_and_empty_failure_are_exercised() {
   1566         let root = tempdir().expect("root");
   1567         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1568             .expect("parent mode");
   1569 
   1570         let empty_path = StagingPath::new(&root.path().join("empty-stage")).expect("empty path");
   1571         let empty = StagingGuard::create(&empty_path, &SystemCaptureOperations).expect("empty");
   1572         assert!(empty.hash_state(&AtomicBool::new(false)).is_err());
   1573         drop(empty);
   1574 
   1575         let path = StagingPath::new(&root.path().join("complete-stage")).expect("path");
   1576         let mut staging =
   1577             StagingGuard::create(&path, &SystemCaptureOperations).expect("complete stage");
   1578         std::fs::write(staging.state_path(), b"captured-state").expect("state bytes");
   1579         staging.validate_inventory().expect("singleton inventory");
   1580         let (length, digest) = staging
   1581             .hash_state(&AtomicBool::new(false))
   1582             .expect("hash state");
   1583         assert_eq!(length, 14);
   1584         assert_eq!(digest, Sha256::digest(b"captured-state").as_slice());
   1585         staging
   1586             .sync_state(&SystemCaptureOperations)
   1587             .expect("sync state");
   1588         staging
   1589             .sync_directories(&SystemCaptureOperations)
   1590             .expect("sync directories");
   1591         staging.commit();
   1592         drop(staging);
   1593         assert!(path.full.exists());
   1594         std::fs::remove_file(path.full.join(STATE_FILE_NAME)).expect("remove state");
   1595         std::fs::remove_dir(path.full).expect("remove stage");
   1596     }
   1597 
   1598     #[test]
   1599     fn staging_inventory_and_sidecar_cleanup_bind_exact_entries() {
   1600         let root = tempdir().expect("root");
   1601         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1602             .expect("parent mode");
   1603         let path = StagingPath::new(&root.path().join("backup-stage")).expect("path");
   1604         let mut staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1605         std::fs::write(staging.state_path(), b"state").expect("state bytes");
   1606         let wal = path.full.join(KNOWN_SIDECARS[0]);
   1607         std::fs::write(&wal, b"sidecar").expect("sidecar");
   1608         std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o600))
   1609             .expect("sidecar mode");
   1610         assert!(staging.validate_inventory().is_err());
   1611         staging.record_sidecars();
   1612         drop(staging);
   1613         assert!(!path.full.exists());
   1614 
   1615         let path = StagingPath::new(&root.path().join("unsafe-sidecar-stage")).expect("path");
   1616         let mut staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1617         std::fs::write(staging.state_path(), b"state").expect("state bytes");
   1618         let wal = path.full.join(KNOWN_SIDECARS[0]);
   1619         std::fs::write(&wal, b"unsafe-sidecar").expect("sidecar");
   1620         std::fs::set_permissions(&wal, std::fs::Permissions::from_mode(0o666))
   1621             .expect("unsafe sidecar mode");
   1622         let outside_wal = root.path().join("unsafe-sidecar-alias");
   1623         std::fs::hard_link(&wal, &outside_wal).expect("make sidecar unsafe by link count");
   1624         staging.record_sidecars();
   1625         drop(staging);
   1626         assert_eq!(
   1627             std::fs::read(wal).expect("unsafe sidecar preserved"),
   1628             b"unsafe-sidecar"
   1629         );
   1630         assert_eq!(
   1631             std::fs::read(outside_wal).expect("outside sidecar link preserved"),
   1632             b"unsafe-sidecar"
   1633         );
   1634     }
   1635 
   1636     #[test]
   1637     fn hash_cancellation_partial_cleanup_and_directory_replacement_are_exact() {
   1638         let root = tempdir().expect("root");
   1639         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1640             .expect("parent mode");
   1641 
   1642         let cancelled_path = StagingPath::new(&root.path().join("cancelled-stage")).expect("path");
   1643         let cancelled =
   1644             StagingGuard::create(&cancelled_path, &SystemCaptureOperations).expect("staging");
   1645         std::fs::write(cancelled.state_path(), b"state").expect("state bytes");
   1646         assert!(cancelled.hash_state(&AtomicBool::new(true)).is_err());
   1647         drop(cancelled);
   1648 
   1649         let exact_path = StagingPath::new(&root.path().join("exact-stage")).expect("path");
   1650         let exact = StagingGuard::create(&exact_path, &SystemCaptureOperations).expect("staging");
   1651         cleanup_partial_staging(
   1652             &exact.parent,
   1653             &exact.path.name,
   1654             Some(exact.directory_identity),
   1655             Some(&exact.directory),
   1656             Some(exact.state_identity),
   1657         );
   1658         assert!(!exact_path.full.exists());
   1659         drop(exact);
   1660 
   1661         let replaced_path = StagingPath::new(&root.path().join("replaced-stage")).expect("path");
   1662         let replaced =
   1663             StagingGuard::create(&replaced_path, &SystemCaptureOperations).expect("staging");
   1664         let retired = root.path().join("retired-stage");
   1665         std::fs::rename(&replaced_path.full, &retired).expect("retire governed directory");
   1666         std::fs::create_dir(&replaced_path.full).expect("replacement directory");
   1667         std::fs::set_permissions(&replaced_path.full, std::fs::Permissions::from_mode(0o700))
   1668             .expect("replacement mode");
   1669         std::fs::write(replaced_path.full.join("foreign"), b"foreign").expect("foreign entry");
   1670         drop(replaced);
   1671         assert_eq!(
   1672             std::fs::read(replaced_path.full.join("foreign")).expect("replacement survives"),
   1673             b"foreign"
   1674         );
   1675     }
   1676 
   1677     #[test]
   1678     fn staging_guard_rejects_collision_and_preserves_replacement() {
   1679         let root = tempdir().expect("root");
   1680         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1681             .expect("parent mode");
   1682         let full = root.path().join("backup-stage");
   1683         std::fs::create_dir(&full).expect("collision");
   1684         assert!(
   1685             StagingGuard::create(
   1686                 &StagingPath::new(&full).expect("path"),
   1687                 &SystemCaptureOperations,
   1688             )
   1689             .is_err()
   1690         );
   1691         std::fs::remove_dir(&full).expect("remove collision");
   1692 
   1693         let path = StagingPath::new(&full).expect("path");
   1694         let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1695         let original = staging.state_path();
   1696         let replacement = full.join("replacement");
   1697         std::fs::write(&replacement, b"foreign").expect("replacement");
   1698         std::fs::set_permissions(&replacement, std::fs::Permissions::from_mode(0o600))
   1699             .expect("replacement mode");
   1700         std::fs::rename(&replacement, &original).expect("replace state");
   1701         drop(staging);
   1702         assert_eq!(std::fs::read(&original).expect("preserved"), b"foreign");
   1703         assert!(full.exists());
   1704     }
   1705 
   1706     #[test]
   1707     fn staging_guard_rejects_insecure_or_symlinked_parent_without_mutation() {
   1708         let root = tempdir().expect("root");
   1709         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1710             .expect("root mode");
   1711         let insecure = root.path().join("insecure");
   1712         std::fs::create_dir(&insecure).expect("insecure parent");
   1713         std::fs::set_permissions(&insecure, std::fs::Permissions::from_mode(0o770))
   1714             .expect("insecure mode");
   1715         let insecure_stage = insecure.join("stage");
   1716         let insecure_error = match StagingGuard::create(
   1717             &StagingPath::new(&insecure_stage).expect("insecure staging path"),
   1718             &SystemCaptureOperations,
   1719         ) {
   1720             Ok(_) => panic!("group-writable parent must be rejected"),
   1721             Err(error) => error,
   1722         };
   1723         assert_eq!(insecure_error.kind(), ServiceSqliteErrorKind::Backup);
   1724         assert!(!insecure_stage.exists());
   1725 
   1726         let parent = root.path().join("real-parent");
   1727         std::fs::create_dir(&parent).expect("real parent");
   1728         std::fs::set_permissions(&parent, std::fs::Permissions::from_mode(0o700))
   1729             .expect("real parent mode");
   1730         let alias = root.path().join("parent-alias");
   1731         symlink(&parent, &alias).expect("parent symlink");
   1732         let symlink_stage = alias.join("stage");
   1733         assert!(
   1734             StagingGuard::create(
   1735                 &StagingPath::new(&symlink_stage).expect("symlink staging path"),
   1736                 &SystemCaptureOperations,
   1737             )
   1738             .is_err()
   1739         );
   1740         assert!(!parent.join("stage").exists());
   1741     }
   1742 
   1743     #[test]
   1744     fn staging_guard_rejects_hardlinked_state_and_preserves_other_link() {
   1745         let root = tempdir().expect("root");
   1746         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1747             .expect("parent mode");
   1748         let path = StagingPath::new(&root.path().join("backup-stage")).expect("path");
   1749         let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1750         std::fs::write(staging.state_path(), b"captured").expect("state bytes");
   1751         let outside = root.path().join("outside-link");
   1752         std::fs::hard_link(staging.state_path(), &outside).expect("hard link");
   1753         assert!(staging.validate().is_err());
   1754         drop(staging);
   1755         assert_eq!(
   1756             std::fs::read(&outside).expect("other link survives"),
   1757             b"captured"
   1758         );
   1759         assert!(!path.full.exists());
   1760     }
   1761 
   1762     #[test]
   1763     fn partial_cleanup_preserves_replaced_directory_and_state_entries() {
   1764         let root = tempdir().expect("root");
   1765         std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
   1766             .expect("root mode");
   1767         let parent = File::from(
   1768             open(
   1769                 root.path(),
   1770                 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1771                 Mode::empty(),
   1772             )
   1773             .expect("open parent"),
   1774         );
   1775         mkdirat(&parent, "stage", Mode::RUSR | Mode::WUSR | Mode::XUSR)
   1776             .expect("create original stage");
   1777         let original_identity =
   1778             created_directory_identity(&parent, OsStr::new("stage")).expect("original identity");
   1779         let original = File::from(
   1780             openat(
   1781                 &parent,
   1782                 "stage",
   1783                 OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
   1784                 Mode::empty(),
   1785             )
   1786             .expect("open original stage"),
   1787         );
   1788         std::fs::rename(root.path().join("stage"), root.path().join("retired"))
   1789             .expect("retire original stage");
   1790         std::fs::create_dir(root.path().join("stage")).expect("replacement stage");
   1791         std::fs::set_permissions(
   1792             root.path().join("stage"),
   1793             std::fs::Permissions::from_mode(0o700),
   1794         )
   1795         .expect("replacement mode");
   1796         std::fs::write(root.path().join("stage/foreign"), b"foreign").expect("replacement member");
   1797         cleanup_partial_staging(
   1798             &parent,
   1799             OsStr::new("stage"),
   1800             Some(original_identity),
   1801             Some(&original),
   1802             None,
   1803         );
   1804         assert_eq!(
   1805             std::fs::read(root.path().join("stage/foreign")).expect("foreign survives"),
   1806             b"foreign"
   1807         );
   1808 
   1809         let path = StagingPath::new(&root.path().join("state-stage")).expect("path");
   1810         let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
   1811         let retired_state = path.full.join("retired-state");
   1812         std::fs::rename(staging.state_path(), &retired_state).expect("retire state");
   1813         std::fs::write(staging.state_path(), b"foreign-state").expect("foreign state");
   1814         std::fs::set_permissions(staging.state_path(), std::fs::Permissions::from_mode(0o600))
   1815             .expect("foreign state mode");
   1816         cleanup_partial_staging(
   1817             &staging.parent,
   1818             &staging.path.name,
   1819             Some(staging.directory_identity),
   1820             Some(&staging.directory),
   1821             Some(staging.state_identity),
   1822         );
   1823         assert_eq!(
   1824             std::fs::read(staging.state_path()).expect("replacement state survives"),
   1825             b"foreign-state"
   1826         );
   1827         drop(staging);
   1828         assert!(path.full.exists());
   1829     }
   1830 
   1831     #[test]
   1832     fn integrity_projection_bounds_corrupt_text_before_semantic_acceptance() {
   1833         assert!(integrity_projection_is_ok(Some(b"ok")));
   1834         assert!(!integrity_projection_is_ok(Some(b"")));
   1835         assert!(!integrity_projection_is_ok(Some(b"not-ok")));
   1836         assert!(!integrity_projection_is_ok(None));
   1837     }
   1838 
   1839     #[test]
   1840     fn capture_database_inventory_rejects_each_independent_projection_drift() {
   1841         assert!(database_inventory_matches(0, "main", false));
   1842         assert!(!database_inventory_matches(1, "main", false));
   1843         assert!(!database_inventory_matches(0, "temp", false));
   1844         assert!(!database_inventory_matches(0, "main", true));
   1845     }
   1846 
   1847     async fn metadata_fixture() -> (SqliteConnection, ServiceDatabaseMetadata) {
   1848         let metadata = ServiceDatabaseMetadata::from_verified_backup(
   1849             ServiceId::new("myc").expect("service"),
   1850             InstanceId::new("primary").expect("instance"),
   1851             SourceGeneration::new([7; 32]).expect("generation"),
   1852             NonZeroU32::new(1).expect("schema"),
   1853             1_234,
   1854             crate::ServiceSqliteApplicationId::new(0x5244_5254).expect("application ID"),
   1855         )
   1856         .expect("metadata");
   1857         let mut connection = SqliteConnection::connect("sqlite::memory:")
   1858             .await
   1859             .expect("database");
   1860         let schema = format!(
   1861             "PRAGMA application_id = {};
   1862                  CREATE TABLE radroots_service_metadata (
   1863                     singleton INTEGER,
   1864                     service_id TEXT,
   1865                     instance_id TEXT,
   1866                     source_generation BLOB,
   1867                     state_schema_version INTEGER,
   1868                     created_at_unix_ms INTEGER
   1869                  );",
   1870             metadata.application_id().get()
   1871         );
   1872         sqlx::raw_sql(sqlx::AssertSqlSafe(schema.as_str()))
   1873             .execute(&mut connection)
   1874             .await
   1875             .expect("metadata schema");
   1876         sqlx::query("INSERT INTO radroots_service_metadata VALUES (1, ?, ?, ?, ?, ?)")
   1877             .bind(metadata.service().as_str())
   1878             .bind(metadata.instance().as_str())
   1879             .bind(metadata.source_generation().as_bytes().as_slice())
   1880             .bind(i64::from(metadata.state_schema_version().get()))
   1881             .bind(i64::try_from(metadata.created_at_unix_ms()).expect("time"))
   1882             .execute(&mut connection)
   1883             .await
   1884             .expect("metadata row");
   1885         (connection, metadata)
   1886     }
   1887 
   1888     #[tokio::test(flavor = "current_thread")]
   1889     async fn capture_database_inventory_metadata_and_integrity_accept_exact_state() {
   1890         let (mut connection, metadata) = metadata_fixture().await;
   1891         verify_database_inventory(&mut connection)
   1892             .await
   1893             .expect("main-only inventory");
   1894         verify_database_metadata(&mut connection, &metadata)
   1895             .await
   1896             .expect("exact metadata");
   1897         verify_integrity(&mut connection)
   1898             .await
   1899             .expect("healthy database");
   1900 
   1901         sqlx::query("ATTACH DATABASE ':memory:' AS extra")
   1902             .execute(&mut connection)
   1903             .await
   1904             .expect("attach extra");
   1905         assert!(verify_database_inventory(&mut connection).await.is_err());
   1906     }
   1907 
   1908     #[tokio::test(flavor = "current_thread")]
   1909     async fn capture_metadata_rejects_every_independent_identity_drift() {
   1910         for statement in [
   1911             "PRAGMA application_id = 1",
   1912             "INSERT INTO radroots_service_metadata SELECT 2, service_id, instance_id, source_generation, state_schema_version, created_at_unix_ms FROM radroots_service_metadata",
   1913             "UPDATE radroots_service_metadata SET service_id = 'rhi'",
   1914             "UPDATE radroots_service_metadata SET instance_id = 'secondary'",
   1915             "UPDATE radroots_service_metadata SET source_generation = zeroblob(32)",
   1916             "UPDATE radroots_service_metadata SET state_schema_version = 2",
   1917             "UPDATE radroots_service_metadata SET created_at_unix_ms = 1235",
   1918             "UPDATE radroots_service_metadata SET service_id = NULL",
   1919         ] {
   1920             let (mut connection, metadata) = metadata_fixture().await;
   1921             sqlx::raw_sql(statement)
   1922                 .execute(&mut connection)
   1923                 .await
   1924                 .expect("apply drift");
   1925             assert!(
   1926                 verify_database_metadata(&mut connection, &metadata)
   1927                     .await
   1928                     .is_err(),
   1929                 "drift must fail: {statement}"
   1930             );
   1931         }
   1932     }
   1933 
   1934     #[tokio::test(flavor = "current_thread")]
   1935     async fn capture_integrity_rejects_foreign_key_violations() {
   1936         let mut connection = SqliteConnection::connect("sqlite::memory:")
   1937             .await
   1938             .expect("database");
   1939         sqlx::raw_sql(
   1940             "PRAGMA foreign_keys = OFF;
   1941                  CREATE TABLE parent(id INTEGER PRIMARY KEY);
   1942                  CREATE TABLE child(parent_id INTEGER REFERENCES parent(id));
   1943                  INSERT INTO child(parent_id) VALUES (41);",
   1944         )
   1945         .execute(&mut connection)
   1946         .await
   1947         .expect("foreign-key violation fixture");
   1948         let error = verify_integrity(&mut connection)
   1949             .await
   1950             .expect_err("foreign-key drift must fail");
   1951         assert_eq!(error.kind(), ServiceSqliteErrorKind::Integrity);
   1952     }
   1953 
   1954     #[test]
   1955     fn active_capture_permit_is_exclusive_and_recoverable() {
   1956         let active = Arc::new(AtomicBool::new(false));
   1957         let first = CapturePermit::acquire(Arc::clone(&active)).expect("first permit");
   1958         assert!(CapturePermit::acquire(Arc::clone(&active)).is_err());
   1959         drop(first);
   1960         assert!(CapturePermit::acquire(active).is_ok());
   1961     }
   1962 }