commit 5f0cb06e5bda866f675567180275220d464a7908
parent 4ffe755de8ae07ad51b3e6b8e1e0d970562a5020
Author: triesap <tyson@radroots.org>
Date: Sun, 2 Aug 2026 19:57:19 +0000
storage-sqlite: implement explicit close and passive status
- share shutdown, integrity, configuration, and writer authority across backend clones
- close both owned pools before releasing writable authority and reporting completion
- expose passive storage and integrity status without initiating hidden maintenance
- prove closing visibility, closed-operation rejection, concurrency, and idempotent release
Diffstat:
5 files changed, 301 insertions(+), 11 deletions(-)
diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs
@@ -13,14 +13,15 @@ use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool};
use std::sync::Arc;
use crate::lock::WriterLock;
+use crate::status::StorageLifecycle;
#[derive(Clone)]
pub struct SqliteStorage {
- pool: SqlitePool,
- private_pool: SqlitePool,
- generation: SourceGeneration,
- mode: EventStoreMode,
- _writer_lock: Option<Arc<WriterLock>>,
+ pub(crate) pool: SqlitePool,
+ pub(crate) private_pool: SqlitePool,
+ pub(crate) generation: SourceGeneration,
+ pub(crate) mode: EventStoreMode,
+ pub(crate) lifecycle: Arc<StorageLifecycle>,
}
struct StoredEventRow {
@@ -41,7 +42,7 @@ impl SqliteStorage {
pool,
generation,
mode,
- _writer_lock: None,
+ lifecycle: Arc::new(StorageLifecycle::scaffold(mode)),
}
}
@@ -57,7 +58,7 @@ impl SqliteStorage {
private_pool,
generation,
mode,
- _writer_lock: None,
+ lifecycle: Arc::new(StorageLifecycle::scaffold(mode)),
}
}
@@ -66,6 +67,8 @@ impl SqliteStorage {
private_pool: SqlitePool,
generation: SourceGeneration,
mode: EventStoreMode,
+ open_mode: crate::OpenMode,
+ busy_timeout: std::time::Duration,
writer_lock: Option<WriterLock>,
) -> Self {
Self {
@@ -73,7 +76,7 @@ impl SqliteStorage {
private_pool,
generation,
mode,
- _writer_lock: writer_lock.map(Arc::new),
+ lifecycle: Arc::new(StorageLifecycle::new(open_mode, busy_timeout, writer_lock)),
}
}
diff --git a/crates/storage_sqlite/src/integrity.rs b/crates/storage_sqlite/src/integrity.rs
@@ -1 +1,20 @@
-//! SQLite integrity validation boundary.
+//! Passive SQLite integrity reporting.
+
+use radroots_storage::{
+ Error,
+ status::{IntegrityHealth, IntegrityStatus},
+};
+
+use crate::SqliteStorage;
+
+impl SqliteStorage {
+ /// Returns the last recorded integrity result without running maintenance,
+ /// querying SQLite pragmas, or mutating either owned database.
+ pub async fn integrity(&self) -> Result<IntegrityStatus, Error> {
+ self.lifecycle.integrity()
+ }
+}
+
+pub(crate) fn unknown() -> Result<IntegrityStatus, Error> {
+ IntegrityStatus::new(IntegrityHealth::Unknown, None, 0, 0)
+}
diff --git a/crates/storage_sqlite/src/lock.rs b/crates/storage_sqlite/src/lock.rs
@@ -37,7 +37,6 @@ impl WriterLock {
/// Explicitly releases the writer lock for the later asynchronous close
/// lifecycle. Dropping the guard remains a fail-safe release path.
- #[allow(dead_code)] // Used by the explicit close lifecycle in its ordered RCL checkpoint.
pub(crate) fn release(self) -> Result<(), Error> {
FileExt::unlock(&self.file).map_err(|source| Error::WriterUnlockFailed {
path: self.path.clone(),
diff --git a/crates/storage_sqlite/src/open.rs b/crates/storage_sqlite/src/open.rs
@@ -480,6 +480,8 @@ impl SqliteStorage {
} else {
EventStoreMode::ReadOnly
},
+ options.mode(),
+ options.busy_timeout(),
writer_lock,
))
}
diff --git a/crates/storage_sqlite/src/status.rs b/crates/storage_sqlite/src/status.rs
@@ -1 +1,268 @@
-//! Passive SQLite storage status boundary.
+//! Passive SQLite storage status and explicit close lifecycle.
+
+use std::{
+ sync::{
+ Mutex, RwLock,
+ atomic::{AtomicU8, Ordering},
+ },
+ time::Duration,
+};
+
+use radroots_storage::{
+ Error,
+ status::{
+ IntegrityStatus, ShutdownState, StorageBackend, StorageOpenMode, StorageStatus,
+ WriterPolicy,
+ },
+};
+
+use crate::{OpenMode, SqliteStorage, integrity, lock::WriterLock};
+
+const OPEN: u8 = 0;
+const CLOSING: u8 = 1;
+const CLOSED: u8 = 2;
+
+pub(crate) struct StorageLifecycle {
+ open_mode: StorageOpenMode,
+ writer_policy: WriterPolicy,
+ wal_enabled: bool,
+ busy_timeout: Duration,
+ shutdown: AtomicU8,
+ integrity: RwLock<Option<IntegrityStatus>>,
+ writer_lock: Mutex<Option<WriterLock>>,
+}
+
+impl StorageLifecycle {
+ pub(crate) fn new(
+ mode: OpenMode,
+ busy_timeout: Duration,
+ writer_lock: Option<WriterLock>,
+ ) -> Self {
+ Self {
+ open_mode: storage_open_mode(mode),
+ writer_policy: if mode.is_writable() {
+ WriterPolicy::AdvisoryProcessLock
+ } else {
+ WriterPolicy::NoWriter
+ },
+ wal_enabled: mode.is_writable(),
+ busy_timeout,
+ shutdown: AtomicU8::new(OPEN),
+ integrity: RwLock::new(None),
+ writer_lock: Mutex::new(writer_lock),
+ }
+ }
+
+ pub(crate) fn scaffold(mode: radroots_storage::status::EventStoreMode) -> Self {
+ let open_mode = match mode {
+ radroots_storage::status::EventStoreMode::ReadOnly => OpenMode::ReadOnly,
+ radroots_storage::status::EventStoreMode::ReadWrite => OpenMode::Create,
+ };
+ Self::new(open_mode, Duration::from_secs(5), None)
+ }
+
+ pub(crate) fn integrity(&self) -> Result<IntegrityStatus, Error> {
+ self.integrity
+ .read()
+ .map_err(|_| Error::BackendUnavailable)?
+ .map_or_else(integrity::unknown, Ok)
+ }
+
+ fn shutdown(&self) -> ShutdownState {
+ match self.shutdown.load(Ordering::Acquire) {
+ OPEN => ShutdownState::Open,
+ CLOSING => ShutdownState::Closing,
+ _ => ShutdownState::Closed,
+ }
+ }
+
+ fn begin_close(&self) {
+ let _ = self
+ .shutdown
+ .compare_exchange(OPEN, CLOSING, Ordering::AcqRel, Ordering::Acquire);
+ }
+
+ fn finish_close(&self) -> Result<(), Error> {
+ let mut writer_lock = self
+ .writer_lock
+ .lock()
+ .map_err(|_| Error::BackendUnavailable)?;
+ let release_result = writer_lock
+ .take()
+ .map(WriterLock::release)
+ .transpose()
+ .map(|_| ())
+ .map_err(|_| Error::BackendUnavailable);
+ self.shutdown.store(CLOSED, Ordering::Release);
+ release_result
+ }
+
+ fn status(&self) -> Result<StorageStatus, Error> {
+ StorageStatus::new(
+ StorageBackend::Sqlite,
+ self.open_mode,
+ self.writer_policy,
+ self.shutdown(),
+ self.integrity()?,
+ self.wal_enabled,
+ u32::try_from(self.busy_timeout.as_millis())
+ .map_err(|_| Error::InvalidStorageStatus)?,
+ )
+ }
+}
+
+impl SqliteStorage {
+ /// Returns backend-level status without opening a connection or initiating
+ /// integrity checks, checkpoints, migrations, or other maintenance.
+ pub async fn storage_status(&self) -> Result<StorageStatus, Error> {
+ self.lifecycle.status()
+ }
+
+ /// Closes both pools, releases writable authority, and returns final
+ /// passive status. Repeated and concurrent calls are idempotent.
+ pub async fn close(&self) -> Result<StorageStatus, Error> {
+ self.lifecycle.begin_close();
+ self.pool.close().await;
+ self.private_pool.close().await;
+ self.lifecycle.finish_close()?;
+ self.lifecycle.status()
+ }
+}
+
+const fn storage_open_mode(mode: OpenMode) -> StorageOpenMode {
+ match mode {
+ OpenMode::ReadOnly => StorageOpenMode::ReadOnly,
+ OpenMode::ReadWriteExisting => StorageOpenMode::ReadWriteExisting,
+ OpenMode::Create => StorageOpenMode::Create,
+ }
+}
+
+#[cfg(test)]
+mod tests {
+ use std::time::Duration;
+
+ use radroots_storage::{
+ EventStore,
+ event::SourceGeneration,
+ status::{IntegrityHealth, ShutdownState, StorageBackend, StorageOpenMode, WriterPolicy},
+ };
+
+ use crate::{OpenOptions, Paths};
+
+ use super::*;
+
+ fn generation(byte: u8) -> SourceGeneration {
+ SourceGeneration::new([byte; 32]).expect("source generation")
+ }
+
+ async fn create(directory: &std::path::Path) -> (Paths, SqliteStorage) {
+ let paths = Paths::from_directory(directory).expect("owned paths");
+ let store = SqliteStorage::open(
+ OpenOptions::new(paths.clone(), OpenMode::Create)
+ .with_busy_timeout(Duration::from_millis(250))
+ .expect("busy timeout")
+ .with_source_generation(generation(73), 7_300)
+ .expect("source generation"),
+ )
+ .await
+ .expect("create storage");
+ (paths, store)
+ }
+
+ #[tokio::test]
+ async fn status_and_integrity_are_passive_and_report_governed_configuration() {
+ let directory = tempfile::tempdir().expect("temporary directory");
+ let (paths, store) = create(directory.path()).await;
+
+ let integrity = store.integrity().await.expect("integrity status");
+ assert_eq!(integrity.health(), IntegrityHealth::Unknown);
+ assert_eq!(integrity.checked_at_unix_ms(), None);
+ assert_eq!(integrity.verified_members(), 0);
+ assert_eq!(integrity.failed_members(), 0);
+
+ let status = store.storage_status().await.expect("storage status");
+ assert_eq!(status.backend(), StorageBackend::Sqlite);
+ assert_eq!(status.open_mode(), StorageOpenMode::Create);
+ assert_eq!(status.writer_policy(), WriterPolicy::AdvisoryProcessLock);
+ assert_eq!(status.shutdown(), ShutdownState::Open);
+ assert_eq!(status.integrity(), integrity);
+ assert!(status.wal_enabled());
+ assert_eq!(status.busy_timeout_ms(), 250);
+
+ let reader = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadOnly))
+ .await
+ .expect("read-only storage");
+ let reader_status = reader.storage_status().await.expect("reader status");
+ assert_eq!(reader_status.open_mode(), StorageOpenMode::ReadOnly);
+ assert_eq!(reader_status.writer_policy(), WriterPolicy::NoWriter);
+ assert!(!reader_status.wal_enabled());
+ assert_eq!(reader_status.busy_timeout_ms(), 5_000);
+ }
+
+ #[tokio::test]
+ async fn close_is_observable_shared_idempotent_and_releases_writable_authority() {
+ let directory = tempfile::tempdir().expect("temporary directory");
+ let (paths, store) = create(directory.path()).await;
+ let clone = store.clone();
+ let held_connection = store.pool.acquire().await.expect("held connection");
+ let mut close = Box::pin(clone.close());
+
+ tokio::select! {
+ biased;
+ result = &mut close => panic!("close completed before the checked-out connection was returned: {result:?}"),
+ () = tokio::task::yield_now() => {}
+ }
+ assert_eq!(
+ store
+ .storage_status()
+ .await
+ .expect("closing status")
+ .shutdown(),
+ ShutdownState::Closing
+ );
+
+ drop(held_connection);
+ assert_eq!(
+ close.await.expect("first close").shutdown(),
+ ShutdownState::Closed
+ );
+ assert_eq!(
+ store
+ .storage_status()
+ .await
+ .expect("shared closed status")
+ .shutdown(),
+ ShutdownState::Closed
+ );
+ assert_eq!(
+ store.close().await.expect("idempotent close").shutdown(),
+ ShutdownState::Closed
+ );
+ assert_eq!(
+ EventStore::status(&store).await,
+ Err(Error::BackendUnavailable)
+ );
+
+ let reopened = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
+ .await
+ .expect("writer authority released before final clone drop");
+ assert_eq!(
+ reopened
+ .storage_status()
+ .await
+ .expect("reopened status")
+ .shutdown(),
+ ShutdownState::Open
+ );
+ let reopened_clone = reopened.clone();
+ let (first, second) = tokio::join!(reopened.close(), reopened_clone.close());
+ assert_eq!(
+ first.expect("concurrent close one").shutdown(),
+ ShutdownState::Closed
+ );
+ assert_eq!(
+ second.expect("concurrent close two").shutdown(),
+ ShutdownState::Closed
+ );
+ }
+}