commit 4ffe755de8ae07ad51b3e6b8e1e0d970562a5020
parent 2a94866a7135c6d2a9ca61e58f62060e8727dab0
Author: triesap <tyson@radroots.org>
Date: Sun, 2 Aug 2026 19:47:15 +0000
storage-sqlite: implement the public open lifecycle
- open distinct runtime and private pools under one governed connection policy
- retain writer locking and forward migrations for the complete backend lifetime
- require explicit source-generation entropy and time before fresh file creation
- prove WAL policy, concurrent readers, clone-safe locking, reuse, and fail-closed input
Diffstat:
8 files changed, 554 insertions(+), 8 deletions(-)
diff --git a/contracts/storage/connection_policy_v1.toml b/contracts/storage/connection_policy_v1.toml
@@ -0,0 +1,12 @@
+schema_version = 1
+databases = ["runtime.sqlite", "private.sqlite"]
+max_connections_per_database = 4
+foreign_keys = true
+journal_mode = "wal"
+synchronous = "full"
+busy_timeout_min_ms = 1
+busy_timeout_default_ms = 5000
+busy_timeout_max_ms = 60000
+fresh_source_generation = "host_supplied_entropy_and_timestamp"
+read_only_migrations = false
+raw_handles_public = false
diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md
@@ -2,6 +2,24 @@
SQLite storage backend for Radroots.
-This package root is established for the Release V1 refactor. Its backend
-implementation and migration authority are introduced in the subsequent
-SQLite storage checkpoints.
+The backend owns separate `runtime.sqlite` and `private.sqlite` files. Writable
+opens hold a process advisory lock, use WAL with bounded busy handling, and
+apply only the governed forward migrations. Fresh stores require a
+host-supplied `SourceGeneration` and creation timestamp; the crate never reads
+hidden entropy or a wall clock.
+
+```rust,no_run
+use radroots_storage::event::SourceGeneration;
+use radroots_storage_sqlite::{OpenMode, OpenOptions, Paths, SqliteStorage};
+
+# async fn open(directory: &std::path::Path) -> Result<(), radroots_storage_sqlite::Error> {
+let paths = Paths::from_directory(directory)?;
+let generation = SourceGeneration::new([7; 32]).expect("non-zero generation");
+let storage = SqliteStorage::open(
+ OpenOptions::new(paths, OpenMode::Create)
+ .with_source_generation(generation, 1_700_000_000_000)?,
+).await?;
+# drop(storage);
+# Ok(())
+# }
+```
diff --git a/crates/storage_sqlite/src/config.rs b/crates/storage_sqlite/src/config.rs
@@ -3,7 +3,7 @@
use std::time::Duration;
use crate::{Error, OpenMode, Paths};
-use radroots_storage::status::WriterPolicy;
+use radroots_storage::{event::SourceGeneration, status::WriterPolicy};
const DEFAULT_BUSY_TIMEOUT: Duration = Duration::from_secs(5);
const MIN_BUSY_TIMEOUT: Duration = Duration::from_millis(1);
@@ -18,6 +18,7 @@ pub struct OpenOptions {
paths: Paths,
mode: OpenMode,
busy_timeout: Duration,
+ source_generation: Option<(SourceGeneration, u64)>,
}
impl OpenOptions {
@@ -27,6 +28,7 @@ impl OpenOptions {
paths,
mode,
busy_timeout: DEFAULT_BUSY_TIMEOUT,
+ source_generation: None,
}
}
@@ -43,6 +45,22 @@ impl OpenOptions {
Ok(self)
}
+ /// Supplies the expected active source generation, or bootstraps it for a
+ /// fresh writable store, without reading hidden entropy or a wall clock.
+ pub fn with_source_generation(
+ mut self,
+ generation: SourceGeneration,
+ created_at_unix_ms: u64,
+ ) -> Result<Self, Error> {
+ if created_at_unix_ms == 0 || i64::try_from(created_at_unix_ms).is_err() {
+ return Err(Error::InvalidSourceGenerationTimestamp {
+ actual: created_at_unix_ms,
+ });
+ }
+ self.source_generation = Some((generation, created_at_unix_ms));
+ Ok(self)
+ }
+
/// Returns the two database paths owned by this backend instance.
pub fn paths(&self) -> &Paths {
&self.paths
@@ -77,6 +95,20 @@ impl OpenOptions {
}
}
+ /// Returns the optional host-supplied source generation expectation.
+ pub fn source_generation(&self) -> Option<SourceGeneration> {
+ self.source_generation.map(|(generation, _)| generation)
+ }
+
+ /// Returns the host-supplied creation time paired with the generation.
+ pub fn source_generation_created_at_unix_ms(&self) -> Option<u64> {
+ self.source_generation.map(|(_, created_at)| created_at)
+ }
+
+ pub(crate) const fn source_generation_bootstrap(&self) -> Option<(SourceGeneration, u64)> {
+ self.source_generation
+ }
+
/// Validates current filesystem state without creating or modifying files.
pub fn validate_filesystem(&self) -> Result<(), Error> {
self.paths.validate_filesystem(self.mode)
diff --git a/crates/storage_sqlite/src/event/mod.rs b/crates/storage_sqlite/src/event/mod.rs
@@ -10,6 +10,9 @@ use radroots_storage::{
status::{EventStoreHealth, EventStoreMode, EventStoreStatus},
};
use sqlx::{QueryBuilder, Row, Sqlite, SqlitePool};
+use std::sync::Arc;
+
+use crate::lock::WriterLock;
#[derive(Clone)]
pub struct SqliteStorage {
@@ -17,6 +20,7 @@ pub struct SqliteStorage {
private_pool: SqlitePool,
generation: SourceGeneration,
mode: EventStoreMode,
+ _writer_lock: Option<Arc<WriterLock>>,
}
struct StoredEventRow {
@@ -26,7 +30,7 @@ struct StoredEventRow {
}
impl SqliteStorage {
- #[allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint.
+ #[allow(dead_code)] // Single-pool in-memory scaffold retained for focused backend tests.
pub(crate) fn new(
pool: SqlitePool,
generation: SourceGeneration,
@@ -37,10 +41,11 @@ impl SqliteStorage {
pool,
generation,
mode,
+ _writer_lock: None,
}
}
- #[allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint.
+ #[allow(dead_code)] // Two-pool in-memory scaffold retained for focused backend tests.
pub(crate) fn with_private_pool(
pool: SqlitePool,
private_pool: SqlitePool,
@@ -52,6 +57,23 @@ impl SqliteStorage {
private_pool,
generation,
mode,
+ _writer_lock: None,
+ }
+ }
+
+ pub(crate) fn from_opened(
+ pool: SqlitePool,
+ private_pool: SqlitePool,
+ generation: SourceGeneration,
+ mode: EventStoreMode,
+ writer_lock: Option<WriterLock>,
+ ) -> Self {
+ Self {
+ pool,
+ private_pool,
+ generation,
+ mode,
+ _writer_lock: writer_lock.map(Arc::new),
}
}
diff --git a/crates/storage_sqlite/src/lock.rs b/crates/storage_sqlite/src/lock.rs
@@ -1,7 +1,5 @@
//! SQLite process and writer locking boundary.
-#![allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint.
-
use crate::{Error, OpenMode, Paths};
use fs2::FileExt;
use std::fs::{self, File, OpenOptions};
@@ -39,6 +37,7 @@ 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
@@ -5,8 +5,16 @@ use std::fmt;
use std::path::{Component, Path, PathBuf};
use std::time::Duration;
+use crate::{OpenOptions, event::SqliteStorage, lock::WriterLock, migration};
+use radroots_storage::{event::SourceGeneration, status::EventStoreMode};
+use sqlx::{
+ ConnectOptions, Connection, Row, SqliteConnection, SqlitePool,
+ sqlite::{SqliteConnectOptions, SqliteJournalMode, SqlitePoolOptions, SqliteSynchronous},
+};
+
const RUNTIME_DATABASE_NAME: &str = "runtime.sqlite";
const PRIVATE_DATABASE_NAME: &str = "private.sqlite";
+const MAX_CONNECTIONS_PER_DATABASE: u32 = 4;
/// Explicit behavior for opening owned SQLite files.
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
@@ -166,6 +174,9 @@ pub enum Error {
maximum: Duration,
actual: Duration,
},
+ InvalidSourceGenerationTimestamp {
+ actual: u64,
+ },
WriterLockOpen {
path: PathBuf,
source: std::io::Error,
@@ -215,6 +226,19 @@ pub enum Error {
database: &'static str,
target_version: u32,
},
+ DatabaseOpenFailed {
+ database: &'static str,
+ },
+ DatabaseCloseFailed {
+ database: &'static str,
+ },
+ ConnectionPolicyMismatch {
+ database: &'static str,
+ },
+ SourceGenerationRequired,
+ SourceGenerationUnavailable,
+ SourceGenerationMismatch,
+ CorruptSourceGeneration,
}
impl fmt::Display for Error {
@@ -269,6 +293,10 @@ impl fmt::Display for Error {
formatter,
"SQLite busy timeout {actual:?} must be within {minimum:?}..={maximum:?}"
),
+ Self::InvalidSourceGenerationTimestamp { actual } => write!(
+ formatter,
+ "source generation creation time {actual} must fit a positive SQLite integer"
+ ),
Self::WriterLockOpen { path, .. } => write!(
formatter,
"failed to open SQLite writer lock: {}",
@@ -341,7 +369,284 @@ impl fmt::Display for Error {
formatter,
"{database} migration to schema version {target_version} failed"
),
+ Self::DatabaseOpenFailed { database } => {
+ write!(
+ formatter,
+ "failed to open governed SQLite database {database}"
+ )
+ }
+ Self::DatabaseCloseFailed { database } => write!(
+ formatter,
+ "failed to close migration connection for {database}"
+ ),
+ Self::ConnectionPolicyMismatch { database } => write!(
+ formatter,
+ "{database} does not satisfy the governed SQLite connection policy"
+ ),
+ Self::SourceGenerationRequired => formatter.write_str(
+ "a fresh writable SQLite store requires a host-supplied source generation",
+ ),
+ Self::SourceGenerationUnavailable => {
+ formatter.write_str("SQLite storage has no active source generation")
+ }
+ Self::SourceGenerationMismatch => formatter
+ .write_str("host-supplied source generation does not match durable storage"),
+ Self::CorruptSourceGeneration => {
+ formatter.write_str("SQLite storage source generation is corrupt")
+ }
+ }
+ }
+}
+
+impl SqliteStorage {
+ /// Opens both governed databases, applying only authorized forward
+ /// migrations and retaining the writer guard for the backend lifetime.
+ pub async fn open(options: OpenOptions) -> Result<Self, Error> {
+ options.validate_filesystem()?;
+ let runtime_exists =
+ options
+ .paths()
+ .runtime()
+ .try_exists()
+ .map_err(|source| Error::Inspect {
+ path: options.paths().runtime().to_path_buf(),
+ source,
+ })?;
+ if options.mode().may_create() && !runtime_exists && options.source_generation().is_none() {
+ return Err(Error::SourceGenerationRequired);
+ }
+ let writer_lock = WriterLock::acquire(options.paths(), options.mode())?;
+ options.validate_filesystem()?;
+
+ let runtime_options = connect_options(
+ options.paths().runtime(),
+ options.mode(),
+ options.busy_timeout(),
+ );
+ let private_options = connect_options(
+ options.paths().private(),
+ options.mode(),
+ options.busy_timeout(),
+ );
+ let mut runtime_connection =
+ connect(runtime_options.clone(), RUNTIME_DATABASE_NAME).await?;
+ let mut private_connection =
+ connect(private_options.clone(), PRIVATE_DATABASE_NAME).await?;
+
+ migration::migrate_runtime(&mut runtime_connection, options.mode()).await?;
+ migration::migrate_private(&mut private_connection, options.mode()).await?;
+ let generation = active_source_generation(
+ &mut runtime_connection,
+ options.mode(),
+ options.source_generation_bootstrap(),
+ )
+ .await?;
+ verify_connection(
+ &mut runtime_connection,
+ RUNTIME_DATABASE_NAME,
+ options.busy_timeout(),
+ )
+ .await?;
+ verify_connection(
+ &mut private_connection,
+ PRIVATE_DATABASE_NAME,
+ options.busy_timeout(),
+ )
+ .await?;
+ runtime_connection
+ .close()
+ .await
+ .map_err(|_| Error::DatabaseCloseFailed {
+ database: RUNTIME_DATABASE_NAME,
+ })?;
+ private_connection
+ .close()
+ .await
+ .map_err(|_| Error::DatabaseCloseFailed {
+ database: PRIVATE_DATABASE_NAME,
+ })?;
+
+ let runtime_pool = pool(runtime_options, RUNTIME_DATABASE_NAME).await?;
+ let private_pool = pool(private_options, PRIVATE_DATABASE_NAME).await?;
+ verify_pool(&runtime_pool, RUNTIME_DATABASE_NAME, options.busy_timeout()).await?;
+ verify_pool(&private_pool, PRIVATE_DATABASE_NAME, options.busy_timeout()).await?;
+
+ Ok(Self::from_opened(
+ runtime_pool,
+ private_pool,
+ generation,
+ if options.mode().is_writable() {
+ EventStoreMode::ReadWrite
+ } else {
+ EventStoreMode::ReadOnly
+ },
+ writer_lock,
+ ))
+ }
+}
+
+fn connect_options(path: &Path, mode: OpenMode, busy_timeout: Duration) -> SqliteConnectOptions {
+ let mut options = SqliteConnectOptions::new()
+ .filename(path)
+ .read_only(!mode.is_writable())
+ .create_if_missing(mode.may_create())
+ .foreign_keys(true)
+ .busy_timeout(busy_timeout)
+ .synchronous(SqliteSynchronous::Full)
+ .disable_statement_logging();
+ if mode.is_writable() {
+ options = options.journal_mode(SqliteJournalMode::Wal);
+ }
+ options
+}
+
+async fn connect(
+ options: SqliteConnectOptions,
+ database: &'static str,
+) -> Result<SqliteConnection, Error> {
+ SqliteConnection::connect_with(&options)
+ .await
+ .map_err(|_| Error::DatabaseOpenFailed { database })
+}
+
+async fn pool(options: SqliteConnectOptions, database: &'static str) -> Result<SqlitePool, Error> {
+ SqlitePoolOptions::new()
+ .max_connections(MAX_CONNECTIONS_PER_DATABASE)
+ .min_connections(1)
+ .connect_with(options)
+ .await
+ .map_err(|_| Error::DatabaseOpenFailed { database })
+}
+
+async fn active_source_generation(
+ connection: &mut SqliteConnection,
+ mode: OpenMode,
+ expected: Option<(SourceGeneration, u64)>,
+) -> Result<SourceGeneration, Error> {
+ if !mode.is_writable() {
+ let rows = active_generation_rows(connection).await?;
+ return existing_source_generation(rows.as_slice(), expected);
+ }
+ let mut transaction = connection
+ .begin_with("BEGIN IMMEDIATE")
+ .await
+ .map_err(|_| Error::SourceGenerationUnavailable)?;
+ let rows = active_generation_rows(&mut transaction).await?;
+ let generation = match rows.as_slice() {
+ [] => {
+ let (generation, created_at) = expected.ok_or(Error::SourceGenerationRequired)?;
+ sqlx::query(
+ "INSERT INTO radroots_runtime_source_generations (
+ generation, sequence_head, state, created_at_unix_ms
+ ) VALUES (?, 0, 'active', ?)",
+ )
+ .bind(generation.as_bytes().as_slice())
+ .bind(i64::try_from(created_at).map_err(|_| Error::CorruptSourceGeneration)?)
+ .execute(&mut *transaction)
+ .await
+ .map_err(|_| Error::SourceGenerationUnavailable)?;
+ generation
}
+ [_] => existing_source_generation(rows.as_slice(), expected)?,
+ _ => return Err(Error::CorruptSourceGeneration),
+ };
+ transaction
+ .commit()
+ .await
+ .map_err(|_| Error::SourceGenerationUnavailable)?;
+ Ok(generation)
+}
+
+async fn active_generation_rows(
+ connection: &mut SqliteConnection,
+) -> Result<Vec<sqlx::sqlite::SqliteRow>, Error> {
+ sqlx::query(
+ "SELECT generation, created_at_unix_ms
+ FROM radroots_runtime_source_generations
+ WHERE state = 'active' ORDER BY generation",
+ )
+ .fetch_all(connection)
+ .await
+ .map_err(|_| Error::SourceGenerationUnavailable)
+}
+
+fn existing_source_generation(
+ rows: &[sqlx::sqlite::SqliteRow],
+ expected: Option<(SourceGeneration, u64)>,
+) -> Result<SourceGeneration, Error> {
+ let [row] = rows else {
+ return if rows.is_empty() {
+ Err(Error::SourceGenerationUnavailable)
+ } else {
+ Err(Error::CorruptSourceGeneration)
+ };
+ };
+ let durable = decode_source_generation(row)?;
+ let created_at = u64::try_from(
+ row.try_get::<i64, _>("created_at_unix_ms")
+ .map_err(|_| Error::CorruptSourceGeneration)?,
+ )
+ .map_err(|_| Error::CorruptSourceGeneration)?;
+ if expected.is_some_and(|candidate| candidate != (durable, created_at)) {
+ Err(Error::SourceGenerationMismatch)
+ } else {
+ Ok(durable)
+ }
+}
+
+fn decode_source_generation(row: &sqlx::sqlite::SqliteRow) -> Result<SourceGeneration, Error> {
+ SourceGeneration::new(
+ row.try_get::<Vec<u8>, _>("generation")
+ .map_err(|_| Error::CorruptSourceGeneration)?
+ .try_into()
+ .map_err(|_| Error::CorruptSourceGeneration)?,
+ )
+ .map_err(|_| Error::CorruptSourceGeneration)
+}
+
+async fn verify_pool(
+ pool: &SqlitePool,
+ database: &'static str,
+ busy_timeout: Duration,
+) -> Result<(), Error> {
+ let mut connection = pool
+ .acquire()
+ .await
+ .map_err(|_| Error::DatabaseOpenFailed { database })?;
+ verify_connection(&mut connection, database, busy_timeout).await
+}
+
+async fn verify_connection(
+ connection: &mut SqliteConnection,
+ database: &'static str,
+ busy_timeout: Duration,
+) -> Result<(), Error> {
+ let foreign_keys = sqlx::query_scalar::<_, i64>("PRAGMA foreign_keys")
+ .fetch_one(&mut *connection)
+ .await
+ .map_err(|_| Error::ConnectionPolicyMismatch { database })?;
+ let journal_mode = sqlx::query_scalar::<_, String>("PRAGMA journal_mode")
+ .fetch_one(&mut *connection)
+ .await
+ .map_err(|_| Error::ConnectionPolicyMismatch { database })?;
+ let configured_busy_timeout = sqlx::query_scalar::<_, i64>("PRAGMA busy_timeout")
+ .fetch_one(&mut *connection)
+ .await
+ .map_err(|_| Error::ConnectionPolicyMismatch { database })?;
+ let synchronous = sqlx::query_scalar::<_, i64>("PRAGMA synchronous")
+ .fetch_one(&mut *connection)
+ .await
+ .map_err(|_| Error::ConnectionPolicyMismatch { database })?;
+ let expected_busy_timeout = i64::try_from(busy_timeout.as_millis())
+ .map_err(|_| Error::ConnectionPolicyMismatch { database })?;
+ if foreign_keys == 1
+ && journal_mode.eq_ignore_ascii_case("wal")
+ && configured_busy_timeout == expected_busy_timeout
+ && synchronous == 2
+ {
+ Ok(())
+ } else {
+ Err(Error::ConnectionPolicyMismatch { database })
}
}
@@ -356,3 +661,53 @@ impl StdError for Error {
}
}
}
+
+#[cfg(test)]
+mod policy_tests {
+ use super::*;
+ use serde::Deserialize;
+
+ const POLICY: &str = include_str!("../../../contracts/storage/connection_policy_v1.toml");
+
+ #[derive(Deserialize)]
+ struct Policy {
+ schema_version: u32,
+ databases: Vec<String>,
+ max_connections_per_database: u32,
+ foreign_keys: bool,
+ journal_mode: String,
+ synchronous: String,
+ busy_timeout_min_ms: u64,
+ busy_timeout_default_ms: u64,
+ busy_timeout_max_ms: u64,
+ fresh_source_generation: String,
+ read_only_migrations: bool,
+ raw_handles_public: bool,
+ }
+
+ #[test]
+ fn implementation_matches_the_governed_connection_policy() {
+ let policy = toml::from_str::<Policy>(POLICY).expect("connection policy");
+ assert_eq!(policy.schema_version, 1);
+ assert_eq!(
+ policy.databases,
+ [RUNTIME_DATABASE_NAME, PRIVATE_DATABASE_NAME]
+ );
+ assert_eq!(
+ policy.max_connections_per_database,
+ MAX_CONNECTIONS_PER_DATABASE
+ );
+ assert!(policy.foreign_keys);
+ assert_eq!(policy.journal_mode, "wal");
+ assert_eq!(policy.synchronous, "full");
+ assert_eq!(policy.busy_timeout_min_ms, 1);
+ assert_eq!(policy.busy_timeout_default_ms, 5_000);
+ assert_eq!(policy.busy_timeout_max_ms, 60_000);
+ assert_eq!(
+ policy.fresh_source_generation,
+ "host_supplied_entropy_and_timestamp"
+ );
+ assert!(!policy.read_only_migrations);
+ assert!(!policy.raw_handles_public);
+ }
+}
diff --git a/crates/storage_sqlite/tests/open_lifecycle.rs b/crates/storage_sqlite/tests/open_lifecycle.rs
@@ -0,0 +1,92 @@
+use std::time::Duration;
+
+use radroots_storage::{EventStore, event::SourceGeneration, status::EventStoreMode};
+use radroots_storage_sqlite::{Error, OpenMode, OpenOptions, Paths, SqliteStorage};
+
+fn generation(byte: u8) -> SourceGeneration {
+ SourceGeneration::new([byte; 32]).expect("source generation")
+}
+
+#[tokio::test]
+async fn public_open_creates_both_databases_and_reuses_the_durable_generation() {
+ let directory = tempfile::tempdir().expect("temporary directory");
+ let paths = Paths::from_directory(directory.path()).expect("owned paths");
+ let expected = generation(31);
+ let store = SqliteStorage::open(
+ OpenOptions::new(paths.clone(), OpenMode::Create)
+ .with_busy_timeout(Duration::from_millis(250))
+ .expect("busy timeout")
+ .with_source_generation(expected, 1_000)
+ .expect("source generation"),
+ )
+ .await
+ .expect("create storage");
+ assert!(paths.runtime().is_file());
+ assert!(paths.private().is_file());
+ let status = EventStore::status(&store).await.expect("event status");
+ assert_eq!(status.generation(), expected);
+ assert_eq!(status.mode(), EventStoreMode::ReadWrite);
+
+ let reader = SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadOnly))
+ .await
+ .expect("concurrent reader");
+ let reader_status = EventStore::status(&reader).await.expect("reader status");
+ assert_eq!(reader_status.generation(), expected);
+ assert_eq!(reader_status.mode(), EventStoreMode::ReadOnly);
+ drop(reader);
+
+ assert!(matches!(
+ SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting,)).await,
+ Err(Error::WriterAlreadyActive { .. })
+ ));
+ let cloned = store.clone();
+ drop(store);
+ assert!(matches!(
+ SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting,)).await,
+ Err(Error::WriterAlreadyActive { .. })
+ ));
+ drop(cloned);
+
+ let reopened = SqliteStorage::open(OpenOptions::new(paths, OpenMode::ReadWriteExisting))
+ .await
+ .expect("reopen after guard release");
+ assert_eq!(
+ EventStore::status(&reopened)
+ .await
+ .expect("reopened status")
+ .generation(),
+ expected
+ );
+}
+
+#[tokio::test]
+async fn fresh_store_requires_explicit_generation_and_exact_expectations() {
+ let directory = tempfile::tempdir().expect("temporary directory");
+ let paths = Paths::from_directory(directory.path()).expect("owned paths");
+ assert!(matches!(
+ SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::Create)).await,
+ Err(Error::SourceGenerationRequired)
+ ));
+ assert!(!paths.runtime().exists());
+ assert!(!paths.private().exists());
+
+ let expected = generation(41);
+ let store = SqliteStorage::open(
+ OpenOptions::new(paths.clone(), OpenMode::Create)
+ .with_source_generation(expected, 2_000)
+ .expect("source generation"),
+ )
+ .await
+ .expect("complete fresh store");
+ drop(store);
+
+ assert!(matches!(
+ SqliteStorage::open(
+ OpenOptions::new(paths, OpenMode::ReadWriteExisting)
+ .with_source_generation(generation(42), 2_000)
+ .expect("wrong expectation"),
+ )
+ .await,
+ Err(Error::SourceGenerationMismatch)
+ ));
+}
diff --git a/crates/storage_sqlite/tests/open_options.rs b/crates/storage_sqlite/tests/open_options.rs
@@ -1,6 +1,7 @@
use std::fs;
use std::time::Duration;
+use radroots_storage::event::SourceGeneration;
use radroots_storage::status::WriterPolicy;
use radroots_storage_sqlite::{Error, OpenMode, OpenOptions, Paths};
@@ -98,4 +99,19 @@ fn options_fix_connection_invariants_and_bound_busy_timeout() {
Err(Error::InvalidBusyTimeout { .. })
));
}
+
+ let generation = SourceGeneration::new([9; 32]).expect("source generation");
+ let bootstrapped = OpenOptions::new(create.paths().clone(), OpenMode::Create)
+ .with_source_generation(generation, 42)
+ .expect("source generation bootstrap");
+ assert_eq!(bootstrapped.source_generation(), Some(generation));
+ assert_eq!(
+ bootstrapped.source_generation_created_at_unix_ms(),
+ Some(42)
+ );
+ assert!(matches!(
+ OpenOptions::new(create.paths().clone(), OpenMode::Create)
+ .with_source_generation(generation, 0),
+ Err(Error::InvalidSourceGenerationTimestamp { actual: 0 })
+ ));
}