commit 18620a5d289190e36353e32643587d2636d9c34e
parent 8a82c10898114d2d244e351fd82c585a2d1e9365
Author: triesap <tyson@radroots.org>
Date: Tue, 11 Aug 2026 22:53:08 +0000
service-sqlite: finalize restore atomically
- Bind live and staged artifacts before durable recovery evidence.
- Transfer commit ownership across cancellation without detached mutation.
- Install restores with no-replace renames and phase-synced markers.
- Reject unresolved recovery evidence before every database open.
Diffstat:
10 files changed, 1445 insertions(+), 19 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -216,6 +216,22 @@ Before editing code:
remain evidence that later admission rejects. Staging must not create a
marker, rename live state, retain an old live database, or install a
replacement; those are later finalization and recovery boundaries.
+- Atomic restore finalization consumes only a sealed `StagedServiceRestore`.
+ Staging must bind the exact live inode, length, and digest that finalization
+ will retain. The owned blocking worker creates and synchronizes `prepared`
+ before disarming stage cleanup, then uses descriptor-relative no-replace
+ renames and parent synchronization for live-to-backup and staged-to-live,
+ advancing the marker only after each durable rename. Cancellation observed
+ before the worker atomically claims commit ownership may cleanly stop;
+ caller loss after that handoff is an unknown immediate outcome, including
+ the interval before `prepared` is durable. Once `prepared` is durable, stage
+ cleanup must remain disarmed after every later error so the marker never
+ loses a bound artifact.
+ Successful finalization returns no host and leaves the old live database and
+ `replacement_installed` marker for open-time recovery. Until that recovery
+ exists, every pool open must reject marker or marker-scratch evidence as
+ `Recovery`. Finalization must not roll back, delete recovery evidence, reopen
+ SQLite, or expose paths, descriptors, marker controls, or rename controls.
- Runtime-management flows consume a sealed `RuntimeContext` for every service
instance. They must not reconstruct service paths from raw identifiers,
ambient selectors, or manager-owned roots, and registries must not persist
diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md
@@ -149,6 +149,30 @@ The operation has no hidden timeout. It does not create or advance a recovery
marker, rename or retain live state, install a replacement, or authorize reopen;
those operations remain the finalization and recovery checkpoints.
+`finalize_staged_restore` consumes that sealed stage in an owned blocking
+worker. The stage has already bound the exact live inode, length, and digest
+that will be retained. Finalization revalidates both retained descriptors,
+creates and synchronizes the `prepared` marker, and only then disarms automatic
+stage cleanup. It renames live to `state.restore-backup.sqlite` and staged to
+live with descriptor-relative no-replace operations. Each rename is followed
+by exact inode and hash verification, state-directory synchronization, and the
+corresponding marker advance to `live_retained` or
+`replacement_installed`.
+
+Cancellation observed before the worker's atomic commit-ownership handoff
+leaves live state untouched and attempts exact stage cleanup. Caller-task loss
+after that handoff has an unknown immediate outcome, including the short
+interval before `prepared` becomes durable; the worker retains writer
+authority until it either fails before durability or establishes recovery
+evidence and continues. Once `prepared` is durable, staged-artifact cleanup is
+disarmed and the bound stage remains available after every later error.
+Success returns no database host, retains the old live database and final
+marker, and requires a new open. Until open-time recovery is implemented,
+every initialized, read-write-existing, or read-only-inspection pool open
+refuses a marker or marker-scratch as `Recovery` before opening SQLite.
+Finalization does not delete evidence, roll back, reconcile an interruption,
+or reopen the database; those remain the next recovery checkpoint.
+
The crate owns mechanics only. Service-specific tables, SQL, repositories,
backup content policy, identity material, process lifecycle, and readiness
policy remain with the consuming service. The crate does not provide callers
diff --git a/crates/service_sqlite/src/initialize.rs b/crates/service_sqlite/src/initialize.rs
@@ -45,6 +45,14 @@ where
#[cfg(any(target_os = "linux", target_os = "macos"))]
{
+ authority.validate_for(paths)?;
+ let recovery = crate::restore::refuse_unresolved_recovery(authority.directory());
+ authority.validate_for(paths)?;
+ recovery?;
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ {
initialize_with_ops(
paths,
authority,
diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs
@@ -47,5 +47,5 @@ pub use migration::{
MigrationKind, MigrationName, MigrationTransactionExecutor,
};
pub use open::{OpenMode, ServiceSqlitePathError, ServiceSqlitePaths};
-pub use restore::{StagedServiceRestore, stage_verified_restore};
+pub use restore::{StagedServiceRestore, finalize_staged_restore, stage_verified_restore};
pub use status::{StorageHealth, StorageIntegrity, StorageStatus};
diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs
@@ -674,6 +674,25 @@ async fn open_connection_pool(
if !schema_catalog.matches_migrations(catalog) {
return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
}
+ let recovery_guard = match (mode, authority.as_ref(), inspection_guard.as_ref()) {
+ (OpenMode::Initialize | OpenMode::ReadWriteExisting, Some(authority), None) => {
+ authority.validate_for(paths)?;
+ let result = crate::restore::refuse_unresolved_recovery(authority.directory());
+ authority.validate_for(paths)?;
+ result
+ }
+ (OpenMode::ReadOnlyInspection, None, Some(inspection_guard)) => {
+ inspection_guard.validate_for(paths)?;
+ let result = crate::restore::refuse_unresolved_recovery(&inspection_guard.directory);
+ inspection_guard.validate_for(paths)?;
+ result
+ }
+ _ => Err(connection_error(
+ ServiceSqliteErrorKind::Authority,
+ ConnectionFailureKind::AuthorityMismatch,
+ )),
+ };
+ recovery_guard?;
let binding = match (mode, authority.as_ref(), inspection_guard.as_ref()) {
(OpenMode::Initialize | OpenMode::ReadWriteExisting, Some(authority), None) => {
authority.validate_for(paths)?;
@@ -1138,6 +1157,7 @@ impl ReadOnlyInspectionGuard {
ConnectionFailureKind::InspectionUnavailable,
));
}
+ let directory = File::from(directory);
let lock = openat(
&directory,
radroots_runtime_paths::SERVICE_STATE_LOCK_FILE_NAME,
@@ -1164,6 +1184,7 @@ impl ReadOnlyInspectionGuard {
inspection_error(ConnectionFailureKind::InspectionUnavailable)
}
})?;
+ crate::restore::refuse_unresolved_recovery(&directory)?;
for sidecar in [WAL_FILE_NAME, SHARED_MEMORY_FILE_NAME] {
match statat(&directory, sidecar, AtFlags::SYMLINK_NOFOLLOW) {
Err(error) if error == rustix::io::Errno::NOENT => {}
@@ -1226,7 +1247,7 @@ impl ReadOnlyInspectionGuard {
lock_device: u64::try_from(lock_status.st_dev)
.map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?,
lock_inode: lock_status.st_ino,
- directory: File::from(directory),
+ directory,
directory_device: u64::try_from(directory_status.st_dev)
.map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?,
directory_inode: directory_status.st_ino,
diff --git a/crates/service_sqlite/src/restore/finalize.rs b/crates/service_sqlite/src/restore/finalize.rs
@@ -0,0 +1,719 @@
+//! Atomic installation of one completely verified offline restore stage.
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+use core::fmt;
+
+use crate::{ServiceSqliteError, ServiceSqliteErrorKind, StagedServiceRestore};
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+use {
+ super::{
+ RestoreArtifactExpectation, RestoreMarkerBinding, RestoreRecoveryMarker,
+ RestoreRecoveryPhase,
+ marker::{BACKUP_FILE_NAME, LIVE_FILE_NAME, STAGED_FILE_NAME},
+ stage::NativeStagedServiceRestore,
+ },
+ rustix::{
+ fs::{FileType, Mode, OFlags, RenameFlags, fstat, openat, renameat_with},
+ process::geteuid,
+ },
+ sha2::{Digest, Sha256},
+ std::{
+ error::Error,
+ fs::File,
+ os::unix::fs::FileExt,
+ sync::{
+ Arc,
+ atomic::{AtomicBool, AtomicU8, Ordering},
+ },
+ },
+};
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+use super::marker::MARKER_NEXT_FILE_NAME;
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const HASH_BUFFER_BYTES: usize = 64 * 1_024;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const PHASE_BEFORE_PREPARED: u8 = 1;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const PHASE_COMMIT_OWNED: u8 = 2;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const PHASE_AFTER_PREPARED: u8 = 3;
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const CANCELLABLE: u8 = 0;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const CANCELLED: u8 = 1;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+const COMMIT_OWNED: u8 = 2;
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) static TEST_FINALIZE_PHASE: AtomicU8 = AtomicU8::new(0);
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) static TEST_FINALIZE_BLOCK_PHASE: AtomicU8 = AtomicU8::new(0);
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+static TEST_FINALIZE_FAILURE: AtomicU8 = AtomicU8::new(0);
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) const TEST_PHASE_BEFORE_PREPARED: u8 = PHASE_BEFORE_PREPARED;
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) const TEST_PHASE_COMMIT_OWNED: u8 = PHASE_COMMIT_OWNED;
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) const TEST_PHASE_AFTER_PREPARED: u8 = PHASE_AFTER_PREPARED;
+
+/// Atomically installs a completely verified adjacent restore stage.
+///
+/// Cancellation observed before the worker atomically claims commit ownership
+/// leaves the live database untouched and attempts exact stage cleanup. Caller
+/// loss after that in-memory handoff has an unknown immediate outcome, even if
+/// the durable `prepared` marker has not appeared yet: the owned worker retains
+/// writer authority until it either fails before durability or establishes
+/// recovery evidence and continues. Once `prepared` is durable, the staged
+/// artifact is retained unconditionally for recovery. A successful return
+/// provides no open database handle; the next open must reconcile and retire
+/// the retained marker and old live database.
+pub async fn finalize_staged_restore(
+ staged: StagedServiceRestore,
+) -> Result<(), ServiceSqliteError> {
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ {
+ let cancellation = Arc::new(AtomicU8::new(CANCELLABLE));
+ let cancellation_on_drop = CancellationOnDrop::new(Arc::clone(&cancellation));
+ let native = staged.into_native();
+ let result = tokio::task::spawn_blocking(move || {
+ finalize_native(native, &cancellation, &SystemFinalizeOperations)
+ })
+ .await
+ .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?;
+ cancellation_on_drop.disarm();
+ result
+ }
+
+ #[cfg(not(any(target_os = "linux", target_os = "macos")))]
+ {
+ let _ = staged;
+ Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Restore))
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn finalize_native(
+ staged: NativeStagedServiceRestore,
+ cancellation: &AtomicU8,
+ operations: &dyn FinalizeOperations,
+) -> Result<(), ServiceSqliteError> {
+ staged.validate()?;
+ check_cancel(cancellation)?;
+ let live = authority_checked(&staged, || open_live(staged.directory()))??;
+ let live_artifact = staged.live_artifact();
+ authority_checked(&staged, || {
+ verify_named_artifact(
+ staged.directory(),
+ LIVE_FILE_NAME,
+ &live,
+ live_artifact,
+ Some(cancellation),
+ )
+ })??;
+ authority_checked(&staged, || {
+ verify_named_artifact(
+ staged.directory(),
+ STAGED_FILE_NAME,
+ staged.staged_file(),
+ staged.artifact(),
+ Some(cancellation),
+ )
+ })??;
+ authority_checked(&staged, || {
+ live.sync_all()
+ .map_err(|source| finalize_source(FinalizeFailureKind::SyncLive, source))
+ })??;
+ authority_checked(&staged, || {
+ staged
+ .staged_file()
+ .sync_all()
+ .map_err(|source| finalize_source(FinalizeFailureKind::SyncStaged, source))
+ })??;
+ test_phase(PHASE_BEFORE_PREPARED, cancellation, true)?;
+ staged.validate()?;
+ claim_commit_ownership(cancellation)?;
+ test_phase(PHASE_COMMIT_OWNED, cancellation, false)?;
+
+ let marker = RestoreRecoveryMarker::prepared(
+ staged.metadata(),
+ staged.manifest_digest(),
+ live_artifact,
+ staged.artifact(),
+ )
+ .map_err(|source| ServiceSqliteError::with_source(ServiceSqliteErrorKind::Restore, source))?;
+ let on_durable = || {
+ staged.disarm_cleanup();
+ let _ = test_phase(PHASE_AFTER_PREPARED, cancellation, false);
+ };
+ #[cfg(test)]
+ let marker_result = if operations.drift_authority_during_marker_sync() {
+ RestoreMarkerBinding::test_create_with_durable_authority_drift(
+ staged.paths(),
+ staged.authority(),
+ &marker,
+ on_durable,
+ )
+ } else {
+ RestoreMarkerBinding::create_with_durable_callback(
+ staged.paths(),
+ staged.authority(),
+ &marker,
+ on_durable,
+ )
+ };
+ #[cfg(not(test))]
+ let marker_result = RestoreMarkerBinding::create_with_durable_callback(
+ staged.paths(),
+ staged.authority(),
+ &marker,
+ on_durable,
+ );
+ let mut marker = marker_result?;
+
+ rename_and_sync(
+ &staged,
+ operations,
+ RenameStep::RetainLive,
+ LIVE_FILE_NAME,
+ BACKUP_FILE_NAME,
+ &live,
+ live_artifact,
+ )?;
+ authority_checked(&staged, || {
+ operations.after_directory_sync(staged.directory(), RenameStep::RetainLive)
+ })??;
+ marker = marker.advance(
+ staged.paths(),
+ staged.authority(),
+ RestoreRecoveryPhase::LiveRetained,
+ )?;
+
+ rename_and_sync(
+ &staged,
+ operations,
+ RenameStep::InstallStage,
+ STAGED_FILE_NAME,
+ LIVE_FILE_NAME,
+ staged.staged_file(),
+ staged.artifact(),
+ )?;
+ authority_checked(&staged, || {
+ operations.after_directory_sync(staged.directory(), RenameStep::InstallStage)
+ })??;
+ marker = marker.advance(
+ staged.paths(),
+ staged.authority(),
+ RestoreRecoveryPhase::ReplacementInstalled,
+ )?;
+ if marker.marker().phase() != RestoreRecoveryPhase::ReplacementInstalled {
+ return Err(finalize_error(FinalizeFailureKind::Marker));
+ }
+ staged.validate_finalization_authority()?;
+ drop(marker);
+ drop(staged);
+ Ok(())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn rename_and_sync(
+ staged: &NativeStagedServiceRestore,
+ operations: &dyn FinalizeOperations,
+ step: RenameStep,
+ source_name: &str,
+ destination_name: &str,
+ held: &File,
+ expected: RestoreArtifactExpectation,
+) -> Result<(), ServiceSqliteError> {
+ authority_checked(staged, || {
+ verify_named_artifact(staged.directory(), source_name, held, expected, None)
+ })??;
+ let rename = authority_checked(staged, || {
+ operations
+ .rename(staged.directory(), source_name, destination_name, step)
+ .map_err(|source| finalize_source(step.rename_failure(), source))
+ })?;
+ rename?;
+ authority_checked(staged, || {
+ verify_named_artifact(staged.directory(), destination_name, held, expected, None)
+ })??;
+ let sync = authority_checked(staged, || {
+ operations
+ .sync_directory(staged.directory(), step)
+ .map_err(|source| finalize_source(step.sync_failure(), source))
+ })?;
+ sync?;
+ authority_checked(staged, || {
+ verify_named_artifact(staged.directory(), destination_name, held, expected, None)
+ })??;
+ Ok(())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn authority_checked<T>(
+ staged: &NativeStagedServiceRestore,
+ operation: impl FnOnce() -> T,
+) -> Result<T, ServiceSqliteError> {
+ staged.validate_finalization_authority()?;
+ let result = operation();
+ staged.validate_finalization_authority()?;
+ Ok(result)
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn open_live(directory: &File) -> Result<File, ServiceSqliteError> {
+ openat(
+ directory,
+ LIVE_FILE_NAME,
+ OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
+ Mode::empty(),
+ )
+ .map(File::from)
+ .map_err(|source| finalize_source(FinalizeFailureKind::Live, source))
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn verify_named_artifact(
+ directory: &File,
+ name: &str,
+ held: &File,
+ expected: RestoreArtifactExpectation,
+ cancellation: Option<&AtomicU8>,
+) -> Result<(), ServiceSqliteError> {
+ let current = openat(
+ directory,
+ name,
+ OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
+ Mode::empty(),
+ )
+ .map(File::from)
+ .map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
+ let held_status =
+ fstat(held).map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
+ let current_status =
+ fstat(¤t).map_err(|source| finalize_source(FinalizeFailureKind::Artifact, source))?;
+ validate_status(&held_status, Some(expected.byte_length()))?;
+ validate_status(¤t_status, Some(expected.byte_length()))?;
+ let held_identity = (
+ u64::try_from(held_status.st_dev)
+ .map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?,
+ held_status.st_ino,
+ );
+ let current_identity = (
+ u64::try_from(current_status.st_dev)
+ .map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?,
+ current_status.st_ino,
+ );
+ if held_identity != (expected.device(), expected.inode())
+ || current_identity != held_identity
+ || hash_exact(held, expected.byte_length(), cancellation)? != expected.sha256()
+ {
+ return Err(finalize_error(FinalizeFailureKind::Artifact));
+ }
+ Ok(())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn validate_status(
+ status: &rustix::fs::Stat,
+ expected_length: Option<u64>,
+) -> Result<(), ServiceSqliteError> {
+ let length =
+ u64::try_from(status.st_size).map_err(|_| finalize_error(FinalizeFailureKind::Artifact))?;
+ if !FileType::from_raw_mode(status.st_mode).is_file()
+ || u64::from(status.st_nlink) != 1
+ || status.st_uid != geteuid().as_raw()
+ || u32::from(status.st_mode) & 0o777 != 0o600
+ || length == 0
+ || length > i64::MAX as u64
+ || expected_length.is_some_and(|expected| length != expected)
+ {
+ return Err(finalize_error(FinalizeFailureKind::Artifact));
+ }
+ Ok(())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn hash_exact(
+ file: &File,
+ expected_length: u64,
+ cancellation: Option<&AtomicU8>,
+) -> Result<[u8; 32], ServiceSqliteError> {
+ let mut offset = 0_u64;
+ let mut buffer = [0_u8; HASH_BUFFER_BYTES];
+ let mut hasher = Sha256::new();
+ while offset < expected_length {
+ if cancellation.is_some_and(|state| state.load(Ordering::Acquire) == CANCELLED) {
+ return Err(finalize_error(FinalizeFailureKind::Cancelled));
+ }
+ let requested = usize::try_from((expected_length - offset).min(HASH_BUFFER_BYTES as u64))
+ .map_err(|_| finalize_error(FinalizeFailureKind::Hash))?;
+ let read = file
+ .read_at(&mut buffer[..requested], offset)
+ .map_err(|source| finalize_source(FinalizeFailureKind::Hash, source))?;
+ if read == 0 {
+ return Err(finalize_error(FinalizeFailureKind::Hash));
+ }
+ hasher.update(&buffer[..read]);
+ offset = offset
+ .checked_add(
+ u64::try_from(read).map_err(|_| finalize_error(FinalizeFailureKind::Hash))?,
+ )
+ .ok_or_else(|| finalize_error(FinalizeFailureKind::Hash))?;
+ }
+ if cancellation.is_some_and(|state| state.load(Ordering::Acquire) == CANCELLED) {
+ return Err(finalize_error(FinalizeFailureKind::Cancelled));
+ }
+ let mut extra = [0_u8; 1];
+ if file
+ .read_at(&mut extra, expected_length)
+ .map_err(|source| finalize_source(FinalizeFailureKind::Hash, source))?
+ != 0
+ {
+ return Err(finalize_error(FinalizeFailureKind::Hash));
+ }
+ Ok(hasher.finalize().into())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn check_cancel(cancellation: &AtomicU8) -> Result<(), ServiceSqliteError> {
+ if cancellation.load(Ordering::Acquire) == CANCELLED {
+ Err(finalize_error(FinalizeFailureKind::Cancelled))
+ } else {
+ Ok(())
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn claim_commit_ownership(cancellation: &AtomicU8) -> Result<(), ServiceSqliteError> {
+ cancellation
+ .compare_exchange(
+ CANCELLABLE,
+ COMMIT_OWNED,
+ Ordering::AcqRel,
+ Ordering::Acquire,
+ )
+ .map(|_| ())
+ .map_err(|_| finalize_error(FinalizeFailureKind::Cancelled))
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn test_phase(
+ phase: u8,
+ cancellation: &AtomicU8,
+ cancellable: bool,
+) -> Result<(), ServiceSqliteError> {
+ #[cfg(test)]
+ {
+ TEST_FINALIZE_PHASE.store(phase, Ordering::Release);
+ while TEST_FINALIZE_BLOCK_PHASE.load(Ordering::Acquire) == phase {
+ if cancellable && cancellation.load(Ordering::Acquire) == CANCELLED {
+ return Err(finalize_error(FinalizeFailureKind::Cancelled));
+ }
+ std::thread::yield_now();
+ }
+ }
+ #[cfg(not(test))]
+ let _ = (phase, cancellation, cancellable);
+ Ok(())
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum RenameStep {
+ RetainLive,
+ InstallStage,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl RenameStep {
+ const fn rename_failure(self) -> FinalizeFailureKind {
+ match self {
+ Self::RetainLive => FinalizeFailureKind::RetainLive,
+ Self::InstallStage => FinalizeFailureKind::InstallStage,
+ }
+ }
+
+ const fn sync_failure(self) -> FinalizeFailureKind {
+ match self {
+ Self::RetainLive => FinalizeFailureKind::SyncRetained,
+ Self::InstallStage => FinalizeFailureKind::SyncInstalled,
+ }
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+trait FinalizeOperations: Send + Sync {
+ fn rename(
+ &self,
+ directory: &File,
+ source: &str,
+ destination: &str,
+ step: RenameStep,
+ ) -> std::io::Result<()>;
+
+ fn sync_directory(&self, directory: &File, step: RenameStep) -> std::io::Result<()>;
+
+ fn after_directory_sync(
+ &self,
+ _directory: &File,
+ _step: RenameStep,
+ ) -> Result<(), ServiceSqliteError> {
+ Ok(())
+ }
+
+ #[cfg(test)]
+ fn drift_authority_during_marker_sync(&self) -> bool {
+ false
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+struct SystemFinalizeOperations;
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl FinalizeOperations for SystemFinalizeOperations {
+ fn rename(
+ &self,
+ directory: &File,
+ source: &str,
+ destination: &str,
+ _step: RenameStep,
+ ) -> std::io::Result<()> {
+ renameat_with(
+ directory,
+ source,
+ directory,
+ destination,
+ RenameFlags::NOREPLACE,
+ )
+ .map_err(std::io::Error::from)
+ }
+
+ fn sync_directory(&self, directory: &File, _step: RenameStep) -> std::io::Result<()> {
+ directory.sync_all()
+ }
+}
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+struct FailingFinalizeOperations;
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+impl FinalizeOperations for FailingFinalizeOperations {
+ fn drift_authority_during_marker_sync(&self) -> bool {
+ TEST_FINALIZE_FAILURE.load(Ordering::Acquire) == 9
+ }
+
+ fn rename(
+ &self,
+ directory: &File,
+ source: &str,
+ destination: &str,
+ step: RenameStep,
+ ) -> std::io::Result<()> {
+ let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
+ let (before, after) = match step {
+ RenameStep::RetainLive => (1, 2),
+ RenameStep::InstallStage => (5, 6),
+ };
+ if failure == before {
+ return Err(std::io::Error::other("injected pre-rename failure"));
+ }
+ SystemFinalizeOperations.rename(directory, source, destination, step)?;
+ if failure == after {
+ return Err(std::io::Error::other("injected post-rename failure"));
+ }
+ Ok(())
+ }
+
+ fn sync_directory(&self, directory: &File, step: RenameStep) -> std::io::Result<()> {
+ let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
+ let target = match step {
+ RenameStep::RetainLive => 3,
+ RenameStep::InstallStage => 7,
+ };
+ if failure == target {
+ return Err(std::io::Error::other("injected directory sync failure"));
+ }
+ SystemFinalizeOperations.sync_directory(directory, step)
+ }
+
+ fn after_directory_sync(
+ &self,
+ directory: &File,
+ step: RenameStep,
+ ) -> Result<(), ServiceSqliteError> {
+ let failure = TEST_FINALIZE_FAILURE.load(Ordering::Acquire);
+ let target = match step {
+ RenameStep::RetainLive => 4,
+ RenameStep::InstallStage => 8,
+ };
+ if failure != target {
+ return Ok(());
+ }
+ openat(
+ directory,
+ MARKER_NEXT_FILE_NAME,
+ OFlags::RDWR
+ | OFlags::CREATE
+ | OFlags::EXCL
+ | OFlags::NOFOLLOW
+ | OFlags::CLOEXEC
+ | OFlags::NONBLOCK,
+ Mode::RUSR | Mode::WUSR,
+ )
+ .map(drop)
+ .map_err(|source| finalize_source(FinalizeFailureKind::Marker, source))
+ }
+}
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) async fn test_finalize_with_failure(
+ staged: StagedServiceRestore,
+ failure: u8,
+) -> Result<(), ServiceSqliteError> {
+ TEST_FINALIZE_FAILURE.store(failure, Ordering::Release);
+ let native = staged.into_native();
+ let result = tokio::task::spawn_blocking(move || {
+ finalize_native(
+ native,
+ &AtomicU8::new(CANCELLABLE),
+ &FailingFinalizeOperations,
+ )
+ })
+ .await
+ .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?;
+ TEST_FINALIZE_FAILURE.store(0, Ordering::Release);
+ result
+}
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) fn reset_test_controls() {
+ TEST_FINALIZE_PHASE.store(0, Ordering::Release);
+ TEST_FINALIZE_BLOCK_PHASE.store(0, Ordering::Release);
+ TEST_FINALIZE_FAILURE.store(0, Ordering::Release);
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+struct CancellationOnDrop {
+ cancellation: Arc<AtomicU8>,
+ armed: AtomicBool,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl CancellationOnDrop {
+ fn new(cancellation: Arc<AtomicU8>) -> Self {
+ Self {
+ cancellation,
+ armed: AtomicBool::new(true),
+ }
+ }
+
+ fn disarm(&self) {
+ self.armed.store(false, Ordering::Release);
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl Drop for CancellationOnDrop {
+ fn drop(&mut self) {
+ if self.armed.load(Ordering::Acquire) {
+ let _ = self.cancellation.compare_exchange(
+ CANCELLABLE,
+ CANCELLED,
+ Ordering::AcqRel,
+ Ordering::Acquire,
+ );
+ }
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum FinalizeFailureKind {
+ Live,
+ Artifact,
+ Hash,
+ SyncLive,
+ SyncStaged,
+ Marker,
+ RetainLive,
+ SyncRetained,
+ InstallStage,
+ SyncInstalled,
+ Cancelled,
+ Join,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+struct FinalizeFailure {
+ kind: FinalizeFailureKind,
+ source: Option<Box<dyn Error + Send + Sync + 'static>>,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl fmt::Debug for FinalizeFailure {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("FinalizeFailure")
+ .field("kind", &self.kind)
+ .field("source", &self.source.as_ref().map(|_| "[redacted]"))
+ .finish()
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl fmt::Display for FinalizeFailure {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ FinalizeFailureKind::Live => "live restore source is invalid",
+ FinalizeFailureKind::Artifact => "restore artifact binding changed",
+ FinalizeFailureKind::Hash => "restore artifact hash failed",
+ FinalizeFailureKind::SyncLive => "live restore source sync failed",
+ FinalizeFailureKind::SyncStaged => "staged restore sync failed",
+ FinalizeFailureKind::Marker => "restore marker transition failed",
+ FinalizeFailureKind::RetainLive => "live restore retention failed",
+ FinalizeFailureKind::SyncRetained => "retained restore sync failed",
+ FinalizeFailureKind::InstallStage => "restore installation failed",
+ FinalizeFailureKind::SyncInstalled => "installed restore sync failed",
+ FinalizeFailureKind::Cancelled => "restore finalization was cancelled",
+ FinalizeFailureKind::Join => "restore finalization worker failed",
+ })
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl Error for FinalizeFailure {
+ fn source(&self) -> Option<&(dyn Error + 'static)> {
+ self.source
+ .as_deref()
+ .map(|source| source as &(dyn Error + 'static))
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn finalize_error(kind: FinalizeFailureKind) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Restore,
+ FinalizeFailure { kind, source: None },
+ )
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn finalize_source(
+ kind: FinalizeFailureKind,
+ source: impl Error + Send + Sync + 'static,
+) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Restore,
+ FinalizeFailure {
+ kind,
+ source: Some(Box::new(source)),
+ },
+ )
+}
diff --git a/crates/service_sqlite/src/restore/marker.rs b/crates/service_sqlite/src/restore/marker.rs
@@ -263,6 +263,18 @@ impl RestoreRecoveryMarker {
self.phase
}
+ pub(crate) const fn live(&self) -> RestoreArtifactExpectation {
+ self.live
+ }
+
+ pub(crate) const fn staged(&self) -> RestoreArtifactExpectation {
+ self.staged
+ }
+
+ pub(crate) const fn backup(&self) -> RestoreArtifactExpectation {
+ self.backup
+ }
+
pub(crate) fn transitioned_to(
&self,
next: RestoreRecoveryPhase,
@@ -569,11 +581,53 @@ mod store {
}
impl RestoreMarkerBinding {
+ #[cfg(test)]
pub(crate) fn create(
paths: &ServiceSqlitePaths,
authority: &WriterAuthority,
marker: &RestoreRecoveryMarker,
) -> Result<Self, ServiceSqliteError> {
+ Self::create_with_durable_callback(paths, authority, marker, || {})
+ }
+
+ pub(crate) fn create_with_durable_callback(
+ paths: &ServiceSqlitePaths,
+ authority: &WriterAuthority,
+ marker: &RestoreRecoveryMarker,
+ on_durable: impl FnOnce(),
+ ) -> Result<Self, ServiceSqliteError> {
+ Self::create_with_operations(
+ paths,
+ authority,
+ marker,
+ &SystemStoreOperations,
+ on_durable,
+ )
+ }
+
+ #[cfg(test)]
+ pub(crate) fn test_create_with_durable_authority_drift(
+ paths: &ServiceSqlitePaths,
+ authority: &WriterAuthority,
+ marker: &RestoreRecoveryMarker,
+ on_durable: impl FnOnce(),
+ ) -> Result<Self, ServiceSqliteError> {
+ Self::create_with_operations(
+ paths,
+ authority,
+ marker,
+ &AuthorityDriftAfterSyncStoreOperations,
+ on_durable,
+ )
+ }
+
+ fn create_with_operations(
+ paths: &ServiceSqlitePaths,
+ authority: &WriterAuthority,
+ marker: &RestoreRecoveryMarker,
+ operations: &dyn StoreOperations,
+ on_durable: impl FnOnce(),
+ ) -> Result<Self, ServiceSqliteError> {
authority.validate_for(paths)?;
if !marker.matches_paths(paths) {
return Err(restore_contract(
@@ -606,13 +660,12 @@ mod store {
Ok::<_, StoreFailure>((file, identity))
})?
.map_err(restore_store)?;
- let operations = SystemStoreOperations;
let write_result = authority_checked(authority, paths, || {
write_and_sync(
&marker_file,
marker.canonical_bytes(),
&directory,
- &operations,
+ operations,
)
})?;
if let Err(cause) = write_result {
@@ -625,12 +678,9 @@ mod store {
)?;
return Err(restore_store(cause));
}
- let parent_sync = authority_checked(authority, paths, || {
- operations
- .sync_directory(&directory)
- .map_err(|_| StoreFailure::Sync)
- })?;
- if parent_sync.is_err() {
+ authority.validate_for(paths)?;
+ if operations.sync_directory(&directory).is_err() {
+ authority.validate_for(paths)?;
cleanup_with_authority(
authority,
paths,
@@ -640,6 +690,11 @@ mod store {
)?;
return Err(restore_store(StoreFailure::Sync));
}
+ // The marker contents and its directory entry are durable from
+ // this point. The caller must transfer ownership of every bound
+ // artifact before any subsequent fallible validation.
+ on_durable();
+ authority.validate_for(paths)?;
let binding = Self {
directory,
directory_identity,
@@ -1017,6 +1072,29 @@ mod store {
}
#[cfg(test)]
+ struct AuthorityDriftAfterSyncStoreOperations;
+
+ #[cfg(test)]
+ impl StoreOperations for AuthorityDriftAfterSyncStoreOperations {
+ fn sync_file(&self, file: &File, _directory: &File) -> std::io::Result<()> {
+ file.sync_all()
+ }
+
+ fn sync_directory(&self, directory: &File) -> std::io::Result<()> {
+ directory.sync_all()?;
+ fchmod(
+ directory,
+ Mode::RUSR | Mode::WUSR | Mode::XUSR | Mode::RGRP | Mode::WGRP | Mode::XGRP,
+ )
+ .map_err(std::io::Error::from)
+ }
+
+ fn replace_marker(&self, directory: &File) -> std::io::Result<()> {
+ SystemStoreOperations.replace_marker(directory)
+ }
+ }
+
+ #[cfg(test)]
#[derive(Clone, Copy, Debug, PartialEq, Eq)]
pub(crate) enum TestStoreFailure {
ScratchSync,
diff --git a/crates/service_sqlite/src/restore/mod.rs b/crates/service_sqlite/src/restore/mod.rs
@@ -1,10 +1,34 @@
//! Private crash-recovery marker mechanics for governed restore.
+mod finalize;
mod marker;
mod stage;
+pub use finalize::finalize_staged_restore;
pub use stage::{StagedServiceRestore, stage_verified_restore};
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+pub(crate) fn refuse_unresolved_recovery(
+ directory: &impl std::os::fd::AsFd,
+) -> Result<(), crate::ServiceSqliteError> {
+ use rustix::{
+ fs::{AtFlags, statat},
+ io::Errno,
+ };
+
+ for name in [marker::MARKER_FILE_NAME, marker::MARKER_NEXT_FILE_NAME] {
+ match statat(directory, name, AtFlags::SYMLINK_NOFOLLOW) {
+ Err(Errno::NOENT) => {}
+ Ok(_) | Err(_) => {
+ return Err(crate::ServiceSqliteError::new(
+ crate::ServiceSqliteErrorKind::Recovery,
+ ));
+ }
+ }
+ }
+ Ok(())
+}
+
#[allow(unused_imports)]
pub(crate) use marker::{
RestoreArtifactExpectation, RestoreMarkerContractError, RestoreRecoveryLayout,
diff --git a/crates/service_sqlite/src/restore/stage.rs b/crates/service_sqlite/src/restore/stage.rs
@@ -380,10 +380,11 @@ pub(crate) struct NativeStagedServiceRestore {
directory_identity: FileIdentity,
staged: File,
staged_identity: FileIdentity,
+ live_artifact: RestoreArtifactExpectation,
metadata: ServiceDatabaseMetadata,
manifest_digest: crate::BackupManifestSha256,
artifact: RestoreArtifactExpectation,
- armed: bool,
+ armed: AtomicBool,
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -412,7 +413,7 @@ impl NativeStagedServiceRestore {
let directory_identity = directory_identity_result?;
let closed_live = validate_closed_live(&authority, &paths, &directory);
authority.validate_for(&paths)?;
- closed_live?;
+ let live_artifact = closed_live?;
let before_create = test_blocking_phase(TEST_PHASE_BEFORE_CREATE, cancellation);
authority.validate_for(&paths)?;
before_create?;
@@ -447,10 +448,11 @@ impl NativeStagedServiceRestore {
directory_identity,
staged,
staged_identity,
+ live_artifact,
metadata,
manifest_digest,
artifact,
- armed: true,
+ armed: AtomicBool::new(true),
};
result.validate_created()?;
let copy = copy_exact(
@@ -541,6 +543,7 @@ impl NativeStagedServiceRestore {
self.staged_identity,
self.artifact.byte_length(),
)?;
+ validate_live_binding(&self.directory, self.live_artifact)?;
require_absent(&self.directory, BACKUP_FILE_NAME)?;
require_absent(&self.directory, MARKER_FILE_NAME)?;
require_absent(&self.directory, MARKER_NEXT_FILE_NAME)?;
@@ -553,6 +556,21 @@ impl NativeStagedServiceRestore {
result
}
+ pub(super) fn validate_finalization_authority(&self) -> Result<(), ServiceSqliteError> {
+ let authority = self
+ .authority
+ .as_ref()
+ .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Authority))?;
+ authority.validate_for(&self.paths)?;
+ let valid_directory = directory_identity(&self.directory)
+ .is_ok_and(|identity| identity == self.directory_identity);
+ authority.validate_for(&self.paths)?;
+ if !valid_directory {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Authority));
+ }
+ Ok(())
+ }
+
fn validate_created(&self) -> Result<(), ServiceSqliteError> {
let authority = self
.authority
@@ -589,6 +607,10 @@ impl NativeStagedServiceRestore {
self.artifact
}
+ pub(super) const fn live_artifact(&self) -> RestoreArtifactExpectation {
+ self.live_artifact
+ }
+
#[allow(dead_code)]
pub(crate) fn authority(&self) -> &WriterAuthority {
self.authority
@@ -596,9 +618,17 @@ impl NativeStagedServiceRestore {
.expect("staged restore retains authority until consumed")
}
+ pub(super) fn directory(&self) -> &File {
+ &self.directory
+ }
+
+ pub(super) fn staged_file(&self) -> &File {
+ &self.staged
+ }
+
#[allow(dead_code)]
- pub(crate) fn disarm_cleanup(&mut self) {
- self.armed = false;
+ pub(crate) fn disarm_cleanup(&self) {
+ self.armed.store(false, Ordering::Release);
}
}
@@ -612,7 +642,7 @@ impl fmt::Debug for NativeStagedServiceRestore {
#[cfg(any(target_os = "linux", target_os = "macos"))]
impl Drop for NativeStagedServiceRestore {
fn drop(&mut self) {
- if self.armed {
+ if self.armed.load(Ordering::Acquire) {
let _ = cleanup_exact_stage(&self.directory, &self.staged, self.staged_identity);
}
if let Some(authority) = self.authority.as_mut() {
@@ -759,7 +789,7 @@ fn validate_closed_live(
authority: &WriterAuthority,
paths: &ServiceSqlitePaths,
directory: &File,
-) -> Result<(), ServiceSqliteError> {
+) -> Result<RestoreArtifactExpectation, ServiceSqliteError> {
authority.validate_for(paths)?;
let live = openat(
directory,
@@ -777,6 +807,20 @@ fn validate_closed_live(
{
return Err(restore_error(RestoreFailureKind::LiveState));
}
+ let length =
+ u64::try_from(status.st_size).map_err(|_| restore_error(RestoreFailureKind::LiveState))?;
+ if length == 0 || length > i64::MAX as u64 {
+ return Err(restore_error(RestoreFailureKind::LiveState));
+ }
+ let live = File::from(live);
+ let digest = hash_exact(&live, length)?;
+ let artifact = RestoreArtifactExpectation::new(
+ u64::try_from(status.st_dev).map_err(|_| restore_error(RestoreFailureKind::LiveState))?,
+ status.st_ino,
+ length,
+ digest,
+ )
+ .map_err(|_| restore_error(RestoreFailureKind::LiveState))?;
for name in [
"state.sqlite-wal",
"state.sqlite-shm",
@@ -787,6 +831,36 @@ fn validate_closed_live(
] {
require_absent(directory, name)?;
}
+ Ok(artifact)
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn validate_live_binding(
+ directory: &File,
+ expected: RestoreArtifactExpectation,
+) -> Result<(), ServiceSqliteError> {
+ let current = openat(
+ directory,
+ radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME,
+ OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC | OFlags::NONBLOCK,
+ Mode::empty(),
+ )
+ .map_err(|source| restore_source(RestoreFailureKind::LiveState, source))?;
+ let status =
+ fstat(¤t).map_err(|source| restore_source(RestoreFailureKind::LiveState, source))?;
+ let device =
+ u64::try_from(status.st_dev).map_err(|_| restore_error(RestoreFailureKind::LiveState))?;
+ let length =
+ u64::try_from(status.st_size).map_err(|_| restore_error(RestoreFailureKind::LiveState))?;
+ if !FileType::from_raw_mode(status.st_mode).is_file()
+ || u64::from(status.st_nlink) != 1
+ || status.st_uid != geteuid().as_raw()
+ || u32::from(status.st_mode) & 0o777 != 0o600
+ || (device, status.st_ino) != (expected.device(), expected.inode())
+ || length != expected.byte_length()
+ {
+ return Err(restore_error(RestoreFailureKind::LiveState));
+ }
Ok(())
}
@@ -1167,11 +1241,17 @@ mod tests {
use sqlx::{ConnectOptions, Connection, sqlite::SqliteConnectOptions};
use super::*;
+ use crate::restore::finalize::{
+ TEST_FINALIZE_BLOCK_PHASE, TEST_FINALIZE_PHASE, TEST_PHASE_AFTER_PREPARED,
+ TEST_PHASE_BEFORE_PREPARED, TEST_PHASE_COMMIT_OWNED,
+ reset_test_controls as reset_finalize_controls, test_finalize_with_failure,
+ };
+ use crate::restore::{RestoreMarkerBinding, RestoreRecoveryMarker, RestoreRecoveryPhase};
use crate::{
BackupCreatedAtUnixMs, BackupMemberSha256, MigrationChecksum, MigrationDescriptor,
SchemaObject, SchemaObjectKind, SchemaVersionCatalog, ServiceBackupManifest,
ServiceDatabaseMetadata, ServiceSqliteApplicationId, ServiceSqliteConnectionOptions,
- ServiceSqliteHost, verify_backup_bundle,
+ ServiceSqliteHost, finalize_staged_restore, verify_backup_bundle,
};
static STAGE_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
@@ -1344,6 +1424,46 @@ mod tests {
panic!("restore staging cleanup did not release authority");
}
+ async fn wait_for_finalize_phase(phase: u8) {
+ for _ in 0..100_000 {
+ if TEST_FINALIZE_PHASE.load(Ordering::Acquire) == phase {
+ return;
+ }
+ tokio::task::yield_now().await;
+ }
+ panic!("restore finalization did not reach phase {phase}");
+ }
+
+ fn recovery_path(fixture: &Fixture, name: &str) -> PathBuf {
+ fixture.staged_path().with_file_name(name)
+ }
+
+ fn marker_phase(fixture: &Fixture) -> RestoreRecoveryPhase {
+ let bytes = fs::read(recovery_path(fixture, MARKER_FILE_NAME)).expect("marker bytes");
+ RestoreRecoveryMarker::from_canonical_bytes(&bytes)
+ .expect("canonical marker")
+ .phase()
+ }
+
+ async fn wait_for_replacement_install(fixture: &Fixture) {
+ for _ in 0..100_000 {
+ if marker_phase_if_present(fixture) == Some(RestoreRecoveryPhase::ReplacementInstalled)
+ && WriterAuthority::acquire(&fixture.paths, OpenMode::ReadWriteExisting).is_ok()
+ {
+ return;
+ }
+ tokio::task::yield_now().await;
+ }
+ panic!("restore finalization worker did not complete");
+ }
+
+ fn marker_phase_if_present(fixture: &Fixture) -> Option<RestoreRecoveryPhase> {
+ let bytes = fs::read(recovery_path(fixture, MARKER_FILE_NAME)).ok()?;
+ RestoreRecoveryMarker::from_canonical_bytes(&bytes)
+ .ok()
+ .map(|marker| marker.phase())
+ }
+
fn spawn_stage(
fixture: &Fixture,
) -> tokio::task::JoinHandle<Result<StagedServiceRestore, ServiceSqliteError>> {
@@ -1882,6 +2002,351 @@ mod tests {
assert!(!ledger_drift.staged_path().exists());
}
+ #[tokio::test(flavor = "current_thread")]
+ async fn atomic_finalization_retains_old_live_installs_stage_and_requires_recovery_open() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+ reset_finalize_controls();
+ let fixture = Fixture::new().await;
+ let old_live = fs::read(fixture.paths.state_database()).expect("old live");
+ let replacement =
+ fs::read(fixture.bundle.join(crate::BACKUP_STATE_MEMBER_NAME)).expect("replacement");
+ let old_metadata = fs::metadata(fixture.paths.state_database()).expect("old metadata");
+ let staged = stage_verified_restore(
+ &fixture.paths,
+ &fixture.identity,
+ &fixture.migrations,
+ &fixture.schema,
+ fixture.proof(),
+ )
+ .await
+ .expect("stage");
+
+ finalize_staged_restore(staged).await.expect("finalize");
+
+ let backup = recovery_path(&fixture, BACKUP_FILE_NAME);
+ assert_eq!(fs::read(&backup).expect("retained old live"), old_live);
+ assert_eq!(
+ fs::read(fixture.paths.state_database()).expect("installed live"),
+ replacement
+ );
+ let backup_metadata = fs::metadata(&backup).expect("backup metadata");
+ assert_eq!(
+ (old_metadata.dev(), old_metadata.ino()),
+ (backup_metadata.dev(), backup_metadata.ino())
+ );
+ assert!(!fixture.staged_path().exists());
+ assert!(!recovery_path(&fixture, MARKER_NEXT_FILE_NAME).exists());
+
+ let binding = RestoreMarkerBinding::load(&fixture.paths)
+ .expect("load marker")
+ .expect("marker present");
+ assert_eq!(
+ binding.marker().phase(),
+ RestoreRecoveryPhase::ReplacementInstalled
+ );
+ assert_eq!(binding.marker().live(), binding.marker().backup());
+ assert_ne!(
+ (
+ binding.marker().live().device(),
+ binding.marker().live().inode()
+ ),
+ (
+ binding.marker().staged().device(),
+ binding.marker().staged().inode()
+ )
+ );
+
+ for mode in [OpenMode::ReadWriteExisting, OpenMode::ReadOnlyInspection] {
+ let error = crate::open::open_existing_connection_pool(
+ &fixture.paths,
+ &fixture.identity,
+ &fixture.migrations,
+ &fixture.schema,
+ mode,
+ ServiceSqliteConnectionOptions::reviewed(),
+ )
+ .await
+ .err()
+ .expect("unresolved recovery must reject open");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Recovery);
+ }
+ let initialize = crate::initialize_database(
+ &fixture.paths,
+ OpenMode::Initialize,
+ &fixture.metadata,
+ &fixture.schema,
+ |_| async { Ok::<(), std::io::Error>(()) },
+ )
+ .await
+ .expect_err("unresolved recovery must reject initialization");
+ assert_eq!(initialize.kind(), ServiceSqliteErrorKind::Recovery);
+ reset_finalize_controls();
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn finalization_cancellation_is_clean_before_handoff_and_owned_after_handoff() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+
+ reset_finalize_controls();
+ let before = Fixture::new().await;
+ let old_live = fs::read(before.paths.state_database()).expect("old live");
+ let staged = stage_verified_restore(
+ &before.paths,
+ &before.identity,
+ &before.migrations,
+ &before.schema,
+ before.proof(),
+ )
+ .await
+ .expect("stage");
+ TEST_FINALIZE_BLOCK_PHASE.store(TEST_PHASE_BEFORE_PREPARED, Ordering::Release);
+ let caller = tokio::spawn(finalize_staged_restore(staged));
+ wait_for_finalize_phase(TEST_PHASE_BEFORE_PREPARED).await;
+ caller.abort();
+ let _ = caller.await;
+ wait_for_cleanup(&before.paths, &before.staged_path()).await;
+ assert_eq!(
+ fs::read(before.paths.state_database()).expect("preserved live"),
+ old_live
+ );
+ assert!(!recovery_path(&before, MARKER_FILE_NAME).exists());
+
+ for phase in [TEST_PHASE_COMMIT_OWNED, TEST_PHASE_AFTER_PREPARED] {
+ reset_finalize_controls();
+ let after = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &after.paths,
+ &after.identity,
+ &after.migrations,
+ &after.schema,
+ after.proof(),
+ )
+ .await
+ .expect("stage");
+ TEST_FINALIZE_BLOCK_PHASE.store(phase, Ordering::Release);
+ let caller = tokio::spawn(finalize_staged_restore(staged));
+ wait_for_finalize_phase(phase).await;
+ if phase == TEST_PHASE_AFTER_PREPARED {
+ assert_eq!(marker_phase(&after), RestoreRecoveryPhase::Prepared);
+ } else {
+ assert!(!recovery_path(&after, MARKER_FILE_NAME).exists());
+ }
+ caller.abort();
+ let _ = caller.await;
+ assert!(WriterAuthority::acquire(&after.paths, OpenMode::ReadWriteExisting).is_err());
+ TEST_FINALIZE_BLOCK_PHASE.store(0, Ordering::Release);
+ wait_for_replacement_install(&after).await;
+ assert_eq!(
+ marker_phase(&after),
+ RestoreRecoveryPhase::ReplacementInstalled
+ );
+ }
+ reset_finalize_controls();
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn every_rename_sync_and_marker_advance_failure_leaves_exact_recovery_evidence() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+ for failure in 1..=8 {
+ reset_finalize_controls();
+ let fixture = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &fixture.paths,
+ &fixture.identity,
+ &fixture.migrations,
+ &fixture.schema,
+ fixture.proof(),
+ )
+ .await
+ .expect("stage");
+ let error = test_finalize_with_failure(staged, failure)
+ .await
+ .expect_err("injected finalization failure");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Restore);
+ assert!(recovery_path(&fixture, MARKER_FILE_NAME).exists());
+ let live = fixture.paths.state_database().exists();
+ let staged = fixture.staged_path().exists();
+ let backup = recovery_path(&fixture, BACKUP_FILE_NAME).exists();
+ match failure {
+ 1 => {
+ assert_eq!((live, staged, backup), (true, true, false));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::Prepared);
+ }
+ 2..=4 => {
+ assert_eq!((live, staged, backup), (false, true, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::Prepared);
+ }
+ 5 => {
+ assert_eq!((live, staged, backup), (false, true, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::LiveRetained);
+ }
+ 6..=8 => {
+ assert_eq!((live, staged, backup), (true, false, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::LiveRetained);
+ }
+ _ => unreachable!("complete injected failure inventory"),
+ }
+ assert_eq!(
+ recovery_path(&fixture, MARKER_NEXT_FILE_NAME).exists(),
+ matches!(failure, 4 | 8)
+ );
+ }
+ reset_finalize_controls();
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn finalization_rejects_live_stage_and_destination_replacement_without_clobbering() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+
+ let live_replaced = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &live_replaced.paths,
+ &live_replaced.identity,
+ &live_replaced.migrations,
+ &live_replaced.schema,
+ live_replaced.proof(),
+ )
+ .await
+ .expect("stage");
+ let foreign = recovery_path(&live_replaced, "foreign-live");
+ fs::write(&foreign, b"foreign-live").expect("foreign live");
+ fs::set_permissions(&foreign, fs::Permissions::from_mode(0o600)).expect("foreign mode");
+ fs::rename(&foreign, live_replaced.paths.state_database()).expect("replace live");
+ let error = finalize_staged_restore(staged)
+ .await
+ .expect_err("live replacement");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Restore);
+ assert_eq!(
+ fs::read(live_replaced.paths.state_database()).expect("preserved foreign live"),
+ b"foreign-live"
+ );
+ assert!(!recovery_path(&live_replaced, MARKER_FILE_NAME).exists());
+
+ let stage_replaced = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &stage_replaced.paths,
+ &stage_replaced.identity,
+ &stage_replaced.migrations,
+ &stage_replaced.schema,
+ stage_replaced.proof(),
+ )
+ .await
+ .expect("stage");
+ let foreign = recovery_path(&stage_replaced, "foreign-stage-finalize");
+ fs::write(&foreign, b"foreign-stage").expect("foreign stage");
+ fs::set_permissions(&foreign, fs::Permissions::from_mode(0o600)).expect("foreign mode");
+ fs::rename(&foreign, stage_replaced.staged_path()).expect("replace stage");
+ let error = finalize_staged_restore(staged)
+ .await
+ .expect_err("stage replacement");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Restore);
+ assert_eq!(
+ fs::read(stage_replaced.staged_path()).expect("preserved foreign stage"),
+ b"foreign-stage"
+ );
+ assert!(!recovery_path(&stage_replaced, MARKER_FILE_NAME).exists());
+
+ let collision = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &collision.paths,
+ &collision.identity,
+ &collision.migrations,
+ &collision.schema,
+ collision.proof(),
+ )
+ .await
+ .expect("stage");
+ let backup = recovery_path(&collision, BACKUP_FILE_NAME);
+ fs::write(&backup, b"foreign-backup").expect("backup collision");
+ fs::set_permissions(&backup, fs::Permissions::from_mode(0o600)).expect("backup mode");
+ let error = finalize_staged_restore(staged)
+ .await
+ .expect_err("backup collision");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Restore);
+ assert_eq!(
+ fs::read(backup).expect("preserved backup"),
+ b"foreign-backup"
+ );
+ assert!(!recovery_path(&collision, MARKER_FILE_NAME).exists());
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn finalization_preserves_authority_precedence_before_and_after_prepared() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+ for phase in [TEST_PHASE_BEFORE_PREPARED, TEST_PHASE_AFTER_PREPARED] {
+ reset_finalize_controls();
+ let fixture = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &fixture.paths,
+ &fixture.identity,
+ &fixture.migrations,
+ &fixture.schema,
+ fixture.proof(),
+ )
+ .await
+ .expect("stage");
+ TEST_FINALIZE_BLOCK_PHASE.store(phase, Ordering::Release);
+ let caller = tokio::spawn(finalize_staged_restore(staged));
+ wait_for_finalize_phase(phase).await;
+ let old_lock = recovery_path(&fixture, "state.lock.finalize-old");
+ fs::rename(fixture.paths.state_lock(), &old_lock).expect("retain old lock");
+ fs::write(fixture.paths.state_lock(), b"").expect("replacement lock");
+ fs::set_permissions(
+ fixture.paths.state_lock(),
+ fs::Permissions::from_mode(0o600),
+ )
+ .expect("lock mode");
+ TEST_FINALIZE_BLOCK_PHASE.store(0, Ordering::Release);
+ let error = caller.await.expect("caller").expect_err("authority drift");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
+ assert_eq!(
+ recovery_path(&fixture, MARKER_FILE_NAME).exists(),
+ phase == TEST_PHASE_AFTER_PREPARED
+ );
+ assert_eq!(
+ fixture.staged_path().exists(),
+ phase == TEST_PHASE_AFTER_PREPARED
+ );
+ fs::remove_file(old_lock).expect("remove old lock");
+ }
+ reset_finalize_controls();
+ }
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn marker_sync_authority_drift_retains_durable_prepared_stage() {
+ let _serial = STAGE_TEST_LOCK.lock().await;
+ reset_finalize_controls();
+ let fixture = Fixture::new().await;
+ let staged = stage_verified_restore(
+ &fixture.paths,
+ &fixture.identity,
+ &fixture.migrations,
+ &fixture.schema,
+ fixture.proof(),
+ )
+ .await
+ .expect("stage");
+
+ let error = test_finalize_with_failure(staged, 9)
+ .await
+ .expect_err("authority drift after durable marker sync");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::Prepared);
+ assert!(fixture.paths.state_database().exists());
+ assert!(fixture.staged_path().exists());
+ assert!(!recovery_path(&fixture, BACKUP_FILE_NAME).exists());
+ fs::set_permissions(
+ fixture
+ .paths
+ .state_database()
+ .parent()
+ .expect("state directory"),
+ fs::Permissions::from_mode(0o700),
+ )
+ .expect("restore directory mode");
+ reset_finalize_controls();
+ }
+
#[test]
fn public_debug_and_errors_do_not_disclose_paths_or_digests() {
let error = restore_error(RestoreFailureKind::StagedChanged);
diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs
@@ -17,6 +17,7 @@ const METADATA_SOURCE: &str = include_str!("../src/metadata.rs");
const MIGRATION_SOURCE: &str = include_str!("../src/migration.rs");
const OPEN_SOURCE: &str = include_str!("../src/open.rs");
const RESTORE_MARKER_SOURCE: &str = include_str!("../src/restore/marker.rs");
+const RESTORE_FINALIZE_SOURCE: &str = include_str!("../src/restore/finalize.rs");
const RESTORE_ROOT_SOURCE: &str = include_str!("../src/restore/mod.rs");
const RESTORE_STAGE_SOURCE: &str = include_str!("../src/restore/stage.rs");
const STATUS_SOURCE: &str = include_str!("../src/status.rs");
@@ -168,6 +169,19 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"cleanup failure leaves staging or recovery evidence",
"detached work retains authority and exact cleanup ownership",
"does not create or advance a recovery marker",
+ "`finalize_staged_restore` consumes that sealed stage",
+ "bound the exact live inode, length, and digest",
+ "creates and synchronizes the `prepared` marker",
+ "descriptor-relative no-replace operations",
+ "marker advance to `live_retained` or `replacement_installed`",
+ "Cancellation observed before the worker's atomic commit-ownership handoff",
+ "interval before `prepared` becomes durable",
+ "Once `prepared` is durable, staged-artifact cleanup is disarmed",
+ "unknown immediate outcome",
+ "retains the old live database and final marker",
+ "every initialized, read-write-existing, or read-only-inspection pool open refuses",
+ "marker or marker-scratch as `Recovery` before opening SQLite",
+ "does not delete evidence, roll back, reconcile an interruption, or reopen the database",
] {
assert!(
readme_words.contains(required),
@@ -618,7 +632,64 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"Step 067 restore staging source contains deferred or raw authority `{forbidden}`"
);
}
- assert!(ROOT.contains("pub use restore::{StagedServiceRestore, stage_verified_restore};"));
+ let restore_finalize_production = RESTORE_FINALIZE_SOURCE
+ .split_once(
+ "#[cfg(all(test, any(target_os = \"linux\", target_os = \"macos\")))]\nstruct FailingFinalizeOperations",
+ )
+ .map(|(production, _)| production)
+ .expect("restore finalization source must keep failure injection separated");
+ for required in [
+ "pub async fn finalize_staged_restore(",
+ "tokio::task::spawn_blocking",
+ "RestoreRecoveryMarker::prepared(",
+ "staged.disarm_cleanup()",
+ "renameat_with(",
+ "RenameFlags::NOREPLACE",
+ "directory.sync_all()",
+ "RestoreRecoveryPhase::LiveRetained",
+ "RestoreRecoveryPhase::ReplacementInstalled",
+ "verify_named_artifact(",
+ "CancellationOnDrop",
+ ] {
+ assert!(
+ restore_finalize_production.contains(required),
+ "Step 068 restore finalization source is missing `{required}`"
+ );
+ }
+ for forbidden in [
+ "pub struct RestoreRecovery",
+ "pub enum RestoreRecovery",
+ "pub fn marker",
+ "pub fn directory",
+ "pub fn path",
+ "sqlx::",
+ "rusqlite::",
+ "remove_file",
+ "unlinkat",
+ "remove_dir_all",
+ "tokio::time::timeout",
+ "ServiceSqliteHost",
+ ] {
+ assert!(
+ !restore_finalize_production.contains(forbidden),
+ "Step 068 restore finalization source contains deferred or raw authority `{forbidden}`"
+ );
+ }
+ for required in [
+ "refuse_unresolved_recovery",
+ "ServiceSqliteErrorKind::Recovery",
+ "MARKER_FILE_NAME",
+ "MARKER_NEXT_FILE_NAME",
+ ] {
+ assert!(
+ RESTORE_ROOT_SOURCE.contains(required),
+ "Step 068 recovery-open guard is missing `{required}`"
+ );
+ }
+ assert!(OPEN_SOURCE.contains("crate::restore::refuse_unresolved_recovery"));
+ assert!(ROOT.contains(
+ "pub use restore::{StagedServiceRestore, finalize_staged_restore, stage_verified_restore};"
+ ));
for required in [
"self.pool.close().await",