lib

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

finalize.rs (33701B)


      1 //! Atomic installation of one completely verified offline restore stage.
      2 
      3 #[cfg(any(target_os = "linux", target_os = "macos"))]
      4 use core::fmt;
      5 
      6 use crate::{ServiceSqliteError, ServiceSqliteErrorKind, StagedServiceRestore};
      7 
      8 #[cfg(any(target_os = "linux", target_os = "macos"))]
      9 use {
     10     super::{
     11         RestoreArtifactExpectation, RestoreMarkerBinding, RestoreRecoveryMarker,
     12         RestoreRecoveryPhase,
     13         marker::{BACKUP_FILE_NAME, LIVE_FILE_NAME, STAGED_FILE_NAME},
     14         stage::NativeStagedServiceRestore,
     15     },
     16     rustix::{
     17         fs::{FileType, Mode, OFlags, RenameFlags, fstat, openat, renameat_with},
     18         process::geteuid,
     19     },
     20     sha2::{Digest, Sha256},
     21     std::{
     22         error::Error,
     23         fs::File,
     24         os::unix::fs::FileExt,
     25         sync::{
     26             Arc,
     27             atomic::{AtomicBool, AtomicU8, Ordering},
     28         },
     29     },
     30 };
     31 
     32 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     33 use super::marker::MARKER_NEXT_FILE_NAME;
     34 
     35 #[cfg(any(target_os = "linux", target_os = "macos"))]
     36 const HASH_BUFFER_BYTES: usize = 64 * 1_024;
     37 #[cfg(any(target_os = "linux", target_os = "macos"))]
     38 const PHASE_BEFORE_PREPARED: u8 = 1;
     39 #[cfg(any(target_os = "linux", target_os = "macos"))]
     40 const PHASE_COMMIT_OWNED: u8 = 2;
     41 #[cfg(any(target_os = "linux", target_os = "macos"))]
     42 const PHASE_AFTER_PREPARED: u8 = 3;
     43 
     44 #[cfg(any(target_os = "linux", target_os = "macos"))]
     45 const CANCELLABLE: u8 = 0;
     46 #[cfg(any(target_os = "linux", target_os = "macos"))]
     47 const CANCELLED: u8 = 1;
     48 #[cfg(any(target_os = "linux", target_os = "macos"))]
     49 const COMMIT_OWNED: u8 = 2;
     50 
     51 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     52 pub(crate) static TEST_FINALIZE_PHASE: AtomicU8 = AtomicU8::new(0);
     53 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     54 pub(crate) static TEST_FINALIZE_BLOCK_PHASE: AtomicU8 = AtomicU8::new(0);
     55 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     56 static TEST_FINALIZE_FAILURE: AtomicU8 = AtomicU8::new(0);
     57 
     58 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     59 pub(crate) const TEST_PHASE_BEFORE_PREPARED: u8 = PHASE_BEFORE_PREPARED;
     60 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     61 pub(crate) const TEST_PHASE_COMMIT_OWNED: u8 = PHASE_COMMIT_OWNED;
     62 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
     63 pub(crate) const TEST_PHASE_AFTER_PREPARED: u8 = PHASE_AFTER_PREPARED;
     64 
     65 /// Atomically installs a completely verified adjacent restore stage.
     66 ///
     67 /// Cancellation observed before the worker atomically claims commit ownership
     68 /// leaves the live database untouched and attempts exact stage cleanup. Caller
     69 /// loss after that in-memory handoff has an unknown immediate outcome, even if
     70 /// the durable `prepared` marker has not appeared yet: the owned worker retains
     71 /// writer authority until it either fails before durability or establishes
     72 /// recovery evidence and continues. Once `prepared` is durable, the staged
     73 /// artifact is retained unconditionally for recovery. A successful return
     74 /// provides no open database handle; the next open must reconcile and retire
     75 /// the retained marker and old live database.
     76 pub async fn finalize_staged_restore(
     77     staged: StagedServiceRestore,
     78 ) -> Result<(), ServiceSqliteError> {
     79     #[cfg(any(target_os = "linux", target_os = "macos"))]
     80     {
     81         let failpoints = crate::failpoint::DurabilityFailpoints::default();
     82         let cancellation = Arc::new(AtomicU8::new(CANCELLABLE));
     83         let cancellation_on_drop = CancellationOnDrop::new(Arc::clone(&cancellation));
     84         let native = staged.into_native();
     85         let result = tokio::task::spawn_blocking(move || {
     86             finalize_native(
     87                 native,
     88                 &cancellation,
     89                 &SystemFinalizeOperations,
     90                 &failpoints,
     91             )
     92         })
     93         .await
     94         .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?;
     95         cancellation_on_drop.disarm();
     96         result
     97     }
     98 
     99     #[cfg(not(any(target_os = "linux", target_os = "macos")))]
    100     {
    101         let _ = staged;
    102         Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Restore))
    103     }
    104 }
    105 
    106 #[cfg(any(target_os = "linux", target_os = "macos"))]
    107 fn finalize_native(
    108     staged: NativeStagedServiceRestore,
    109     cancellation: &AtomicU8,
    110     operations: &dyn FinalizeOperations,
    111     failpoints: &crate::failpoint::DurabilityFailpoints,
    112 ) -> Result<(), ServiceSqliteError> {
    113     staged.validate()?;
    114     check_cancel(cancellation)?;
    115     let live = authority_checked(&staged, || open_live(staged.directory()))??;
    116     let live_artifact = staged.live_artifact();
    117     authority_checked(&staged, || {
    118         verify_named_artifact(
    119             staged.directory(),
    120             LIVE_FILE_NAME,
    121             &live,
    122             live_artifact,
    123             Some(cancellation),
    124         )
    125     })??;
    126     authority_checked(&staged, || {
    127         verify_named_artifact(
    128             staged.directory(),
    129             STAGED_FILE_NAME,
    130             staged.staged_file(),
    131             staged.artifact(),
    132             Some(cancellation),
    133         )
    134     })??;
    135     authority_checked(&staged, || {
    136         live.sync_all()
    137             .map_err(|source| finalize_source(FinalizeFailureKind::SyncLive, source))
    138     })??;
    139     authority_checked(&staged, || {
    140         staged
    141             .staged_file()
    142             .sync_all()
    143             .map_err(|source| finalize_source(FinalizeFailureKind::SyncStaged, source))
    144     })??;
    145     test_phase(PHASE_BEFORE_PREPARED, cancellation, true)?;
    146     staged.validate()?;
    147     claim_commit_ownership(cancellation)?;
    148     test_phase(PHASE_COMMIT_OWNED, cancellation, false)?;
    149 
    150     let marker = RestoreRecoveryMarker::prepared(
    151         staged.metadata(),
    152         staged.manifest_digest(),
    153         live_artifact,
    154         staged.artifact(),
    155     )
    156     .map_err(|source| ServiceSqliteError::with_source(ServiceSqliteErrorKind::Restore, source))?;
    157     let on_durable = || {
    158         staged.disarm_cleanup();
    159         let _ = test_phase(PHASE_AFTER_PREPARED, cancellation, false);
    160     };
    161     #[cfg(test)]
    162     let marker_result = if operations.drift_authority_during_marker_sync() {
    163         RestoreMarkerBinding::test_create_with_durable_authority_drift(
    164             staged.paths(),
    165             staged.authority(),
    166             &marker,
    167             on_durable,
    168         )
    169     } else {
    170         RestoreMarkerBinding::create_with_durable_callback_and_failpoints(
    171             staged.paths(),
    172             staged.authority(),
    173             &marker,
    174             failpoints,
    175             on_durable,
    176         )
    177     };
    178     #[cfg(not(test))]
    179     let marker_result = RestoreMarkerBinding::create_with_durable_callback_and_failpoints(
    180         staged.paths(),
    181         staged.authority(),
    182         &marker,
    183         failpoints,
    184         on_durable,
    185     );
    186     let mut marker = marker_result?;
    187 
    188     rename_and_sync(
    189         &staged,
    190         operations,
    191         RenameStep::RetainLive,
    192         RenameArtifact {
    193             source_name: LIVE_FILE_NAME,
    194             destination_name: BACKUP_FILE_NAME,
    195             held: &live,
    196             expected: live_artifact,
    197         },
    198         failpoints,
    199     )?;
    200     authority_checked(&staged, || {
    201         operations.after_directory_sync(staged.directory(), RenameStep::RetainLive)
    202     })??;
    203     marker = marker.advance_with_failpoints(
    204         staged.paths(),
    205         staged.authority(),
    206         RestoreRecoveryPhase::LiveRetained,
    207         failpoints,
    208     )?;
    209 
    210     rename_and_sync(
    211         &staged,
    212         operations,
    213         RenameStep::InstallStage,
    214         RenameArtifact {
    215             source_name: STAGED_FILE_NAME,
    216             destination_name: LIVE_FILE_NAME,
    217             held: staged.staged_file(),
    218             expected: staged.artifact(),
    219         },
    220         failpoints,
    221     )?;
    222     authority_checked(&staged, || {
    223         operations.after_directory_sync(staged.directory(), RenameStep::InstallStage)
    224     })??;
    225     marker = marker.advance_with_failpoints(
    226         staged.paths(),
    227         staged.authority(),
    228         RestoreRecoveryPhase::ReplacementInstalled,
    229         failpoints,
    230     )?;
    231     require_finalize_condition(
    232         marker.marker().phase() == RestoreRecoveryPhase::ReplacementInstalled,
    233         FinalizeFailureKind::Marker,
    234     )?;
    235     staged.validate_finalization_authority()?;
    236     drop(marker);
    237     drop(staged);
    238     Ok(())
    239 }
    240 
    241 #[cfg(any(target_os = "linux", target_os = "macos"))]
    242 fn rename_and_sync(
    243     staged: &NativeStagedServiceRestore,
    244     operations: &dyn FinalizeOperations,
    245     step: RenameStep,
    246     artifact: RenameArtifact<'_>,
    247     failpoints: &crate::failpoint::DurabilityFailpoints,
    248 ) -> Result<(), ServiceSqliteError> {
    249     authority_checked(staged, || {
    250         verify_named_artifact(
    251             staged.directory(),
    252             artifact.source_name,
    253             artifact.held,
    254             artifact.expected,
    255             None,
    256         )
    257     })??;
    258     authority_checked(staged, || {
    259         hit(
    260             failpoints,
    261             step.before_rename_failpoint(),
    262             step.rename_failure(),
    263         )
    264     })??;
    265     let rename = authority_checked(staged, || {
    266         operations
    267             .rename(
    268                 staged.directory(),
    269                 artifact.source_name,
    270                 artifact.destination_name,
    271                 step,
    272             )
    273             .map_err(|source| finalize_source(step.rename_failure(), source))
    274     })?;
    275     rename?;
    276     authority_checked(staged, || {
    277         hit(
    278             failpoints,
    279             step.after_rename_failpoint(),
    280             step.rename_failure(),
    281         )
    282     })??;
    283     authority_checked(staged, || {
    284         verify_named_artifact(
    285             staged.directory(),
    286             artifact.destination_name,
    287             artifact.held,
    288             artifact.expected,
    289             None,
    290         )
    291     })??;
    292     authority_checked(staged, || {
    293         hit(
    294             failpoints,
    295             step.before_sync_failpoint(),
    296             step.sync_failure(),
    297         )
    298     })??;
    299     let sync = authority_checked(staged, || {
    300         operations
    301             .sync_directory(staged.directory(), step)
    302             .map_err(|source| finalize_source(step.sync_failure(), source))
    303     })?;
    304     sync?;
    305     authority_checked(staged, || {
    306         hit(failpoints, step.after_sync_failpoint(), step.sync_failure())
    307     })??;
    308     authority_checked(staged, || {
    309         verify_named_artifact(
    310             staged.directory(),
    311             artifact.destination_name,
    312             artifact.held,
    313             artifact.expected,
    314             None,
    315         )
    316     })??;
    317     Ok(())
    318 }
    319 
    320 #[cfg(any(target_os = "linux", target_os = "macos"))]
    321 fn authority_checked<T>(
    322     staged: &NativeStagedServiceRestore,
    323     operation: impl FnOnce() -> T,
    324 ) -> Result<T, ServiceSqliteError> {
    325     staged.validate_finalization_authority()?;
    326     let result = operation();
    327     staged.validate_finalization_authority()?;
    328     Ok(result)
    329 }
    330 
    331 #[cfg(any(target_os = "linux", target_os = "macos"))]
    332 fn open_live(directory: &File) -> Result<File, ServiceSqliteError> {
    333     openat(
    334         directory,
    335         LIVE_FILE_NAME,
    336         OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
    337         Mode::empty(),
    338     )
    339     .map(File::from)
    340     .map_err(|source| finalize_source(FinalizeFailureKind::Live, source))
    341 }
    342 
    343 #[cfg(any(target_os = "linux", target_os = "macos"))]
    344 fn verify_named_artifact(
    345     directory: &File,
    346     name: &str,
    347     held: &File,
    348     expected: RestoreArtifactExpectation,
    349     cancellation: Option<&AtomicU8>,
    350 ) -> Result<(), ServiceSqliteError> {
    351     let current = openat(
    352         directory,
    353         name,
    354         OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
    355         Mode::empty(),
    356     )
    357     .map(File::from)
    358     .map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
    359     let held_status =
    360         fstat(held).map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
    361     let current_status =
    362         fstat(&current).map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
    363     validate_status(&held_status, Some(expected.byte_length()))?;
    364     validate_status(&current_status, Some(expected.byte_length()))?;
    365     let held_identity = (
    366         crate::native_metadata::device(held_status.st_dev)
    367             .map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?,
    368         held_status.st_ino,
    369     );
    370     let current_identity = (
    371         crate::native_metadata::device(current_status.st_dev)
    372             .map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?,
    373         current_status.st_ino,
    374     );
    375     require_finalize_condition(
    376         crate::native_metadata::identity_pair_matches(
    377             held_identity.0,
    378             held_identity.1,
    379             current_identity.0,
    380             current_identity.1,
    381             expected.device(),
    382             expected.inode(),
    383         ),
    384         FinalizeFailureKind::Artifact,
    385     )?;
    386     require_finalize_condition(
    387         hash_exact(held, expected.byte_length(), cancellation)? == expected.sha256(),
    388         FinalizeFailureKind::Artifact,
    389     )?;
    390     Ok(())
    391 }
    392 
    393 #[cfg(any(target_os = "linux", target_os = "macos"))]
    394 fn validate_status(
    395     status: &rustix::fs::Stat,
    396     expected_length: Option<u64>,
    397 ) -> Result<(), ServiceSqliteError> {
    398     let length =
    399         u64::try_from(status.st_size).map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?;
    400     require_finalize_condition(
    401         crate::all_constraints([
    402             crate::native_metadata::exact_regular_file(
    403                 FileType::from_raw_mode(status.st_mode).is_file(),
    404                 crate::native_metadata::link_count(status.st_nlink),
    405                 status.st_uid,
    406                 geteuid().as_raw(),
    407                 crate::native_metadata::mode(status.st_mode),
    408             ),
    409             crate::native_metadata::valid_artifact_length(length, expected_length),
    410         ]),
    411         FinalizeFailureKind::Artifact,
    412     )?;
    413     Ok(())
    414 }
    415 
    416 #[cfg(any(target_os = "linux", target_os = "macos"))]
    417 fn hash_exact(
    418     file: &File,
    419     expected_length: u64,
    420     cancellation: Option<&AtomicU8>,
    421 ) -> Result<[u8; 32], ServiceSqliteError> {
    422     let mut offset = 0_u64;
    423     let mut buffer = [0_u8; HASH_BUFFER_BYTES];
    424     let mut hasher = Sha256::new();
    425     while offset < expected_length {
    426         if cancellation.is_some_and(|state| state.load(Ordering::Acquire) == CANCELLED) {
    427             return Err(finalize_error(FinalizeFailureKind::Cancelled));
    428         }
    429         let requested = usize::try_from((expected_length - offset).min(HASH_BUFFER_BYTES as u64))
    430             .map_err(|_| finalize_error(FinalizeFailureKind::Hash))?;
    431         let read = file
    432             .read_at(&mut buffer[..requested], offset)
    433             .map_err(|source| finalize_source(FinalizeFailureKind::Hash, source))?;
    434         if read == 0 {
    435             return Err(finalize_error(FinalizeFailureKind::Hash));
    436         }
    437         hasher.update(&buffer[..read]);
    438         offset = offset
    439             .checked_add(
    440                 u64::try_from(read).map_err(|_| finalize_error(FinalizeFailureKind::Hash))?,
    441             )
    442             .ok_or_else(|| finalize_error(FinalizeFailureKind::Hash))?;
    443     }
    444     require_finalize_condition(
    445         !cancellation.is_some_and(|state| state.load(Ordering::Acquire) == CANCELLED),
    446         FinalizeFailureKind::Cancelled,
    447     )?;
    448     let mut extra = [0_u8; 1];
    449     if file
    450         .read_at(&mut extra, expected_length)
    451         .map_err(|source| finalize_source(FinalizeFailureKind::Hash, source))?
    452         != 0
    453     {
    454         return Err(finalize_error(FinalizeFailureKind::Hash));
    455     }
    456     Ok(hasher.finalize().into())
    457 }
    458 
    459 #[cfg(any(target_os = "linux", target_os = "macos"))]
    460 fn check_cancel(cancellation: &AtomicU8) -> Result<(), ServiceSqliteError> {
    461     if cancellation.load(Ordering::Acquire) == CANCELLED {
    462         Err(finalize_error(FinalizeFailureKind::Cancelled))
    463     } else {
    464         Ok(())
    465     }
    466 }
    467 
    468 #[cfg(any(target_os = "linux", target_os = "macos"))]
    469 fn require_finalize_condition(
    470     condition: bool,
    471     kind: FinalizeFailureKind,
    472 ) -> Result<(), ServiceSqliteError> {
    473     if condition {
    474         Ok(())
    475     } else {
    476         Err(finalize_error(kind))
    477     }
    478 }
    479 
    480 #[cfg(any(target_os = "linux", target_os = "macos"))]
    481 fn claim_commit_ownership(cancellation: &AtomicU8) -> Result<(), ServiceSqliteError> {
    482     cancellation
    483         .compare_exchange(
    484             CANCELLABLE,
    485             COMMIT_OWNED,
    486             Ordering::AcqRel,
    487             Ordering::Acquire,
    488         )
    489         .map(|_| ())
    490         .map_err(|_| finalize_error(FinalizeFailureKind::Cancelled))
    491 }
    492 
    493 #[cfg(any(target_os = "linux", target_os = "macos"))]
    494 fn test_phase(
    495     phase: u8,
    496     cancellation: &AtomicU8,
    497     cancellable: bool,
    498 ) -> Result<(), ServiceSqliteError> {
    499     #[cfg(test)]
    500     {
    501         TEST_FINALIZE_PHASE.store(phase, Ordering::Release);
    502         while TEST_FINALIZE_BLOCK_PHASE.load(Ordering::Acquire) == phase {
    503             if cancellable && cancellation.load(Ordering::Acquire) == CANCELLED {
    504                 return Err(finalize_error(FinalizeFailureKind::Cancelled));
    505             }
    506             std::thread::yield_now();
    507         }
    508     }
    509     #[cfg(not(test))]
    510     let _ = (phase, cancellation, cancellable);
    511     Ok(())
    512 }
    513 
    514 #[cfg(any(target_os = "linux", target_os = "macos"))]
    515 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    516 enum RenameStep {
    517     RetainLive,
    518     InstallStage,
    519 }
    520 
    521 #[cfg(any(target_os = "linux", target_os = "macos"))]
    522 struct RenameArtifact<'a> {
    523     source_name: &'static str,
    524     destination_name: &'static str,
    525     held: &'a File,
    526     expected: RestoreArtifactExpectation,
    527 }
    528 
    529 #[cfg(any(target_os = "linux", target_os = "macos"))]
    530 impl RenameStep {
    531     const fn rename_failure(self) -> FinalizeFailureKind {
    532         match self {
    533             Self::RetainLive => FinalizeFailureKind::RetainLive,
    534             Self::InstallStage => FinalizeFailureKind::InstallStage,
    535         }
    536     }
    537 
    538     const fn sync_failure(self) -> FinalizeFailureKind {
    539         match self {
    540             Self::RetainLive => FinalizeFailureKind::SyncRetained,
    541             Self::InstallStage => FinalizeFailureKind::SyncInstalled,
    542         }
    543     }
    544 
    545     const fn before_rename_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
    546         match self {
    547             Self::RetainLive => {
    548                 crate::failpoint::DurabilityFailpoint::RestoreBeforeRetainLiveRename
    549             }
    550             Self::InstallStage => {
    551                 crate::failpoint::DurabilityFailpoint::RestoreBeforeInstallStageRename
    552             }
    553         }
    554     }
    555 
    556     const fn after_rename_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
    557         match self {
    558             Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreAfterRetainLiveRename,
    559             Self::InstallStage => {
    560                 crate::failpoint::DurabilityFailpoint::RestoreAfterInstallStageRename
    561             }
    562         }
    563     }
    564 
    565     const fn before_sync_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
    566         match self {
    567             Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreBeforeRetainLiveSync,
    568             Self::InstallStage => {
    569                 crate::failpoint::DurabilityFailpoint::RestoreBeforeInstallStageSync
    570             }
    571         }
    572     }
    573 
    574     const fn after_sync_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
    575         match self {
    576             Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreAfterRetainLiveSync,
    577             Self::InstallStage => {
    578                 crate::failpoint::DurabilityFailpoint::RestoreAfterInstallStageSync
    579             }
    580         }
    581     }
    582 }
    583 
    584 #[cfg(any(target_os = "linux", target_os = "macos"))]
    585 trait FinalizeOperations: Send + Sync {
    586     fn rename(
    587         &self,
    588         directory: &File,
    589         source: &str,
    590         destination: &str,
    591         step: RenameStep,
    592     ) -> std::io::Result<()>;
    593 
    594     fn sync_directory(&self, directory: &File, step: RenameStep) -> std::io::Result<()>;
    595 
    596     fn after_directory_sync(
    597         &self,
    598         _directory: &File,
    599         _step: RenameStep,
    600     ) -> Result<(), ServiceSqliteError> {
    601         Ok(())
    602     }
    603 
    604     #[cfg(test)]
    605     fn drift_authority_during_marker_sync(&self) -> bool {
    606         false
    607     }
    608 }
    609 
    610 #[cfg(any(target_os = "linux", target_os = "macos"))]
    611 struct SystemFinalizeOperations;
    612 
    613 #[cfg(any(target_os = "linux", target_os = "macos"))]
    614 impl FinalizeOperations for SystemFinalizeOperations {
    615     #[cfg_attr(coverage_nightly, coverage(off))]
    616     fn rename(
    617         &self,
    618         directory: &File,
    619         source: &str,
    620         destination: &str,
    621         _step: RenameStep,
    622     ) -> std::io::Result<()> {
    623         renameat_with(
    624             directory,
    625             source,
    626             directory,
    627             destination,
    628             RenameFlags::NOREPLACE,
    629         )
    630         .map_err(std::io::Error::from)
    631     }
    632 
    633     #[cfg_attr(coverage_nightly, coverage(off))]
    634     fn sync_directory(&self, directory: &File, _step: RenameStep) -> std::io::Result<()> {
    635         directory.sync_all()
    636     }
    637 }
    638 
    639 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    640 struct FailingFinalizeOperations;
    641 
    642 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    643 impl FinalizeOperations for FailingFinalizeOperations {
    644     fn drift_authority_during_marker_sync(&self) -> bool {
    645         TEST_FINALIZE_FAILURE.load(Ordering::Acquire) == 9
    646     }
    647 
    648     fn rename(
    649         &self,
    650         directory: &File,
    651         source: &str,
    652         destination: &str,
    653         step: RenameStep,
    654     ) -> std::io::Result<()> {
    655         let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
    656         let (before, after) = match step {
    657             RenameStep::RetainLive => (1, 2),
    658             RenameStep::InstallStage => (5, 6),
    659         };
    660         if failure == before {
    661             return Err(std::io::Error::other("injected pre-rename failure"));
    662         }
    663         SystemFinalizeOperations.rename(directory, source, destination, step)?;
    664         if failure == after {
    665             return Err(std::io::Error::other("injected post-rename failure"));
    666         }
    667         Ok(())
    668     }
    669 
    670     fn sync_directory(&self, directory: &File, step: RenameStep) -> std::io::Result<()> {
    671         let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
    672         let target = match step {
    673             RenameStep::RetainLive => 3,
    674             RenameStep::InstallStage => 7,
    675         };
    676         if failure == target {
    677             return Err(crate::failpoint::storage_full_error());
    678         }
    679         SystemFinalizeOperations.sync_directory(directory, step)
    680     }
    681 
    682     fn after_directory_sync(
    683         &self,
    684         directory: &File,
    685         step: RenameStep,
    686     ) -> Result<(), ServiceSqliteError> {
    687         let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
    688         let target = match step {
    689             RenameStep::RetainLive => 4,
    690             RenameStep::InstallStage => 8,
    691         };
    692         if failure != target {
    693             return Ok(());
    694         }
    695         openat(
    696             directory,
    697             MARKER_NEXT_FILE_NAME,
    698             OFlags::RDWR
    699                 | OFlags::CREATE
    700                 | OFlags::EXCL
    701                 | OFlags::NOFOLLOW
    702                 | OFlags::CLOEXEC
    703                 | OFlags::NONBLOCK,
    704             Mode::RUSR | Mode::WUSR,
    705         )
    706         .map(drop)
    707         .map_err(|source| finalize_source(FinalizeFailureKind::Marker, source))
    708     }
    709 }
    710 
    711 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    712 pub(crate) async fn test_finalize_with_failure(
    713     staged: StagedServiceRestore,
    714     failure: u8,
    715 ) -> Result<(), ServiceSqliteError> {
    716     TEST_FINALIZE_FAILURE.store(failure, Ordering::Release);
    717     let native = staged.into_native();
    718     let result = tokio::task::spawn_blocking(move || {
    719         finalize_native(
    720             native,
    721             &AtomicU8::new(CANCELLABLE),
    722             &FailingFinalizeOperations,
    723             &crate::failpoint::DurabilityFailpoints::default(),
    724         )
    725     })
    726     .await
    727     .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?;
    728     TEST_FINALIZE_FAILURE.store(0, Ordering::Release);
    729     result
    730 }
    731 
    732 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    733 pub(crate) async fn test_finalize_with_failpoint(
    734     staged: StagedServiceRestore,
    735     failpoints: crate::failpoint::DurabilityFailpoints,
    736 ) -> Result<(), ServiceSqliteError> {
    737     let native = staged.into_native();
    738     tokio::task::spawn_blocking(move || {
    739         finalize_native(
    740             native,
    741             &AtomicU8::new(CANCELLABLE),
    742             &SystemFinalizeOperations,
    743             &failpoints,
    744         )
    745     })
    746     .await
    747     .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?
    748 }
    749 
    750 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    751 pub(crate) fn reset_test_controls() {
    752     TEST_FINALIZE_PHASE.store(0, Ordering::Release);
    753     TEST_FINALIZE_BLOCK_PHASE.store(0, Ordering::Release);
    754     TEST_FINALIZE_FAILURE.store(0, Ordering::Release);
    755 }
    756 
    757 #[cfg(any(target_os = "linux", target_os = "macos"))]
    758 struct CancellationOnDrop {
    759     cancellation: Arc<AtomicU8>,
    760     armed: AtomicBool,
    761 }
    762 
    763 #[cfg(any(target_os = "linux", target_os = "macos"))]
    764 impl CancellationOnDrop {
    765     fn new(cancellation: Arc<AtomicU8>) -> Self {
    766         Self {
    767             cancellation,
    768             armed: AtomicBool::new(true),
    769         }
    770     }
    771 
    772     fn disarm(&self) {
    773         self.armed.store(false, Ordering::Release);
    774     }
    775 }
    776 
    777 #[cfg(any(target_os = "linux", target_os = "macos"))]
    778 impl Drop for CancellationOnDrop {
    779     fn drop(&mut self) {
    780         if self.armed.load(Ordering::Acquire) {
    781             let _ = self.cancellation.compare_exchange(
    782                 CANCELLABLE,
    783                 CANCELLED,
    784                 Ordering::AcqRel,
    785                 Ordering::Acquire,
    786             );
    787         }
    788     }
    789 }
    790 
    791 #[cfg(any(target_os = "linux", target_os = "macos"))]
    792 #[derive(Clone, Copy, Debug, PartialEq, Eq)]
    793 enum FinalizeFailureKind {
    794     Live,
    795     Artifact,
    796     Hash,
    797     SyncLive,
    798     SyncStaged,
    799     Marker,
    800     RetainLive,
    801     SyncRetained,
    802     InstallStage,
    803     SyncInstalled,
    804     Cancelled,
    805     Join,
    806 }
    807 
    808 #[cfg(any(target_os = "linux", target_os = "macos"))]
    809 struct FinalizeFailure {
    810     kind: FinalizeFailureKind,
    811     source: Option<Box<dyn Error + Send + Sync + 'static>>,
    812 }
    813 
    814 #[cfg(any(target_os = "linux", target_os = "macos"))]
    815 impl fmt::Debug for FinalizeFailure {
    816     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    817         formatter
    818             .debug_struct("FinalizeFailure")
    819             .field("kind", &self.kind)
    820             .field("source", &self.source.as_ref().map(|_| "[redacted]"))
    821             .finish()
    822     }
    823 }
    824 
    825 #[cfg(any(target_os = "linux", target_os = "macos"))]
    826 impl fmt::Display for FinalizeFailure {
    827     fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
    828         formatter.write_str(match self.kind {
    829             FinalizeFailureKind::Live => "live restore source is invalid",
    830             FinalizeFailureKind::Artifact => "restore artifact binding changed",
    831             FinalizeFailureKind::Hash => "restore artifact hash failed",
    832             FinalizeFailureKind::SyncLive => "live restore source sync failed",
    833             FinalizeFailureKind::SyncStaged => "staged restore sync failed",
    834             FinalizeFailureKind::Marker => "restore marker transition failed",
    835             FinalizeFailureKind::RetainLive => "live restore retention failed",
    836             FinalizeFailureKind::SyncRetained => "retained restore sync failed",
    837             FinalizeFailureKind::InstallStage => "restore installation failed",
    838             FinalizeFailureKind::SyncInstalled => "installed restore sync failed",
    839             FinalizeFailureKind::Cancelled => "restore finalization was cancelled",
    840             FinalizeFailureKind::Join => "restore finalization worker failed",
    841         })
    842     }
    843 }
    844 
    845 #[cfg(any(target_os = "linux", target_os = "macos"))]
    846 impl Error for FinalizeFailure {
    847     fn source(&self) -> Option<&(dyn Error + 'static)> {
    848         self.source
    849             .as_deref()
    850             .map(|source| source as &(dyn Error + 'static))
    851     }
    852 }
    853 
    854 #[cfg(any(target_os = "linux", target_os = "macos"))]
    855 fn finalize_error(kind: FinalizeFailureKind) -> ServiceSqliteError {
    856     ServiceSqliteError::with_source(
    857         ServiceSqliteErrorKind::Restore,
    858         FinalizeFailure { kind, source: None },
    859     )
    860 }
    861 
    862 #[cfg(any(target_os = "linux", target_os = "macos"))]
    863 fn finalize_source(
    864     kind: FinalizeFailureKind,
    865     source: impl Error + Send + Sync + 'static,
    866 ) -> ServiceSqliteError {
    867     ServiceSqliteError::with_source(
    868         ServiceSqliteErrorKind::Restore,
    869         FinalizeFailure {
    870             kind,
    871             source: Some(Box::new(source)),
    872         },
    873     )
    874 }
    875 
    876 #[cfg(any(target_os = "linux", target_os = "macos"))]
    877 fn hit(
    878     failpoints: &crate::failpoint::DurabilityFailpoints,
    879     point: crate::failpoint::DurabilityFailpoint,
    880     kind: FinalizeFailureKind,
    881 ) -> Result<(), ServiceSqliteError> {
    882     failpoints
    883         .hit(point)
    884         .map_err(|source| finalize_source(kind, source))
    885 }
    886 
    887 #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
    888 mod tests {
    889     use std::io::Write;
    890 
    891     use sha2::{Digest, Sha256};
    892 
    893     use super::*;
    894 
    895     #[test]
    896     fn finalize_failure_inventory_is_complete_and_source_aware() {
    897         let cases = [
    898             (FinalizeFailureKind::Live, "live restore source is invalid"),
    899             (
    900                 FinalizeFailureKind::Artifact,
    901                 "restore artifact binding changed",
    902             ),
    903             (FinalizeFailureKind::Hash, "restore artifact hash failed"),
    904             (
    905                 FinalizeFailureKind::SyncLive,
    906                 "live restore source sync failed",
    907             ),
    908             (
    909                 FinalizeFailureKind::SyncStaged,
    910                 "staged restore sync failed",
    911             ),
    912             (
    913                 FinalizeFailureKind::Marker,
    914                 "restore marker transition failed",
    915             ),
    916             (
    917                 FinalizeFailureKind::RetainLive,
    918                 "live restore retention failed",
    919             ),
    920             (
    921                 FinalizeFailureKind::SyncRetained,
    922                 "retained restore sync failed",
    923             ),
    924             (
    925                 FinalizeFailureKind::InstallStage,
    926                 "restore installation failed",
    927             ),
    928             (
    929                 FinalizeFailureKind::SyncInstalled,
    930                 "installed restore sync failed",
    931             ),
    932             (
    933                 FinalizeFailureKind::Cancelled,
    934                 "restore finalization was cancelled",
    935             ),
    936             (
    937                 FinalizeFailureKind::Join,
    938                 "restore finalization worker failed",
    939             ),
    940         ];
    941         for (kind, message) in cases {
    942             let plain = FinalizeFailure { kind, source: None };
    943             assert_eq!(plain.to_string(), message);
    944             assert!(plain.source().is_none());
    945             let sourced = FinalizeFailure {
    946                 kind,
    947                 source: Some(Box::new(std::io::Error::other("private-cause"))),
    948             };
    949             assert_eq!(sourced.to_string(), message);
    950             assert!(sourced.source().is_some());
    951             assert!(format!("{sourced:?}").contains("[redacted]"));
    952             assert!(require_finalize_condition(true, kind).is_ok());
    953             assert_eq!(
    954                 require_finalize_condition(false, kind)
    955                     .expect_err("false condition")
    956                     .kind(),
    957                 ServiceSqliteErrorKind::Restore
    958             );
    959         }
    960     }
    961 
    962     #[test]
    963     fn hash_and_cancellation_boundaries_fail_closed() {
    964         let root = tempfile::tempdir().expect("root");
    965         let path = root.path().join("artifact.sqlite");
    966         let payload = b"finalize-artifact";
    967         std::fs::write(&path, payload).expect("artifact");
    968         let file = File::open(&path).expect("open artifact");
    969         let length = u64::try_from(payload.len()).expect("length");
    970         let digest: [u8; 32] = Sha256::digest(payload).into();
    971         assert_eq!(hash_exact(&file, length, None).expect("exact hash"), digest);
    972         assert_eq!(
    973             hash_exact(&file, length + 1, None)
    974                 .expect_err("short artifact")
    975                 .kind(),
    976             ServiceSqliteErrorKind::Restore
    977         );
    978         assert_eq!(
    979             hash_exact(&file, length - 1, None)
    980                 .expect_err("long artifact")
    981                 .kind(),
    982             ServiceSqliteErrorKind::Restore
    983         );
    984 
    985         let cancelled = AtomicU8::new(CANCELLED);
    986         assert_eq!(
    987             hash_exact(&file, length, Some(&cancelled))
    988                 .expect_err("pre-read cancellation")
    989                 .kind(),
    990             ServiceSqliteErrorKind::Restore
    991         );
    992         assert_eq!(
    993             check_cancel(&cancelled).expect_err("cancelled").kind(),
    994             ServiceSqliteErrorKind::Restore
    995         );
    996         assert_eq!(
    997             claim_commit_ownership(&cancelled)
    998                 .expect_err("cancelled ownership handoff")
    999                 .kind(),
   1000             ServiceSqliteErrorKind::Restore
   1001         );
   1002         let active = AtomicU8::new(CANCELLABLE);
   1003         check_cancel(&active).expect("active");
   1004         claim_commit_ownership(&active).expect("commit ownership");
   1005         assert_eq!(active.load(Ordering::Acquire), COMMIT_OWNED);
   1006 
   1007         let empty_path = root.path().join("empty.sqlite");
   1008         let mut empty = File::create(&empty_path).expect("empty artifact");
   1009         empty.flush().expect("flush empty artifact");
   1010         let empty = File::open(empty_path).expect("open empty artifact");
   1011         assert_eq!(
   1012             hash_exact(&empty, 1, None)
   1013                 .expect_err("zero-byte read")
   1014                 .kind(),
   1015             ServiceSqliteErrorKind::Restore
   1016         );
   1017     }
   1018 }