commit 0f50fc836019e3435e1a8324bbaa6befe7f4efa1
parent a1204b53aea99527c401575e85b82b93d830c0f5
Author: triesap <tyson@radroots.org>
Date: Tue, 11 Aug 2026 19:17:23 +0000
service-sqlite: capture online backups
- add cancellation-aware online backup capture for writable service hosts
- secure exact staging identity modes cleanup integrity hashing and durability
- retain pool authority through capture completion cancellation and close drain
- cover WAL consistency races failures portability and package boundaries
Diffstat:
9 files changed, 2217 insertions(+), 13 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -177,9 +177,14 @@ Before editing code:
exact 1,024-byte compact canonical JSON, raw canonical-byte SHA-256, typed
service/instance/source-generation/schema/time binding, singleton
`state.sqlite` inventory, exact `ok` integrity projection, and mandatory
- protected-material exclusion. Parsing is structural only; do not let the
- model perform filesystem capture, SQLite backup, file/digest verification,
- restore, ambient clock access, or public dependency-owned error exposure.
+ protected-material exclusion. Parsing is structural only and must not perform
+ filesystem or SQLite work. Online capture belongs only to the writable
+ `ServiceSqliteHost`: admit one capture at a time, use SQLite's incremental
+ online-backup API, create a caller-selected new owner-only staging directory,
+ return the manifest in memory, and retain host authority until success or
+ exact-artifact cancellation cleanup completes. Do not expose raw backup
+ handles, capture credentials, invent a manifest filename, read an ambient
+ clock, or fold untrusted verification or restore behavior into capture.
- 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/Cargo.lock b/Cargo.lock
@@ -3913,6 +3913,7 @@ dependencies = [
"futures",
"radroots_runtime_paths",
"radroots_storage",
+ "rusqlite",
"rustix 1.1.4",
"serde",
"serde_json",
diff --git a/crates/service_sqlite/Cargo.toml b/crates/service_sqlite/Cargo.toml
@@ -16,12 +16,13 @@ fs2 = { workspace = true }
futures = { workspace = true }
radroots_runtime_paths = { workspace = true }
radroots_storage = { workspace = true }
+rusqlite = { workspace = true, features = ["backup", "bundled"] }
rustix = { workspace = true }
serde = { workspace = true, features = ["derive", "std"] }
serde_json = { workspace = true }
sha2 = { workspace = true }
sqlx = { workspace = true, features = ["runtime-tokio", "sqlite-bundled"] }
-tokio = { workspace = true, features = ["sync"] }
+tokio = { workspace = true, features = ["rt", "sync"] }
[dev-dependencies]
tempfile = { workspace = true }
diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md
@@ -48,8 +48,31 @@ integrity are exactly `ok`, and protected material is always excluded.
Parsing proves only the strict structural and canonical contract. It rejects
unknown, duplicate, null, reordered, whitespace-altered, or version-drifted
input; member bytes, digest, SQLite identity, and actual integrity remain the
-separate backup-verification boundary. This crate does no backup filesystem or
-SQLite capture work while constructing or parsing the manifest model.
+separate backup-verification boundary. Constructing or parsing the manifest
+model performs no filesystem or SQLite work.
+
+Writable hosts provide `ServiceSqliteHost::capture_online_backup` for one
+incremental, point-in-time SQLite capture at a time. The caller supplies an
+injected creation time and the exact new absolute staging-directory path under
+an existing owner-controlled parent; capture creates that directory with mode
+`0700` and its sole `state.sqlite` member with mode `0600`. It uses SQLite's
+online-backup API without checkpointing or copying the live source file, then
+requires exact service metadata, bounded `integrity_check`, an empty
+`foreign_key_check`, a singleton member inventory, SHA-256, and file, staging,
+and parent synchronization before returning the canonical manifest in memory.
+No manifest file, bundle identifier, credential, or protected material is
+written to the staging directory.
+
+Capture rejects read-only, closing, unsupported, colliding, or concurrent
+admission before publishing a result. Dropping the capture future requests
+cancellation; the blocking worker retains its checked-out pool admission,
+writer authority, and exact staging identities until SQLite handles are closed
+and cleanup completes. Host close therefore drains capture and cancellation
+cleanup before it checkpoints or releases authority. Capture has no hidden
+timeout: callers own any deadline by cancelling the future. A completed capture
+is still untrusted backup input until the separate verifier binds its manifest,
+member bytes, expected intent, application metadata, and integrity; capture
+does not provide restore or replacement behavior.
The crate owns mechanics only. Service-specific tables, SQL, repositories,
backup content policy, identity material, process lifecycle, and readiness
diff --git a/crates/service_sqlite/src/backup/capture.rs b/crates/service_sqlite/src/backup/capture.rs
@@ -0,0 +1,1465 @@
+//! Cancellation-aware online capture into a caller-owned staging directory.
+
+use core::fmt;
+use std::{
+ error::Error,
+ ffi::{OsStr, OsString},
+ fs::File,
+ io::{self, Read, Seek, SeekFrom},
+ os::unix::ffi::OsStrExt,
+ path::{Component, Path, PathBuf},
+ sync::{
+ Arc,
+ atomic::{AtomicBool, Ordering},
+ },
+ thread,
+};
+
+#[cfg(test)]
+use core::sync::atomic::AtomicU8;
+
+use rusqlite::{Connection, OpenFlags, OptionalExtension, backup::StepResult, types::ValueRef};
+use rustix::{
+ fs::{AtFlags, FileType, Mode, OFlags, fchmod, fstat, mkdirat, open, openat, statat, unlinkat},
+ process::geteuid,
+};
+use sha2::{Digest, Sha256};
+use sqlx::{Sqlite, pool::PoolConnection};
+
+use crate::{
+ BackupCreatedAtUnixMs, BackupMemberSha256, OpenMode, ServiceBackupManifest,
+ ServiceDatabaseMetadata, ServiceSqliteError, ServiceSqliteErrorKind,
+ open::{BackupSourceValidator, PrivateConnectionPool},
+};
+
+const MAX_STAGING_PATH_BYTES: usize = 4_096;
+const BACKUP_PAGES_PER_STEP: i32 = 64;
+const HASH_BUFFER_BYTES: usize = 16 * 1_024;
+const MAX_ID_UTF8_BYTES: i64 = 128;
+const MAX_INTEGRITY_RESULT_UTF8_BYTES: usize = 64;
+const STATE_FILE_NAME: &str = radroots_runtime_paths::SERVICE_STATE_DATABASE_FILE_NAME;
+const KNOWN_SIDECARS: [&str; 3] = [
+ "state.sqlite-wal",
+ "state.sqlite-shm",
+ "state.sqlite-journal",
+];
+
+pub(crate) const TEST_CAPTURE_PHASE_BEFORE_CREATE: u8 = 1;
+pub(crate) const TEST_CAPTURE_PHASE_STAGING_CREATED: u8 = 2;
+pub(crate) const TEST_CAPTURE_PHASE_BACKUP_STEPPED: u8 = 3;
+pub(crate) const TEST_CAPTURE_PHASE_POST_COPY: u8 = 4;
+pub(crate) const TEST_CAPTURE_PHASE_PRE_FINAL_SYNC: u8 = 5;
+pub(crate) const TEST_CAPTURE_PHASE_METADATA_AWAITED: u8 = 10;
+pub(crate) const TEST_CAPTURE_PHASE_JOIN_AWAITED: u8 = 11;
+#[cfg(test)]
+static TEST_CAPTURE_PHASE: AtomicU8 = AtomicU8::new(0);
+#[cfg(test)]
+static TEST_CAPTURE_BLOCK_PHASE: AtomicU8 = AtomicU8::new(0);
+#[cfg(test)]
+static TEST_CAPTURE_INJECT_METADATA_FAILURE: AtomicBool = AtomicBool::new(false);
+#[cfg(test)]
+static TEST_CAPTURE_PANIC_WORKER: AtomicBool = AtomicBool::new(false);
+
+#[cfg(test)]
+pub(crate) fn test_capture_phase() -> u8 {
+ TEST_CAPTURE_PHASE.load(Ordering::Acquire)
+}
+
+#[cfg(test)]
+pub(crate) fn test_capture_block_phase(phase: u8) {
+ TEST_CAPTURE_BLOCK_PHASE.store(phase, Ordering::Release);
+}
+
+#[cfg(test)]
+pub(crate) fn test_capture_reset() {
+ TEST_CAPTURE_BLOCK_PHASE.store(0, Ordering::Release);
+ TEST_CAPTURE_PHASE.store(0, Ordering::Release);
+ TEST_CAPTURE_INJECT_METADATA_FAILURE.store(false, Ordering::Release);
+ TEST_CAPTURE_PANIC_WORKER.store(false, Ordering::Release);
+}
+
+#[cfg(test)]
+pub(crate) fn test_capture_inject_metadata_failure(enabled: bool) {
+ TEST_CAPTURE_INJECT_METADATA_FAILURE.store(enabled, Ordering::Release);
+}
+
+#[cfg(test)]
+pub(crate) fn test_capture_panic_worker(enabled: bool) {
+ TEST_CAPTURE_PANIC_WORKER.store(enabled, Ordering::Release);
+}
+
+async fn test_async_phase(phase: u8) {
+ #[cfg(test)]
+ {
+ TEST_CAPTURE_PHASE.store(phase, Ordering::Release);
+ while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase {
+ tokio::task::yield_now().await;
+ }
+ }
+ #[cfg(not(test))]
+ let _ = phase;
+}
+
+trait CaptureOperations: Send + Sync {
+ fn sync_state(&self, state: &File) -> io::Result<()>;
+ fn sync_staging(&self, staging: &File) -> io::Result<()>;
+ fn sync_parent(&self, parent: &File) -> io::Result<()>;
+}
+
+struct SystemCaptureOperations;
+
+impl CaptureOperations for SystemCaptureOperations {
+ fn sync_state(&self, state: &File) -> io::Result<()> {
+ state.sync_all()
+ }
+
+ fn sync_staging(&self, staging: &File) -> io::Result<()> {
+ staging.sync_all()
+ }
+
+ fn sync_parent(&self, parent: &File) -> io::Result<()> {
+ parent.sync_all()
+ }
+}
+
+#[cfg(test)]
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) enum TestCaptureSyncFailure {
+ State,
+ Staging,
+ FinalParent,
+}
+
+#[cfg(test)]
+struct FailingCaptureOperations {
+ failure: TestCaptureSyncFailure,
+ parent_syncs: core::sync::atomic::AtomicU8,
+}
+
+#[cfg(test)]
+impl CaptureOperations for FailingCaptureOperations {
+ fn sync_state(&self, state: &File) -> io::Result<()> {
+ if self.failure == TestCaptureSyncFailure::State {
+ Err(io::Error::other("injected state sync failure"))
+ } else {
+ state.sync_all()
+ }
+ }
+
+ fn sync_staging(&self, staging: &File) -> io::Result<()> {
+ if self.failure == TestCaptureSyncFailure::Staging {
+ Err(io::Error::other("injected staging sync failure"))
+ } else {
+ staging.sync_all()
+ }
+ }
+
+ fn sync_parent(&self, parent: &File) -> io::Result<()> {
+ let occurrence = self.parent_syncs.fetch_add(1, Ordering::AcqRel);
+ if self.failure == TestCaptureSyncFailure::FinalParent && occurrence == 1 {
+ Err(io::Error::other("injected parent sync failure"))
+ } else {
+ parent.sync_all()
+ }
+ }
+}
+
+pub(crate) async fn capture_online_backup(
+ pool: &PrivateConnectionPool,
+ closing: &AtomicBool,
+ active: &Arc<AtomicBool>,
+ staging_directory: &Path,
+ created_at_unix_ms: BackupCreatedAtUnixMs,
+) -> Result<ServiceBackupManifest, ServiceSqliteError> {
+ capture_online_backup_with_operations(
+ pool,
+ closing,
+ active,
+ staging_directory,
+ created_at_unix_ms,
+ Arc::new(SystemCaptureOperations),
+ )
+ .await
+}
+
+#[cfg(test)]
+pub(crate) async fn test_capture_online_backup_with_sync_failure(
+ pool: &PrivateConnectionPool,
+ closing: &AtomicBool,
+ active: &Arc<AtomicBool>,
+ staging_directory: &Path,
+ created_at_unix_ms: BackupCreatedAtUnixMs,
+ failure: TestCaptureSyncFailure,
+) -> Result<ServiceBackupManifest, ServiceSqliteError> {
+ capture_online_backup_with_operations(
+ pool,
+ closing,
+ active,
+ staging_directory,
+ created_at_unix_ms,
+ Arc::new(FailingCaptureOperations {
+ failure,
+ parent_syncs: core::sync::atomic::AtomicU8::new(0),
+ }),
+ )
+ .await
+}
+
+async fn capture_online_backup_with_operations(
+ pool: &PrivateConnectionPool,
+ closing: &AtomicBool,
+ active: &Arc<AtomicBool>,
+ staging_directory: &Path,
+ created_at_unix_ms: BackupCreatedAtUnixMs,
+ operations: Arc<dyn CaptureOperations>,
+) -> Result<ServiceBackupManifest, ServiceSqliteError> {
+ if !matches!(
+ pool.mode(),
+ OpenMode::Initialize | OpenMode::ReadWriteExisting
+ ) || closing.load(Ordering::Acquire)
+ {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ }
+ let staging = StagingPath::new(staging_directory)?;
+ let permit = CapturePermit::acquire(Arc::clone(active))?;
+ pool.validate()?;
+ let mut admission = pool.acquire().await?;
+ pool.validate()?;
+ if closing.load(Ordering::Acquire) {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ }
+ let metadata = crate::metadata::verify_database_metadata(&mut admission, pool.identity()).await;
+ #[cfg(test)]
+ let metadata = if TEST_CAPTURE_INJECT_METADATA_FAILURE.load(Ordering::Acquire) {
+ Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata))
+ } else {
+ metadata
+ };
+ test_async_phase(TEST_CAPTURE_PHASE_METADATA_AWAITED).await;
+ pool.validate()?;
+ let metadata = metadata?;
+ let validator = pool.backup_source_validator();
+ validator.validate()?;
+
+ let cancellation = Arc::new(AtomicBool::new(false));
+ let cancellation_guard = CaptureCancellation::new(Arc::clone(&cancellation));
+ let worker = CaptureWorker {
+ _admission: admission,
+ _permit: permit,
+ validator,
+ metadata,
+ staging,
+ created_at_unix_ms,
+ cancellation,
+ operations,
+ };
+ let joined = tokio::task::spawn_blocking(move || worker.run()).await;
+ test_async_phase(TEST_CAPTURE_PHASE_JOIN_AWAITED).await;
+ cancellation_guard.complete();
+ pool.validate()?;
+ let result = joined.map_err(|source| backup_source(BackupFailureKind::Join, source))?;
+ result.map(PendingCapture::commit)
+}
+
+struct CapturePermit {
+ active: Arc<AtomicBool>,
+}
+
+impl CapturePermit {
+ fn acquire(active: Arc<AtomicBool>) -> Result<Self, ServiceSqliteError> {
+ active
+ .compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
+ .map_err(|_| backup_error(BackupFailureKind::AlreadyActive))?;
+ Ok(Self { active })
+ }
+}
+
+impl Drop for CapturePermit {
+ fn drop(&mut self) {
+ self.active.store(false, Ordering::Release);
+ }
+}
+
+struct CaptureCancellation {
+ cancelled: Arc<AtomicBool>,
+ complete: AtomicBool,
+}
+
+impl CaptureCancellation {
+ fn new(cancelled: Arc<AtomicBool>) -> Self {
+ Self {
+ cancelled,
+ complete: AtomicBool::new(false),
+ }
+ }
+
+ fn complete(&self) {
+ self.complete.store(true, Ordering::Release);
+ }
+}
+
+impl Drop for CaptureCancellation {
+ fn drop(&mut self) {
+ if !self.complete.load(Ordering::Acquire) {
+ self.cancelled.store(true, Ordering::Release);
+ }
+ }
+}
+
+struct CaptureWorker {
+ _admission: PoolConnection<Sqlite>,
+ _permit: CapturePermit,
+ validator: BackupSourceValidator,
+ metadata: ServiceDatabaseMetadata,
+ staging: StagingPath,
+ created_at_unix_ms: BackupCreatedAtUnixMs,
+ cancellation: Arc<AtomicBool>,
+ operations: Arc<dyn CaptureOperations>,
+}
+
+impl CaptureWorker {
+ fn run(self) -> Result<PendingCapture, ServiceSqliteError> {
+ self.check_cancelled()?;
+ self.validator.validate()?;
+ self.test_phase(TEST_CAPTURE_PHASE_BEFORE_CREATE);
+ self.check_cancelled()?;
+ let mut staging = StagingGuard::create(&self.staging, self.operations.as_ref())?;
+ self.test_phase(TEST_CAPTURE_PHASE_STAGING_CREATED);
+ #[cfg(test)]
+ if TEST_CAPTURE_PANIC_WORKER.load(Ordering::Acquire) {
+ panic!("injected backup worker failure");
+ }
+ self.check_cancelled()?;
+ self.validator.validate()?;
+ staging.validate()?;
+
+ let source = self.open_source()?;
+ let mut destination = self.open_destination(&staging)?;
+ staging.record_sidecars();
+ verify_database_inventory(&source)?;
+ verify_database_metadata(&source, &self.metadata)?;
+ self.validator.validate()?;
+ staging.validate()?;
+
+ {
+ let backup = match rusqlite::backup::Backup::new(&source, &mut destination) {
+ Ok(backup) => backup,
+ Err(source) => {
+ staging.record_sidecars();
+ return Err(backup_source(BackupFailureKind::Capture, source));
+ }
+ };
+ loop {
+ self.check_cancelled()?;
+ self.validator.validate()?;
+ staging.validate()?;
+ let step = backup.step(BACKUP_PAGES_PER_STEP);
+ staging.record_sidecars();
+ self.test_phase(TEST_CAPTURE_PHASE_BACKUP_STEPPED);
+ let step =
+ step.map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ self.validator.validate()?;
+ staging.validate()?;
+ match step {
+ StepResult::Done => break,
+ StepResult::More => {}
+ StepResult::Busy | StepResult::Locked => thread::yield_now(),
+ _ => return Err(backup_error(BackupFailureKind::Capture)),
+ }
+ }
+ }
+ staging.record_sidecars();
+
+ self.test_phase(TEST_CAPTURE_PHASE_POST_COPY);
+ self.check_cancelled()?;
+ self.validator.validate()?;
+ staging.validate()?;
+ verify_database_inventory(&destination)?;
+ verify_database_metadata(&destination, &self.metadata)?;
+ verify_integrity(&destination)?;
+ staging.record_sidecars();
+ self.check_cancelled()?;
+ self.validator.validate()?;
+ staging.validate()?;
+
+ destination
+ .close()
+ .map_err(|(_, source)| backup_source(BackupFailureKind::Capture, source))?;
+ staging.record_sidecars();
+ source
+ .close()
+ .map_err(|(_, source)| backup_source(BackupFailureKind::Capture, source))?;
+ self.validator.validate()?;
+ staging.validate()?;
+ staging.validate_inventory()?;
+ self.check_cancelled()?;
+
+ staging.sync_state(self.operations.as_ref())?;
+ let (byte_length, digest) = staging.hash_state(&self.cancellation)?;
+ self.test_phase(TEST_CAPTURE_PHASE_PRE_FINAL_SYNC);
+ self.check_cancelled()?;
+ staging.sync_directories(self.operations.as_ref())?;
+ self.validator.validate()?;
+ staging.validate()?;
+ staging.validate_inventory()?;
+ self.check_cancelled()?;
+
+ let manifest = ServiceBackupManifest::from_capture(
+ &self.metadata,
+ self.created_at_unix_ms,
+ byte_length,
+ BackupMemberSha256::from_bytes(digest),
+ )
+ .map_err(|source| backup_source(BackupFailureKind::Manifest, source))?;
+ Ok(PendingCapture {
+ staging,
+ _admission: self._admission,
+ _permit: self._permit,
+ manifest,
+ })
+ }
+
+ fn open_source(&self) -> Result<Connection, ServiceSqliteError> {
+ self.validator.validate()?;
+ let result = Connection::open_with_flags(
+ self.validator.database_path(),
+ OpenFlags::SQLITE_OPEN_READ_ONLY
+ | OpenFlags::SQLITE_OPEN_NO_MUTEX
+ | OpenFlags::SQLITE_OPEN_NOFOLLOW,
+ );
+ self.validator.validate()?;
+ let connection =
+ result.map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ connection
+ .pragma_update(None, "query_only", true)
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ Ok(connection)
+ }
+
+ fn open_destination(&self, staging: &StagingGuard) -> Result<Connection, ServiceSqliteError> {
+ staging.validate()?;
+ let result = Connection::open_with_flags(
+ staging.state_path(),
+ OpenFlags::SQLITE_OPEN_READ_WRITE
+ | OpenFlags::SQLITE_OPEN_NO_MUTEX
+ | OpenFlags::SQLITE_OPEN_NOFOLLOW,
+ );
+ staging.validate()?;
+ result.map_err(|source| backup_source(BackupFailureKind::Capture, source))
+ }
+
+ fn check_cancelled(&self) -> Result<(), ServiceSqliteError> {
+ if self.cancellation.load(Ordering::Acquire) {
+ Err(backup_error(BackupFailureKind::Cancelled))
+ } else {
+ Ok(())
+ }
+ }
+
+ fn test_phase(&self, phase: u8) {
+ #[cfg(test)]
+ {
+ TEST_CAPTURE_PHASE.store(phase, Ordering::Release);
+ while TEST_CAPTURE_BLOCK_PHASE.load(Ordering::Acquire) == phase
+ && !self.cancellation.load(Ordering::Acquire)
+ {
+ thread::yield_now();
+ }
+ }
+ #[cfg(not(test))]
+ let _ = phase;
+ }
+}
+
+struct PendingCapture {
+ staging: StagingGuard,
+ _admission: PoolConnection<Sqlite>,
+ _permit: CapturePermit,
+ manifest: ServiceBackupManifest,
+}
+
+impl PendingCapture {
+ fn commit(mut self) -> ServiceBackupManifest {
+ self.staging.commit();
+ self.manifest
+ }
+}
+
+#[derive(Clone)]
+struct StagingPath {
+ full: PathBuf,
+ parent: PathBuf,
+ name: OsString,
+}
+
+impl StagingPath {
+ fn new(path: &Path) -> Result<Self, ServiceSqliteError> {
+ if path.as_os_str().as_bytes().len() > MAX_STAGING_PATH_BYTES || !path.is_absolute() {
+ return Err(backup_error(BackupFailureKind::InvalidStagingPath));
+ }
+ if path
+ .components()
+ .any(|component| matches!(component, Component::CurDir | Component::ParentDir))
+ {
+ return Err(backup_error(BackupFailureKind::InvalidStagingPath));
+ }
+ let name = match path.components().next_back() {
+ Some(Component::Normal(name)) if !name.as_bytes().is_empty() => name.to_os_string(),
+ _ => return Err(backup_error(BackupFailureKind::InvalidStagingPath)),
+ };
+ let parent = path
+ .parent()
+ .filter(|parent| parent.is_absolute())
+ .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingPath))?
+ .to_path_buf();
+ Ok(Self {
+ full: path.to_path_buf(),
+ parent,
+ name,
+ })
+ }
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+struct FileIdentity {
+ device: u64,
+ inode: u64,
+}
+
+struct StagingGuard {
+ path: StagingPath,
+ parent: File,
+ parent_identity: FileIdentity,
+ directory: File,
+ directory_identity: FileIdentity,
+ state: File,
+ state_identity: FileIdentity,
+ sidecar_identities: [Option<FileIdentity>; KNOWN_SIDECARS.len()],
+ committed: bool,
+}
+
+impl StagingGuard {
+ fn create(
+ path: &StagingPath,
+ operations: &dyn CaptureOperations,
+ ) -> Result<Self, ServiceSqliteError> {
+ let parent = File::from(
+ open(
+ &path.parent,
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?,
+ );
+ let parent_identity = validate_directory_descriptor(&parent, false)?;
+ mkdirat(&parent, &path.name, Mode::RUSR | Mode::WUSR | Mode::XUSR).map_err(|source| {
+ backup_source(
+ if source == rustix::io::Errno::EXIST {
+ BackupFailureKind::StagingCollision
+ } else {
+ BackupFailureKind::CreateStaging
+ },
+ source,
+ )
+ })?;
+ let created_directory_identity = match created_directory_identity(&parent, &path.name) {
+ Ok(identity) => identity,
+ Err(error) => {
+ // Without a proven identity, the current entry must be preserved.
+ let _ = parent.sync_all();
+ return Err(error);
+ }
+ };
+ let directory = match openat(
+ &parent,
+ &path.name,
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ ) {
+ Ok(directory) => File::from(directory),
+ Err(source) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ None,
+ None,
+ );
+ return Err(backup_source(BackupFailureKind::CreateStaging, source));
+ }
+ };
+ let opened_directory_identity = match descriptor_identity(&directory) {
+ Ok(identity) => identity,
+ Err(error) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(error);
+ }
+ };
+ if opened_directory_identity != created_directory_identity {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(backup_error(BackupFailureKind::StagingReplaced));
+ }
+ if let Err(source) = fchmod(&directory, Mode::RUSR | Mode::WUSR | Mode::XUSR) {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(backup_source(BackupFailureKind::CreateStaging, source));
+ }
+ let directory_identity = match validate_directory_descriptor(&directory, true) {
+ Ok(identity) => identity,
+ Err(error) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(error);
+ }
+ };
+ let state = match openat(
+ &directory,
+ STATE_FILE_NAME,
+ OFlags::RDWR | OFlags::CREATE | OFlags::EXCL | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::RUSR | Mode::WUSR,
+ ) {
+ Ok(state) => File::from(state),
+ Err(source) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(backup_source(BackupFailureKind::CreateState, source));
+ }
+ };
+ let created_state_identity = match descriptor_identity(&state) {
+ Ok(identity) => identity,
+ Err(error) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ None,
+ );
+ return Err(error);
+ }
+ };
+ if let Err(source) = fchmod(&state, Mode::RUSR | Mode::WUSR) {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ Some(created_state_identity),
+ );
+ return Err(backup_source(BackupFailureKind::CreateState, source));
+ }
+ let state_identity = match validate_file_descriptor(&state) {
+ Ok(identity) => identity,
+ Err(error) => {
+ cleanup_partial_staging(
+ &parent,
+ &path.name,
+ Some(created_directory_identity),
+ Some(&directory),
+ Some(created_state_identity),
+ );
+ return Err(error);
+ }
+ };
+ let staging = Self {
+ path: path.clone(),
+ parent,
+ parent_identity,
+ directory,
+ directory_identity,
+ state,
+ state_identity,
+ sidecar_identities: [None; KNOWN_SIDECARS.len()],
+ committed: false,
+ };
+ staging.validate()?;
+ operations
+ .sync_parent(&staging.parent)
+ .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?;
+ Ok(staging)
+ }
+
+ fn state_path(&self) -> PathBuf {
+ self.path.full.join(STATE_FILE_NAME)
+ }
+
+ fn validate(&self) -> Result<(), ServiceSqliteError> {
+ validate_reopened_directory(&self.path.parent, &self.parent, self.parent_identity, false)?;
+ validate_directory_entry(
+ &self.parent,
+ &self.path.name,
+ &self.directory,
+ self.directory_identity,
+ )?;
+ validate_file_entry(
+ &self.directory,
+ OsStr::new(STATE_FILE_NAME),
+ &self.state,
+ self.state_identity,
+ )
+ }
+
+ fn validate_inventory(&self) -> Result<(), ServiceSqliteError> {
+ self.validate()?;
+ let mut entries = std::fs::read_dir(&self.path.full)
+ .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?;
+ let first = entries
+ .next()
+ .transpose()
+ .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?
+ .ok_or_else(|| backup_error(BackupFailureKind::InvalidStagingInventory))?;
+ if first.file_name() != OsStr::new(STATE_FILE_NAME) || entries.next().is_some() {
+ return Err(backup_error(BackupFailureKind::InvalidStagingInventory));
+ }
+ Ok(())
+ }
+
+ fn sync_state(&self, operations: &dyn CaptureOperations) -> Result<(), ServiceSqliteError> {
+ self.validate()?;
+ operations
+ .sync_state(&self.state)
+ .map_err(|source| backup_source(BackupFailureKind::SyncState, source))?;
+ self.validate()
+ }
+
+ fn hash_state(&self, cancellation: &AtomicBool) -> Result<(u64, [u8; 32]), ServiceSqliteError> {
+ self.validate()?;
+ let mut state = self
+ .state
+ .try_clone()
+ .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
+ state
+ .seek(SeekFrom::Start(0))
+ .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
+ let mut hasher = Sha256::new();
+ let mut length = 0_u64;
+ let mut buffer = [0_u8; HASH_BUFFER_BYTES];
+ loop {
+ if cancellation.load(Ordering::Acquire) {
+ return Err(backup_error(BackupFailureKind::Cancelled));
+ }
+ let count = state
+ .read(&mut buffer)
+ .map_err(|source| backup_source(BackupFailureKind::HashState, source))?;
+ if count == 0 {
+ break;
+ }
+ length = length
+ .checked_add(
+ u64::try_from(count).map_err(|_| backup_error(BackupFailureKind::HashState))?,
+ )
+ .ok_or_else(|| backup_error(BackupFailureKind::HashState))?;
+ if length > i64::MAX as u64 {
+ return Err(backup_error(BackupFailureKind::HashState));
+ }
+ hasher.update(&buffer[..count]);
+ }
+ if length == 0 {
+ return Err(backup_error(BackupFailureKind::HashState));
+ }
+ self.validate()?;
+ Ok((length, hasher.finalize().into()))
+ }
+
+ fn sync_directories(
+ &self,
+ operations: &dyn CaptureOperations,
+ ) -> Result<(), ServiceSqliteError> {
+ self.validate()?;
+ operations
+ .sync_staging(&self.directory)
+ .map_err(|source| backup_source(BackupFailureKind::SyncStaging, source))?;
+ self.validate()?;
+ operations
+ .sync_parent(&self.parent)
+ .map_err(|source| backup_source(BackupFailureKind::SyncParent, source))?;
+ self.validate()
+ }
+
+ fn record_sidecars(&mut self) {
+ for (index, name) in KNOWN_SIDECARS.iter().enumerate() {
+ if self.sidecar_identities[index].is_none() {
+ self.sidecar_identities[index] = safe_sidecar_identity(&self.directory, name);
+ }
+ }
+ }
+
+ fn commit(&mut self) {
+ self.committed = true;
+ }
+
+ fn cleanup(&mut self) {
+ if self.committed {
+ return;
+ }
+ if validate_directory_entry(
+ &self.parent,
+ &self.path.name,
+ &self.directory,
+ self.directory_identity,
+ )
+ .is_err()
+ {
+ return;
+ }
+ if current_entry_identity(&self.directory, OsStr::new(STATE_FILE_NAME))
+ == Some(self.state_identity)
+ {
+ let _ = unlinkat(&self.directory, STATE_FILE_NAME, AtFlags::empty());
+ }
+ for (sidecar, identity) in KNOWN_SIDECARS.iter().zip(self.sidecar_identities) {
+ if identity.is_some() && safe_sidecar_identity(&self.directory, sidecar) == identity {
+ let _ = unlinkat(&self.directory, *sidecar, AtFlags::empty());
+ }
+ }
+ if current_entry_identity(&self.parent, &self.path.name) == Some(self.directory_identity) {
+ let _ = unlinkat(&self.parent, &self.path.name, AtFlags::REMOVEDIR);
+ }
+ let _ = self.parent.sync_all();
+ self.committed = true;
+ }
+}
+
+fn cleanup_partial_staging(
+ parent: &File,
+ name: &OsStr,
+ directory_identity: Option<FileIdentity>,
+ directory: Option<&File>,
+ state_identity: Option<FileIdentity>,
+) {
+ if let (Some(directory), Some(state_identity)) = (directory, state_identity)
+ && current_entry_identity(directory, OsStr::new(STATE_FILE_NAME)) == Some(state_identity)
+ {
+ let _ = unlinkat(directory, STATE_FILE_NAME, AtFlags::empty());
+ }
+ if directory_identity.is_some() && current_entry_identity(parent, name) == directory_identity {
+ let _ = unlinkat(parent, name, AtFlags::REMOVEDIR);
+ }
+ let _ = parent.sync_all();
+}
+
+impl Drop for StagingGuard {
+ fn drop(&mut self) {
+ self.cleanup();
+ }
+}
+
+fn created_directory_identity(
+ parent: &File,
+ name: &OsStr,
+) -> Result<FileIdentity, ServiceSqliteError> {
+ let status = statat(parent, name, AtFlags::SYMLINK_NOFOLLOW)
+ .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?;
+ if !FileType::from_raw_mode(status.st_mode).is_dir()
+ || status.st_uid != geteuid().as_raw()
+ || u32::from(status.st_mode) & 0o022 != 0
+ {
+ return Err(backup_error(BackupFailureKind::StagingReplaced));
+ }
+ identity(&status)
+}
+
+fn descriptor_identity(descriptor: &File) -> Result<FileIdentity, ServiceSqliteError> {
+ let status = fstat(descriptor)
+ .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?;
+ identity(&status)
+}
+
+fn validate_directory_descriptor(
+ directory: &File,
+ exact_owner_mode: bool,
+) -> Result<FileIdentity, ServiceSqliteError> {
+ let status = fstat(directory)
+ .map_err(|source| backup_source(BackupFailureKind::InvalidStagingParent, source))?;
+ let mode = u32::from(status.st_mode) & 0o777;
+ if !FileType::from_raw_mode(status.st_mode).is_dir()
+ || status.st_uid != geteuid().as_raw()
+ || if exact_owner_mode {
+ mode != 0o700
+ } else {
+ mode & 0o022 != 0
+ }
+ {
+ return Err(backup_error(BackupFailureKind::InvalidStagingParent));
+ }
+ identity(&status)
+}
+
+fn validate_file_descriptor(file: &File) -> Result<FileIdentity, ServiceSqliteError> {
+ let status = fstat(file)
+ .map_err(|source| backup_source(BackupFailureKind::InvalidStagingInventory, source))?;
+ 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
+ {
+ return Err(backup_error(BackupFailureKind::InvalidStagingInventory));
+ }
+ identity(&status)
+}
+
+fn validate_reopened_directory(
+ path: &Path,
+ held: &File,
+ expected: FileIdentity,
+ exact_owner_mode: bool,
+) -> Result<(), ServiceSqliteError> {
+ let current = File::from(
+ open(
+ path,
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
+ );
+ if validate_directory_descriptor(¤t, exact_owner_mode)? != expected
+ || validate_directory_descriptor(held, exact_owner_mode)? != expected
+ {
+ return Err(backup_error(BackupFailureKind::StagingReplaced));
+ }
+ Ok(())
+}
+
+fn validate_directory_entry(
+ parent: &File,
+ name: &OsStr,
+ held: &File,
+ expected: FileIdentity,
+) -> Result<(), ServiceSqliteError> {
+ let current = File::from(
+ openat(
+ parent,
+ name,
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
+ );
+ if validate_directory_descriptor(¤t, true)? != expected
+ || validate_directory_descriptor(held, true)? != expected
+ {
+ return Err(backup_error(BackupFailureKind::StagingReplaced));
+ }
+ Ok(())
+}
+
+fn validate_file_entry(
+ directory: &File,
+ name: &OsStr,
+ held: &File,
+ expected: FileIdentity,
+) -> Result<(), ServiceSqliteError> {
+ let current = File::from(
+ openat(
+ directory,
+ name,
+ OFlags::RDONLY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .map_err(|source| backup_source(BackupFailureKind::StagingReplaced, source))?,
+ );
+ if validate_file_descriptor(¤t)? != expected
+ || validate_file_descriptor(held)? != expected
+ {
+ return Err(backup_error(BackupFailureKind::StagingReplaced));
+ }
+ Ok(())
+}
+
+fn current_entry_identity(directory: &File, name: &OsStr) -> Option<FileIdentity> {
+ let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?;
+ identity(&status).ok()
+}
+
+fn safe_sidecar_identity(directory: &File, name: &str) -> Option<FileIdentity> {
+ let status = statat(directory, name, AtFlags::SYMLINK_NOFOLLOW).ok()?;
+ if !FileType::from_raw_mode(status.st_mode).is_file()
+ || u64::from(status.st_nlink) != 1
+ || status.st_uid != geteuid().as_raw()
+ {
+ return None;
+ }
+ identity(&status).ok()
+}
+
+fn identity(status: &rustix::fs::Stat) -> Result<FileIdentity, ServiceSqliteError> {
+ Ok(FileIdentity {
+ device: u64::try_from(status.st_dev)
+ .map_err(|_| backup_error(BackupFailureKind::StagingReplaced))?,
+ inode: status.st_ino,
+ })
+}
+
+fn verify_database_inventory(connection: &Connection) -> Result<(), ServiceSqliteError> {
+ let mut statement = connection
+ .prepare("PRAGMA database_list")
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ let mut rows = statement
+ .query([])
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ let first = rows
+ .next()
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?
+ .ok_or_else(|| backup_error(BackupFailureKind::Capture))?;
+ let sequence: i64 = first
+ .get(0)
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ let name: String = first
+ .get(1)
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?;
+ if sequence != 0 || name != "main" {
+ return Err(backup_error(BackupFailureKind::Capture));
+ }
+ if rows
+ .next()
+ .map_err(|source| backup_source(BackupFailureKind::Capture, source))?
+ .is_some()
+ {
+ return Err(backup_error(BackupFailureKind::Capture));
+ }
+ Ok(())
+}
+
+fn verify_database_metadata(
+ connection: &Connection,
+ expected: &ServiceDatabaseMetadata,
+) -> Result<(), ServiceSqliteError> {
+ let application_id: i64 = connection
+ .pragma_query_value(None, "application_id", |row| row.get(0))
+ .map_err(metadata_source)?;
+ let row_count: i64 = connection
+ .query_row(
+ "SELECT COUNT(*) FROM (SELECT 1 FROM radroots_service_metadata LIMIT 2)",
+ [],
+ |row| row.get(0),
+ )
+ .map_err(metadata_source)?;
+ if row_count != 1 || application_id != i64::from(expected.application_id().get()) {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
+ }
+ let row = connection
+ .query_row(
+ "SELECT
+ CASE WHEN typeof(service_id) = 'text'
+ AND length(CAST(service_id AS BLOB)) BETWEEN 1 AND ?1
+ THEN service_id END,
+ CASE WHEN typeof(instance_id) = 'text'
+ AND length(CAST(instance_id AS BLOB)) BETWEEN 1 AND ?1
+ THEN instance_id END,
+ CASE WHEN typeof(source_generation) = 'blob'
+ AND length(source_generation) = 32
+ THEN source_generation END,
+ CASE WHEN typeof(state_schema_version) = 'integer'
+ THEN state_schema_version END,
+ CASE WHEN typeof(created_at_unix_ms) = 'integer'
+ THEN created_at_unix_ms END
+ FROM radroots_service_metadata
+ WHERE singleton = 1
+ LIMIT 1",
+ [MAX_ID_UTF8_BYTES],
+ |row| {
+ Ok((
+ row.get::<_, Option<String>>(0)?,
+ row.get::<_, Option<String>>(1)?,
+ row.get::<_, Option<Vec<u8>>>(2)?,
+ row.get::<_, Option<i64>>(3)?,
+ row.get::<_, Option<i64>>(4)?,
+ ))
+ },
+ )
+ .optional()
+ .map_err(metadata_source)?;
+ let Some((Some(service), Some(instance), Some(generation), Some(schema), Some(created_at))) =
+ row
+ else {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
+ };
+ if service != expected.service().as_str()
+ || instance != expected.instance().as_str()
+ || generation.as_slice() != expected.source_generation().as_bytes()
+ || schema != i64::from(expected.state_schema_version().get())
+ || created_at != i64::try_from(expected.created_at_unix_ms()).unwrap_or(-1)
+ {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Metadata));
+ }
+ Ok(())
+}
+
+fn verify_integrity(connection: &Connection) -> Result<(), ServiceSqliteError> {
+ let mut statement = connection
+ .prepare("PRAGMA integrity_check(1)")
+ .map_err(integrity_source)?;
+ let mut rows = statement.query([]).map_err(integrity_source)?;
+ let row = rows
+ .next()
+ .map_err(integrity_source)?
+ .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
+ let result = row.get_ref(0).map_err(integrity_source)?;
+ if !integrity_projection_is_ok(result) || rows.next().map_err(integrity_source)?.is_some() {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
+ }
+ let mut statement = connection
+ .prepare("PRAGMA foreign_key_check")
+ .map_err(integrity_source)?;
+ if statement
+ .query([])
+ .map_err(integrity_source)?
+ .next()
+ .map_err(integrity_source)?
+ .is_some()
+ {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
+ }
+ Ok(())
+}
+
+fn integrity_projection_is_ok(value: ValueRef<'_>) -> bool {
+ matches!(
+ value,
+ ValueRef::Text(bytes)
+ if !bytes.is_empty()
+ && bytes.len() <= MAX_INTEGRITY_RESULT_UTF8_BYTES
+ && bytes == b"ok"
+ )
+}
+
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum BackupFailureKind {
+ InvalidStagingPath,
+ InvalidStagingParent,
+ StagingCollision,
+ CreateStaging,
+ CreateState,
+ StagingReplaced,
+ InvalidStagingInventory,
+ AlreadyActive,
+ Capture,
+ Cancelled,
+ HashState,
+ SyncState,
+ SyncStaging,
+ SyncParent,
+ Manifest,
+ Join,
+}
+
+struct BackupFailure {
+ kind: BackupFailureKind,
+ source: Option<Box<dyn Error + Send + Sync + 'static>>,
+}
+
+impl fmt::Debug for BackupFailure {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter
+ .debug_struct("BackupFailure")
+ .field("kind", &self.kind)
+ .field("source", &self.source.as_ref().map(|_| "[redacted]"))
+ .finish()
+ }
+}
+
+impl fmt::Display for BackupFailure {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str(match self.kind {
+ BackupFailureKind::InvalidStagingPath => "backup staging path is invalid",
+ BackupFailureKind::InvalidStagingParent => "backup staging parent is invalid",
+ BackupFailureKind::StagingCollision => "backup staging destination already exists",
+ BackupFailureKind::CreateStaging => "backup staging directory could not be created",
+ BackupFailureKind::CreateState => "backup state member could not be created",
+ BackupFailureKind::StagingReplaced => "backup staging identity changed",
+ BackupFailureKind::InvalidStagingInventory => "backup staging inventory is invalid",
+ BackupFailureKind::AlreadyActive => "another backup capture is active",
+ BackupFailureKind::Capture => "online backup capture failed",
+ BackupFailureKind::Cancelled => "online backup capture was cancelled",
+ BackupFailureKind::HashState => "backup state member could not be hashed",
+ BackupFailureKind::SyncState => "backup state member could not be synchronized",
+ BackupFailureKind::SyncStaging => "backup staging directory could not be synchronized",
+ BackupFailureKind::SyncParent => "backup staging parent could not be synchronized",
+ BackupFailureKind::Manifest => "backup manifest could not be constructed",
+ BackupFailureKind::Join => "backup worker could not be joined",
+ })
+ }
+}
+
+impl Error for BackupFailure {
+ fn source(&self) -> Option<&(dyn Error + 'static)> {
+ self.source
+ .as_deref()
+ .map(|source| source as &(dyn Error + 'static))
+ }
+}
+
+fn backup_error(kind: BackupFailureKind) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Backup,
+ BackupFailure { kind, source: None },
+ )
+}
+
+fn backup_source(
+ kind: BackupFailureKind,
+ source: impl Error + Send + Sync + 'static,
+) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Backup,
+ BackupFailure {
+ kind,
+ source: Some(Box::new(source)),
+ },
+ )
+}
+
+fn integrity_source(source: rusqlite::Error) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Integrity, source)
+}
+
+fn metadata_source(source: rusqlite::Error) -> ServiceSqliteError {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Metadata, source)
+}
+
+#[cfg(test)]
+mod tests {
+ use std::os::unix::fs::{PermissionsExt, symlink};
+
+ use tempfile::tempdir;
+
+ use super::*;
+
+ #[test]
+ fn staging_path_rejects_relative_parent_and_oversize_inputs() {
+ assert!(StagingPath::new(Path::new("relative/stage")).is_err());
+ assert!(StagingPath::new(Path::new("/tmp/../stage")).is_err());
+ assert!(StagingPath::new(Path::new("/")).is_err());
+ let large = format!("/tmp/{}", "x".repeat(MAX_STAGING_PATH_BYTES));
+ assert!(StagingPath::new(Path::new(&large)).is_err());
+ }
+
+ #[test]
+ fn staging_guard_creates_exact_modes_and_cleans_uncommitted_state() {
+ let root = tempdir().expect("root");
+ std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
+ .expect("parent mode");
+ let path = StagingPath::new(&root.path().join("backup-stage")).expect("path");
+ {
+ let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
+ staging.validate().expect("validate");
+ assert_eq!(
+ std::fs::metadata(&path.full)
+ .expect("directory metadata")
+ .permissions()
+ .mode()
+ & 0o777,
+ 0o700
+ );
+ assert_eq!(
+ std::fs::metadata(staging.state_path())
+ .expect("state metadata")
+ .permissions()
+ .mode()
+ & 0o777,
+ 0o600
+ );
+ }
+ assert!(!path.full.exists());
+ }
+
+ #[test]
+ fn staging_guard_rejects_collision_and_preserves_replacement() {
+ let root = tempdir().expect("root");
+ std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
+ .expect("parent mode");
+ let full = root.path().join("backup-stage");
+ std::fs::create_dir(&full).expect("collision");
+ assert!(
+ StagingGuard::create(
+ &StagingPath::new(&full).expect("path"),
+ &SystemCaptureOperations,
+ )
+ .is_err()
+ );
+ std::fs::remove_dir(&full).expect("remove collision");
+
+ let path = StagingPath::new(&full).expect("path");
+ let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
+ let original = staging.state_path();
+ let replacement = full.join("replacement");
+ std::fs::write(&replacement, b"foreign").expect("replacement");
+ std::fs::set_permissions(&replacement, std::fs::Permissions::from_mode(0o600))
+ .expect("replacement mode");
+ std::fs::rename(&replacement, &original).expect("replace state");
+ drop(staging);
+ assert_eq!(std::fs::read(&original).expect("preserved"), b"foreign");
+ assert!(full.exists());
+ }
+
+ #[test]
+ fn staging_guard_rejects_insecure_or_symlinked_parent_without_mutation() {
+ let root = tempdir().expect("root");
+ std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
+ .expect("root mode");
+ let insecure = root.path().join("insecure");
+ std::fs::create_dir(&insecure).expect("insecure parent");
+ std::fs::set_permissions(&insecure, std::fs::Permissions::from_mode(0o770))
+ .expect("insecure mode");
+ let insecure_stage = insecure.join("stage");
+ let insecure_error = match StagingGuard::create(
+ &StagingPath::new(&insecure_stage).expect("insecure staging path"),
+ &SystemCaptureOperations,
+ ) {
+ Ok(_) => panic!("group-writable parent must be rejected"),
+ Err(error) => error,
+ };
+ assert_eq!(insecure_error.kind(), ServiceSqliteErrorKind::Backup);
+ assert!(!insecure_stage.exists());
+
+ let parent = root.path().join("real-parent");
+ std::fs::create_dir(&parent).expect("real parent");
+ std::fs::set_permissions(&parent, std::fs::Permissions::from_mode(0o700))
+ .expect("real parent mode");
+ let alias = root.path().join("parent-alias");
+ symlink(&parent, &alias).expect("parent symlink");
+ let symlink_stage = alias.join("stage");
+ assert!(
+ StagingGuard::create(
+ &StagingPath::new(&symlink_stage).expect("symlink staging path"),
+ &SystemCaptureOperations,
+ )
+ .is_err()
+ );
+ assert!(!parent.join("stage").exists());
+ }
+
+ #[test]
+ fn staging_guard_rejects_hardlinked_state_and_preserves_other_link() {
+ let root = tempdir().expect("root");
+ std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
+ .expect("parent mode");
+ let path = StagingPath::new(&root.path().join("backup-stage")).expect("path");
+ let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
+ std::fs::write(staging.state_path(), b"captured").expect("state bytes");
+ let outside = root.path().join("outside-link");
+ std::fs::hard_link(staging.state_path(), &outside).expect("hard link");
+ assert!(staging.validate().is_err());
+ drop(staging);
+ assert_eq!(
+ std::fs::read(&outside).expect("other link survives"),
+ b"captured"
+ );
+ assert!(!path.full.exists());
+ }
+
+ #[test]
+ fn partial_cleanup_preserves_replaced_directory_and_state_entries() {
+ let root = tempdir().expect("root");
+ std::fs::set_permissions(root.path(), std::fs::Permissions::from_mode(0o700))
+ .expect("root mode");
+ let parent = File::from(
+ open(
+ root.path(),
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .expect("open parent"),
+ );
+ mkdirat(&parent, "stage", Mode::RUSR | Mode::WUSR | Mode::XUSR)
+ .expect("create original stage");
+ let original_identity =
+ created_directory_identity(&parent, OsStr::new("stage")).expect("original identity");
+ let original = File::from(
+ openat(
+ &parent,
+ "stage",
+ OFlags::RDONLY | OFlags::DIRECTORY | OFlags::NOFOLLOW | OFlags::CLOEXEC,
+ Mode::empty(),
+ )
+ .expect("open original stage"),
+ );
+ std::fs::rename(root.path().join("stage"), root.path().join("retired"))
+ .expect("retire original stage");
+ std::fs::create_dir(root.path().join("stage")).expect("replacement stage");
+ std::fs::set_permissions(
+ root.path().join("stage"),
+ std::fs::Permissions::from_mode(0o700),
+ )
+ .expect("replacement mode");
+ std::fs::write(root.path().join("stage/foreign"), b"foreign").expect("replacement member");
+ cleanup_partial_staging(
+ &parent,
+ OsStr::new("stage"),
+ Some(original_identity),
+ Some(&original),
+ None,
+ );
+ assert_eq!(
+ std::fs::read(root.path().join("stage/foreign")).expect("foreign survives"),
+ b"foreign"
+ );
+
+ let path = StagingPath::new(&root.path().join("state-stage")).expect("path");
+ let staging = StagingGuard::create(&path, &SystemCaptureOperations).expect("staging");
+ let retired_state = path.full.join("retired-state");
+ std::fs::rename(staging.state_path(), &retired_state).expect("retire state");
+ std::fs::write(staging.state_path(), b"foreign-state").expect("foreign state");
+ std::fs::set_permissions(staging.state_path(), std::fs::Permissions::from_mode(0o600))
+ .expect("foreign state mode");
+ cleanup_partial_staging(
+ &staging.parent,
+ &staging.path.name,
+ Some(staging.directory_identity),
+ Some(&staging.directory),
+ Some(staging.state_identity),
+ );
+ assert_eq!(
+ std::fs::read(staging.state_path()).expect("replacement state survives"),
+ b"foreign-state"
+ );
+ drop(staging);
+ assert!(path.full.exists());
+ }
+
+ #[test]
+ fn integrity_projection_bounds_corrupt_text_before_semantic_acceptance() {
+ assert!(integrity_projection_is_ok(ValueRef::Text(b"ok")));
+ let maximum = vec![b'x'; MAX_INTEGRITY_RESULT_UTF8_BYTES];
+ assert!(!integrity_projection_is_ok(ValueRef::Text(&maximum)));
+ let over_maximum = vec![b'x'; MAX_INTEGRITY_RESULT_UTF8_BYTES + 1];
+ assert!(!integrity_projection_is_ok(ValueRef::Text(&over_maximum)));
+ assert!(!integrity_projection_is_ok(ValueRef::Null));
+ }
+
+ #[test]
+ fn active_capture_permit_is_exclusive_and_recoverable() {
+ let active = Arc::new(AtomicBool::new(false));
+ let first = CapturePermit::acquire(Arc::clone(&active)).expect("first permit");
+ assert!(CapturePermit::acquire(Arc::clone(&active)).is_err());
+ drop(first);
+ assert!(CapturePermit::acquire(active).is_ok());
+ }
+}
diff --git a/crates/service_sqlite/src/backup/mod.rs b/crates/service_sqlite/src/backup/mod.rs
@@ -1,7 +1,21 @@
//! Stable service backup manifest identity.
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+mod capture;
mod manifest;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+pub(crate) use capture::capture_online_backup;
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) use capture::{
+ TEST_CAPTURE_PHASE_BACKUP_STEPPED, TEST_CAPTURE_PHASE_BEFORE_CREATE,
+ TEST_CAPTURE_PHASE_JOIN_AWAITED, TEST_CAPTURE_PHASE_METADATA_AWAITED,
+ TEST_CAPTURE_PHASE_POST_COPY, TEST_CAPTURE_PHASE_PRE_FINAL_SYNC,
+ TEST_CAPTURE_PHASE_STAGING_CREATED, TestCaptureSyncFailure, test_capture_block_phase,
+ test_capture_inject_metadata_failure, test_capture_online_backup_with_sync_failure,
+ test_capture_panic_worker, test_capture_phase, test_capture_reset,
+};
+
pub use manifest::{
BACKUP_MANIFEST_CANONICAL_MAX_BYTES, BACKUP_MANIFEST_SCHEMA, BACKUP_MANIFEST_SCHEMA_VERSION,
BACKUP_STATE_MEMBER_NAME, BackupCreatedAtUnixMs, BackupManifestContractError,
diff --git a/crates/service_sqlite/src/connection.rs b/crates/service_sqlite/src/connection.rs
@@ -4,6 +4,7 @@ use core::fmt;
use std::{
error::Error,
future::Future,
+ path::Path,
pin::Pin,
sync::{
Arc,
@@ -46,6 +47,8 @@ pub struct ServiceSqliteHost {
closing: AtomicBool,
#[cfg(any(target_os = "linux", target_os = "macos"))]
close_state: tokio::sync::Mutex<ServiceSqliteHostCloseState>,
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ backup_active: Arc<AtomicBool>,
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -201,6 +204,43 @@ impl ServiceSqliteHost {
}
}
+ /// Captures one point-in-time SQLite backup into a new staging directory.
+ ///
+ /// The caller supplies the exact new absolute staging-directory path and an
+ /// injected creation time. Capture is available only on writable hosts and
+ /// admits at most one active capture per host. The canonical manifest and a
+ /// successful result are returned only after the visible staging directory's
+ /// sole `state.sqlite` member has passed metadata, integrity, digest, and
+ /// durability checks. The manifest remains in memory and is not written into
+ /// the staging directory.
+ ///
+ /// Dropping this future requests cancellation. The admitted worker retains
+ /// host authority until it has closed SQLite handles and either completed or
+ /// cleaned the exact staging artifacts, so `close` drains that work before
+ /// releasing writer authority.
+ pub async fn capture_online_backup(
+ &self,
+ staging_directory: &Path,
+ created_at_unix_ms: crate::BackupCreatedAtUnixMs,
+ ) -> Result<crate::ServiceBackupManifest, ServiceSqliteError> {
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ {
+ crate::backup::capture_online_backup(
+ &self.pool,
+ &self.closing,
+ &self.backup_active,
+ staging_directory,
+ created_at_unix_ms,
+ )
+ .await
+ }
+ #[cfg(not(any(target_os = "linux", target_os = "macos")))]
+ {
+ let _ = (staging_directory, created_at_unix_ms);
+ Err(unsupported_host())
+ }
+ }
+
/// Executes one runner-owned transaction without exposing its connection or pool.
///
/// Dropping this future before the runner enables its outer commit quarantines
@@ -436,6 +476,7 @@ impl ServiceSqliteHost {
pool,
closing: AtomicBool::new(false),
close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending),
+ backup_active: Arc::new(AtomicBool::new(false)),
}
}
@@ -836,6 +877,9 @@ mod tests {
#[cfg(any(target_os = "linux", target_os = "macos"))]
use tokio::sync::Notify;
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ static CAPTURE_TEST_LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
+
use super::*;
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -998,6 +1042,546 @@ mod tests {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn online_backup_captures_exact_member_manifest_and_preserves_source() {
+ use sha2::Digest;
+
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ crate::backup::test_capture_reset();
+ let (root, paths, identity, migrations, schema, host) = initialized_host().await;
+ host.transaction(|transaction| {
+ Box::pin(async move {
+ sqlx::query("INSERT INTO host_probe (value) VALUES (41), (42)")
+ .execute(&mut *transaction)
+ .await
+ .map(|_| ())
+ })
+ })
+ .await
+ .expect("seed live WAL state");
+
+ let output = root.path().join("backup-output");
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+ let collision = output.join("collision");
+ fs::create_dir(&collision).expect("preexisting collision");
+ fs::write(collision.join("foreign"), b"preserve").expect("foreign collision member");
+ let collision_error = host
+ .capture_online_backup(
+ &collision,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_100).expect("backup creation time"),
+ )
+ .await
+ .expect_err("existing destination is rejected");
+ assert_eq!(collision_error.kind(), ServiceSqliteErrorKind::Backup);
+ assert_eq!(
+ fs::read(collision.join("foreign")).expect("preserved collision"),
+ b"preserve"
+ );
+
+ let source_database_before = fs::read(paths.state_database()).expect("source bytes");
+ let source_inventory_before =
+ fs::read_dir(paths.state_database().parent().expect("source directory"))
+ .expect("source inventory")
+ .map(|entry| entry.expect("source entry").file_name())
+ .collect::<std::collections::BTreeSet<_>>();
+ let stage = output.join("successful");
+ let created_at =
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_101).expect("backup creation time");
+ let manifest = host
+ .capture_online_backup(&stage, created_at)
+ .await
+ .expect("online backup");
+
+ assert_eq!(manifest.service(), paths.service());
+ assert_eq!(manifest.instance(), paths.instance());
+ assert_eq!(manifest.source_generation(), identity.source_generation());
+ assert_eq!(
+ manifest.state_schema_version(),
+ identity.supported_state_schema_version()
+ );
+ assert_eq!(manifest.created_at_unix_ms(), created_at);
+ assert_eq!(manifest.members().len(), 1);
+ assert!(!manifest.protected_material_included());
+ assert_eq!(manifest.integrity().sqlite(), "ok");
+ assert_eq!(manifest.integrity().foreign_keys(), "ok");
+
+ let state = stage.join("state.sqlite");
+ let bytes = fs::read(&state).expect("captured state bytes");
+ let digest: [u8; 32] = sha2::Sha256::digest(&bytes).into();
+ assert_eq!(manifest.members()[0].byte_length(), bytes.len() as u64);
+ assert_eq!(manifest.members()[0].sha256().as_bytes(), &digest);
+ assert_eq!(
+ fs::metadata(&stage)
+ .expect("stage metadata")
+ .permissions()
+ .mode()
+ & 0o777,
+ 0o700
+ );
+ assert_eq!(
+ fs::metadata(&state)
+ .expect("state metadata")
+ .permissions()
+ .mode()
+ & 0o777,
+ 0o600
+ );
+ assert_eq!(
+ fs::read_dir(&stage)
+ .expect("stage inventory")
+ .map(|entry| entry.expect("stage entry").file_name())
+ .collect::<Vec<_>>(),
+ vec![std::ffi::OsString::from("state.sqlite")]
+ );
+
+ let backup = rusqlite::Connection::open_with_flags(
+ &state,
+ rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NOFOLLOW,
+ )
+ .expect("open captured database");
+ assert_eq!(
+ backup
+ .query_row("SELECT COUNT(*) FROM host_probe", [], |row| row
+ .get::<_, i64>(0))
+ .expect("captured row count"),
+ 2
+ );
+ backup.close().expect("close captured database");
+
+ assert_eq!(
+ fs::read(paths.state_database()).expect("source bytes after capture"),
+ source_database_before
+ );
+ assert_eq!(
+ fs::read_dir(paths.state_database().parent().expect("source directory"))
+ .expect("source inventory after")
+ .map(|entry| entry.expect("source entry").file_name())
+ .collect::<std::collections::BTreeSet<_>>(),
+ source_inventory_before
+ );
+
+ host.close().await.expect("close writable host");
+ let inspection = ServiceSqliteHost::open_read_only_inspection(
+ &paths,
+ &identity,
+ &migrations,
+ &schema,
+ ServiceSqliteConnectionOptions::reviewed(),
+ )
+ .await
+ .expect("open inspection");
+ let forbidden = output.join("read-only-forbidden");
+ let error = inspection
+ .capture_online_backup(&forbidden, created_at)
+ .await
+ .expect_err("read-only capture is unavailable");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
+ assert!(!forbidden.exists());
+ inspection.close().await.expect("close inspection");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn cancelled_capture_cleans_up_close_drains_and_next_capture_recovers() {
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ crate::backup::test_capture_reset();
+ let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+ let output = root.path().join("cancel-output");
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+ let cancelled_stage = output.join("cancelled");
+ crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED);
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = cancelled_stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_200)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase()
+ != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED
+ {
+ tokio::task::yield_now().await;
+ }
+
+ let concurrent_stage = output.join("concurrent");
+ let concurrent = host
+ .capture_online_backup(
+ &concurrent_stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_201).expect("backup creation time"),
+ )
+ .await
+ .expect_err("second capture is rejected");
+ assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Backup);
+ assert!(!concurrent_stage.exists());
+
+ capture.abort();
+ assert!(
+ capture
+ .await
+ .expect_err("capture task is cancelled")
+ .is_cancelled()
+ );
+ while host.backup_active.load(Ordering::Acquire) {
+ tokio::task::yield_now().await;
+ }
+ crate::backup::test_capture_reset();
+ assert!(!cancelled_stage.exists());
+
+ let recovered_stage = output.join("recovered");
+ host.capture_online_backup(
+ &recovered_stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_202).expect("backup creation time"),
+ )
+ .await
+ .expect("capture recovers after cancellation");
+ assert!(recovered_stage.join("state.sqlite").exists());
+
+ crate::backup::test_capture_reset();
+ crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED);
+ let close_drained_stage = output.join("close-drained");
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = close_drained_stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_203)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase()
+ != crate::backup::TEST_CAPTURE_PHASE_STAGING_CREATED
+ {
+ tokio::task::yield_now().await;
+ }
+ capture.abort();
+ assert!(
+ capture
+ .await
+ .expect_err("capture task is cancelled")
+ .is_cancelled()
+ );
+ host.close()
+ .await
+ .expect("close drains every admitted capture");
+ crate::backup::test_capture_reset();
+ assert!(!close_drained_stage.exists());
+ assert!(!host.backup_active.load(Ordering::Acquire));
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn capture_cancellation_is_cleanup_safe_at_every_governed_phase() {
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+
+ for (index, phase) in [
+ crate::backup::TEST_CAPTURE_PHASE_BEFORE_CREATE,
+ crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED,
+ crate::backup::TEST_CAPTURE_PHASE_POST_COPY,
+ crate::backup::TEST_CAPTURE_PHASE_PRE_FINAL_SYNC,
+ ]
+ .into_iter()
+ .enumerate()
+ {
+ crate::backup::test_capture_reset();
+ let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+ if phase == crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED {
+ host.transaction(|transaction| {
+ Box::pin(async move {
+ sqlx::raw_sql(
+ "WITH RECURSIVE sequence(value) AS (
+ VALUES(1)
+ UNION ALL
+ SELECT value + 1 FROM sequence WHERE value < 100000
+ )
+ INSERT INTO host_probe(value) SELECT 0 FROM sequence",
+ )
+ .execute(&mut *transaction)
+ .await
+ .map(|_| ())
+ })
+ })
+ .await
+ .expect("seed multi-batch cancellation fixture");
+ }
+
+ let output = root.path().join(format!("phase-cancel-output-{index}"));
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+ let stage = output.join("cancelled");
+ crate::backup::test_capture_block_phase(phase);
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_220 + index as u64)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase() != phase {
+ tokio::task::yield_now().await;
+ }
+
+ capture.abort();
+ assert!(
+ capture
+ .await
+ .expect_err("capture task is cancelled")
+ .is_cancelled()
+ );
+ host.close()
+ .await
+ .expect("close drains phase-cancelled capture cleanup");
+ crate::backup::test_capture_reset();
+
+ assert!(!stage.exists(), "phase {phase} must leave no staging tree");
+ assert!(!host.backup_active.load(Ordering::Acquire));
+ let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
+ .expect("writer authority can be reacquired")
+ .expect("writable mode returns authority");
+ authority.release().expect("release reacquired authority");
+ }
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn online_backup_remains_consistent_with_a_concurrent_wal_writer() {
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ crate::backup::test_capture_reset();
+ let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+ host.transaction(|transaction| {
+ Box::pin(async move {
+ sqlx::raw_sql(
+ "WITH RECURSIVE sequence(value) AS (
+ VALUES(1) UNION ALL SELECT value + 1 FROM sequence WHERE value < 100000
+ )
+ INSERT INTO host_probe(value) SELECT 0 FROM sequence",
+ )
+ .execute(&mut *transaction)
+ .await
+ .map(|_| ())
+ })
+ })
+ .await
+ .expect("seed a multi-batch database");
+
+ let output = root.path().join("concurrent-output");
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+ let stage = output.join("consistent");
+ crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED);
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_300)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase()
+ != crate::backup::TEST_CAPTURE_PHASE_BACKUP_STEPPED
+ {
+ tokio::task::yield_now().await;
+ }
+ host.transaction(|transaction| {
+ Box::pin(async move {
+ sqlx::query("UPDATE host_probe SET value = 1 WHERE rowid IN (1, 2)")
+ .execute(&mut *transaction)
+ .await
+ .map(|_| ())
+ })
+ })
+ .await
+ .expect("commit concurrent WAL transaction");
+ crate::backup::test_capture_block_phase(0);
+ capture
+ .await
+ .expect("capture joins")
+ .expect("capture remains consistent");
+ crate::backup::test_capture_reset();
+
+ let backup = rusqlite::Connection::open_with_flags(
+ stage.join("state.sqlite"),
+ rusqlite::OpenFlags::SQLITE_OPEN_READ_ONLY | rusqlite::OpenFlags::SQLITE_OPEN_NOFOLLOW,
+ )
+ .expect("open backup");
+ let updated: i64 = backup
+ .query_row(
+ "SELECT COUNT(*) FROM host_probe WHERE value = 1",
+ [],
+ |row| row.get(0),
+ )
+ .expect("count transaction projection");
+ assert!(
+ updated == 0 || updated == 2,
+ "backup must not tear a transaction"
+ );
+ assert_eq!(
+ backup
+ .query_row("PRAGMA integrity_check(1)", [], |row| row
+ .get::<_, String>(0))
+ .expect("backup integrity"),
+ "ok"
+ );
+ backup.close().expect("close backup");
+ host.close().await.expect("close writer host");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn backup_sync_failures_cleanup_and_leave_host_recoverable() {
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ crate::backup::test_capture_reset();
+ let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let output = root.path().join("sync-failure-output");
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+
+ for (index, failure) in [
+ crate::backup::TestCaptureSyncFailure::State,
+ crate::backup::TestCaptureSyncFailure::Staging,
+ crate::backup::TestCaptureSyncFailure::FinalParent,
+ ]
+ .into_iter()
+ .enumerate()
+ {
+ let stage = output.join(format!("failure-{index}"));
+ let error = crate::backup::test_capture_online_backup_with_sync_failure(
+ &host.pool,
+ &host.closing,
+ &host.backup_active,
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_400 + index as u64)
+ .expect("backup creation time"),
+ failure,
+ )
+ .await
+ .expect_err("injected synchronization failure must reject capture");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup);
+ assert!(!stage.exists());
+ assert!(!host.backup_active.load(Ordering::Acquire));
+ assert_eq!(row_count(&host).await, 0);
+ }
+
+ let recovered = output.join("recovered");
+ host.capture_online_backup(
+ &recovered,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_410).expect("backup creation time"),
+ )
+ .await
+ .expect("host remains usable after sync failures");
+ assert!(recovered.join("state.sqlite").exists());
+ host.close().await.expect("close host");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn backup_await_boundaries_preserve_authority_precedence() {
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ crate::backup::test_capture_reset();
+ let (root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+ let output = root.path().join("precedence-output");
+ fs::create_dir(&output).expect("backup parent");
+ fs::set_permissions(&output, fs::Permissions::from_mode(0o700))
+ .expect("backup parent mode");
+
+ crate::backup::test_capture_inject_metadata_failure(true);
+ crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED);
+ let metadata_stage = output.join("metadata");
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = metadata_stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_420)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase()
+ != crate::backup::TEST_CAPTURE_PHASE_METADATA_AWAITED
+ {
+ tokio::task::yield_now().await;
+ }
+ let retired_lock = paths
+ .state_lock()
+ .parent()
+ .expect("state directory")
+ .join("retired-backup-precedence.lock");
+ fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
+ crate::backup::test_capture_block_phase(0);
+ let metadata_error = capture
+ .await
+ .expect("metadata capture joins")
+ .expect_err("authority overrides metadata failure");
+ assert_eq!(metadata_error.kind(), ServiceSqliteErrorKind::Authority);
+ fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
+ crate::backup::test_capture_reset();
+ assert!(!metadata_stage.exists());
+ assert_eq!(row_count(&host).await, 0);
+
+ crate::backup::test_capture_panic_worker(true);
+ crate::backup::test_capture_block_phase(crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED);
+ let join_stage = output.join("join");
+ let capture = tokio::spawn({
+ let host = Arc::clone(&host);
+ let stage = join_stage.clone();
+ async move {
+ host.capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_000_421)
+ .expect("backup creation time"),
+ )
+ .await
+ }
+ });
+ while crate::backup::test_capture_phase() != crate::backup::TEST_CAPTURE_PHASE_JOIN_AWAITED
+ {
+ tokio::task::yield_now().await;
+ }
+ fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock again");
+ crate::backup::test_capture_block_phase(0);
+ let join_error = capture
+ .await
+ .expect("join-failure capture joins")
+ .expect_err("authority overrides worker join failure");
+ assert_eq!(join_error.kind(), ServiceSqliteErrorKind::Authority);
+ fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock again");
+ crate::backup::test_capture_reset();
+ assert!(!join_stage.exists());
+ assert_eq!(row_count(&host).await, 0);
+ host.close().await.expect("close host");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
#[derive(Debug, PartialEq, Eq)]
struct StateFileSnapshot {
bytes: Vec<u8>,
@@ -1777,7 +2361,7 @@ mod tests {
#[test]
fn host_and_transaction_errors_are_redacted_and_source_free() {
- let error = ServiceSqliteTransactionError::rollback_failed(
+ let error = ServiceSqliteTransactionError::not_committed_with_operation(
Some("operation-secret"),
ServiceSqliteError::new(ServiceSqliteErrorKind::Open),
);
diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs
@@ -193,7 +193,7 @@ pub(crate) struct PrivateConnectionPool {
schema_catalog: SchemaCatalog,
mode: OpenMode,
policy: ServiceSqliteConnectionOptions,
- resources: Mutex<PrivateConnectionResources>,
+ resources: Arc<Mutex<PrivateConnectionResources>>,
close_driver: tokio::sync::Mutex<PrivateCloseDriver>,
#[cfg(test)]
close_phase: std::sync::atomic::AtomicU8,
@@ -211,6 +211,41 @@ struct PrivateConnectionResources {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
+#[derive(Clone)]
+pub(crate) struct BackupSourceValidator {
+ binding: DirectoryBinding,
+ paths: ServiceSqlitePaths,
+ resources: Arc<Mutex<PrivateConnectionResources>>,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl BackupSourceValidator {
+ pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> {
+ let resources = self.resources.lock().map_err(|_| {
+ connection_error(
+ ServiceSqliteErrorKind::Authority,
+ ConnectionFailureKind::AuthorityMismatch,
+ )
+ })?;
+ resources
+ .authority
+ .as_ref()
+ .ok_or_else(|| {
+ connection_error(
+ ServiceSqliteErrorKind::Authority,
+ ConnectionFailureKind::AuthorityMismatch,
+ )
+ })?
+ .validate_for(&self.paths)?;
+ self.binding.validate(&self.paths)
+ }
+
+ pub(crate) fn database_path(&self) -> &Path {
+ self.paths.state_database()
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
enum PrivateCloseDriver {
Pending,
Connecting(BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>>),
@@ -223,7 +258,7 @@ enum PrivateCloseDriver {
Complete(Option<ServiceSqliteErrorKind>),
}
-#[cfg(test)]
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
pub(crate) const TEST_CLOSE_PHASE_CHECKPOINT: u8 = 4;
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -273,6 +308,14 @@ impl PrivateConnectionPool {
&self.identity
}
+ pub(crate) fn backup_source_validator(&self) -> BackupSourceValidator {
+ BackupSourceValidator {
+ binding: self.binding.clone(),
+ paths: self.paths.clone(),
+ resources: Arc::clone(&self.resources),
+ }
+ }
+
pub(crate) fn catalog(&self) -> &MigrationCatalog {
&self.catalog
}
@@ -351,7 +394,7 @@ impl PrivateConnectionPool {
pub(crate) async fn close(self) -> Option<WriterAuthority> {
self.pool.close().await;
- let mut resources = match self.resources.into_inner() {
+ let mut resources = match self.resources.lock() {
Ok(resources) => resources,
Err(poisoned) => poisoned.into_inner(),
};
@@ -891,10 +934,10 @@ async fn open_connection_pool(
schema_catalog: retained_schema_catalog,
mode,
policy,
- resources: Mutex::new(PrivateConnectionResources {
+ resources: Arc::new(Mutex::new(PrivateConnectionResources {
authority,
inspection_guard,
- }),
+ })),
close_driver: tokio::sync::Mutex::new(PrivateCloseDriver::Pending),
#[cfg(test)]
close_phase: std::sync::atomic::AtomicU8::new(0),
diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs
@@ -5,6 +5,7 @@ const README: &str = include_str!("../README.md");
const ROOT: &str = include_str!("../src/lib.rs");
const AUTHORITY_SOURCE: &str = include_str!("../src/authority.rs");
const BACKUP_SOURCE: &str = include_str!("../src/backup/manifest.rs");
+const BACKUP_CAPTURE_SOURCE: &str = include_str!("../src/backup/capture.rs");
const CONFIG_SOURCE: &str = include_str!("../src/config.rs");
const CONNECTION_SOURCE: &str = include_str!("../src/connection.rs");
const ERROR_SOURCE: &str = include_str!("../src/error.rs");
@@ -39,6 +40,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"futures",
"radroots_runtime_paths",
"radroots_storage",
+ "rusqlite",
"rustix",
"serde",
"serde_json",
@@ -105,7 +107,22 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"integrity are exactly `ok`, and protected material is always excluded",
"Parsing proves only the strict structural and canonical contract",
"unknown, duplicate, null, reordered, whitespace-altered, or version-drifted input",
- "does no backup filesystem or SQLite capture work",
+ "Constructing or parsing the manifest model performs no filesystem or SQLite work",
+ "Writable hosts provide `ServiceSqliteHost::capture_online_backup`",
+ "one incremental, point-in-time SQLite capture at a time",
+ "exact new absolute staging-directory path",
+ "directory with mode `0700` and its sole `state.sqlite` member with mode `0600`",
+ "online-backup API without checkpointing or copying the live source file",
+ "requires exact service metadata, bounded `integrity_check`, an empty `foreign_key_check`",
+ "SHA-256, and file, staging, and parent synchronization",
+ "returning the canonical manifest in memory",
+ "No manifest file, bundle identifier, credential, or protected material",
+ "Dropping the capture future requests cancellation",
+ "retains its checked-out pool admission, writer authority, and exact staging identities",
+ "Host close therefore drains capture and cancellation cleanup",
+ "Capture has no hidden timeout",
+ "callers own any deadline by cancelling the future",
+ "does not provide restore or replacement behavior",
] {
assert!(
readme_words.contains(required),
@@ -132,6 +149,10 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
.split_once("#[cfg(test)]")
.map(|(production, _)| production)
.expect("backup source must keep tests separated");
+ let backup_capture_production = BACKUP_CAPTURE_SOURCE
+ .split_once("#[cfg(test)]\nmod tests")
+ .map(|(production, _)| production)
+ .expect("backup capture source must keep tests separated");
for required in [
"ServiceSqliteErrorCode",
@@ -282,6 +303,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"cancelling the future yields no result",
"must reread authoritative state before any idempotent retry",
"pub async fn close(&self)",
+ "pub async fn capture_online_backup(",
"closing.store(true, Ordering::Release)",
"close_state.lock().await",
"ServiceSqliteHostCloseState::Complete",
@@ -331,6 +353,52 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
}
for required in [
+ "rusqlite::backup::Backup::new",
+ "tokio::task::spawn_blocking",
+ "BACKUP_PAGES_PER_STEP",
+ "HASH_BUFFER_BYTES",
+ "OFlags::CREATE | OFlags::EXCL | OFlags::NOFOLLOW | OFlags::CLOEXEC",
+ "Mode::RUSR | Mode::WUSR | Mode::XUSR",
+ "Mode::RUSR | Mode::WUSR",
+ "PRAGMA integrity_check(1)",
+ "ValueRef::Text",
+ "MAX_INTEGRITY_RESULT_UTF8_BYTES",
+ "PRAGMA foreign_key_check",
+ "BackupSourceValidator",
+ "PoolConnection<Sqlite>",
+ "CaptureCancellation",
+ "CapturePermit",
+ "ServiceBackupManifest::from_capture",
+ "sync_state",
+ "sync_directories",
+ "hash_state",
+ "validate_inventory",
+ ] {
+ assert!(
+ backup_capture_production.contains(required),
+ "Step 064 backup capture source is missing `{required}`"
+ );
+ }
+ for forbidden in [
+ "pub use rusqlite",
+ "pub fn restore",
+ "pub async fn restore",
+ "pub fn verify_backup",
+ "pub async fn verify_backup",
+ "VACUUM INTO",
+ "SystemTime::now",
+ "tokio::time::timeout",
+ "manifest.json",
+ "create_dir_all",
+ "std::fs::copy",
+ ] {
+ assert!(
+ !backup_capture_production.contains(forbidden) && !ROOT.contains(forbidden),
+ "Step 064 backup capture source contains deferred or public authority `{forbidden}`"
+ );
+ }
+
+ for required in [
"self.pool.close().await",
"PRAGMA wal_checkpoint(TRUNCATE)",
"if busy == 0",