lib

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

commit 61b7527d16b78f464194853d5d53c514d9a84075
parent 7bfb2347ce1ac26a9880403b2be158e1b77c3ccf
Author: triesap <tyson@radroots.org>
Date:   Tue, 11 Aug 2026 17:20:05 +0000

service-sqlite: close state explicitly

- add serialized explicit host close and permanent transaction admission fencing
- retain checkpoint connection work across cancellation and release authority explicitly
- enforce fixed unblocked WAL truncation with stable redacted failure precedence
- cover concurrent, cancelled, read-only, busy, drift, and reopen close paths

Diffstat:
MAGENTS.md | 9++++++++-
Mcrates/service_sqlite/Cargo.toml | 1+
Mcrates/service_sqlite/README.md | 15+++++++++++++++
Mcrates/service_sqlite/src/connection.rs | 557+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----
Mcrates/service_sqlite/src/open.rs | 386+++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++--------
Mcrates/service_sqlite/tests/package_boundary.rs | 71++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++-----
6 files changed, 972 insertions(+), 67 deletions(-)

diff --git a/AGENTS.md b/AGENTS.md @@ -165,7 +165,14 @@ Before editing code: accessors. Do not attach or detach secondary SQLite databases through the transaction executor. Writable host construction must finish governed migrations before returning, while read-only inspection must require current - migration and schema state. + migration and schema state. Every host owner must explicitly await + `ServiceSqliteHost::close`: close permanently stops admission, drains admitted + work, applies the fixed unblocked `TRUNCATE` WAL checkpoint for writable + hosts only, closes its private checkpoint connection, and explicitly releases + writer or inspection authority. A cancelled close retains authority and must + be resumed through the host-owned connect/checkpoint/connection-close driver; + Drop is not an asynchronous close or completion proof. Do not add public + checkpoint knobs, background close tasks, or Drop-based async cleanup. - 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/Cargo.toml b/crates/service_sqlite/Cargo.toml @@ -20,6 +20,7 @@ rustix = { workspace = true } serde = { workspace = true, features = ["derive", "std"] } sha2 = { workspace = true } sqlx = { workspace = true, features = ["runtime-tokio", "sqlite-bundled"] } +tokio = { workspace = true, features = ["sync"] } [dev-dependencies] serde_json = { workspace = true } diff --git a/crates/service_sqlite/README.md b/crates/service_sqlite/README.md @@ -23,6 +23,21 @@ commit begins yields no result and must be treated as an unknown commit outcome. Both that case and `CommitOutcomeUnknown` require rereading authoritative state before an idempotent retry. +Every host must be closed explicitly with `ServiceSqliteHost::close`. Close +permanently stops new transaction admission, drains transactions that were +already admitted, and is safe to call sequentially or concurrently. Writable +close applies the fixed `PRAGMA wal_checkpoint(TRUNCATE)` policy, requires an +unblocked checkpoint, closes the private checkpoint connection, and explicitly +releases writer authority. Read-only inspection close drains its pool and +releases its shared inspection guard without checkpointing or mutating the +database or filesystem. Cancelling close before terminal completion leaves the +host non-admitting and retains authority; the private connect, checkpoint, and +explicit connection-close driver remains host-owned so a later call resumes +close without losing the SQLite handle or its close proof. Once authority +release is proven, the stable outer result is cached for every later call. +Dropping a host performs no asynchronous close work and is not proof that the +governed checkpoint and authority-release sequence completed. + 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 @@ -42,6 +42,16 @@ pub struct ServiceSqliteHost { mode: OpenMode, #[cfg(any(target_os = "linux", target_os = "macos"))] pool: crate::open::PrivateConnectionPool, + #[cfg(any(target_os = "linux", target_os = "macos"))] + closing: AtomicBool, + #[cfg(any(target_os = "linux", target_os = "macos"))] + close_state: tokio::sync::Mutex<ServiceSqliteHostCloseState>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +enum ServiceSqliteHostCloseState { + Pending, + Complete(Option<ServiceSqliteErrorKind>), } impl ServiceSqliteHost { @@ -69,13 +79,7 @@ impl ServiceSqliteHost { ) .await?; match pool.apply_migrations(applied_at, build, callbacks).await { - Ok(outcome) => Ok(( - Self { - mode: OpenMode::ReadWriteExisting, - pool, - }, - outcome, - )), + Ok(outcome) => Ok((Self::from_pool(OpenMode::ReadWriteExisting, pool), outcome)), Err(error) => { drop(pool.close().await); Err(error) @@ -111,13 +115,7 @@ impl ServiceSqliteHost { ) .await?; match pool.apply_migrations(applied_at, build, callbacks).await { - Ok(outcome) => Ok(( - Self { - mode: OpenMode::Initialize, - pool, - }, - outcome, - )), + Ok(outcome) => Ok((Self::from_pool(OpenMode::Initialize, pool), outcome)), Err(error) => { drop(pool.close().await); Err(error) @@ -153,10 +151,7 @@ impl ServiceSqliteHost { options, ) .await?; - Ok(Self { - mode: OpenMode::ReadOnlyInspection, - pool, - }) + Ok(Self::from_pool(OpenMode::ReadOnlyInspection, pool)) } #[cfg(not(any(target_os = "linux", target_os = "macos")))] { @@ -171,6 +166,41 @@ impl ServiceSqliteHost { self.mode } + /// Closes all connections and explicitly releases retained instance authority. + /// + /// Close rejects new transactions as soon as it starts and waits for already + /// admitted transactions to finish. Writable hosts then perform the fixed + /// governed `TRUNCATE` WAL checkpoint before releasing writer authority; + /// read-only inspection performs no checkpoint or filesystem mutation. + /// + /// Cancelling this future leaves the host permanently non-admitting and retains + /// authority until a later call resumes close. A completed result is cached, so + /// sequential or concurrent later calls return the same stable outer outcome. + pub async fn close(&self) -> Result<(), ServiceSqliteError> { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + self.closing.store(true, Ordering::Release); + let mut state = self.close_state.lock().await; + if let ServiceSqliteHostCloseState::Complete(kind) = *state { + return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind))); + } + let close = self.pool.close_explicit().await; + match close { + Err(retryable) => Err(retryable), + Ok(terminal) => { + *state = ServiceSqliteHostCloseState::Complete( + terminal.as_ref().err().map(ServiceSqliteError::kind), + ); + terminal + } + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + 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 @@ -218,6 +248,11 @@ impl ServiceSqliteHost { ) -> ServiceSqliteTransactionFuture<'a, T, E> + Send, { + if self.closing.load(Ordering::Acquire) { + return Err(ServiceSqliteTransactionError::not_committed( + ServiceSqliteError::new(ServiceSqliteErrorKind::Open), + )); + } self.pool .validate() .map_err(ServiceSqliteTransactionError::not_committed)?; @@ -395,6 +430,31 @@ impl ServiceSqliteHost { } #[cfg(any(target_os = "linux", target_os = "macos"))] + fn from_pool(mode: OpenMode, pool: crate::open::PrivateConnectionPool) -> Self { + Self { + mode, + pool, + closing: AtomicBool::new(false), + close_state: tokio::sync::Mutex::new(ServiceSqliteHostCloseState::Pending), + } + } + + fn lifecycle_state(&self) -> &'static str { + #[cfg(any(target_os = "linux", target_os = "macos"))] + { + if self.closing.load(Ordering::Acquire) { + "closing_or_closed" + } else { + "open" + } + } + #[cfg(not(any(target_os = "linux", target_os = "macos")))] + { + "closing_or_closed" + } + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] async fn verify_before_commit( &self, connection: &mut SqliteConnection, @@ -428,11 +488,6 @@ impl ServiceSqliteHost { } Ok(()) } - - #[cfg(all(test, any(target_os = "linux", target_os = "macos")))] - async fn close_pool(&self) { - self.pool.close_pool().await; - } } impl fmt::Debug for ServiceSqliteHost { @@ -441,6 +496,7 @@ impl fmt::Debug for ServiceSqliteHost { .debug_struct("ServiceSqliteHost") .field("mode", &self.mode) .field("pool", &"[redacted]") + .field("lifecycle", &self.lifecycle_state()) .finish() } } @@ -755,14 +811,17 @@ fn unsupported_host() -> ServiceSqliteError { mod tests { #[cfg(any(target_os = "linux", target_os = "macos"))] use std::{ + collections::BTreeMap, convert::Infallible, fs, num::NonZeroU32, os::unix::fs::PermissionsExt, + path::Path, sync::{ Arc, atomic::{AtomicBool, Ordering}, }, + time::{Duration, SystemTime}, }; #[cfg(any(target_os = "linux", target_os = "macos"))] @@ -774,6 +833,8 @@ mod tests { use radroots_storage::event::SourceGeneration; #[cfg(any(target_os = "linux", target_os = "macos"))] use sqlx::{Connection, sqlite::SqliteConnectOptions}; + #[cfg(any(target_os = "linux", target_os = "macos"))] + use tokio::sync::Notify; use super::*; @@ -852,6 +913,20 @@ mod tests { SchemaCatalog, ServiceSqliteHost, ) { + initialized_host_with_options(ServiceSqliteConnectionOptions::reviewed()).await + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + async fn initialized_host_with_options( + options: ServiceSqliteConnectionOptions, + ) -> ( + tempfile::TempDir, + ServiceSqlitePaths, + ServiceDatabaseIdentity, + MigrationCatalog, + SchemaCatalog, + ServiceSqliteHost, + ) { let root = tempfile::tempdir().expect("temporary root"); let paths = ServiceSqlitePaths::from_runtime_context(&runtime_context(root.path())) .expect("SQLite paths"); @@ -895,7 +970,7 @@ mod tests { &identity, &migrations, &schema, - ServiceSqliteConnectionOptions::reviewed(), + options, authority, MigrationAppliedAtUnixSeconds::new(1_700_000_000).expect("migration time"), &build_identity(), @@ -923,6 +998,436 @@ mod tests { } #[cfg(any(target_os = "linux", target_os = "macos"))] + #[derive(Debug, PartialEq, Eq)] + struct StateFileSnapshot { + bytes: Vec<u8>, + length: u64, + modified: SystemTime, + mode: u32, + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + fn state_directory_snapshot(directory: &Path) -> BTreeMap<String, StateFileSnapshot> { + fs::read_dir(directory) + .expect("read state directory") + .map(|entry| { + let entry = entry.expect("state entry"); + let name = entry.file_name().into_string().expect("UTF-8 state entry"); + let metadata = entry.metadata().expect("state entry metadata"); + ( + name, + StateFileSnapshot { + bytes: fs::read(entry.path()).expect("state entry bytes"), + length: metadata.len(), + modified: metadata.modified().expect("state entry mtime"), + mode: metadata.permissions().mode() & 0o777, + }, + ) + }) + .collect() + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn close_drains_admitted_work_rejects_new_work_and_is_idempotent() { + let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; + let host = Arc::new(host); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let transaction = tokio::spawn({ + let host = Arc::clone(&host); + let entered = Arc::clone(&entered); + let release = Arc::clone(&release); + async move { + host.transaction::<i64, Infallible, _>(|transaction| { + Box::pin(async move { + let count = sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *transaction) + .await + .expect("read in admitted transaction"); + entered.notify_one(); + release.notified().await; + Ok(count) + }) + }) + .await + } + }); + entered.notified().await; + + let first_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 second_close = tokio::spawn({ + let host = Arc::clone(&host); + async move { host.close().await } + }); + let rejected = host + .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) + .await + .expect_err("close admission must reject new work"); + assert_eq!( + rejected.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + assert_eq!( + rejected.sqlite_error().map(ServiceSqliteError::kind), + Some(ServiceSqliteErrorKind::Open) + ); + let contended = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); + assert!(matches!( + contended, + Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority + )); + + release.notify_one(); + assert_eq!( + transaction + .await + .expect("transaction task joins") + .expect("admitted transaction"), + 0 + ); + first_close + .await + .expect("first close task joins") + .expect("first close succeeds"); + second_close + .await + .expect("second close task joins") + .expect("concurrent close is idempotent"); + host.close().await.expect("sequential close is idempotent"); + + let mut next = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("authority can be reacquired") + .expect("writer mode yields authority"); + next.release().expect("release reacquired authority"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn cancelled_close_retains_authority_and_retry_finishes() { + let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; + let host = Arc::new(host); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let transaction = tokio::spawn({ + let host = Arc::clone(&host); + let entered = Arc::clone(&entered); + let release = Arc::clone(&release); + async move { + host.transaction::<(), Infallible, _>(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *transaction) + .await + .expect("admitted read"); + entered.notify_one(); + release.notified().await; + Ok(()) + }) + }) + .await + } + }); + entered.notified().await; + let close_task = tokio::spawn({ + let host = Arc::clone(&host); + async move { host.close().await } + }); + while !host.closing.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + close_task.abort(); + assert!( + close_task + .await + .expect_err("close task is cancelled") + .is_cancelled() + ); + let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); + assert!(matches!( + retained, + Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority + )); + let rejected = host + .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) + .await + .expect_err("cancelled close remains non-admitting"); + assert_eq!( + rejected.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + + release.notify_one(); + transaction + .await + .expect("transaction task joins") + .expect("admitted transaction finishes"); + host.close().await.expect("close retry succeeds"); + assert!( + WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("authority reacquisition after retry") + .is_some() + ); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn writable_close_checkpoints_and_read_only_close_is_side_effect_free() { + let (_root, paths, identity, migrations, schema, host) = initialized_host().await; + host.transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (41)") + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect("write WAL frame"); + host.close().await.expect("writable close checkpoints"); + let state_directory = paths.state_database().parent().expect("state directory"); + assert!(!state_directory.join("state.sqlite-wal").exists()); + assert!(!state_directory.join("state.sqlite-shm").exists()); + let before = state_directory_snapshot(state_directory); + + let inspection = ServiceSqliteHost::open_read_only_inspection( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + ) + .await + .expect("open read-only inspection"); + assert_eq!(row_count(&inspection).await, 1); + inspection + .close() + .await + .expect("close read-only inspection"); + assert_eq!(state_directory_snapshot(state_directory), before); + + let (reopened, 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("reopen writer after read-only close"); + assert_eq!(outcome.applied_count(), 0); + assert_eq!(row_count(&reopened).await, 1); + reopened.close().await.expect("close reopened writer"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn read_only_close_revalidates_after_drain_before_releasing_stale_authority() { + let (_root, paths, identity, migrations, schema, writer) = initialized_host().await; + writer.close().await.expect("close writer host"); + let inspection = Arc::new( + ServiceSqliteHost::open_read_only_inspection( + &paths, + &identity, + &migrations, + &schema, + ServiceSqliteConnectionOptions::reviewed(), + ) + .await + .expect("open read-only inspection"), + ); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let transaction = tokio::spawn({ + let inspection = Arc::clone(&inspection); + let entered = Arc::clone(&entered); + let release = Arc::clone(&release); + async move { + inspection + .transaction::<(), Infallible, _>(|transaction| { + Box::pin(async move { + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *transaction) + .await + .expect("read through retained inspection"); + entered.notify_one(); + release.notified().await; + Ok(()) + }) + }) + .await + } + }); + entered.notified().await; + let close_task = tokio::spawn({ + let inspection = Arc::clone(&inspection); + async move { inspection.close().await } + }); + while !inspection.closing.load(Ordering::Acquire) { + tokio::task::yield_now().await; + } + + let retired_lock = paths + .state_lock() + .parent() + .expect("state directory") + .join("retired-inspection-close.lock"); + fs::rename(paths.state_lock(), retired_lock).expect("replace inspection lock"); + let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("replacement acquisition") + .expect("replacement authority"); + release.notify_one(); + let transaction_error = transaction + .await + .expect("inspection transaction task joins") + .expect_err("binding drift revokes admitted inspection"); + assert_eq!( + transaction_error.kind(), + ServiceSqliteTransactionErrorKind::NotCommitted + ); + let close_error = close_task + .await + .expect("inspection close task joins") + .expect_err("close must report stale inspection authority"); + assert_eq!(close_error.kind(), ServiceSqliteErrorKind::Authority); + assert_eq!( + inspection + .close() + .await + .expect_err("terminal authority result is cached") + .kind(), + ServiceSqliteErrorKind::Authority + ); + replacement + .release() + .expect("release replacement authority"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn cancelled_checkpoint_resumes_to_terminal_error_and_releases_authority() { + let options = ServiceSqliteConnectionOptions::new(Duration::from_secs(1), 1) + .expect("short reviewed limits"); + let (_root, paths, _identity, _migrations, _schema, host) = + initialized_host_with_options(options).await; + let host = Arc::new(host); + host.transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (1)") + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect("seed reader snapshot"); + + let mut reader = SqliteConnection::connect_with( + &SqliteConnectOptions::new() + .filename(paths.state_database()) + .read_only(true) + .create_if_missing(false), + ) + .await + .expect("open external reader"); + let mut reader_transaction = reader.begin().await.expect("begin external read"); + assert_eq!( + sqlx::query_scalar::<_, i64>("SELECT COUNT(*) FROM host_probe") + .fetch_one(&mut *reader_transaction) + .await + .expect("establish reader snapshot"), + 1 + ); + host.transaction(|transaction| { + Box::pin(async move { + sqlx::query("INSERT INTO host_probe (value) VALUES (2)") + .execute(&mut *transaction) + .await + .map(|_| ()) + }) + }) + .await + .expect("append frame after reader snapshot"); + + let close_task = tokio::spawn({ + let host = Arc::clone(&host); + async move { host.close().await } + }); + while host.pool.close_phase() != crate::open::TEST_CLOSE_PHASE_CHECKPOINT { + tokio::task::yield_now().await; + } + close_task.abort(); + assert!( + close_task + .await + .expect_err("checkpoint close task is cancelled") + .is_cancelled() + ); + let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); + assert!(matches!( + retained, + Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority + )); + + let first = host + .close() + .await + .expect_err("active reader prevents TRUNCATE"); + assert_eq!(first.kind(), ServiceSqliteErrorKind::Pragma); + assert!(!first.to_string().contains("state.sqlite")); + let repeated = host.close().await.expect_err("terminal result is cached"); + assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Pragma); + let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("close releases authority despite checkpoint failure") + .expect("writer mode yields authority"); + replacement + .release() + .expect("release replacement authority"); + + reader_transaction + .rollback() + .await + .expect("release external snapshot"); + reader.close().await.expect("close external reader"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn close_authority_drift_is_cached_and_releases_the_stale_lock() { + let (_root, paths, _identity, _migrations, _schema, host) = initialized_host().await; + let retired_lock = paths + .state_lock() + .parent() + .expect("state directory") + .join("retired-close-state.lock"); + fs::rename(paths.state_lock(), retired_lock).expect("replace canonical lock"); + let mut replacement = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("replacement acquisition") + .expect("replacement authority"); + + let first = host + .close() + .await + .expect_err("close detects authority drift"); + assert_eq!(first.kind(), ServiceSqliteErrorKind::Authority); + let repeated = host.close().await.expect_err("authority result is cached"); + assert_eq!(repeated.kind(), ServiceSqliteErrorKind::Authority); + assert!(!format!("{host:?}").contains(paths.state_database().to_string_lossy().as_ref())); + replacement + .release() + .expect("release replacement authority"); + } + + #[cfg(any(target_os = "linux", target_os = "macos"))] #[tokio::test] async fn typed_execution_commits_and_operation_failure_rolls_back() { let (_root, _paths, _identity, _migrations, _schema, host) = initialized_host().await; @@ -1147,7 +1652,7 @@ mod tests { #[tokio::test] async fn inspection_lock_replacement_and_new_writer_revoke_live_inspection() { let (_root, paths, identity, migrations, schema, host) = initialized_host().await; - host.close_pool().await; + host.close().await.expect("close writer host"); drop(host); let inspection = ServiceSqliteHost::open_read_only_inspection( &paths, @@ -1230,7 +1735,7 @@ mod tests { #[tokio::test] async fn read_only_writes_roll_back_and_internal_pool_close_refuses_new_work() { let (root, paths, identity, migrations, schema, host) = initialized_host().await; - host.close_pool().await; + host.close().await.expect("close writer host"); let closed = host .transaction(|_| Box::pin(async { Ok::<_, Infallible>(()) })) .await diff --git a/crates/service_sqlite/src/open.rs b/crates/service_sqlite/src/open.rs @@ -10,13 +10,15 @@ use std::{ use std::{ fs::File, sync::{ - Arc, + Arc, Mutex, atomic::{AtomicBool, Ordering}, }, }; #[cfg(any(target_os = "linux", target_os = "macos"))] use fs2::FileExt; +#[cfg(any(target_os = "linux", target_os = "macos"))] +use futures::future::BoxFuture; use radroots_runtime_paths::{ InstanceId, RuntimeContext, ServiceId, default_service_instance_artifacts, }; @@ -190,8 +192,11 @@ pub(crate) struct PrivateConnectionPool { catalog: MigrationCatalog, schema_catalog: SchemaCatalog, mode: OpenMode, - authority: Option<WriterAuthority>, - inspection_guard: Option<ReadOnlyInspectionGuard>, + policy: ServiceSqliteConnectionOptions, + resources: Mutex<PrivateConnectionResources>, + close_driver: tokio::sync::Mutex<PrivateCloseDriver>, + #[cfg(test)] + close_phase: std::sync::atomic::AtomicU8, authority_failure: Arc<AtomicBool>, metadata_failure: Arc<AtomicBool>, migration_failure: Arc<AtomicBool>, @@ -200,6 +205,28 @@ pub(crate) struct PrivateConnectionPool { } #[cfg(any(target_os = "linux", target_os = "macos"))] +struct PrivateConnectionResources { + authority: Option<WriterAuthority>, + inspection_guard: Option<ReadOnlyInspectionGuard>, +} + +#[cfg(any(target_os = "linux", target_os = "macos"))] +enum PrivateCloseDriver { + Pending, + Connecting(BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>>), + Connected(SqliteConnection), + Closing { + future: BoxFuture<'static, Result<(), ServiceSqliteError>>, + authority_error: Option<ServiceSqliteError>, + checkpoint_error: Option<ServiceSqliteError>, + }, + Complete(Option<ServiceSqliteErrorKind>), +} + +#[cfg(test)] +pub(crate) const TEST_CLOSE_PHASE_CHECKPOINT: u8 = 4; + +#[cfg(any(target_os = "linux", target_os = "macos"))] impl PrivateConnectionPool { fn connection_failure_kind(&self) -> ServiceSqliteErrorKind { connection_failure_kind( @@ -212,11 +239,28 @@ impl PrivateConnectionPool { } pub(crate) fn validate(&self) -> Result<(), ServiceSqliteError> { - if let Some(authority) = self.authority.as_ref() { - authority.validate_for(&self.paths)?; - } - if let Some(inspection_guard) = self.inspection_guard.as_ref() { - inspection_guard.validate_for(&self.paths)?; + let resources = self.resources.lock().map_err(|_| { + connection_error( + ServiceSqliteErrorKind::Authority, + ConnectionFailureKind::AuthorityMismatch, + ) + })?; + match self.mode { + OpenMode::Initialize | OpenMode::ReadWriteExisting => resources + .authority + .as_ref() + .ok_or_else(|| { + connection_error( + ServiceSqliteErrorKind::Authority, + ConnectionFailureKind::AuthorityMismatch, + ) + })? + .validate_for(&self.paths)?, + OpenMode::ReadOnlyInspection => resources + .inspection_guard + .as_ref() + .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))? + .validate_for(&self.paths)?, } self.binding.validate(&self.paths) } @@ -237,6 +281,11 @@ impl PrivateConnectionPool { &self.schema_catalog } + #[cfg(test)] + pub(crate) fn close_phase(&self) -> u8 { + self.close_phase.load(Ordering::Acquire) + } + pub(crate) async fn acquire(&self) -> Result<PoolConnection<Sqlite>, ServiceSqliteError> { self.validate()?; let result = self.pool.acquire().await; @@ -261,25 +310,16 @@ impl PrivateConnectionPool { build: &MigrationBuildIdentity, callbacks: &[crate::migration::MigrationCallbackBinding], ) -> Result<crate::migration::MigrationApplicationOutcome, ServiceSqliteError> { - let authority = self - .authority - .as_ref() - .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Authority))?; - authority.validate_for(&self.paths)?; - self.binding.validate(&self.paths)?; + self.validate_writer_authority()?; let acquired = self.pool.acquire().await; - authority.validate_for(&self.paths)?; - self.binding.validate(&self.paths)?; + self.validate_writer_authority()?; let mut connection = acquired.map_err(|source| connection_source(self.connection_failure_kind(), source))?; // Migration execution installs connection-local fail-closed guards. Always // discard this one-time connection so cancellation cannot return a guarded // or callback-altered handle to the pool. connection.close_on_drop(); - let mut validate_authority = || { - authority.validate_for(&self.paths)?; - self.binding.validate(&self.paths) - }; + let mut validate_authority = || self.validate_writer_authority(); let result = crate::migration::apply_governed_migrations( &mut connection, &self.catalog, @@ -290,20 +330,204 @@ impl PrivateConnectionPool { &mut validate_authority, ) .await; - authority.validate_for(&self.paths)?; - self.binding.validate(&self.paths)?; + self.validate_writer_authority()?; result } - pub(crate) async fn close(mut self) -> Option<WriterAuthority> { + fn validate_writer_authority(&self) -> Result<(), ServiceSqliteError> { + let resources = self.resources.lock().map_err(|_| { + connection_error( + ServiceSqliteErrorKind::Authority, + ConnectionFailureKind::AuthorityMismatch, + ) + })?; + resources + .authority + .as_ref() + .ok_or_else(|| ServiceSqliteError::new(ServiceSqliteErrorKind::Authority))? + .validate_for(&self.paths)?; + self.binding.validate(&self.paths) + } + + pub(crate) async fn close(self) -> Option<WriterAuthority> { self.pool.close().await; - self.inspection_guard.take(); - self.authority.take() + let mut resources = match self.resources.into_inner() { + Ok(resources) => resources, + Err(poisoned) => poisoned.into_inner(), + }; + resources.inspection_guard.take(); + resources.authority.take() } - #[cfg(test)] - pub(crate) async fn close_pool(&self) { + /// The outer error means authority release was not proven and close must retry. + /// The inner result is terminal and may be cached by the host. + pub(crate) async fn close_explicit( + &self, + ) -> Result<Result<(), ServiceSqliteError>, ServiceSqliteError> { self.pool.close().await; + let terminal = if self.mode.requires_writer_authority() { + self.drive_writable_close().await + } else { + self.validate() + }; + self.release_resources()?; + Ok(terminal) + } + + async fn drive_writable_close(&self) -> Result<(), ServiceSqliteError> { + let mut driver = self.close_driver.lock().await; + loop { + match &mut *driver { + PrivateCloseDriver::Pending => { + if let Err(error) = self.validate_writer_authority() { + #[cfg(test)] + self.close_phase.store(6, Ordering::Release); + *driver = PrivateCloseDriver::Complete(Some(error.kind())); + return Err(error); + } + let options = sqlite_connect_options(&self.paths, self.mode, self.policy); + let connect: BoxFuture<'static, Result<SqliteConnection, ServiceSqliteError>> = + Box::pin(async move { + SqliteConnection::connect_with(&options) + .await + .map_err(|source| { + connection_source(ServiceSqliteErrorKind::Open, source) + }) + }); + #[cfg(test)] + self.close_phase.store(1, Ordering::Release); + *driver = PrivateCloseDriver::Connecting(connect); + } + PrivateCloseDriver::Connecting(connect) => { + let connected = connect.as_mut().await; + match connected { + Ok(connection) => { + #[cfg(test)] + self.close_phase.store(2, Ordering::Release); + *driver = PrivateCloseDriver::Connected(connection); + } + Err(error) => { + let authority_error = self.validate_writer_authority().err(); + let error = authority_error.unwrap_or(error); + #[cfg(test)] + self.close_phase.store(6, Ordering::Release); + *driver = PrivateCloseDriver::Complete(Some(error.kind())); + return Err(error); + } + } + } + PrivateCloseDriver::Connected(connection) => { + let mut authority_error = self.validate_writer_authority().err(); + let mut checkpoint_error = None; + if authority_error.is_none() { + #[cfg(test)] + self.close_phase.store(3, Ordering::Release); + checkpoint_error = + verify_connection_policy(connection, self.mode, self.policy) + .await + .map_err(|source| { + connection_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() { + #[cfg(test)] + self.close_phase + .store(TEST_CLOSE_PHASE_CHECKPOINT, Ordering::Release); + checkpoint_error = + sqlx::query_as::<_, (i64, i64, i64)>("PRAGMA wal_checkpoint(TRUNCATE)") + .fetch_one(&mut *connection) + .await + .map_err(|source| { + connection_source(ServiceSqliteErrorKind::Pragma, source) + }) + .and_then(|(busy, _log_frames, _checkpointed_frames)| { + if busy == 0 { + Ok(()) + } else { + Err(connection_error( + ServiceSqliteErrorKind::Pragma, + ConnectionFailureKind::CheckpointBusy, + )) + } + }) + .err(); + authority_error = + authority_error.or_else(|| self.validate_writer_authority().err()); + } + + let connected = core::mem::replace(&mut *driver, PrivateCloseDriver::Pending); + let PrivateCloseDriver::Connected(connection) = connected else { + unreachable!("close driver retains its connected phase") + }; + let close: BoxFuture<'static, Result<(), ServiceSqliteError>> = + Box::pin(async move { + connection.close().await.map_err(|source| { + connection_source(ServiceSqliteErrorKind::Open, source) + }) + }); + #[cfg(test)] + self.close_phase.store(5, Ordering::Release); + *driver = PrivateCloseDriver::Closing { + future: close, + authority_error, + checkpoint_error, + }; + } + PrivateCloseDriver::Closing { + future, + authority_error, + checkpoint_error, + } => { + let close_error = future.as_mut().await.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); + #[cfg(test)] + self.close_phase.store(6, Ordering::Release); + *driver = + PrivateCloseDriver::Complete(error.as_ref().map(ServiceSqliteError::kind)); + return error.map_or(Ok(()), Err); + } + PrivateCloseDriver::Complete(kind) => { + return kind.map_or(Ok(()), |kind| Err(ServiceSqliteError::new(kind))); + } + } + } + } + + fn release_resources(&self) -> Result<(), ServiceSqliteError> { + let mut resources = self.resources.lock().map_err(|_| { + connection_error( + ServiceSqliteErrorKind::Authority, + ConnectionFailureKind::AuthorityMismatch, + ) + })?; + match self.mode { + OpenMode::Initialize | OpenMode::ReadWriteExisting => { + let authority = resources.authority.as_mut().ok_or_else(|| { + connection_error( + ServiceSqliteErrorKind::Authority, + ConnectionFailureKind::AuthorityMismatch, + ) + })?; + authority.release()?; + resources.authority.take(); + } + OpenMode::ReadOnlyInspection => { + let inspection = resources.inspection_guard.as_mut().ok_or_else(|| { + inspection_error(ConnectionFailureKind::InspectionUnavailable) + })?; + inspection.release()?; + resources.inspection_guard.take(); + } + } + Ok(()) } } @@ -666,8 +890,14 @@ async fn open_connection_pool( catalog: retained_catalog, schema_catalog: retained_schema_catalog, mode, - authority, - inspection_guard, + policy, + resources: Mutex::new(PrivateConnectionResources { + authority, + inspection_guard, + }), + close_driver: tokio::sync::Mutex::new(PrivateCloseDriver::Pending), + #[cfg(test)] + close_phase: std::sync::atomic::AtomicU8::new(0), authority_failure, metadata_failure, migration_failure, @@ -787,6 +1017,7 @@ enum ConnectionFailureKind { AuthorityMismatch, InspectionUnavailable, InspectionContended, + CheckpointBusy, } #[cfg(any(target_os = "linux", target_os = "macos"))] @@ -797,6 +1028,7 @@ impl fmt::Display for ConnectionFailureKind { Self::AuthorityMismatch => "SQLite writer authority is missing or mismatched", Self::InspectionUnavailable => "SQLite inspection authority is unavailable", Self::InspectionContended => "SQLite inspection requires an offline writer", + Self::CheckpointBusy => "SQLite close checkpoint could not drain active readers", }) } } @@ -827,7 +1059,7 @@ fn connection_source(kind: ServiceSqliteErrorKind, cause: sqlx::Error) -> Servic #[cfg(any(target_os = "linux", target_os = "macos"))] struct ReadOnlyInspectionGuard { - lock: File, + lock: Option<File>, lock_device: u64, lock_inode: u64, directory: File, @@ -947,7 +1179,7 @@ impl ReadOnlyInspectionGuard { )); } Ok(Self { - lock, + lock: Some(lock), lock_device: u64::try_from(lock_status.st_dev) .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?, lock_inode: lock_status.st_ino, @@ -1009,7 +1241,11 @@ impl ReadOnlyInspectionGuard { .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; let lock_status = fstat(&lock) .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; - let held_lock_status = fstat(&self.lock) + let held_lock = self + .lock + .as_ref() + .ok_or_else(|| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + let held_lock_status = fstat(held_lock) .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; let lock_device = u64::try_from(lock_status.st_dev) .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; @@ -1044,12 +1280,22 @@ impl ReadOnlyInspectionGuard { } Ok(()) } + + fn release(&mut self) -> Result<(), ServiceSqliteError> { + let Some(lock) = self.lock.as_ref() else { + return Ok(()); + }; + FileExt::unlock(lock) + .map_err(|_| inspection_error(ConnectionFailureKind::InspectionUnavailable))?; + self.lock.take(); + Ok(()) + } } #[cfg(any(target_os = "linux", target_os = "macos"))] impl Drop for ReadOnlyInspectionGuard { fn drop(&mut self) { - let _ = FileExt::unlock(&self.lock); + let _ = self.release(); } } @@ -1280,7 +1526,10 @@ mod tests { fs, num::NonZeroU32, os::unix::fs::{PermissionsExt, symlink}, - sync::atomic::{AtomicUsize, Ordering as AtomicOrdering}, + sync::{ + Arc, + atomic::{AtomicUsize, Ordering as AtomicOrdering}, + }, time::{Duration, SystemTime}, }; @@ -1293,6 +1542,8 @@ mod tests { #[cfg(any(target_os = "linux", target_os = "macos"))] use crate::{ServiceDatabaseMetadata, ServiceSqliteApplicationId}; + #[cfg(any(target_os = "linux", target_os = "macos"))] + use tokio::sync::Notify; use super::*; @@ -2439,7 +2690,72 @@ mod tests { ) .await .expect("read-only inspection"); - assert!(read_only.authority.is_none()); + assert_eq!(read_only.mode(), OpenMode::ReadOnlyInspection); assert!(read_only.close().await.is_none()); } + + #[cfg(any(target_os = "linux", target_os = "macos"))] + #[tokio::test] + async fn cancellation_during_explicit_connection_close_retains_the_close_driver() { + let directory = tempfile::tempdir().expect("temporary directory"); + let policy = ServiceSqliteConnectionOptions::reviewed(); + let (paths, pool) = initialized_pool(directory.path(), policy).await; + let connection = SqliteConnection::connect_with(&sqlite_connect_options( + &paths, + OpenMode::Initialize, + policy, + )) + .await + .expect("open private checkpoint connection"); + let entered = Arc::new(Notify::new()); + let release = Arc::new(Notify::new()); + let close: BoxFuture<'static, Result<(), ServiceSqliteError>> = Box::pin({ + let entered = Arc::clone(&entered); + let release = Arc::clone(&release); + async move { + entered.notify_one(); + release.notified().await; + connection + .close() + .await + .map_err(|source| connection_source(ServiceSqliteErrorKind::Open, source)) + } + }); + { + let mut driver = pool.close_driver.lock().await; + *driver = PrivateCloseDriver::Closing { + future: close, + authority_error: None, + checkpoint_error: None, + }; + } + let pool = Arc::new(pool); + let close_task = tokio::spawn({ + let pool = Arc::clone(&pool); + async move { pool.close_explicit().await } + }); + entered.notified().await; + close_task.abort(); + assert!( + close_task + .await + .expect_err("explicit close task is cancelled") + .is_cancelled() + ); + let retained = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting); + assert!(matches!( + retained, + Err(ref error) if error.kind() == ServiceSqliteErrorKind::Authority + )); + + release.notify_one(); + pool.close_explicit() + .await + .expect("authority release is proven") + .expect("retained connection closes explicitly"); + let mut reacquired = WriterAuthority::acquire(&paths, OpenMode::ReadWriteExisting) + .expect("authority reacquisition") + .expect("writer mode yields authority"); + reacquired.release().expect("release reacquired authority"); + } } diff --git a/crates/service_sqlite/tests/package_boundary.rs b/crates/service_sqlite/tests/package_boundary.rs @@ -18,6 +18,7 @@ const TRANSACTION_CONTROL_SOURCE: &str = include_str!("../src/transaction_contro #[test] fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { + let readme_words = README.split_whitespace().collect::<Vec<_>>().join(" "); for required in [ "name = \"radroots_service_sqlite\"", "publish = false", @@ -40,7 +41,8 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "rustix", "serde", "sha2", - "sqlx" + "sqlx", + "tokio" ]) ); assert_eq!( @@ -80,9 +82,22 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "must be treated as an unknown commit outcome", "require rereading authoritative state", "before an idempotent retry", + "Every host must be closed explicitly with `ServiceSqliteHost::close`", + "permanently stops new transaction admission", + "drains transactions that were", + "safe to call sequentially or concurrently", + "fixed `PRAGMA wal_checkpoint(TRUNCATE)` policy", + "requires an unblocked checkpoint", + "releases its shared inspection guard without checkpointing or mutating", + "Cancelling close before terminal completion", + "private connect, checkpoint, and explicit connection-close driver remains host-owned", + "without losing the SQLite handle or its close proof", + "a later call resumes close", + "stable outer result is cached", + "Dropping a host performs no asynchronous close work", ] { assert!( - README.contains(required), + readme_words.contains(required), "Step 061 README contract is missing `{required}`" ); } @@ -91,13 +106,17 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { .map(|(production, _)| production) .expect("authority source must keep tests separated"); let open_production = OPEN_SOURCE - .split_once("#[cfg(test)]") + .split_once("#[cfg(test)]\nmod tests") .map(|(production, _)| production) .expect("open source must keep tests separated"); let migration_production = MIGRATION_SOURCE .split_once("#[cfg(test)]") .map(|(production, _)| production) .expect("migration source must keep tests separated"); + let connection_production = CONNECTION_SOURCE + .split_once("#[cfg(test)]\nmod tests") + .map(|(production, _)| production) + .expect("connection source must keep tests separated"); for required in [ "ServiceSqliteErrorCode", @@ -199,6 +218,10 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "CommitOutcomeUnknown", "cancelling the future yields no result", "must reread authoritative state before any idempotent retry", + "pub async fn close(&self)", + "closing.store(true, Ordering::Release)", + "close_state.lock().await", + "ServiceSqliteHostCloseState::Complete", ] { assert!( CONNECTION_SOURCE.contains(required), @@ -233,14 +256,51 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { "pub async fn commit(", "pub async fn rollback(", "pub use sqlx", + "impl Drop for ServiceSqliteHost", + "pub async fn checkpoint", + "pub fn checkpoint", + "checkpoint_mode", ] { assert!( - !CONNECTION_SOURCE.contains(forbidden) && !ROOT.contains(forbidden), + !connection_production.contains(forbidden) && !ROOT.contains(forbidden), "Step 061 source exposes forbidden raw authority `{forbidden}`" ); } for required in [ + "self.pool.close().await", + "PRAGMA wal_checkpoint(TRUNCATE)", + "if busy == 0", + ".close()", + "authority.release()?", + "inspection.release()?", + "release_resources", + "CheckpointBusy", + "close_driver: tokio::sync::Mutex<PrivateCloseDriver>", + "PrivateCloseDriver::Connecting", + "PrivateCloseDriver::Connected", + "PrivateCloseDriver::Closing", + "connect.as_mut().await", + "future.as_mut().await", + ] { + assert!( + open_production.contains(required), + "Step 062 close source is missing `{required}`" + ); + } + for forbidden in [ + "tokio::spawn", + "Runtime::new", + "pub enum Checkpoint", + "pub struct Checkpoint", + ] { + assert!( + !connection_production.contains(forbidden) && !open_production.contains(forbidden), + "Step 062 close source contains forbidden authority `{forbidden}`" + ); + } + + for required in [ "radroots.service_sqlite.schema_object.v1\\0", "radroots.service_sqlite.schema_snapshot.v1\\0", "radroots.service_sqlite.schema_catalog.v1\\0", @@ -483,7 +543,8 @@ fn service_sqlite_is_unpublished_lint_governed_and_dependency_bounded() { } for required in [ - "inspection_guard.validate_for(&self.paths)", + ".inspection_guard", + ".validate_for(&self.paths)?", "fn validate_for(&self, paths: &ServiceSqlitePaths)", "held_lock_status", "lock_device",