commit 657970b028e7ce0b61a2b5879a0d9011f3aa725a
parent 3225d43c15f5124ce7e990aaeb6817b1713b847b
Author: triesap <tyson@radroots.org>
Date: Wed, 12 Aug 2026 00:45:17 +0000
service-sqlite: add integrity inspection
- add bounded typed active integrity reports for every host mode
- retain cancelled SQLite work in a host-owned close driver
- preserve authority and close-error precedence across cleanup
- bind documentation, package boundaries, and portability checks
Diffstat:
7 files changed, 1337 insertions(+), 4 deletions(-)
diff --git a/AGENTS.md b/AGENTS.md
@@ -251,6 +251,19 @@ Before editing code:
`Recovery` evidence. Do not expose recovery controls, add a background task
or hidden timeout, repair without writer authority, or fold Step 070
integrity/status APIs and the later process failpoint harness into recovery.
+- Explicit active integrity inspection belongs only to
+ `ServiceSqliteHost::inspect_integrity`. It admits at most one check per host,
+ uses one governed read snapshot, accepts an injected positive wall-clock
+ timestamp, and returns only the closed SQLite/foreign-key outcomes plus at
+ most two stable diagnostic codes in canonical order. Preserve authority
+ precedence after every await. Do not expose raw SQLite diagnostics, paths,
+ SQL, pool handles, or dependency errors; persist or cache the result; read an
+ ambient clock; create a timer/task; or weaken the strict restore/backup
+ integrity verifier. Callers own the monotonic deadline by cancelling the
+ future. The host-owned integrity driver must retain a cancelled in-flight
+ connection and its explicit close future until the SQLx worker terminates;
+ retry and host close resume that cleanup before proceeding. A retry must
+ inject a new timestamp.
- 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
@@ -199,6 +199,30 @@ cancelled. A later cancelled SQLite open is retried by rereading the already
durable, marker-free state. Finalization itself does not reconcile or reopen the
database.
+`ServiceSqliteHost::inspect_integrity` is the explicit active operator check.
+It is available on initialized, writable-existing, and read-only inspection
+hosts, admits at most one check per host, and uses one deferred read transaction
+as the SQLite snapshot. The caller injects a positive wall-clock
+`IntegrityCheckedAtUnixMs`; the library does not read an ambient clock or
+create a timer. The completed report contains only `verified` or `failed` for
+SQLite integrity and foreign keys, plus at most the fixed
+`sqlite_integrity_failed` and `foreign_key_violation` diagnostic codes in that
+canonical order. It can be projected to the passive `StorageIntegrity`
+vocabulary, but the library does not persist or cache the report.
+
+The check never publishes raw SQLite diagnostics, table or row identity,
+filesystem paths, SQL, or dependency errors. Inability to execute, decode, or
+finish either bounded check is an `Integrity` error rather than a fabricated
+completed result. Authority is revalidated after every await and has precedence
+over integrity classification. The operation has no hidden timeout or task.
+Callers own a positive monotonic deadline by dropping the future; cancellation
+returns no report, writes nothing, quarantines the checked-out connection, and
+leaves it in a host-owned close driver. Retry or host close explicitly awaits
+that retained close future until the prior SQLite worker terminates before any
+new check or authority release. A retry uses a newly injected wall-clock time.
+The strict backup and restore integrity verifier remains a separate fail-closed
+boundary.
+
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/connection.rs b/crates/service_sqlite/src/connection.rs
@@ -21,8 +21,8 @@ use sqlx::{
use crate::{
MigrationApplicationOutcome, MigrationAppliedAtUnixSeconds, MigrationBuildIdentity,
MigrationCallbackBinding, MigrationCatalog, OpenMode, SchemaCatalog, ServiceDatabaseIdentity,
- ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind, ServiceSqlitePaths,
- WriterAuthority,
+ ServiceSqliteConnectionOptions, ServiceSqliteError, ServiceSqliteErrorKind,
+ ServiceSqliteIntegrityReport, ServiceSqlitePaths, WriterAuthority,
};
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -49,6 +49,8 @@ pub struct ServiceSqliteHost {
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"))]
+ integrity_driver: tokio::sync::Mutex<IntegrityInspectionDriver>,
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -57,6 +59,78 @@ enum ServiceSqliteHostCloseState {
Complete(Option<ServiceSqliteErrorKind>),
}
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+enum IntegrityInspectionDriver {
+ Idle,
+ Connected(QuarantinedConnection),
+ Closing(BoxFuture<'static, Result<(), sqlx::Error>>),
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+#[derive(Clone, Copy, Debug, PartialEq, Eq)]
+enum IntegrityInspectionDriverFailure {
+ Invariant,
+ ConnectionClose,
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+impl IntegrityInspectionDriver {
+ async fn close_retained(&mut self) -> Result<(), IntegrityInspectionDriverFailure> {
+ loop {
+ match self {
+ Self::Idle => return Ok(()),
+ Self::Connected(_) => {
+ let Self::Connected(connection) = core::mem::replace(self, Self::Idle) else {
+ return Err(IntegrityInspectionDriverFailure::Invariant);
+ };
+ *self = Self::Closing(
+ connection
+ .into_close_future()
+ .ok_or(IntegrityInspectionDriverFailure::Invariant)?,
+ );
+ }
+ Self::Closing(close) => {
+ #[cfg(test)]
+ crate::integrity::integrity_test_seam::pause(
+ crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING,
+ )
+ .await;
+ let result = close.await;
+ *self = Self::Idle;
+ #[cfg(test)]
+ if crate::integrity::integrity_test_seam::take_connection_close_failure() {
+ return Err(IntegrityInspectionDriverFailure::ConnectionClose);
+ }
+ return result.map_err(|_| IntegrityInspectionDriverFailure::ConnectionClose);
+ }
+ }
+ }
+ }
+
+ fn connection_mut(&mut self) -> Result<&mut SqliteConnection, ServiceSqliteError> {
+ match self {
+ Self::Connected(connection) => Ok(connection),
+ Self::Idle | Self::Closing(_) => {
+ Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))
+ }
+ }
+ }
+
+ fn return_to_pool(&mut self) -> Result<(), ServiceSqliteError> {
+ let Self::Connected(mut connection) = core::mem::replace(self, Self::Idle) else {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity));
+ };
+ connection.trust();
+ drop(connection);
+ Ok(())
+ }
+
+ #[cfg(test)]
+ const fn is_idle(&self) -> bool {
+ matches!(self, Self::Idle)
+ }
+}
+
impl ServiceSqliteHost {
/// Opens existing writable state and finishes every pending governed migration.
///
@@ -193,10 +267,20 @@ impl ServiceSqliteHost {
if let ServiceSqliteHostCloseState::Complete(kind) = *state {
return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind)));
}
+ let (integrity_cleanup, integrity_validation) = {
+ let mut driver = self.integrity_driver.lock().await;
+ let cleanup = driver
+ .close_retained()
+ .await
+ .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ let validation = self.pool.validate();
+ (cleanup, validation)
+ };
let close = self.pool.close_explicit().await;
match close {
Err(retryable) => Err(retryable),
Ok(terminal) => {
+ let terminal = integrity_validation.and(terminal).and(integrity_cleanup);
*state = ServiceSqliteHostCloseState::Complete(
terminal.as_ref().err().map(ServiceSqliteError::kind),
);
@@ -247,6 +331,32 @@ impl ServiceSqliteHost {
}
}
+ /// Runs one explicit bounded integrity inspection over a single read snapshot.
+ ///
+ /// The caller injects the wall-clock completion time and owns any monotonic
+ /// deadline. The host admits at most one inspection at a time. Dropping this
+ /// future before it returns publishes no report, persists no status, and
+ /// leaves the checked-out connection in a host-owned explicit-close driver.
+ /// Retry or host close finishes that close before another check or authority
+ /// release; a retry must inject a new time.
+ /// Completed SQLite and foreign-key failures are returned only as fixed safe
+ /// diagnostic codes. An inability to execute or decode either check is an
+ /// `Integrity` error.
+ pub async fn inspect_integrity(
+ &self,
+ checked_at: crate::IntegrityCheckedAtUnixMs,
+ ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ {
+ self.inspect_integrity_supported(checked_at).await
+ }
+ #[cfg(not(any(target_os = "linux", target_os = "macos")))]
+ {
+ let _ = checked_at;
+ 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
@@ -476,6 +586,62 @@ impl ServiceSqliteHost {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
+ async fn inspect_integrity_supported(
+ &self,
+ checked_at: crate::IntegrityCheckedAtUnixMs,
+ ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
+ if self.closing.load(Ordering::Acquire) {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ }
+ let mut driver = self
+ .integrity_driver
+ .try_lock()
+ .map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
+ let cleanup = driver.close_retained().await;
+ self.pool.validate()?;
+ cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
+ if self.closing.load(Ordering::Acquire) {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ }
+ self.pool.validate()?;
+ let connection = self.pool.acquire().await;
+ self.pool.validate()?;
+ *driver = IntegrityInspectionDriver::Connected(QuarantinedConnection::new(connection?));
+ if self.closing.load(Ordering::Acquire) {
+ return Err(ServiceSqliteError::new(ServiceSqliteErrorKind::Open));
+ }
+ let report = crate::integrity::inspect_database_integrity(
+ driver.connection_mut()?,
+ checked_at,
+ || self.pool.validate(),
+ )
+ .await;
+ let validation = self.pool.validate();
+ match (validation, report) {
+ (Err(error), _) => {
+ let cleanup = driver.close_retained().await;
+ let validation = self.pool.validate();
+ if error.kind() == ServiceSqliteErrorKind::Authority {
+ return Err(error);
+ }
+ validation?;
+ cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
+ Err(error)
+ }
+ (Ok(()), Err(error)) => {
+ let cleanup = driver.close_retained().await;
+ self.pool.validate()?;
+ cleanup.map_err(|_| ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity))?;
+ Err(error)
+ }
+ (Ok(()), Ok(report)) => {
+ driver.return_to_pool()?;
+ Ok(report)
+ }
+ }
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
fn from_pool(mode: OpenMode, pool: crate::open::PrivateConnectionPool) -> Self {
Self {
mode,
@@ -483,6 +649,7 @@ impl ServiceSqliteHost {
closing: AtomicBool::new(false),
close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending),
backup_active: Arc::new(AtomicBool::new(false)),
+ integrity_driver: tokio::sync::Mutex::new(IntegrityInspectionDriver::Idle),
}
}
@@ -813,6 +980,11 @@ impl QuarantinedConnection {
fn trust(&mut self) {
self.trusted = true;
}
+
+ fn into_close_future(mut self) -> Option<BoxFuture<'static, Result<(), sqlx::Error>>> {
+ let connection = self.connection.take()?;
+ Some(Box::pin(async move { connection.close().await }))
+ }
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
@@ -1048,6 +1220,514 @@ mod tests {
}
#[cfg(any(target_os = "linux", target_os = "macos"))]
+ fn integrity_checked_at(value: u64) -> crate::IntegrityCheckedAtUnixMs {
+ crate::IntegrityCheckedAtUnixMs::new(value).expect("integrity inspection time")
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn integrity_inspection_is_explicit_safe_and_available_in_every_host_mode() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ let (_root, paths, identity, migrations, schema, initialized) = initialized_host().await;
+
+ let initialized_report = initialized
+ .inspect_integrity(integrity_checked_at(1_700_000_000_500))
+ .await
+ .expect("inspect initialized host");
+ assert_eq!(initialized.mode(), OpenMode::Initialize);
+ assert_eq!(
+ initialized_report.sqlite(),
+ crate::IntegrityCheckOutcome::Verified
+ );
+ assert_eq!(
+ initialized_report.foreign_keys(),
+ crate::IntegrityCheckOutcome::Verified
+ );
+ assert!(initialized_report.diagnostics().is_empty());
+ assert_eq!(
+ initialized_report.storage_integrity(),
+ crate::StorageIntegrity::Verified
+ );
+ initialized.close().await.expect("close initialized host");
+
+ let (writable, outcome) = ServiceSqliteHost::open_read_write_existing(
+ &paths,
+ &identity,
+ &migrations,
+ &schema,
+ ServiceSqliteConnectionOptions::reviewed(),
+ MigrationAppliedAtUnixSeconds::new(1_700_000_001).expect("migration time"),
+ &build_identity(),
+ &[],
+ )
+ .await
+ .expect("open writable host");
+ assert_eq!(outcome.applied_count(), 0);
+ let writable_report = writable
+ .inspect_integrity(integrity_checked_at(1_700_000_000_501))
+ .await
+ .expect("inspect writable host");
+ assert_eq!(writable.mode(), OpenMode::ReadWriteExisting);
+ assert_eq!(
+ writable_report.storage_integrity(),
+ crate::StorageIntegrity::Verified
+ );
+ writable.close().await.expect("close writable host");
+
+ let read_only = ServiceSqliteHost::open_read_only_inspection(
+ &paths,
+ &identity,
+ &migrations,
+ &schema,
+ ServiceSqliteConnectionOptions::reviewed(),
+ )
+ .await
+ .expect("open read-only host");
+ let read_only_report = read_only
+ .inspect_integrity(integrity_checked_at(1_700_000_000_502))
+ .await
+ .expect("inspect read-only host");
+ assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection);
+ assert_eq!(
+ read_only_report.storage_integrity(),
+ crate::StorageIntegrity::Verified
+ );
+ read_only.close().await.expect("close read-only host");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn integrity_inspection_is_single_admission_cancel_safe_and_close_drained() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE,
+ );
+ let first = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_510))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_BEFORE_SQLITE
+ {
+ tokio::task::yield_now().await;
+ }
+ let concurrent = host
+ .inspect_integrity(integrity_checked_at(1_700_000_000_511))
+ .await
+ .expect_err("second integrity inspection is rejected");
+ assert_eq!(concurrent.kind(), ServiceSqliteErrorKind::Integrity);
+ first.abort();
+ assert!(
+ first
+ .await
+ .expect_err("inspection is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::release();
+ assert!(
+ !host.integrity_driver.lock().await.is_idle(),
+ "cancelled SQLx connection remains host-owned until explicit cleanup"
+ );
+
+ let recovered = host
+ .inspect_integrity(integrity_checked_at(1_700_000_000_512))
+ .await
+ .expect("inspection recovers after cancellation");
+ assert_eq!(
+ recovered.storage_integrity(),
+ crate::StorageIntegrity::Verified
+ );
+
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
+ );
+ let timed_out = tokio::time::timeout(
+ Duration::from_millis(20),
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_513)),
+ )
+ .await;
+ assert!(timed_out.is_err(), "caller deadline cancels inspection");
+ crate::integrity::integrity_test_seam::release();
+ assert!(!host.integrity_driver.lock().await.is_idle());
+
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK,
+ );
+ let pre_rollback = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_514))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK
+ {
+ tokio::task::yield_now().await;
+ }
+ pre_rollback.abort();
+ assert!(
+ pre_rollback
+ .await
+ .expect_err("pre-rollback inspection is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::release();
+ assert!(!host.integrity_driver.lock().await.is_idle());
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_515))
+ .await
+ .expect("inspection recovers after pre-rollback cancellation");
+
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK,
+ );
+ let admitted = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_516))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_BEFORE_ROLLBACK
+ {
+ tokio::task::yield_now().await;
+ }
+ let close = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move { host.close().await }
+ });
+ while !host.closing.load(Ordering::Acquire) {
+ tokio::task::yield_now().await;
+ }
+ let rejected = host
+ .inspect_integrity(integrity_checked_at(1_700_000_000_517))
+ .await
+ .expect_err("closing host rejects inspection");
+ assert_eq!(rejected.kind(), ServiceSqliteErrorKind::Open);
+ assert!(!close.is_finished());
+ crate::integrity::integrity_test_seam::release();
+ admitted
+ .await
+ .expect("admitted inspection task joins")
+ .expect("admitted inspection completes");
+ close
+ .await
+ .expect("close task joins")
+ .expect("close drains admitted inspection");
+ let closed = host
+ .inspect_integrity(integrity_checked_at(1_700_000_000_518))
+ .await
+ .expect_err("closed host rejects inspection");
+ assert_eq!(closed.kind(), ServiceSqliteErrorKind::Open);
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn integrity_inspection_preserves_authority_precedence_after_await() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
+ );
+ let inspection = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_520))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS
+ {
+ tokio::task::yield_now().await;
+ }
+ let retired_lock = paths
+ .state_lock()
+ .parent()
+ .expect("state directory")
+ .join("retired-integrity-state.lock");
+ fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
+ crate::integrity::integrity_test_seam::release();
+ let error = inspection
+ .await
+ .expect("inspection task joins")
+ .expect_err("authority drift rejects inspection");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
+ fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
+ host.close().await.expect("close host after restored lock");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn observed_authority_drift_precedes_transient_restore_and_close_failure() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
+ let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS,
+ );
+ let inspection = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_525))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_BEFORE_FOREIGN_KEYS
+ {
+ tokio::task::yield_now().await;
+ }
+
+ let retired_lock = paths
+ .state_lock()
+ .parent()
+ .expect("state directory")
+ .join("transient-integrity-state.lock");
+ fs::rename(paths.state_lock(), &retired_lock).expect("retire writer lock");
+ crate::integrity::integrity_test_seam::block(
+ crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING,
+ );
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+
+ fs::rename(&retired_lock, paths.state_lock()).expect("restore writer lock");
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(true);
+ crate::integrity::integrity_test_seam::release();
+ let error = inspection
+ .await
+ .expect("inspection task joins")
+ .expect_err("observed authority drift remains terminal");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
+ assert!(host.integrity_driver.lock().await.is_idle());
+ host.close().await.expect("close host after restored lock");
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn cancelled_real_sqlite_work_is_explicitly_closed_before_retry_or_host_close() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
+ let cancelled = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_530))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ assert!(!cancelled.is_finished(), "SQLite probe remains in flight");
+ cancelled.abort();
+ assert!(
+ cancelled
+ .await
+ .expect_err("real SQLite inspection is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
+ assert!(!host.integrity_driver.lock().await.is_idle());
+
+ let cancelled_cleanup = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_531))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ cancelled_cleanup.abort();
+ assert!(
+ cancelled_cleanup
+ .await
+ .expect_err("retained close retry is cancelled")
+ .is_cancelled()
+ );
+ assert!(!host.integrity_driver.lock().await.is_idle());
+
+ let recovered = host
+ .inspect_integrity(integrity_checked_at(1_700_000_000_532))
+ .await
+ .expect("retry explicitly closes prior SQLite worker before inspecting");
+ assert_eq!(
+ recovered.storage_integrity(),
+ crate::StorageIntegrity::Verified
+ );
+ assert!(host.integrity_driver.lock().await.is_idle());
+
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
+ let cancelled = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_533))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ cancelled.abort();
+ assert!(
+ cancelled
+ .await
+ .expect_err("second real SQLite inspection is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
+ assert!(!host.integrity_driver.lock().await.is_idle());
+ host.close()
+ .await
+ .expect("host close explicitly terminates retained SQLite worker");
+ assert!(host.integrity_driver.lock().await.is_idle());
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn retained_integrity_close_releases_authority_and_caches_concurrent_lock_drift() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
+ let inspection = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_540))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ inspection.abort();
+ assert!(
+ inspection
+ .await
+ .expect_err("real integrity work is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
+
+ let close = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move { host.close().await }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_CONNECTION_CLOSE_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ let retired_lock = paths
+ .state_lock()
+ .parent()
+ .expect("state directory")
+ .join("retired-integrity-close-state.lock");
+ fs::rename(paths.state_lock(), &retired_lock).expect("retire held writer lock");
+ fs::write(paths.state_lock(), b"").expect("create replacement writer lock");
+ fs::set_permissions(paths.state_lock(), fs::Permissions::from_mode(0o600))
+ .expect("replacement writer lock mode");
+
+ let error = close
+ .await
+ .expect("close task joins")
+ .expect_err("authority drift is terminal");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Authority);
+ let repeated = host.close().await.expect_err("terminal result is cached");
+ assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority);
+ assert!(host.integrity_driver.lock().await.is_idle());
+
+ fs::remove_file(paths.state_lock()).expect("remove replacement writer lock");
+ fs::rename(&retired_lock, paths.state_lock()).expect("restore original writer lock");
+ let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
+ .expect("authority reacquisition")
+ .expect("writer authority");
+ authority.release().expect("release reacquired authority");
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test]
+ async fn retained_integrity_close_failure_is_cached_as_open_and_releases_authority() {
+ let _serial = crate::integrity::integrity_test_seam::LOCK.lock().await;
+ crate::integrity::integrity_test_seam::release();
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
+ let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await;
+ let host = Arc::new(host);
+
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(true);
+ let inspection = tokio::spawn({
+ let host = Arc::clone(&host);
+ async move {
+ host.inspect_integrity(integrity_checked_at(1_700_000_000_550))
+ .await
+ }
+ });
+ while crate::integrity::integrity_test_seam::reached()
+ != crate::integrity::integrity_test_seam::PHASE_SQLITE_EXECUTION_AWAITING
+ {
+ tokio::task::yield_now().await;
+ }
+ inspection.abort();
+ assert!(
+ inspection
+ .await
+ .expect_err("real integrity work is cancelled")
+ .is_cancelled()
+ );
+ crate::integrity::integrity_test_seam::enable_real_sqlite_probe(false);
+ assert!(!host.integrity_driver.lock().await.is_idle());
+
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(true);
+ let error = host
+ .close()
+ .await
+ .expect_err("retained connection close failure is terminal");
+ assert_eq!(error.kind(), ServiceSqliteErrorKind::Open);
+ let repeated = host.close().await.expect_err("terminal result is cached");
+ assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Open);
+ assert!(host.integrity_driver.lock().await.is_idle());
+
+ let mut authority = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting)
+ .expect("authority reacquisition")
+ .expect("writer authority");
+ authority.release().expect("release reacquired authority");
+ crate::integrity::integrity_test_seam::inject_connection_close_failure(false);
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
#[tokio::test]
async fn online_backup_captures_exact_member_manifest_and_preserves_source() {
use sha2::Digest;
diff --git a/crates/service_sqlite/src/integrity/inspection.rs b/crates/service_sqlite/src/integrity/inspection.rs
@@ -0,0 +1,515 @@
+//! Explicit, bounded integrity inspection over one governed SQLite snapshot.
+
+use serde::Serialize;
+
+use crate::StorageIntegrity;
+
+/// Caller-injected wall-clock time for one completed integrity inspection.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, PartialOrd, Ord, Serialize)]
+#[serde(transparent)]
+pub struct IntegrityCheckedAtUnixMs(u64);
+
+impl IntegrityCheckedAtUnixMs {
+ /// Constructs a positive timestamp that SQLite can represent exactly.
+ #[must_use]
+ pub const fn new(value: u64) -> Option<Self> {
+ if value > 0 && value <= i64::MAX as u64 {
+ Some(Self(value))
+ } else {
+ None
+ }
+ }
+
+ /// Returns the validated Unix timestamp in milliseconds.
+ #[must_use]
+ pub const fn get(self) -> u64 {
+ self.0
+ }
+}
+
+/// Closed result of one completed bounded database check.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
+#[serde(rename_all = "snake_case")]
+pub enum IntegrityCheckOutcome {
+ Verified,
+ Failed,
+}
+
+/// Stable, content-free diagnostic code for a completed failed check.
+#[derive(Clone, Copy, Debug, PartialEq, Eq, Serialize)]
+#[serde(rename_all = "snake_case")]
+pub enum IntegrityDiagnosticCode {
+ SqliteIntegrityFailed,
+ ForeignKeyViolation,
+}
+
+/// Safe bounded result of an explicit host integrity inspection.
+#[derive(Clone, Debug, PartialEq, Eq, Serialize)]
+pub struct ServiceSqliteIntegrityReport {
+ checked_at_unix_ms: IntegrityCheckedAtUnixMs,
+ sqlite: IntegrityCheckOutcome,
+ foreign_keys: IntegrityCheckOutcome,
+ diagnostics: Box<[IntegrityDiagnosticCode]>,
+}
+
+impl ServiceSqliteIntegrityReport {
+ #[cfg(any(test, target_os = "linux", target_os = "macos"))]
+ pub(crate) fn new(
+ checked_at_unix_ms: IntegrityCheckedAtUnixMs,
+ sqlite: IntegrityCheckOutcome,
+ foreign_keys: IntegrityCheckOutcome,
+ ) -> Self {
+ let diagnostics: Box<[IntegrityDiagnosticCode]> = match (sqlite, foreign_keys) {
+ (IntegrityCheckOutcome::Verified, IntegrityCheckOutcome::Verified) => Box::new([]),
+ (IntegrityCheckOutcome::Failed, IntegrityCheckOutcome::Verified) => {
+ Box::new([IntegrityDiagnosticCode::SqliteIntegrityFailed])
+ }
+ (IntegrityCheckOutcome::Verified, IntegrityCheckOutcome::Failed) => {
+ Box::new([IntegrityDiagnosticCode::ForeignKeyViolation])
+ }
+ (IntegrityCheckOutcome::Failed, IntegrityCheckOutcome::Failed) => Box::new([
+ IntegrityDiagnosticCode::SqliteIntegrityFailed,
+ IntegrityDiagnosticCode::ForeignKeyViolation,
+ ]),
+ };
+ Self {
+ checked_at_unix_ms,
+ sqlite,
+ foreign_keys,
+ diagnostics,
+ }
+ }
+
+ /// Returns the caller-injected completion time.
+ #[must_use]
+ pub const fn checked_at_unix_ms(&self) -> IntegrityCheckedAtUnixMs {
+ self.checked_at_unix_ms
+ }
+
+ /// Returns the completed SQLite integrity-check outcome.
+ #[must_use]
+ pub const fn sqlite(&self) -> IntegrityCheckOutcome {
+ self.sqlite
+ }
+
+ /// Returns the completed foreign-key-check outcome.
+ #[must_use]
+ pub const fn foreign_keys(&self) -> IntegrityCheckOutcome {
+ self.foreign_keys
+ }
+
+ /// Returns zero to two stable diagnostic codes in canonical order.
+ #[must_use]
+ pub fn diagnostics(&self) -> &[IntegrityDiagnosticCode] {
+ &self.diagnostics
+ }
+
+ /// Projects this active result into the passive storage-status vocabulary.
+ #[must_use]
+ pub const fn storage_integrity(&self) -> StorageIntegrity {
+ if matches!(self.sqlite, IntegrityCheckOutcome::Verified)
+ && matches!(self.foreign_keys, IntegrityCheckOutcome::Verified)
+ {
+ StorageIntegrity::Verified
+ } else {
+ StorageIntegrity::Failed
+ }
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+mod native {
+ use sqlx::{Connection, Row, SqliteConnection};
+
+ use super::{IntegrityCheckOutcome, IntegrityCheckedAtUnixMs, ServiceSqliteIntegrityReport};
+ use crate::{ServiceSqliteError, ServiceSqliteErrorKind};
+
+ const SQLITE_INTEGRITY_SQL: &str = "PRAGMA integrity_check(1)";
+ const FOREIGN_KEY_SQL: &str = "SELECT 1 FROM pragma_foreign_key_check LIMIT 1";
+
+ pub(crate) async fn inspect_database_integrity(
+ connection: &mut SqliteConnection,
+ checked_at: IntegrityCheckedAtUnixMs,
+ mut validate: impl FnMut() -> Result<(), ServiceSqliteError>,
+ ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
+ validate()?;
+ let transaction = connection.begin().await;
+ validate()?;
+ let mut transaction = transaction.map_err(|_| integrity_error())?;
+
+ #[cfg(test)]
+ if super::test_seam::real_sqlite_probe_enabled() {
+ super::test_seam::observe(super::test_seam::PHASE_SQLITE_EXECUTION_AWAITING);
+ let probe = sqlx::query_scalar::<_, i64>(
+ "WITH RECURSIVE counter(value) AS (
+ VALUES(0) UNION ALL SELECT value + 1 FROM counter WHERE value < 5000000
+ ) SELECT sum(value) FROM counter",
+ )
+ .fetch_one(&mut *transaction)
+ .await;
+ validate()?;
+ if probe.is_err() {
+ return rollback_error(transaction, &mut validate).await;
+ }
+ }
+
+ #[cfg(test)]
+ super::test_seam::pause(super::test_seam::PHASE_BEFORE_SQLITE).await;
+ let sqlite_rows = sqlx::query(SQLITE_INTEGRITY_SQL)
+ .fetch_all(&mut *transaction)
+ .await;
+ validate()?;
+ let sqlite = match sqlite_rows {
+ Ok(rows) if rows.len() == 1 => match rows[0].try_get::<&str, _>(0) {
+ Ok(value) => classify_integrity_value(value),
+ Err(_) => return rollback_error(transaction, &mut validate).await,
+ },
+ Ok(_) | Err(_) => return rollback_error(transaction, &mut validate).await,
+ };
+
+ #[cfg(test)]
+ super::test_seam::pause(super::test_seam::PHASE_BEFORE_FOREIGN_KEYS).await;
+ let foreign_key_row = sqlx::query_scalar::<_, i64>(FOREIGN_KEY_SQL)
+ .fetch_optional(&mut *transaction)
+ .await;
+ validate()?;
+ let foreign_keys = match foreign_key_row {
+ Ok(None) => IntegrityCheckOutcome::Verified,
+ Ok(Some(1)) => IntegrityCheckOutcome::Failed,
+ Ok(Some(_)) | Err(_) => return rollback_error(transaction, &mut validate).await,
+ };
+
+ #[cfg(test)]
+ super::test_seam::pause(super::test_seam::PHASE_BEFORE_ROLLBACK).await;
+ let rollback = transaction.rollback().await;
+ validate()?;
+ rollback.map_err(|_| integrity_error())?;
+ Ok(ServiceSqliteIntegrityReport::new(
+ checked_at,
+ sqlite,
+ foreign_keys,
+ ))
+ }
+
+ async fn rollback_error(
+ transaction: sqlx::Transaction<'_, sqlx::Sqlite>,
+ validate: &mut impl FnMut() -> Result<(), ServiceSqliteError>,
+ ) -> Result<ServiceSqliteIntegrityReport, ServiceSqliteError> {
+ let rollback = transaction.rollback().await;
+ validate()?;
+ rollback.map_err(|_| integrity_error())?;
+ Err(integrity_error())
+ }
+
+ fn integrity_error() -> ServiceSqliteError {
+ ServiceSqliteError::new(ServiceSqliteErrorKind::Integrity)
+ }
+
+ fn classify_integrity_value(value: &str) -> IntegrityCheckOutcome {
+ if value == "ok" {
+ IntegrityCheckOutcome::Verified
+ } else {
+ IntegrityCheckOutcome::Failed
+ }
+ }
+
+ #[cfg(test)]
+ pub(super) fn classify_test_value(value: &str) -> IntegrityCheckOutcome {
+ classify_integrity_value(value)
+ }
+}
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+pub(crate) use native::inspect_database_integrity;
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) mod test_seam {
+ use std::sync::atomic::{AtomicBool, AtomicU8, Ordering};
+
+ pub(crate) const PHASE_BEFORE_SQLITE: u8 = 1;
+ pub(crate) const PHASE_BEFORE_FOREIGN_KEYS: u8 = 2;
+ pub(crate) const PHASE_BEFORE_ROLLBACK: u8 = 3;
+ pub(crate) const PHASE_SQLITE_EXECUTION_AWAITING: u8 = 4;
+ pub(crate) const PHASE_CONNECTION_CLOSE_AWAITING: u8 = 5;
+
+ static BLOCKED: AtomicU8 = AtomicU8::new(0);
+ static REACHED: AtomicU8 = AtomicU8::new(0);
+ static RELEASED: AtomicBool = AtomicBool::new(true);
+ static REAL_SQLITE_PROBE: AtomicBool = AtomicBool::new(false);
+ static CONNECTION_CLOSE_FAILURE: AtomicBool = AtomicBool::new(false);
+ pub(crate) static LOCK: tokio::sync::Mutex<()> = tokio::sync::Mutex::const_new(());
+
+ pub(crate) async fn pause(phase: u8) {
+ REACHED.store(phase, Ordering::Release);
+ while BLOCKED.load(Ordering::Acquire) == phase && !RELEASED.load(Ordering::Acquire) {
+ tokio::task::yield_now().await;
+ }
+ }
+
+ pub(crate) fn block(phase: u8) {
+ REACHED.store(0, Ordering::Release);
+ BLOCKED.store(phase, Ordering::Release);
+ RELEASED.store(false, Ordering::Release);
+ }
+
+ pub(crate) fn observe(phase: u8) {
+ REACHED.store(phase, Ordering::Release);
+ }
+
+ pub(crate) fn enable_real_sqlite_probe(enabled: bool) {
+ REACHED.store(0, Ordering::Release);
+ REAL_SQLITE_PROBE.store(enabled, Ordering::Release);
+ }
+
+ pub(crate) fn real_sqlite_probe_enabled() -> bool {
+ REAL_SQLITE_PROBE.load(Ordering::Acquire)
+ }
+
+ pub(crate) fn inject_connection_close_failure(enabled: bool) {
+ CONNECTION_CLOSE_FAILURE.store(enabled, Ordering::Release);
+ }
+
+ pub(crate) fn take_connection_close_failure() -> bool {
+ CONNECTION_CLOSE_FAILURE.swap(false, Ordering::AcqRel)
+ }
+
+ pub(crate) fn reached() -> u8 {
+ REACHED.load(Ordering::Acquire)
+ }
+
+ pub(crate) fn release() {
+ RELEASED.store(true, Ordering::Release);
+ BLOCKED.store(0, Ordering::Release);
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ use std::{
+ fs::OpenOptions,
+ io::{Seek, SeekFrom, Write},
+ };
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ use sqlx::{Connection, SqliteConnection, sqlite::SqliteConnectOptions};
+
+ #[test]
+ fn timestamp_bounds_and_report_wire_vocabulary_are_exact() {
+ assert!(IntegrityCheckedAtUnixMs::new(0).is_none());
+ let maximum = IntegrityCheckedAtUnixMs::new(i64::MAX as u64).expect("maximum timestamp");
+ assert_eq!(maximum.get(), i64::MAX as u64);
+ assert!(IntegrityCheckedAtUnixMs::new(i64::MAX as u64 + 1).is_none());
+
+ let checked_at = IntegrityCheckedAtUnixMs::new(1_700_000_000_000).unwrap();
+ let verified = ServiceSqliteIntegrityReport::new(
+ checked_at,
+ IntegrityCheckOutcome::Verified,
+ IntegrityCheckOutcome::Verified,
+ );
+ assert!(verified.diagnostics().is_empty());
+ assert_eq!(verified.storage_integrity(), StorageIntegrity::Verified);
+ assert_eq!(
+ serde_json::to_string(&verified).unwrap(),
+ r#"{"checked_at_unix_ms":1700000000000,"sqlite":"verified","foreign_keys":"verified","diagnostics":[]}"#
+ );
+
+ let failed = ServiceSqliteIntegrityReport::new(
+ checked_at,
+ IntegrityCheckOutcome::Failed,
+ IntegrityCheckOutcome::Failed,
+ );
+ assert_eq!(
+ failed.diagnostics(),
+ [
+ IntegrityDiagnosticCode::SqliteIntegrityFailed,
+ IntegrityDiagnosticCode::ForeignKeyViolation,
+ ]
+ );
+ assert_eq!(failed.storage_integrity(), StorageIntegrity::Failed);
+ assert_eq!(
+ serde_json::to_string(&failed).unwrap(),
+ r#"{"checked_at_unix_ms":1700000000000,"sqlite":"failed","foreign_keys":"failed","diagnostics":["sqlite_integrity_failed","foreign_key_violation"]}"#
+ );
+ assert!(!format!("{failed:?}").contains("sqlite_schema"));
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ assert_eq!(
+ native::classify_test_value(
+ "a completed SQLite diagnostic that is intentionally much longer than sixty-four bytes"
+ ),
+ IntegrityCheckOutcome::Failed
+ );
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ async fn in_memory_database() -> SqliteConnection {
+ SqliteConnection::connect_with(&SqliteConnectOptions::new().filename(":memory:"))
+ .await
+ .expect("in-memory database")
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test(flavor = "current_thread")]
+ async fn native_inspection_reports_healthy_and_foreign_key_failure() {
+ let _serial = test_seam::LOCK.lock().await;
+ test_seam::release();
+ let checked_at = IntegrityCheckedAtUnixMs::new(1).unwrap();
+ let mut healthy = in_memory_database().await;
+ let healthy = inspect_database_integrity(&mut healthy, checked_at, || Ok(()))
+ .await
+ .expect("healthy inspection");
+ assert_eq!(healthy.sqlite(), IntegrityCheckOutcome::Verified);
+ assert_eq!(healthy.foreign_keys(), IntegrityCheckOutcome::Verified);
+ assert!(healthy.diagnostics().is_empty());
+
+ let mut foreign_keys = in_memory_database().await;
+ sqlx::raw_sql(
+ "PRAGMA foreign_keys=OFF;
+ CREATE TABLE parent (id INTEGER PRIMARY KEY) STRICT;
+ CREATE TABLE child (parent_id INTEGER REFERENCES parent(id)) STRICT;
+ INSERT INTO child(parent_id) VALUES (99);",
+ )
+ .execute(&mut foreign_keys)
+ .await
+ .expect("seed foreign-key violation");
+ let failed = inspect_database_integrity(&mut foreign_keys, checked_at, || Ok(()))
+ .await
+ .expect("completed foreign-key inspection");
+ assert_eq!(failed.sqlite(), IntegrityCheckOutcome::Verified);
+ assert_eq!(failed.foreign_keys(), IntegrityCheckOutcome::Failed);
+ assert_eq!(
+ failed.diagnostics(),
+ [IntegrityDiagnosticCode::ForeignKeyViolation]
+ );
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test(flavor = "current_thread")]
+ async fn native_inspection_keeps_completed_failure_diagnostics_bounded() {
+ let _serial = test_seam::LOCK.lock().await;
+ test_seam::release();
+ let mut corrupt = in_memory_database().await;
+ sqlx::raw_sql(
+ "PRAGMA foreign_keys=OFF;
+ CREATE TABLE parent (id INTEGER PRIMARY KEY) STRICT;
+ CREATE TABLE child (parent_id INTEGER REFERENCES parent(id)) STRICT;
+ CREATE INDEX parent_index ON parent(id);
+ INSERT INTO child(parent_id) VALUES (99);
+ PRAGMA writable_schema=ON;
+ UPDATE sqlite_schema SET rootpage=0 WHERE name='parent_index';
+ PRAGMA writable_schema=OFF;
+ PRAGMA schema_version=99;",
+ )
+ .execute(&mut corrupt)
+ .await
+ .expect("seed bounded corruption");
+ let report = inspect_database_integrity(
+ &mut corrupt,
+ IntegrityCheckedAtUnixMs::new(2).unwrap(),
+ || Ok(()),
+ )
+ .await
+ .expect("completed corruption inspection");
+ assert_eq!(report.sqlite(), IntegrityCheckOutcome::Failed);
+ assert_eq!(report.foreign_keys(), IntegrityCheckOutcome::Failed);
+ assert_eq!(report.diagnostics().len(), 2);
+ let rendered = format!("{report:?}");
+ assert!(!rendered.contains("parent"));
+ assert!(!rendered.contains("rootpage"));
+ }
+
+ #[cfg(any(target_os = "linux", target_os = "macos"))]
+ #[tokio::test(flavor = "current_thread")]
+ async fn physical_corruption_and_query_failure_remain_redacted_and_typed() {
+ let _serial = test_seam::LOCK.lock().await;
+ test_seam::release();
+ let directory = tempfile::tempdir().expect("temporary database directory");
+ let path = directory.path().join("sensitive-state-name.sqlite");
+ let options = SqliteConnectOptions::new()
+ .filename(&path)
+ .create_if_missing(true);
+ let mut connection = SqliteConnection::connect_with(&options)
+ .await
+ .expect("create database");
+ sqlx::query("CREATE TABLE integrity_probe (value BLOB NOT NULL) STRICT")
+ .execute(&mut connection)
+ .await
+ .expect("create probe table");
+ sqlx::query("INSERT INTO integrity_probe(value) VALUES (zeroblob(4096))")
+ .execute(&mut connection)
+ .await
+ .expect("allocate probe page");
+ let page_size = sqlx::query_scalar::<_, i64>("PRAGMA page_size")
+ .fetch_one(&mut connection)
+ .await
+ .expect("page size");
+ let root_page = sqlx::query_scalar::<_, i64>(
+ "SELECT rootpage FROM sqlite_schema WHERE name='integrity_probe'",
+ )
+ .fetch_one(&mut connection)
+ .await
+ .expect("probe root page");
+ connection.close().await.expect("close database");
+ let offset = u64::try_from(root_page - 1)
+ .ok()
+ .and_then(|page| page.checked_mul(u64::try_from(page_size).ok()?))
+ .expect("corrupt page offset");
+ let mut file = OpenOptions::new()
+ .write(true)
+ .open(&path)
+ .expect("open database bytes");
+ file.seek(SeekFrom::Start(offset)).expect("seek root page");
+ file.write_all(&[0xff]).expect("corrupt page type");
+ file.sync_all().expect("sync corrupt database");
+ drop(file);
+
+ let options = SqliteConnectOptions::new()
+ .filename(&path)
+ .create_if_missing(false);
+ let mut corrupt = SqliteConnection::connect_with(&options)
+ .await
+ .expect("open corrupt database shell");
+ let report = inspect_database_integrity(
+ &mut corrupt,
+ IntegrityCheckedAtUnixMs::new(3).unwrap(),
+ || Ok(()),
+ )
+ .await
+ .expect("completed physical-corruption result");
+ assert_eq!(report.sqlite(), IntegrityCheckOutcome::Failed);
+ assert_eq!(
+ report.diagnostics(),
+ [IntegrityDiagnosticCode::SqliteIntegrityFailed]
+ );
+ let rendered = format!("{report:?}");
+ assert!(!rendered.contains("sensitive-state-name"));
+ assert!(!rendered.contains("integrity_probe"));
+ assert!(!rendered.contains("database disk image"));
+
+ let mut malformed = in_memory_database().await;
+ sqlx::raw_sql(
+ "CREATE TABLE secret_schema_name (value INTEGER) STRICT;
+ PRAGMA writable_schema=ON;
+ UPDATE sqlite_schema SET sql='CREATE TABLE secret_schema_name('
+ WHERE name='secret_schema_name';
+ PRAGMA writable_schema=OFF;
+ PRAGMA schema_version=99;",
+ )
+ .execute(&mut malformed)
+ .await
+ .expect("seed malformed schema");
+ let error = inspect_database_integrity(
+ &mut malformed,
+ IntegrityCheckedAtUnixMs::new(4).unwrap(),
+ || Ok(()),
+ )
+ .await
+ .expect_err("query failure is not a completed report");
+ assert_eq!(error.kind(), crate::ServiceSqliteErrorKind::Integrity);
+ let rendered = format!("{error:?} {}", error);
+ assert!(!rendered.contains("secret_schema_name"));
+ assert!(!rendered.contains("incomplete input"));
+ }
+}
diff --git a/crates/service_sqlite/src/integrity/mod.rs b/crates/service_sqlite/src/integrity/mod.rs
@@ -1,11 +1,22 @@
//! Exact schema-object catalog verification.
pub(crate) mod catalog;
+mod inspection;
pub use catalog::{
SchemaCatalog, SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind,
SchemaVersionCatalog,
};
+pub use inspection::{
+ IntegrityCheckOutcome, IntegrityCheckedAtUnixMs, IntegrityDiagnosticCode,
+ ServiceSqliteIntegrityReport,
+};
+
+#[cfg(any(target_os = "linux", target_os = "macos"))]
+pub(crate) use inspection::inspect_database_integrity;
+
+#[cfg(all(test, any(target_os = "linux", target_os = "macos")))]
+pub(crate) use inspection::test_seam as integrity_test_seam;
#[cfg(any(target_os = "linux", target_os = "macos"))]
use core::fmt;
diff --git a/crates/service_sqlite/src/lib.rs b/crates/service_sqlite/src/lib.rs
@@ -33,8 +33,9 @@ pub use error::{
};
pub use initialize::initialize_database;
pub use integrity::{
- SchemaCatalog, SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind,
- SchemaVersionCatalog,
+ IntegrityCheckOutcome, IntegrityCheckedAtUnixMs, IntegrityDiagnosticCode, SchemaCatalog,
+ SchemaCatalogContractError, SchemaDigest, SchemaObject, SchemaObjectKind, SchemaVersionCatalog,
+ ServiceSqliteIntegrityReport,
};
pub use metadata::{
ServiceDatabaseIdentity, ServiceDatabaseMetadata, ServiceSqliteApplicationId,
diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs
@@ -13,6 +13,7 @@ const ERROR_SOURCE: &str = include_str!("../src/error.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");
+const INTEGRITY_INSPECTION_SOURCE: &str = include_str!("../src/integrity/inspection.rs");
const METADATA_SOURCE: &str = include_str!("../src/metadata.rs");
const MIGRATION_SOURCE: &str = include_str!("../src/migration.rs");
const OPEN_SOURCE: &str = include_str!("../src/open.rs");
@@ -196,6 +197,29 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"Recovery has no await point or hidden task",
"each synchronous filesystem step and its authority checks complete",
"Finalization itself does not reconcile or reopen the database",
+ "`ServiceSqliteHost::inspect_integrity` is the explicit active operator check",
+ "available on initialized, writable-existing, and read-only inspection hosts",
+ "admits at most one check per host",
+ "uses one deferred read transaction as the SQLite snapshot",
+ "caller injects a positive wall-clock `IntegrityCheckedAtUnixMs`",
+ "does not read an ambient clock or create a timer",
+ "only `verified` or `failed` for SQLite integrity and foreign keys",
+ "at most the fixed `sqlite_integrity_failed` and `foreign_key_violation` diagnostic codes",
+ "canonical order",
+ "projected to the passive `StorageIntegrity` vocabulary",
+ "does not persist or cache the report",
+ "never publishes raw SQLite diagnostics, table or row identity",
+ "Inability to execute, decode, or finish either bounded check is an `Integrity` error",
+ "Authority is revalidated after every await and has precedence",
+ "operation has no hidden timeout or task",
+ "Callers own a positive monotonic deadline by dropping the future",
+ "cancellation returns no report, writes nothing, quarantines the checked-out connection",
+ "leaves it in a host-owned close driver",
+ "Retry or host close explicitly awaits that retained close future",
+ "until the prior SQLite worker terminates",
+ "before any new check or authority release",
+ "retry uses a newly injected wall-clock time",
+ "strict backup and restore integrity verifier remains a separate fail-closed boundary",
] {
assert!(
readme_words.contains(required),
@@ -229,6 +253,10 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
let backup_verify_production = BACKUP_VERIFY_SOURCE
.split_once("#[cfg(test)]")
.map_or(BACKUP_VERIFY_SOURCE, |(production, _)| production);
+ let integrity_inspection_production = INTEGRITY_INSPECTION_SOURCE
+ .split_once("#[cfg(all(test, any(target_os = \"linux\", target_os = \"macos\")))]")
+ .map(|(production, _)| production)
+ .expect("integrity inspection source must keep test seams separated");
let restore_marker_production = RESTORE_MARKER_SOURCE
.split_once("#[cfg(test)]\nmod tests")
.map(|(production, _)| production)
@@ -274,6 +302,10 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"MigrationKind",
"MigrationName",
"MigrationTransactionExecutor",
+ "IntegrityCheckOutcome",
+ "IntegrityCheckedAtUnixMs",
+ "IntegrityDiagnosticCode",
+ "ServiceSqliteIntegrityReport",
"SchemaCatalog",
"SchemaCatalogContractError",
"SchemaDigest",
@@ -386,6 +418,7 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
"must reread authoritative state before any idempotent retry",
"pub async fn close(&self)",
"pub async fn capture_online_backup(",
+ "pub async fn inspect_integrity(",
"closing.store(true, Ordering::Release)",
"close_state.lock().await",
"ServiceSqliteHostCloseState::Complete",
@@ -435,6 +468,62 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() {
}
for required in [
+ "pub struct IntegrityCheckedAtUnixMs",
+ "pub enum IntegrityCheckOutcome",
+ "pub enum IntegrityDiagnosticCode",
+ "pub struct ServiceSqliteIntegrityReport",
+ "SqliteIntegrityFailed",
+ "ForeignKeyViolation",
+ "Box<[IntegrityDiagnosticCode]>",
+ "pub const fn storage_integrity",
+ "PRAGMA integrity_check(1)",
+ "SELECT 1 FROM pragma_foreign_key_check LIMIT 1",
+ "connection.begin().await",
+ "transaction.rollback().await",
+ "validate()?",
+ ] {
+ assert!(
+ integrity_inspection_production.contains(required),
+ "Step 070 integrity inspection source is missing `{required}`"
+ );
+ }
+ for required in [
+ "integrity_driver: tokio::sync::Mutex<IntegrityInspectionDriver>",
+ ".try_lock()",
+ "driver.close_retained().await",
+ "QuarantinedConnection::new",
+ "crate::integrity::inspect_database_integrity",
+ "connection.trust()",
+ ] {
+ assert!(
+ connection_production.contains(required),
+ "Step 070 host integration is missing `{required}`"
+ );
+ }
+ for forbidden in [
+ "pub use sqlx",
+ "SqlitePool",
+ "PoolConnection",
+ "SystemTime",
+ "Instant",
+ "tokio::time",
+ "tokio::task",
+ "spawn",
+ "std::fs",
+ "OpenOptions",
+ "write_all",
+ "persist",
+ "cache",
+ "myc_",
+ "rhi_",
+ ] {
+ assert!(
+ !integrity_inspection_production.contains(forbidden),
+ "Step 070 integrity inspection contains forbidden authority `{forbidden}`"
+ );
+ }
+
+ for required in [
"rusqlite::backup::Backup::new",
"tokio::task::spawn_blocking",
"BACKUP_PAGES_PER_STEP",