commit 6f70340ad71fa51c135059f408ca6926ea0821ae
parent 120621ba9381ea1ea06b5f500974da43fa0fed98
Author: triesap <tyson@radroots.org>
Date: Wed, 12 Aug 2026 02:07:13 +0000
service-sqlite: add durability failpoints
- add a private instance-scoped one-shot durability edge inventory
- instrument initialize, transaction, backup, restore, and close boundaries
- preserve authority precedence and exact rollback or recovery evidence
- qualify edge ordering, portability, redaction, and package boundaries
Diffstat:
12 files changed, 1436 insertions(+), 49 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -277,6 +277,14 @@ Before editing code:
the adapter. Keep inspection advisory: do not add a reservation, host/pool or
SQLite dependency, ambient timer, background sampler, service default, or
status persistence to this crate.
+- Durability failpoints are private, instance-scoped test mechanics only. Keep
+ a closed before/after inventory across initialization, transaction commit,
+ online backup, restore-marker persistence, restore rename/synchronization,
+ and explicit close. One armed controller may fail one selected edge once;
+ ordinary controllers have zero behavior. Never export failpoint types, use
+ process-global failpoint state, select a point from environment or service
+ configuration, or add a Cargo feature that alters production behavior.
+ Process-level crash and signal qualification remains a separate layer.
- 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
@@ -248,6 +248,20 @@ or hidden sampling. Service configuration, threshold defaults, cache refresh,
status persistence, admission wiring, and route projection remain consumer
responsibilities.
+Durability fault injection is a private test-only mechanism. A closed
+instance-scoped controller can arm exactly one named before/after boundary and
+returns one injected error the first time that boundary is reached; later hits
+are no-ops. The complete inventory covers database initialization, runner-owned
+transaction begin and commit, online-backup creation/copy/synchronization,
+restore-marker creation and advancement, both restore rename/synchronization
+steps, and explicit host drain/checkpoint/connection-close/authority-release.
+An ordinary controller has zero behavior, and no failpoint type or selector is
+exported from the crate root. There is no process-global failpoint state,
+environment or configuration selector, Cargo feature, hidden task, timer,
+panic, or process-exit behavior. These deterministic in-process edges qualify
+error ordering, rollback, cleanup, recovery evidence, and one-shot retry
+semantics; process crashes and signals remain a separate qualification layer.
+
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/backup/capture.rs b/crates/service_sqlite/src/backup/capture.rs
@@ -170,6 +170,7 @@ pub(crate) async fn capture_online_backup(
active: &Arc<AtomicBool>,
staging_directory: &Path,
created_at_unix_ms: BackupCreatedAtUnixMs,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<ServiceBackupManifest, ServiceSqliteError> {
capture_online_backup_with_operations(
pool,
@@ -178,6 +179,7 @@ pub(crate) async fn capture_online_backup(
staging_directory,
created_at_unix_ms,
Arc::new(SystemCaptureOperations),
+ failpoints,
)
.await
}
@@ -201,6 +203,7 @@ pub(crate) async fn test_capture_online_backup_with_sync_failure(
failure,
parent_syncs: core::sync::atomic::AtomicU8::new(0),
}),
+ &crate::failpoint::DurabilityFailpoints::default(),
)
.await
}
@@ -212,6 +215,7 @@ async fn capture_online_backup_with_operations(
staging_directory: &Path,
created_at_unix_ms: BackupCreatedAtUnixMs,
operations: Arc<dyn CaptureOperations>,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<ServiceBackupManifest, ServiceSqliteError> {
if !matches!(
pool.mode(),
@@ -252,6 +256,7 @@ async fn capture_online_backup_with_operations(
created_at_unix_ms,
cancellation,
operations,
+ failpoints: failpoints.clone(),
};
let joined = tokio::task::spawn_blocking(move || worker.run()).await;
test_async_phase(TEST_CAPTURE_PHASE_JOIN_AWAITED).await;
@@ -315,6 +320,7 @@ struct CaptureWorker {
created_at_unix_ms: BackupCreatedAtUnixMs,
cancellation: Arc<AtomicBool>,
operations: Arc<dyn CaptureOperations>,
+ failpoints: crate::failpoint::DurabilityFailpoints,
}
impl CaptureWorker {
@@ -323,7 +329,17 @@ impl CaptureWorker {
self.validator.validate()?;
self.test_phase(TEST_CAPTURE_PHASE_BEFORE_CREATE);
self.check_cancelled()?;
+ self.hit_checked(
+ None,
+ crate::failpoint::DurabilityFailpoint::BackupBeforeCreate,
+ BackupFailureKind::CreateStaging,
+ )?;
let mut staging = StagingGuard::create(&self.staging, self.operations.as_ref())?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupAfterCreate,
+ BackupFailureKind::CreateStaging,
+ )?;
self.test_phase(TEST_CAPTURE_PHASE_STAGING_CREATED);
#[cfg(test)]
if TEST_CAPTURE_PANIC_WORKER.load(Ordering::Acquire) {
@@ -341,6 +357,11 @@ impl CaptureWorker {
self.validator.validate()?;
staging.validate()?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupBeforeCopy,
+ BackupFailureKind::Capture,
+ )?;
{
let backup = match rusqlite::backup::Backup::new(&source, &mut destination) {
Ok(backup) => backup,
@@ -368,6 +389,11 @@ impl CaptureWorker {
}
}
}
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupAfterCopy,
+ BackupFailureKind::Capture,
+ )?;
staging.record_sidecars();
self.test_phase(TEST_CAPTURE_PHASE_POST_COPY);
@@ -394,11 +420,31 @@ impl CaptureWorker {
staging.validate_inventory()?;
self.check_cancelled()?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupBeforeFileSync,
+ BackupFailureKind::SyncState,
+ )?;
staging.sync_state(self.operations.as_ref())?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupAfterFileSync,
+ BackupFailureKind::SyncState,
+ )?;
let (byte_length, digest) = staging.hash_state(&self.cancellation)?;
self.test_phase(TEST_CAPTURE_PHASE_PRE_FINAL_SYNC);
self.check_cancelled()?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupBeforeDirectorySync,
+ BackupFailureKind::SyncStaging,
+ )?;
staging.sync_directories(self.operations.as_ref())?;
+ self.hit_checked(
+ Some(&staging),
+ crate::failpoint::DurabilityFailpoint::BackupAfterDirectorySync,
+ BackupFailureKind::SyncParent,
+ )?;
self.validator.validate()?;
staging.validate()?;
staging.validate_inventory()?;
@@ -456,6 +502,26 @@ impl CaptureWorker {
}
}
+ fn hit_checked(
+ &self,
+ staging: Option<&StagingGuard>,
+ point: crate::failpoint::DurabilityFailpoint,
+ failure: BackupFailureKind,
+ ) -> Result<(), ServiceSqliteError> {
+ #[cfg(test)]
+ self.failpoints
+ .observe(point, TEST_CAPTURE_PHASE.load(Ordering::Acquire));
+ let injected = self
+ .failpoints
+ .hit(point)
+ .map_err(|source| backup_source(failure, source));
+ self.validator.validate()?;
+ if let Some(staging) = staging {
+ staging.validate()?;
+ }
+ injected
+ }
+
fn test_phase(&self, phase: u8) {
#[cfg(test)]
{
diff --git a/crates/service_sqlite/src/connection.rs b/crates/service_sqlite/src/connection.rs
@@ -51,6 +51,8 @@ pub struct ServiceSqliteHost {
backup_active: Arc<AtomicBool>,
#[cfg(any(target_os = "linux", target_os = "macos"))]
integrity_driver: tokio::sync::Mutex<IntegrityInspectionDriver>,
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ failpoints: crate::failpoint::DurabilityFailpoints,
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -276,7 +278,7 @@ impl ServiceSqliteHost {
let validation = self.pool.validate();
(cleanup, validation)
};
- let close = self.pool.close_explicit().await;
+ let close = self.pool.close_explicit(&self.failpoints).await;
match close {
Err(retryable) => Err(retryable),
Ok(terminal) => {
@@ -321,6 +323,7 @@ impl ServiceSqliteHost {
&self.backup_active,
staging_directory,
created_at_unix_ms,
+ &self.failpoints,
)
.await
}
@@ -435,6 +438,16 @@ impl ServiceSqliteHost {
self.pool
.validate()
.map_err(ServiceSqliteTransactionError::not_committed)?;
+ let before_begin = self
+ .failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeBegin)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ });
+ self.pool
+ .validate()
+ .map_err(ServiceSqliteTransactionError::not_committed)?;
+ before_begin.map_err(ServiceSqliteTransactionError::not_committed)?;
let mut transaction = match match self.pool.mode() {
OpenMode::Initialize | OpenMode::ReadWriteExisting => {
connection.begin_with("BEGIN IMMEDIATE").await
@@ -448,9 +461,33 @@ impl ServiceSqliteHost {
)));
}
};
- self.pool
- .validate()
- .map_err(ServiceSqliteTransactionError::not_committed)?;
+ let injected_after_begin = self
+ .failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterBegin)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ });
+ let after_begin = self.pool.validate().and(injected_after_begin);
+ if let Err(error) = after_begin {
+ let permit = gate.permit_runner_rollback();
+ let rollback = transaction.rollback().await.map_err(sqlite_source);
+ drop(permit);
+ let rollback_was_confirmed =
+ gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
+ let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
+ let authority = self.pool.validate();
+ if let Some(rollback_error) = authority
+ .err()
+ .or_else(|| rollback.err().filter(|_| !rollback_was_confirmed))
+ .or_else(|| remove.err())
+ {
+ return Err(ServiceSqliteTransactionError::rollback_failed(
+ None,
+ rollback_error,
+ ));
+ }
+ return Err(ServiceSqliteTransactionError::not_committed(error));
+ }
let operation_result = {
let database_control_rejected = Arc::new(AtomicBool::new(false));
@@ -518,6 +555,18 @@ impl ServiceSqliteHost {
&database_control_rejected,
)
.await;
+ let precommit = match precommit {
+ Ok(()) => {
+ let injected = self
+ .failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::TransactionBeforeCommit)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ });
+ self.pool.validate().and(injected)
+ }
+ Err(error) => Err(error),
+ };
if let Err(error) = precommit {
let permit = gate.permit_runner_rollback();
let rollback = transaction.rollback().await.map_err(sqlite_source);
@@ -526,11 +575,10 @@ impl ServiceSqliteHost {
gate.rejected_commit_rolled_back() && !connection.is_in_transaction();
let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
let authority = self.pool.validate();
- if let Some(rollback_error) = rollback
+ if let Some(rollback_error) = authority
.err()
- .filter(|_| !rollback_was_confirmed)
+ .or_else(|| rollback.err().filter(|_| !rollback_was_confirmed))
.or_else(|| remove.err())
- .or_else(|| authority.err())
{
return Err(ServiceSqliteTransactionError::rollback_failed(
None,
@@ -543,16 +591,22 @@ impl ServiceSqliteHost {
let permit = gate.permit_outer_commit();
let commit = transaction.commit().await.map_err(sqlite_source);
drop(permit);
- let commit = match commit {
- Ok(()) => Ok(()),
- Err(error) => Err(ServiceSqliteTransactionError::commit_outcome_unknown(error)),
+ let injected_after_commit = if commit.is_ok() {
+ self.failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::TransactionAfterCommit)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ })
+ } else {
+ Ok(())
};
let remove = gate.remove(&mut connection).await.map_err(sqlite_source);
- commit?;
- remove.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
self.pool
.validate()
.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
+ commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
+ remove.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
+ injected_after_commit.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
let final_policy = crate::migration::read_connection_policy(&mut connection)
.await
.map_err(ServiceSqliteTransactionError::commit_outcome_unknown)?;
@@ -650,9 +704,15 @@ impl ServiceSqliteHost {
close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending),
backup_active: Arc::new(AtomicBool::new(false)),
integrity_driver: tokio::sync::Mutex::new(IntegrityInspectionDriver::Idle),
+ failpoints: crate::failpoint::DurabilityFailpoints::default(),
}
}
+ #[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+ fn arm_durability_failpoint(&self, point: crate::failpoint::DurabilityFailpoint) {
+ self.failpoints.arm(point);
+ }
+
fn lifecycle_state(&self) -> &'static str {
#[cfg(any(target_os = "linux", target_os = "macos"))]
{
@@ -2757,6 +2817,263 @@ mod tests {
#[cfg(any(target_os = "linux", target_os = "macos"))]
#[tokio::test]
+ async fn transaction_durability_edges_preserve_exact_commit_semantics() {
+ use crate::failpoint::DurabilityFailpoint;
+
+ for (point, expected_reached) in [
+ (
+ DurabilityFailpoint::TransactionBeforeBegin,
+ &[DurabilityFailpoint::TransactionBeforeBegin][..],
+ ),
+ (
+ DurabilityFailpoint::TransactionAfterBegin,
+ &[
+ DurabilityFailpoint::TransactionBeforeBegin,
+ DurabilityFailpoint::TransactionAfterBegin,
+ ][..],
+ ),
+ (
+ DurabilityFailpoint::TransactionBeforeCommit,
+ &[
+ DurabilityFailpoint::TransactionBeforeBegin,
+ DurabilityFailpoint::TransactionAfterBegin,
+ DurabilityFailpoint::TransactionBeforeCommit,
+ ][..],
+ ),
+ (
+ DurabilityFailpoint::TransactionAfterCommit,
+ &[
+ DurabilityFailpoint::TransactionBeforeBegin,
+ DurabilityFailpoint::TransactionAfterBegin,
+ DurabilityFailpoint::TransactionBeforeCommit,
+ DurabilityFailpoint::TransactionAfterCommit,
+ ][..],
+ ),
+ ] {
+ let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ host.arm_durability_failpoint(point);
+ let error = host
+ .transaction(|transaction| {
+ Box::pin(async move {
+ sqlx::query("INSERT INTO host_probe (value) VALUES (72)")
+ .execute(&mut *transaction)
+ .await
+ .map(|_| ())
+ })
+ })
+ .await
+ .expect_err("injected transaction edge");
+ let reached = host.failpoints.reached();
+ if point == DurabilityFailpoint::TransactionAfterCommit {
+ assert_eq!(
+ error.kind(),
+ ServiceSqliteTransactionErrorKind::CommitOutcomeUnknown
+ );
+ assert_eq!(row_count(&host).await, 1);
+ } else {
+ assert_eq!(
+ error.kind(),
+ ServiceSqliteTransactionErrorKind::NotCommitted
+ );
+ assert_eq!(row_count(&host).await, 0);
+ }
+ assert!(host.failpoints.fired());
+ assert_eq!(reached, expected_reached);
+ host.close()
+ .await
+ .expect("close host after transaction edge");
+ }
+
+ let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ host.arm_durability_failpoint(DurabilityFailpoint::TransactionBeforeCommit);
+ let error = host
+ .transaction(|transaction| {
+ Box::pin(async move {
+ let _ = sqlx::raw_sql(
+ "PRAGMA trusted_schema=ON; INSERT INTO host_probe (value) VALUES (73)",
+ )
+ .execute(&mut *transaction)
+ .await;
+ Ok::<_, Infallible>(())
+ })
+ })
+ .await
+ .expect_err("precommit policy rejection must precede the commit-edge hook");
+ assert_eq!(
+ error.kind(),
+ ServiceSqliteTransactionErrorKind::NotCommitted
+ );
+ assert!(!host.failpoints.fired());
+ assert_eq!(
+ host.failpoints.reached(),
+ [
+ DurabilityFailpoint::TransactionBeforeBegin,
+ DurabilityFailpoint::TransactionAfterBegin,
+ ]
+ );
+ host.failpoints.disarm();
+ assert_eq!(row_count(&host).await, 0);
+ host.close().await.expect("close policy-drift host");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn backup_durability_edges_fail_once_clean_exact_stage_and_recover() {
+ use crate::failpoint::DurabilityFailpoint;
+
+ let _serial = CAPTURE_TEST_LOCK.lock().await;
+ for (index, (point, expected_phase)) in [
+ (DurabilityFailpoint::BackupBeforeCreate, 1),
+ (DurabilityFailpoint::BackupAfterCreate, 1),
+ (DurabilityFailpoint::BackupBeforeCopy, 2),
+ (DurabilityFailpoint::BackupAfterCopy, 3),
+ (DurabilityFailpoint::BackupBeforeFileSync, 4),
+ (DurabilityFailpoint::BackupAfterFileSync, 4),
+ (DurabilityFailpoint::BackupBeforeDirectorySync, 5),
+ (DurabilityFailpoint::BackupAfterDirectorySync, 5),
+ ]
+ .into_iter()
+ .enumerate()
+ {
+ let (root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let stage = root.path().join(format!("failpoint-backup-{index}"));
+ host.arm_durability_failpoint(point);
+ let error = host
+ .capture_online_backup(
+ &stage,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_072_000).expect("capture time"),
+ )
+ .await
+ .expect_err("injected backup edge");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Backup);
+ assert!(host.failpoints.fired());
+ assert_eq!(host.failpoints.reached().last(), Some(&point));
+ assert_eq!(
+ host.failpoints.observation(point),
+ Some(expected_phase),
+ "backup failpoint must fire in its named lifecycle phase"
+ );
+ assert!(!stage.exists(), "owned failed stage is cleaned");
+ let recovery = root.path().join(format!("recovered-backup-{index}"));
+ host.capture_online_backup(
+ &recovery,
+ crate::BackupCreatedAtUnixMs::new(1_700_000_072_001).expect("capture time"),
+ )
+ .await
+ .expect("one-shot failpoint permits retry");
+ assert!(recovery.join(crate::BACKUP_STATE_MEMBER_NAME).is_file());
+ host.close().await.expect("close backup host");
+ }
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn close_durability_edges_are_once_only_retryable_or_terminal() {
+ use crate::failpoint::DurabilityFailpoint;
+
+ let all_close_edges = [
+ DurabilityFailpoint::CloseBeforeDrain,
+ DurabilityFailpoint::CloseAfterDrain,
+ DurabilityFailpoint::CloseBeforeCheckpoint,
+ DurabilityFailpoint::CloseAfterCheckpoint,
+ DurabilityFailpoint::CloseBeforeConnectionClose,
+ DurabilityFailpoint::CloseAfterConnectionClose,
+ DurabilityFailpoint::CloseBeforeAuthorityRelease,
+ DurabilityFailpoint::CloseAfterAuthorityRelease,
+ ];
+ for (point, expected, retryable, expected_phase, expected_reached) in [
+ (
+ DurabilityFailpoint::CloseBeforeDrain,
+ ServiceSqliteErrorKind::Open,
+ true,
+ 0,
+ &all_close_edges[..1],
+ ),
+ (
+ DurabilityFailpoint::CloseAfterDrain,
+ ServiceSqliteErrorKind::Open,
+ true,
+ 0,
+ &all_close_edges[..2],
+ ),
+ (
+ DurabilityFailpoint::CloseBeforeCheckpoint,
+ ServiceSqliteErrorKind::Pragma,
+ false,
+ 6,
+ &[
+ DurabilityFailpoint::CloseBeforeDrain,
+ DurabilityFailpoint::CloseAfterDrain,
+ DurabilityFailpoint::CloseBeforeCheckpoint,
+ DurabilityFailpoint::CloseBeforeConnectionClose,
+ DurabilityFailpoint::CloseAfterConnectionClose,
+ DurabilityFailpoint::CloseBeforeAuthorityRelease,
+ DurabilityFailpoint::CloseAfterAuthorityRelease,
+ ][..],
+ ),
+ (
+ DurabilityFailpoint::CloseAfterCheckpoint,
+ ServiceSqliteErrorKind::Pragma,
+ false,
+ 6,
+ &all_close_edges[..],
+ ),
+ (
+ DurabilityFailpoint::CloseBeforeConnectionClose,
+ ServiceSqliteErrorKind::Open,
+ false,
+ 6,
+ &all_close_edges[..],
+ ),
+ (
+ DurabilityFailpoint::CloseAfterConnectionClose,
+ ServiceSqliteErrorKind::Open,
+ false,
+ 6,
+ &all_close_edges[..],
+ ),
+ (
+ DurabilityFailpoint::CloseBeforeAuthorityRelease,
+ ServiceSqliteErrorKind::Authority,
+ true,
+ 6,
+ &all_close_edges[..7],
+ ),
+ (
+ DurabilityFailpoint::CloseAfterAuthorityRelease,
+ ServiceSqliteErrorKind::Authority,
+ false,
+ 6,
+ &all_close_edges[..],
+ ),
+ ] {
+ let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ host.arm_durability_failpoint(point);
+ let error = host.close().await.expect_err("injected close edge");
+ assert_eq!(error.kind(), expected);
+ assert!(host.failpoints.fired());
+ assert_eq!(host.failpoints.reached(), expected_reached);
+ assert_eq!(
+ host.pool.close_phase(),
+ expected_phase,
+ "close failpoint must fire in its named driver phase"
+ );
+ if retryable {
+ host.close().await.expect("one-shot close edge resumes");
+ } else {
+ assert_eq!(
+ host.close()
+ .await
+ .expect_err("terminal close result is cached")
+ .kind(),
+ expected
+ );
+ }
+ }
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
async fn cancellation_quarantines_connection_and_pool_recovers() {
let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
let host = Arc::new(host);
diff --git a/crates/service_sqlite/src/failpoint.rs b/crates/service_sqlite/src/failpoint.rs
@@ -0,0 +1,255 @@
+//! Private per-operation durability failpoints used by deterministic tests.
+
+use core::fmt;
+
+#[cfg(test)]
+use std::sync::{Arc, Mutex};
+
+/// Closed inventory of durability edges exercised by the crash-boundary harness.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) enum DurabilityFailpoint {
+ InitializeBeforeCreate,
+ InitializeAfterCreate,
+ InitializeBeforeReservationDirectorySync,
+ InitializeAfterReservationDirectorySync,
+ InitializeBeforeFileSync,
+ InitializeAfterFileSync,
+ InitializeBeforeCommitDirectorySync,
+ InitializeAfterCommitDirectorySync,
+ TransactionBeforeBegin,
+ TransactionAfterBegin,
+ TransactionBeforeCommit,
+ TransactionAfterCommit,
+ BackupBeforeCreate,
+ BackupAfterCreate,
+ BackupBeforeCopy,
+ BackupAfterCopy,
+ BackupBeforeFileSync,
+ BackupAfterFileSync,
+ BackupBeforeDirectorySync,
+ BackupAfterDirectorySync,
+ MarkerBeforeCreate,
+ MarkerAfterCreate,
+ MarkerBeforeFileSync,
+ MarkerAfterFileSync,
+ MarkerBeforeDirectorySync,
+ MarkerAfterDirectorySync,
+ MarkerAdvanceBeforeWriteAndFileSync,
+ MarkerAdvanceAfterWriteAndFileSync,
+ MarkerAdvanceBeforeReplace,
+ MarkerAdvanceAfterReplace,
+ MarkerAdvanceBeforeDirectorySync,
+ MarkerAdvanceAfterDirectorySync,
+ RestoreBeforeRetainLiveRename,
+ RestoreAfterRetainLiveRename,
+ RestoreBeforeRetainLiveSync,
+ RestoreAfterRetainLiveSync,
+ RestoreBeforeInstallStageRename,
+ RestoreAfterInstallStageRename,
+ RestoreBeforeInstallStageSync,
+ RestoreAfterInstallStageSync,
+ CloseBeforeDrain,
+ CloseAfterDrain,
+ CloseBeforeCheckpoint,
+ CloseAfterCheckpoint,
+ CloseBeforeConnectionClose,
+ CloseAfterConnectionClose,
+ CloseBeforeAuthorityRelease,
+ CloseAfterAuthorityRelease,
+}
+
+impl DurabilityFailpoint {
+ #[cfg(test)]
+ pub(crate) const ALL: [Self; 48] = [
+ Self::InitializeBeforeCreate,
+ Self::InitializeAfterCreate,
+ Self::InitializeBeforeReservationDirectorySync,
+ Self::InitializeAfterReservationDirectorySync,
+ Self::InitializeBeforeFileSync,
+ Self::InitializeAfterFileSync,
+ Self::InitializeBeforeCommitDirectorySync,
+ Self::InitializeAfterCommitDirectorySync,
+ Self::TransactionBeforeBegin,
+ Self::TransactionAfterBegin,
+ Self::TransactionBeforeCommit,
+ Self::TransactionAfterCommit,
+ Self::BackupBeforeCreate,
+ Self::BackupAfterCreate,
+ Self::BackupBeforeCopy,
+ Self::BackupAfterCopy,
+ Self::BackupBeforeFileSync,
+ Self::BackupAfterFileSync,
+ Self::BackupBeforeDirectorySync,
+ Self::BackupAfterDirectorySync,
+ Self::MarkerBeforeCreate,
+ Self::MarkerAfterCreate,
+ Self::MarkerBeforeFileSync,
+ Self::MarkerAfterFileSync,
+ Self::MarkerBeforeDirectorySync,
+ Self::MarkerAfterDirectorySync,
+ Self::MarkerAdvanceBeforeWriteAndFileSync,
+ Self::MarkerAdvanceAfterWriteAndFileSync,
+ Self::MarkerAdvanceBeforeReplace,
+ Self::MarkerAdvanceAfterReplace,
+ Self::MarkerAdvanceBeforeDirectorySync,
+ Self::MarkerAdvanceAfterDirectorySync,
+ Self::RestoreBeforeRetainLiveRename,
+ Self::RestoreAfterRetainLiveRename,
+ Self::RestoreBeforeRetainLiveSync,
+ Self::RestoreAfterRetainLiveSync,
+ Self::RestoreBeforeInstallStageRename,
+ Self::RestoreAfterInstallStageRename,
+ Self::RestoreBeforeInstallStageSync,
+ Self::RestoreAfterInstallStageSync,
+ Self::CloseBeforeDrain,
+ Self::CloseAfterDrain,
+ Self::CloseBeforeCheckpoint,
+ Self::CloseAfterCheckpoint,
+ Self::CloseBeforeConnectionClose,
+ Self::CloseAfterConnectionClose,
+ Self::CloseBeforeAuthorityRelease,
+ Self::CloseAfterAuthorityRelease,
+ ];
+}
+
+/// Disabled in ordinary builds; tests may arm one edge on one owned controller.
+#[derive(Clone, Default)]
+pub(crate) struct DurabilityFailpoints {
+ #[cfg(test)]
+ state: Arc<Mutex<TestState>>,
+}
+
+#[cfg(test)]
+#[derive(Default)]
+struct TestState {
+ armed: Option<DurabilityFailpoint>,
+ fired: bool,
+ reached: Vec<DurabilityFailpoint>,
+ observations: Vec<(DurabilityFailpoint, u8)>,
+}
+
+impl DurabilityFailpoints {
+ pub(crate) fn hit(&self, point: DurabilityFailpoint) -> Result<(), DurabilityFailpointError> {
+ #[cfg(test)]
+ {
+ let mut state = self.state.lock().map_err(|_| DurabilityFailpointError)?;
+ if state.reached.len() < DurabilityFailpoint::ALL.len() {
+ state.reached.push(point);
+ }
+ if state.armed == Some(point) && !state.fired {
+ state.fired = true;
+ return Err(DurabilityFailpointError);
+ }
+ }
+ #[cfg(not(test))]
+ let _ = point;
+ Ok(())
+ }
+
+ #[cfg(test)]
+ pub(crate) fn armed(point: DurabilityFailpoint) -> Self {
+ Self {
+ state: Arc::new(Mutex::new(TestState {
+ armed: Some(point),
+ fired: false,
+ reached: Vec::new(),
+ observations: Vec::new(),
+ })),
+ }
+ }
+
+ #[cfg(test)]
+ pub(crate) fn arm(&self, point: DurabilityFailpoint) {
+ let mut state = self.state.lock().expect("durability failpoint state");
+ state.armed = Some(point);
+ state.fired = false;
+ state.reached.clear();
+ state.observations.clear();
+ }
+
+ #[cfg(test)]
+ pub(crate) fn disarm(&self) {
+ let mut state = self.state.lock().expect("durability failpoint state");
+ state.armed = None;
+ }
+
+ #[cfg(test)]
+ pub(crate) fn fired(&self) -> bool {
+ self.state.lock().is_ok_and(|state| state.fired)
+ }
+
+ #[cfg(test)]
+ pub(crate) fn reached(&self) -> Vec<DurabilityFailpoint> {
+ self.state
+ .lock()
+ .map_or_else(|_| Vec::new(), |state| state.reached.clone())
+ }
+
+ #[cfg(test)]
+ pub(crate) fn observe(&self, point: DurabilityFailpoint, state_value: u8) {
+ let mut state = self.state.lock().expect("durability failpoint state");
+ state.observations.push((point, state_value));
+ }
+
+ #[cfg(test)]
+ pub(crate) fn observation(&self, point: DurabilityFailpoint) -> Option<u8> {
+ self.state.lock().ok().and_then(|state| {
+ state
+ .observations
+ .iter()
+ .find_map(|(observed, value)| (*observed == point).then_some(*value))
+ })
+ }
+}
+
+impl fmt::Debug for DurabilityFailpoints {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("DurabilityFailpoints([redacted])")
+ }
+}
+
+/// Source-free injected failure; subsystem adapters retain their stable error kind.
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+pub(crate) struct DurabilityFailpointError;
+
+impl fmt::Display for DurabilityFailpointError {
+ fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
+ formatter.write_str("injected durability boundary failure")
+ }
+}
+
+impl std::error::Error for DurabilityFailpointError {}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[test]
+ fn every_closed_point_fires_once_on_its_owned_plan() {
+ for point in DurabilityFailpoint::ALL {
+ let plan = DurabilityFailpoints::armed(point);
+ assert_eq!(plan.hit(point), Err(DurabilityFailpointError));
+ assert_eq!(plan.hit(point), Ok(()));
+ assert!(plan.fired());
+ assert_eq!(plan.reached(), [point, point]);
+ }
+ }
+
+ #[test]
+ fn plans_are_instance_local_and_disabled_plan_never_fails() {
+ let first = DurabilityFailpoints::armed(DurabilityFailpoint::TransactionBeforeCommit);
+ let second = DurabilityFailpoints::armed(DurabilityFailpoint::BackupBeforeCopy);
+ assert_eq!(first.hit(DurabilityFailpoint::BackupBeforeCopy), Ok(()));
+ assert!(!first.fired());
+ assert_eq!(
+ second.hit(DurabilityFailpoint::BackupBeforeCopy),
+ Err(DurabilityFailpointError)
+ );
+ assert!(!first.fired());
+ assert!(second.fired());
+ assert_eq!(
+ DurabilityFailpoints::default().hit(DurabilityFailpoint::CloseBeforeDrain),
+ Ok(())
+ );
+ }
+}
diff --git a/crates/service_sqlite/src/initialize.rs b/crates/service_sqlite/src/initialize.rs
@@ -53,6 +53,7 @@ where
#[cfg(any(target_os = "linux", target_os = "macos"))]
{
+ let failpoints = crate::failpoint::DurabilityFailpoints::default();
initialize_with_ops(
paths,
authority,
@@ -60,6 +61,7 @@ where
schema_catalog,
initialize_schema,
&SystemInitializationOperations,
+ &failpoints,
)
.await
}
@@ -91,6 +93,8 @@ enum InitializationFailureKind {
DirectorySyncFailed,
#[cfg(any(target_os = "linux", target_os = "macos"))]
CleanupFailed,
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ InjectedFailure,
}
struct InitializationCause {
@@ -154,6 +158,10 @@ impl fmt::Display for InitializationCause {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
InitializationFailureKind::CleanupFailed => "SQLite initialization cleanup failed",
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ InitializationFailureKind::InjectedFailure => {
+ "SQLite initialization durability boundary failed"
+ }
})
}
}
@@ -335,11 +343,31 @@ mod supported {
Ok(())
}
- fn commit(&mut self, canonical_path: &std::path::Path) -> Result<(), InitializationCause> {
+ fn commit(
+ &mut self,
+ canonical_path: &std::path::Path,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ ) -> Result<(), InitializationCause> {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeBeforeFileSync,
+ )?;
self.operations.sync_database(&self.database)?;
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeAfterFileSync,
+ )?;
self.validate()?;
self.validate_canonical_path(canonical_path)?;
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeBeforeCommitDirectorySync,
+ )?;
self.operations.sync_directory(self.directory)?;
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeAfterCommitDirectorySync,
+ )?;
self.committed = true;
Ok(())
}
@@ -452,6 +480,7 @@ mod supported {
schema_catalog: &SchemaCatalog,
initialize_schema: F,
operations: &O,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<WriterAuthority, ServiceSqliteError>
where
F: FnOnce(PathBuf) -> Fut,
@@ -459,11 +488,34 @@ mod supported {
E: Error + Send + Sync + 'static,
O: InitializationOperations,
{
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeBeforeCreate,
+ )
+ .map_err(initialization_error)?;
let mut pending = PendingDatabase::create(authority.directory(), operations)
.map_err(initialization_error)?;
+ if let Err(error) = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeAfterCreate,
+ ) {
+ return fail_with_rollback(pending, error).await;
+ }
+ if let Err(error) = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeBeforeReservationDirectorySync,
+ ) {
+ return fail_with_rollback(pending, error).await;
+ }
if let Err(error) = operations.sync_directory(authority.directory()) {
return fail_with_rollback(pending, error).await;
}
+ if let Err(error) = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::InitializeAfterReservationDirectorySync,
+ ) {
+ return fail_with_rollback(pending, error).await;
+ }
if let Err(error) = pending.validate_canonical_path(paths.state_database()) {
return fail_with_rollback(pending, error).await;
}
@@ -520,13 +572,22 @@ mod supported {
if let Err(error) = pending.validate_canonical_path(paths.state_database()) {
return fail_with_rollback(pending, error).await;
}
- if let Err(error) = pending.commit(paths.state_database()) {
+ if let Err(error) = pending.commit(paths.state_database(), failpoints) {
return fail_with_rollback(pending, error).await;
}
drop(pending);
Ok(authority)
}
+ fn hit(
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ point: crate::failpoint::DurabilityFailpoint,
+ ) -> Result<(), InitializationCause> {
+ failpoints.hit(point).map_err(|source| {
+ InitializationCause::with_source(InitializationFailureKind::InjectedFailure, source)
+ })
+ }
+
#[cfg(test)]
mod tests {
use std::{
@@ -1042,6 +1103,7 @@ mod supported {
&schema_catalog,
|_| ready(Err::<(), _>(CallbackFailure)),
&operations,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
.await
} else {
@@ -1052,6 +1114,7 @@ mod supported {
&schema_catalog,
|_| ready(Ok::<(), CallbackFailure>(())),
&operations,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
.await
};
@@ -1094,6 +1157,101 @@ mod supported {
}
}
}
+
+ #[tokio::test(flavor = "current_thread")]
+ async fn every_initialization_durability_edge_fails_once_and_rolls_back() {
+ use crate::failpoint::{DurabilityFailpoint, DurabilityFailpoints};
+
+ for (index, (point, expected_events)) in [
+ (DurabilityFailpoint::InitializeBeforeCreate, &[][..]),
+ (
+ DurabilityFailpoint::InitializeAfterCreate,
+ &["unlink_database", "sync_directory"][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeBeforeReservationDirectorySync,
+ &["unlink_database", "sync_directory"][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeAfterReservationDirectorySync,
+ &["sync_directory", "unlink_database", "sync_directory"][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeBeforeFileSync,
+ &["sync_directory", "unlink_database", "sync_directory"][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeAfterFileSync,
+ &[
+ "sync_directory",
+ "sync_database",
+ "unlink_database",
+ "sync_directory",
+ ][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeBeforeCommitDirectorySync,
+ &[
+ "sync_directory",
+ "sync_database",
+ "unlink_database",
+ "sync_directory",
+ ][..],
+ ),
+ (
+ DurabilityFailpoint::InitializeAfterCommitDirectorySync,
+ &[
+ "sync_directory",
+ "sync_database",
+ "sync_directory",
+ "unlink_database",
+ "sync_directory",
+ ][..],
+ ),
+ ]
+ .into_iter()
+ .enumerate()
+ {
+ let root = tempfile::tempdir().expect("root");
+ let paths = paths(root.path(), &format!("failpoint-{index}"));
+ prepare(&paths);
+ let metadata = metadata(&paths);
+ let schema_catalog = base_schema_catalog();
+ let authority = WriterAuthority::acquire(&paths, OpenMode::Initialize)
+ .expect("authority acquisition")
+ .expect("initialize authority");
+ let failpoints = DurabilityFailpoints::armed(point);
+ let operations = RecordingOperations::default();
+ let error = initialize_with_ops(
+ &paths,
+ authority,
+ &metadata,
+ &schema_catalog,
+ |_| ready(Ok::<(), CallbackFailure>(())),
+ &operations,
+ &failpoints,
+ )
+ .await
+ .expect_err("durability edge must fail");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Create);
+ assert!(failpoints.fired());
+ assert_eq!(
+ failpoints.reached().last(),
+ Some(&point),
+ "named edge must be the last reached boundary"
+ );
+ assert_eq!(
+ operations.events.borrow().as_slice(),
+ expected_events,
+ "named before/after edge must bracket the expected sync operation"
+ );
+ assert!(!paths.state_database().exists());
+ let mut recovered = WriterAuthority::acquire(&paths, OpenMode::Initialize)
+ .expect("reacquire after rollback")
+ .expect("initialize authority");
+ recovered.release().expect("release recovered authority");
+ }
+ }
}
}
diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs
@@ -7,6 +7,8 @@ mod backup;
mod config;
mod connection;
mod error;
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+mod failpoint;
mod initialize;
mod integrity;
mod metadata;
diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs
@@ -254,6 +254,7 @@ enum PrivateCloseDriver {
future: BoxFuture<'static, Result<(), ServiceSqliteError>>,
authority_error: Option<ServiceSqliteError>,
checkpoint_error: Option<ServiceSqliteError>,
+ connection_close_error: Option<ServiceSqliteError>,
},
Complete(Option<ServiceSqliteErrorKind>),
}
@@ -406,18 +407,48 @@ impl PrivateConnectionPool {
/// The inner result is terminal and may be cached by the host.
pub(crate) async fn close_explicit(
&self,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<Result<(), ServiceSqliteError>, ServiceSqliteError> {
+ let mut authority_error = self.validate().err();
+ let before_drain = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeDrain)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ });
+ authority_error = authority_error.or_else(|| self.validate().err());
+ if authority_error.is_none() {
+ before_drain?;
+ }
self.pool.close().await;
- let terminal = if self.mode.requires_writer_authority() {
- self.drive_writable_close().await
- } else {
- self.validate()
+ let after_drain = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseAfterDrain)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ });
+ authority_error = authority_error.or_else(|| self.validate().err());
+ let terminal = match authority_error {
+ Some(authority) => Err(authority),
+ None => {
+ after_drain?;
+ if self.mode.requires_writer_authority() {
+ self.drive_writable_close(failpoints).await
+ } else {
+ self.validate()
+ }
+ }
};
- self.release_resources()?;
- Ok(terminal)
+ let release_error = self.release_resources(failpoints)?;
+ Ok(match (terminal, release_error) {
+ (Err(error), _) if error.kind() == ServiceSqliteErrorKind::Authority => Err(error),
+ (_, Some(release_error)) => Err(release_error),
+ (terminal, None) => terminal,
+ })
}
- async fn drive_writable_close(&self) -> Result<(), ServiceSqliteError> {
+ async fn drive_writable_close(
+ &self,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ ) -> Result<(), ServiceSqliteError> {
let mut driver = self.close_driver.lock().await;
loop {
match &mut *driver {
@@ -479,6 +510,19 @@ impl PrivateConnectionPool {
#[cfg(test)]
self.close_phase
.store(TEST_CLOSE_PHASE_CHECKPOINT, Ordering::Release);
+ checkpoint_error = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeCheckpoint)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Pragma,
+ source,
+ )
+ })
+ .err();
+ authority_error =
+ authority_error.or_else(|| self.validate_writer_authority().err());
+ }
+ if checkpoint_error.is_none() && authority_error.is_none() {
checkpoint_error =
sqlx::query_as::<_, (i64, i64, i64)>("PRAGMA wal_checkpoint(TRUNCATE)")
.fetch_one(&mut *connection)
@@ -497,10 +541,28 @@ impl PrivateConnectionPool {
}
})
.err();
+ if checkpoint_error.is_none() {
+ checkpoint_error = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseAfterCheckpoint)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(
+ ServiceSqliteErrorKind::Pragma,
+ source,
+ )
+ })
+ .err();
+ }
authority_error =
authority_error.or_else(|| self.validate_writer_authority().err());
}
+ let connection_close_error = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeConnectionClose)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ })
+ .err();
+
let connected = core::mem::replace(&mut *driver, PrivateCloseDriver::Pending);
let PrivateCloseDriver::Connected(connection) = connected else {
unreachable!("close driver retains its connected phase")
@@ -517,20 +579,30 @@ impl PrivateConnectionPool {
future: close,
authority_error,
checkpoint_error,
+ connection_close_error,
};
}
PrivateCloseDriver::Closing {
future,
authority_error,
checkpoint_error,
+ connection_close_error,
} => {
let close_error = future.as_mut().await.err();
+ let injected_after_close = failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseAfterConnectionClose)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Open, source)
+ })
+ .err();
let authority_error = authority_error
.take()
.or_else(|| self.validate_writer_authority().err());
let error = authority_error
.or_else(|| checkpoint_error.take())
- .or(close_error);
+ .or_else(|| connection_close_error.take())
+ .or(close_error)
+ .or(injected_after_close);
#[cfg(test)]
self.close_phase.store(6, Ordering::Release);
*driver =
@@ -544,7 +616,15 @@ impl PrivateConnectionPool {
}
}
- fn release_resources(&self) -> Result<(), ServiceSqliteError> {
+ fn release_resources(
+ &self,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ ) -> Result<Option<ServiceSqliteError>, ServiceSqliteError> {
+ failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseBeforeAuthorityRelease)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source)
+ })?;
let mut resources = self.resources.lock().map_err(|_| {
connection_error(
ServiceSqliteErrorKind::Authority,
@@ -570,7 +650,12 @@ impl PrivateConnectionPool {
resources.inspection_guard.take();
}
}
- Ok(())
+ Ok(failpoints
+ .hit(crate::failpoint::DurabilityFailpoint::CloseAfterAuthorityRelease)
+ .map_err(|source| {
+ ServiceSqliteError::with_source(ServiceSqliteErrorKind::Authority, source)
+ })
+ .err())
}
}
@@ -2906,12 +2991,16 @@ mod tests {
future: close,
authority_error: None,
checkpoint_error: None,
+ connection_close_error: None,
};
}
let pool = Arc::new(pool);
let close_task = tokio::spawn({
let pool = Arc::clone(&pool);
- async move { pool.close_explicit().await }
+ async move {
+ pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default())
+ .await
+ }
});
entered.notified().await;
close_task.abort();
@@ -2928,7 +3017,7 @@ mod tests {
));
release.notify_one();
- pool.close_explicit()
+ pool.close_explicit(&crate::failpoint::DurabilityFailpoints::default())
.await
.expect("authority release is proven")
.expect("retained connection closes explicitly");
diff --git a/crates/service_sqlite/src/restore/finalize.rs b/crates/service_sqlite/src/restore/finalize.rs
@@ -78,11 +78,17 @@ pub async fn finalize_staged_restore(
) -> Result<(), ServiceSqliteError> {
#[cfg(any(target_os = "linux", target_os = "macos"))]
{
+ let failpoints = crate::failpoint::DurabilityFailpoints::default();
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)
+ finalize_native(
+ native,
+ &cancellation,
+ &SystemFinalizeOperations,
+ &failpoints,
+ )
})
.await
.map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?;
@@ -102,6 +108,7 @@ fn finalize_native(
staged: NativeStagedServiceRestore,
cancellation: &AtomicU8,
operations: &dyn FinalizeOperations,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<(), ServiceSqliteError> {
staged.validate()?;
check_cancel(cancellation)?;
@@ -160,18 +167,20 @@ fn finalize_native(
on_durable,
)
} else {
- RestoreMarkerBinding::create_with_durable_callback(
+ RestoreMarkerBinding::create_with_durable_callback_and_failpoints(
staged.paths(),
staged.authority(),
&marker,
+ failpoints,
on_durable,
)
};
#[cfg(not(test))]
- let marker_result = RestoreMarkerBinding::create_with_durable_callback(
+ let marker_result = RestoreMarkerBinding::create_with_durable_callback_and_failpoints(
staged.paths(),
staged.authority(),
&marker,
+ failpoints,
on_durable,
);
let mut marker = marker_result?;
@@ -180,36 +189,44 @@ fn finalize_native(
&staged,
operations,
RenameStep::RetainLive,
- LIVE_FILE_NAME,
- BACKUP_FILE_NAME,
- &live,
- live_artifact,
+ RenameArtifact {
+ source_name: LIVE_FILE_NAME,
+ destination_name: BACKUP_FILE_NAME,
+ held: &live,
+ expected: live_artifact,
+ },
+ failpoints,
)?;
authority_checked(&staged, || {
operations.after_directory_sync(staged.directory(), RenameStep::RetainLive)
})??;
- marker = marker.advance(
+ marker = marker.advance_with_failpoints(
staged.paths(),
staged.authority(),
RestoreRecoveryPhase::LiveRetained,
+ failpoints,
)?;
rename_and_sync(
&staged,
operations,
RenameStep::InstallStage,
- STAGED_FILE_NAME,
- LIVE_FILE_NAME,
- staged.staged_file(),
- staged.artifact(),
+ RenameArtifact {
+ source_name: STAGED_FILE_NAME,
+ destination_name: LIVE_FILE_NAME,
+ held: staged.staged_file(),
+ expected: staged.artifact(),
+ },
+ failpoints,
)?;
authority_checked(&staged, || {
operations.after_directory_sync(staged.directory(), RenameStep::InstallStage)
})??;
- marker = marker.advance(
+ marker = marker.advance_with_failpoints(
staged.paths(),
staged.authority(),
RestoreRecoveryPhase::ReplacementInstalled,
+ failpoints,
)?;
if marker.marker().phase() != RestoreRecoveryPhase::ReplacementInstalled {
return Err(finalize_error(FinalizeFailureKind::Marker));
@@ -225,22 +242,58 @@ fn rename_and_sync(
staged: &NativeStagedServiceRestore,
operations: &dyn FinalizeOperations,
step: RenameStep,
- source_name: &str,
- destination_name: &str,
- held: &File,
- expected: RestoreArtifactExpectation,
+ artifact: RenameArtifact<'_>,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<(), ServiceSqliteError> {
authority_checked(staged, || {
- verify_named_artifact(staged.directory(), source_name, held, expected, None)
+ verify_named_artifact(
+ staged.directory(),
+ artifact.source_name,
+ artifact.held,
+ artifact.expected,
+ None,
+ )
+ })??;
+ authority_checked(staged, || {
+ hit(
+ failpoints,
+ step.before_rename_failpoint(),
+ step.rename_failure(),
+ )
})??;
let rename = authority_checked(staged, || {
operations
- .rename(staged.directory(), source_name, destination_name, step)
+ .rename(
+ staged.directory(),
+ artifact.source_name,
+ artifact.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)
+ hit(
+ failpoints,
+ step.after_rename_failpoint(),
+ step.rename_failure(),
+ )
+ })??;
+ authority_checked(staged, || {
+ verify_named_artifact(
+ staged.directory(),
+ artifact.destination_name,
+ artifact.held,
+ artifact.expected,
+ None,
+ )
+ })??;
+ authority_checked(staged, || {
+ hit(
+ failpoints,
+ step.before_sync_failpoint(),
+ step.sync_failure(),
+ )
})??;
let sync = authority_checked(staged, || {
operations
@@ -249,7 +302,16 @@ fn rename_and_sync(
})?;
sync?;
authority_checked(staged, || {
- verify_named_artifact(staged.directory(), destination_name, held, expected, None)
+ hit(failpoints, step.after_sync_failpoint(), step.sync_failure())
+ })??;
+ authority_checked(staged, || {
+ verify_named_artifact(
+ staged.directory(),
+ artifact.destination_name,
+ artifact.held,
+ artifact.expected,
+ None,
+ )
})??;
Ok(())
}
@@ -431,6 +493,14 @@ enum RenameStep {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
+struct RenameArtifact<'a> {
+ source_name: &'static str,
+ destination_name: &'static str,
+ held: &'a File,
+ expected: RestoreArtifactExpectation,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
impl RenameStep {
const fn rename_failure(self) -> FinalizeFailureKind {
match self {
@@ -445,6 +515,44 @@ impl RenameStep {
Self::InstallStage => FinalizeFailureKind::SyncInstalled,
}
}
+
+ const fn before_rename_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
+ match self {
+ Self::RetainLive => {
+ crate::failpoint::DurabilityFailpoint::RestoreBeforeRetainLiveRename
+ }
+ Self::InstallStage => {
+ crate::failpoint::DurabilityFailpoint::RestoreBeforeInstallStageRename
+ }
+ }
+ }
+
+ const fn after_rename_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
+ match self {
+ Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreAfterRetainLiveRename,
+ Self::InstallStage => {
+ crate::failpoint::DurabilityFailpoint::RestoreAfterInstallStageRename
+ }
+ }
+ }
+
+ const fn before_sync_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
+ match self {
+ Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreBeforeRetainLiveSync,
+ Self::InstallStage => {
+ crate::failpoint::DurabilityFailpoint::RestoreBeforeInstallStageSync
+ }
+ }
+ }
+
+ const fn after_sync_failpoint(self) -> crate::failpoint::DurabilityFailpoint {
+ match self {
+ Self::RetainLive => crate::failpoint::DurabilityFailpoint::RestoreAfterRetainLiveSync,
+ Self::InstallStage => {
+ crate::failpoint::DurabilityFailpoint::RestoreAfterInstallStageSync
+ }
+ }
+ }
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -584,6 +692,7 @@ pub(crate) async fn test_finalize_with_failure(
native,
&AtomicU8::new(CANCELLABLE),
&FailingFinalizeOperations,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
})
.await
@@ -593,6 +702,24 @@ pub(crate) async fn test_finalize_with_failure(
}
#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) async fn test_finalize_with_failpoint(
+ staged: StagedServiceRestore,
+ failpoints: crate::failpoint::DurabilityFailpoints,
+) -> Result<(), ServiceSqliteError> {
+ let native = staged.into_native();
+ tokio::task::spawn_blocking(move || {
+ finalize_native(
+ native,
+ &AtomicU8::new(CANCELLABLE),
+ &SystemFinalizeOperations,
+ &failpoints,
+ )
+ })
+ .await
+ .map_err(|source| finalize_source(FinalizeFailureKind::Join, source))?
+}
+
+#[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);
@@ -717,3 +844,14 @@ fn finalize_source(
},
)
}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+fn hit(
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ point: crate::failpoint::DurabilityFailpoint,
+ kind: FinalizeFailureKind,
+) -> Result<(), ServiceSqliteError> {
+ failpoints
+ .hit(point)
+ .map_err(|source| finalize_source(kind, source))
+}
diff --git a/crates/service_sqlite/src/restore/marker.rs b/crates/service_sqlite/src/restore/marker.rs
@@ -605,11 +605,29 @@ mod store {
marker: &RestoreRecoveryMarker,
on_durable: impl FnOnce(),
) -> Result<Self, ServiceSqliteError> {
+ let failpoints = crate::failpoint::DurabilityFailpoints::default();
+ Self::create_with_durable_callback_and_failpoints(
+ paths,
+ authority,
+ marker,
+ &failpoints,
+ on_durable,
+ )
+ }
+
+ pub(crate) fn create_with_durable_callback_and_failpoints(
+ paths: &ServiceSqlitePaths,
+ authority: &WriterAuthority,
+ marker: &RestoreRecoveryMarker,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ on_durable: impl FnOnce(),
+ ) -> Result<Self, ServiceSqliteError> {
Self::create_with_operations(
paths,
authority,
marker,
&SystemStoreOperations,
+ failpoints,
on_durable,
)
}
@@ -626,6 +644,7 @@ mod store {
authority,
marker,
&AuthorityDriftAfterSyncStoreOperations,
+ &crate::failpoint::DurabilityFailpoints::default(),
on_durable,
)
}
@@ -635,6 +654,7 @@ mod store {
authority: &WriterAuthority,
marker: &RestoreRecoveryMarker,
operations: &dyn StoreOperations,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
on_durable: impl FnOnce(),
) -> Result<Self, ServiceSqliteError> {
authority.validate_for(paths)?;
@@ -663,18 +683,41 @@ mod store {
require_absent(&directory, MARKER_NEXT_FILE_NAME)
})?
.map_err(restore_store)?;
+ authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerBeforeCreate,
+ )
+ })?
+ .map_err(restore_store)?;
let (marker_file, marker_identity) = authority_checked(authority, paths, || {
let file = create_marker_file(&directory, MARKER_FILE_NAME)?;
let identity = file_identity(&file)?;
Ok::<_, StoreFailure>((file, identity))
})?
.map_err(restore_store)?;
+ if let Err(cause) = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAfterCreate,
+ ) {
+ cleanup_with_authority(
+ authority,
+ paths,
+ &directory,
+ MARKER_FILE_NAME,
+ marker_identity,
+ )?;
+ return Err(restore_store(cause));
+ }
let write_result = authority_checked(authority, paths, || {
write_and_sync(
&marker_file,
marker.canonical_bytes(),
&directory,
operations,
+ failpoints,
+ Some(crate::failpoint::DurabilityFailpoint::MarkerBeforeFileSync),
+ Some(crate::failpoint::DurabilityFailpoint::MarkerAfterFileSync),
)
})?;
if let Err(cause) = write_result {
@@ -688,6 +731,19 @@ mod store {
return Err(restore_store(cause));
}
authority.validate_for(paths)?;
+ if let Err(cause) = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerBeforeDirectorySync,
+ ) {
+ cleanup_with_authority(
+ authority,
+ paths,
+ &directory,
+ MARKER_FILE_NAME,
+ marker_identity,
+ )?;
+ return Err(restore_store(cause));
+ }
if operations.sync_directory(&directory).is_err() {
authority.validate_for(paths)?;
cleanup_with_authority(
@@ -703,7 +759,12 @@ mod store {
// this point. The caller must transfer ownership of every bound
// artifact before any subsequent fallible validation.
on_durable();
+ let after_directory_sync = hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAfterDirectorySync,
+ );
authority.validate_for(paths)?;
+ after_directory_sync.map_err(restore_store)?;
let binding = Self {
directory,
directory_identity,
@@ -822,12 +883,24 @@ mod store {
authority: &WriterAuthority,
next: RestoreRecoveryPhase,
) -> Result<Self, ServiceSqliteError> {
+ let failpoints = crate::failpoint::DurabilityFailpoints::default();
+ self.advance_with_failpoints(paths, authority, next, &failpoints)
+ }
+
+ pub(crate) fn advance_with_failpoints(
+ self,
+ paths: &ServiceSqlitePaths,
+ authority: &WriterAuthority,
+ next: RestoreRecoveryPhase,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ ) -> Result<Self, ServiceSqliteError> {
self.advance_with_operations(
paths,
authority,
next,
&SystemStoreOperations,
ServiceSqliteErrorKind::Restore,
+ failpoints,
)
}
@@ -843,6 +916,7 @@ mod store {
next,
&SystemStoreOperations,
ServiceSqliteErrorKind::Recovery,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
}
@@ -853,6 +927,7 @@ mod store {
next: RestoreRecoveryPhase,
operations: &dyn StoreOperations,
operation_kind: ServiceSqliteErrorKind,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
) -> Result<Self, ServiceSqliteError> {
authority_checked(authority, paths, || {
self.validate_inner(paths, true, operation_kind)
@@ -879,12 +954,31 @@ mod store {
Ok::<_, StoreFailure>((file, identity))
})?
.map_err(|cause| operation_store(operation_kind, cause))?;
+ let before_write = authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceBeforeWriteAndFileSync,
+ )
+ })?;
+ if let Err(cause) = before_write {
+ cleanup_with_authority(
+ authority,
+ paths,
+ &self.directory,
+ MARKER_NEXT_FILE_NAME,
+ scratch_identity,
+ )?;
+ return Err(operation_store(operation_kind, cause));
+ }
let scratch_write = authority_checked(authority, paths, || {
write_and_sync(
&scratch,
next_marker.canonical_bytes(),
&self.directory,
operations,
+ failpoints,
+ None,
+ None,
)
})?;
if let Err(cause) = scratch_write {
@@ -898,6 +992,13 @@ mod store {
return Err(operation_store(operation_kind, cause));
}
authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceAfterWriteAndFileSync,
+ )
+ })?
+ .map_err(|cause| operation_store(operation_kind, cause))?;
+ authority_checked(authority, paths, || {
self.validate_inner(paths, false, operation_kind)
})??;
let scratch_matches = authority_checked(authority, paths, || {
@@ -918,6 +1019,22 @@ mod store {
)?;
return Err(operation_store(operation_kind, StoreFailure::Conflict));
}
+ let before_replace = authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceBeforeReplace,
+ )
+ })?;
+ if let Err(cause) = before_replace {
+ cleanup_with_authority(
+ authority,
+ paths,
+ &self.directory,
+ MARKER_NEXT_FILE_NAME,
+ scratch_identity,
+ )?;
+ return Err(operation_store(operation_kind, cause));
+ }
let replacement = authority_checked(authority, paths, || {
operations
.replace_marker(&self.directory)
@@ -933,6 +1050,20 @@ mod store {
)?;
return Err(operation_store(operation_kind, StoreFailure::Rename));
}
+ authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceAfterReplace,
+ )
+ })?
+ .map_err(|cause| operation_store(operation_kind, cause))?;
+ authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceBeforeDirectorySync,
+ )
+ })?
+ .map_err(|cause| operation_store(operation_kind, cause))?;
let parent_sync = authority_checked(authority, paths, || {
operations
.sync_directory(&self.directory)
@@ -941,6 +1072,13 @@ mod store {
if parent_sync.is_err() {
return Err(operation_store(operation_kind, StoreFailure::Sync));
}
+ authority_checked(authority, paths, || {
+ hit(
+ failpoints,
+ crate::failpoint::DurabilityFailpoint::MarkerAdvanceAfterDirectorySync,
+ )
+ })?
+ .map_err(|cause| operation_store(operation_kind, cause))?;
let (marker_file, marker_identity, reread) =
authority_checked(authority, paths, || {
let file = open_marker_file(&self.directory, MARKER_FILE_NAME)?;
@@ -979,6 +1117,7 @@ mod store {
next,
&FailingStoreOperations { failure },
ServiceSqliteErrorKind::Restore,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
}
@@ -996,6 +1135,7 @@ mod store {
next,
&FailingStoreOperations { failure },
ServiceSqliteErrorKind::Recovery,
+ &crate::failpoint::DurabilityFailpoints::default(),
)
}
@@ -1362,12 +1502,29 @@ mod store {
bytes: &[u8],
directory: &File,
operations: &dyn StoreOperations,
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ before_sync: Option<crate::failpoint::DurabilityFailpoint>,
+ after_sync: Option<crate::failpoint::DurabilityFailpoint>,
) -> Result<(), StoreFailure> {
let mut file = file.try_clone().map_err(|_| StoreFailure::Write)?;
file.write_all(bytes).map_err(|_| StoreFailure::Write)?;
+ if let Some(before_sync) = before_sync {
+ hit(failpoints, before_sync)?;
+ }
operations
.sync_file(&file, directory)
- .map_err(|_| StoreFailure::Sync)
+ .map_err(|_| StoreFailure::Sync)?;
+ if let Some(after_sync) = after_sync {
+ hit(failpoints, after_sync)?;
+ }
+ Ok(())
+ }
+
+ fn hit(
+ failpoints: &crate::failpoint::DurabilityFailpoints,
+ point: crate::failpoint::DurabilityFailpoint,
+ ) -> Result<(), StoreFailure> {
+ failpoints.hit(point).map_err(|_| StoreFailure::Injected)
}
fn read_marker(file: &File) -> Result<RestoreRecoveryMarker, StoreFailure> {
@@ -1429,6 +1586,7 @@ mod store {
Sync,
Rename,
Conflict,
+ Injected,
Contract(RestoreMarkerContractError),
}
@@ -1445,6 +1603,7 @@ mod store {
Self::Sync => "restore marker durability could not be proven",
Self::Rename => "restore marker replacement failed",
Self::Conflict => "restore marker binding changed",
+ Self::Injected => "restore marker durability boundary failed",
Self::Contract(error) => return error.fmt(formatter),
})
}
diff --git a/crates/service_sqlite/src/restore/stage.rs b/crates/service_sqlite/src/restore/stage.rs
@@ -1244,7 +1244,8 @@ mod tests {
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,
+ reset_test_controls as reset_finalize_controls, test_finalize_with_failpoint,
+ test_finalize_with_failure,
};
use crate::restore::{RestoreMarkerBinding, RestoreRecoveryMarker, RestoreRecoveryPhase};
use crate::{
@@ -2217,6 +2218,106 @@ mod tests {
}
#[tokio::test(flavor = "current_thread")]
+ async fn every_marker_and_restore_durability_edge_is_wired_once() {
+ use crate::failpoint::{DurabilityFailpoint, DurabilityFailpoints};
+
+ let _serial = STAGE_TEST_LOCK.lock().await;
+ for point in [
+ DurabilityFailpoint::MarkerBeforeCreate,
+ DurabilityFailpoint::MarkerAfterCreate,
+ DurabilityFailpoint::MarkerBeforeFileSync,
+ DurabilityFailpoint::MarkerAfterFileSync,
+ DurabilityFailpoint::MarkerBeforeDirectorySync,
+ DurabilityFailpoint::MarkerAfterDirectorySync,
+ DurabilityFailpoint::MarkerAdvanceBeforeWriteAndFileSync,
+ DurabilityFailpoint::MarkerAdvanceAfterWriteAndFileSync,
+ DurabilityFailpoint::MarkerAdvanceBeforeReplace,
+ DurabilityFailpoint::MarkerAdvanceAfterReplace,
+ DurabilityFailpoint::MarkerAdvanceBeforeDirectorySync,
+ DurabilityFailpoint::MarkerAdvanceAfterDirectorySync,
+ DurabilityFailpoint::RestoreBeforeRetainLiveRename,
+ DurabilityFailpoint::RestoreAfterRetainLiveRename,
+ DurabilityFailpoint::RestoreBeforeRetainLiveSync,
+ DurabilityFailpoint::RestoreAfterRetainLiveSync,
+ DurabilityFailpoint::RestoreBeforeInstallStageRename,
+ DurabilityFailpoint::RestoreAfterInstallStageRename,
+ DurabilityFailpoint::RestoreBeforeInstallStageSync,
+ DurabilityFailpoint::RestoreAfterInstallStageSync,
+ ] {
+ 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 failpoints = DurabilityFailpoints::armed(point);
+ let error = test_finalize_with_failpoint(staged, failpoints.clone())
+ .await
+ .expect_err("injected marker or restore edge");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Restore);
+ assert!(failpoints.fired(), "edge was not reached: {point:?}");
+ assert!(
+ failpoints.reached().contains(&point),
+ "named marker/rename edge must be observed"
+ );
+ assert!(
+ fixture.paths.state_database().exists()
+ || recovery_path(&fixture, BACKUP_FILE_NAME).exists(),
+ "live content remains represented by an exact governed artifact"
+ );
+ let topology = (
+ fixture.paths.state_database().exists(),
+ fixture.staged_path().exists(),
+ recovery_path(&fixture, BACKUP_FILE_NAME).exists(),
+ recovery_path(&fixture, MARKER_FILE_NAME).exists(),
+ );
+ match point {
+ DurabilityFailpoint::MarkerBeforeCreate
+ | DurabilityFailpoint::MarkerAfterCreate
+ | DurabilityFailpoint::MarkerBeforeFileSync
+ | DurabilityFailpoint::MarkerAfterFileSync
+ | DurabilityFailpoint::MarkerBeforeDirectorySync => {
+ assert_eq!(topology, (true, false, false, false));
+ }
+ DurabilityFailpoint::MarkerAfterDirectorySync
+ | DurabilityFailpoint::RestoreBeforeRetainLiveRename => {
+ assert_eq!(topology, (true, true, false, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::Prepared);
+ }
+ DurabilityFailpoint::RestoreAfterRetainLiveRename
+ | DurabilityFailpoint::RestoreBeforeRetainLiveSync
+ | DurabilityFailpoint::RestoreAfterRetainLiveSync
+ | DurabilityFailpoint::MarkerAdvanceBeforeWriteAndFileSync
+ | DurabilityFailpoint::MarkerAdvanceAfterWriteAndFileSync
+ | DurabilityFailpoint::MarkerAdvanceBeforeReplace => {
+ assert_eq!(topology, (false, true, true, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::Prepared);
+ }
+ DurabilityFailpoint::MarkerAdvanceAfterReplace
+ | DurabilityFailpoint::MarkerAdvanceBeforeDirectorySync
+ | DurabilityFailpoint::MarkerAdvanceAfterDirectorySync
+ | DurabilityFailpoint::RestoreBeforeInstallStageRename => {
+ assert_eq!(topology, (false, true, true, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::LiveRetained);
+ }
+ DurabilityFailpoint::RestoreAfterInstallStageRename
+ | DurabilityFailpoint::RestoreBeforeInstallStageSync
+ | DurabilityFailpoint::RestoreAfterInstallStageSync => {
+ assert_eq!(topology, (true, false, true, true));
+ assert_eq!(marker_phase(&fixture), RestoreRecoveryPhase::LiveRetained);
+ }
+ _ => unreachable!("complete marker/restore failpoint inventory"),
+ }
+ }
+ 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;
diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs
@@ -10,6 +10,7 @@ const BACKUP_VERIFY_SOURCE: &str = include_str!("../src/backup/verify.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");
+const FAILPOINT_SOURCE: &str = include_str!("../src/failpoint.rs");
const INITIALIZE_SOURCE: &str = include_str!("../src/initialize.rs");
const INTEGRITY_SOURCE: &str = include_str!("../src/integrity/mod.rs");
const INTEGRITY_CATALOG_SOURCE: &str = include_str!("../src/integrity/catalog.rs");
@@ -69,6 +70,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"config",
"connection",
"error",
+ "failpoint",
"initialize",
"integrity",
"metadata",
@@ -234,6 +236,19 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"measurement is advisory rather than a space reservation",
"performs no database open, pool operation, SQLite query, filesystem mutation, ambient time read, timer, task",
"Service configuration, threshold defaults, cache refresh, status persistence, admission wiring, and route projection remain consumer responsibilities",
+ "Durability fault injection is a private test-only mechanism",
+ "closed instance-scoped controller can arm exactly one named before/after boundary",
+ "returns one injected error the first time that boundary is reached; later hits are no-ops",
+ "complete inventory covers database initialization, runner-owned transaction begin and commit",
+ "online-backup creation/copy/synchronization",
+ "restore-marker creation and advancement",
+ "both restore rename/synchronization steps",
+ "explicit host drain/checkpoint/connection-close/authority-release",
+ "ordinary controller has zero behavior",
+ "no failpoint type or selector is exported from the crate root",
+ "no process-global failpoint state, environment or configuration selector, Cargo feature",
+ "hidden task, timer, panic, or process-exit behavior",
+ "process crashes and signals remain a separate qualification layer",
] {
assert!(
readme_words.contains(required),
@@ -403,6 +418,71 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
);
}
+ let failpoint_production = FAILPOINT_SOURCE
+ .split_once("#[cfg(test)]\nmod tests")
+ .map(|(production, _)| production)
+ .expect("failpoint source must keep tests separated");
+ for required in [
+ "pub(crate) enum DurabilityFailpoint",
+ "pub(crate) struct DurabilityFailpoints",
+ "InitializeBeforeCreate",
+ "InitializeAfterReservationDirectorySync",
+ "InitializeAfterCommitDirectorySync",
+ "TransactionBeforeBegin",
+ "TransactionAfterCommit",
+ "BackupBeforeCreate",
+ "BackupAfterDirectorySync",
+ "MarkerBeforeCreate",
+ "MarkerAfterDirectorySync",
+ "MarkerAdvanceBeforeWriteAndFileSync",
+ "MarkerAdvanceAfterDirectorySync",
+ "RestoreBeforeRetainLiveRename",
+ "RestoreAfterInstallStageSync",
+ "CloseBeforeDrain",
+ "CloseAfterAuthorityRelease",
+ "pub(crate) fn hit",
+ "#[cfg(test)]\n pub(crate) fn armed",
+ "DurabilityFailpoints([redacted])",
+ ] {
+ assert!(
+ failpoint_production.contains(required),
+ "Step 072 failpoint source is missing `{required}`"
+ );
+ }
+ for forbidden in [
+ "pub enum DurabilityFailpoint",
+ "pub struct DurabilityFailpoints",
+ "static ",
+ "std::env",
+ "SystemTime",
+ "Instant",
+ "tokio::",
+ "spawn",
+ "sleep",
+ "panic!",
+ "process::exit",
+ ] {
+ assert!(
+ !failpoint_production.contains(forbidden),
+ "Step 072 failpoint source contains forbidden authority `{forbidden}`"
+ );
+ }
+ assert!(!ROOT.contains("pub use failpoint"));
+ assert!(!MANIFEST.contains("[features]"));
+ for (source, required) in [
+ (INITIALIZE_SOURCE, "InitializeBeforeCreate"),
+ (CONNECTION_SOURCE, "TransactionAfterCommit"),
+ (BACKUP_CAPTURE_SOURCE, "BackupAfterDirectorySync"),
+ (RESTORE_MARKER_SOURCE, "MarkerAdvanceAfterDirectorySync"),
+ (RESTORE_FINALIZE_SOURCE, "RestoreAfterInstallStageSync"),
+ (OPEN_SOURCE, "CloseAfterAuthorityRelease"),
+ ] {
+ assert!(
+ source.contains(required),
+ "Step 072 integration source is missing `{required}`"
+ );
+ }
+
for required in [
"radroots.service-backup",
"BACKUP_MANIFEST_SCHEMA_VERSION: u32 = 1",