commit 1cc197c2d48cd537de50e5f2637e00bf2c10fd36
parent e6f296604363e5293d6abc6e18bf5eaea2177b9f
Author: triesap <tyson@radroots.org>
Date: Tue, 22 Sep 2026 00:37:21 +0000
storage: preflight both schemas before migration
- Inspect existing files before writable connection setup
- Preserve incompatible schemas, generations and envelope evidence
- Retain transactional forward migrations and typed errors
- Verify owner coverage, public APIs and workspace consumers
Diffstat:
4 files changed, 481 insertions(+), 52 deletions(-)
diff --git a/crates/storage_sqlite/README.md b/crates/storage_sqlite/README.md
@@ -175,3 +175,11 @@ partial replay rolls back. Complete replay returns immutable historical rows and
does not establish current application authority. No migration or second storage
owner is introduced. Callers may recover an unknown commit result by replaying
the exact pair.
+
+Open inspects both existing database schemas through read-only connections before
+configuring writable connections or applying either pending migration. A future
+version, foreign namespace, invalid catalog, ineligible authored migration or
+incompatible source generation therefore leaves both original files intact. Each accepted migration still
+rechecks metadata and applies its pending suffix atomically in its own database;
+this does not claim a cross-database atomic commit. An interrupted compatible
+upgrade resumes from its committed version on the next open.
diff --git a/crates/storage_sqlite/src/migration.rs b/crates/storage_sqlite/src/migration.rs
@@ -115,12 +115,23 @@ impl MigrationReport {
}
}
-#[allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint.
#[cfg_attr(coverage_nightly, coverage(off))]
pub(crate) async fn migrate_runtime(
connection: &mut SqliteConnection,
mode: OpenMode,
) -> Result<MigrationReport, Error> {
+ migrate(connection, mode, &runtime_plan()?).await
+}
+
+#[cfg_attr(coverage_nightly, coverage(off))]
+pub(crate) async fn migrate_private(
+ connection: &mut SqliteConnection,
+ mode: OpenMode,
+) -> Result<MigrationReport, Error> {
+ migrate(connection, mode, &private_plan()?).await
+}
+
+fn runtime_plan() -> Result<MigrationPlan, Error> {
let steps = runtime::MIGRATIONS
.iter()
.map(|migration| {
@@ -135,27 +146,17 @@ pub(crate) async fn migrate_runtime(
})
})
.collect::<Result<Vec<_>, Error>>()?;
- migrate(
- connection,
- mode,
- &MigrationPlan {
- database: RUNTIME_DATABASE,
- application_id: RUNTIME_APPLICATION_ID,
- set_application_id_sql: SET_RUNTIME_APPLICATION_ID,
- minimum_version: runtime::MINIMUM_VERSION,
- current_version: runtime::CURRENT_VERSION,
- steps,
- },
- )
- .await
+ Ok(MigrationPlan {
+ database: RUNTIME_DATABASE,
+ application_id: RUNTIME_APPLICATION_ID,
+ set_application_id_sql: SET_RUNTIME_APPLICATION_ID,
+ minimum_version: runtime::MINIMUM_VERSION,
+ current_version: runtime::CURRENT_VERSION,
+ steps,
+ })
}
-#[allow(dead_code)] // Wired into the public open lifecycle in its ordered RCL checkpoint.
-#[cfg_attr(coverage_nightly, coverage(off))]
-pub(crate) async fn migrate_private(
- connection: &mut SqliteConnection,
- mode: OpenMode,
-) -> Result<MigrationReport, Error> {
+fn private_plan() -> Result<MigrationPlan, Error> {
let steps = private::MIGRATIONS
.iter()
.map(|migration| {
@@ -170,52 +171,124 @@ pub(crate) async fn migrate_private(
})
})
.collect::<Result<Vec<_>, Error>>()?;
- migrate(
- connection,
- mode,
- &MigrationPlan {
- database: PRIVATE_DATABASE,
- application_id: PRIVATE_APPLICATION_ID,
- set_application_id_sql: SET_PRIVATE_APPLICATION_ID,
- minimum_version: private::MINIMUM_VERSION,
- current_version: private::CURRENT_VERSION,
- steps,
- },
- )
- .await
+ Ok(MigrationPlan {
+ database: PRIVATE_DATABASE,
+ application_id: PRIVATE_APPLICATION_ID,
+ set_application_id_sql: SET_PRIVATE_APPLICATION_ID,
+ minimum_version: private::MINIMUM_VERSION,
+ current_version: private::CURRENT_VERSION,
+ steps,
+ })
}
-#[cfg_attr(coverage_nightly, coverage(off))]
-async fn migrate(
+// Validate both existing members before WAL setup or any forward migration.
+// The caller retains the canonical writer lock; migration still rechecks under
+// its own transaction. This is compatibility preflight, not a cross-file commit.
+pub(crate) async fn preflight_existing(options: &crate::OpenOptions) -> Result<(), Error> {
+ for (path, plan) in [
+ (options.paths().runtime(), runtime_plan()?),
+ (options.paths().private(), private_plan()?),
+ ] {
+ if !path.try_exists().map_err(|source| Error::Inspect {
+ path: path.to_path_buf(),
+ source,
+ })? {
+ if options.mode().may_create() {
+ continue;
+ }
+ return Err(Error::MissingFile(path.to_path_buf()));
+ }
+ let mut connection = SqliteConnection::connect_with(
+ &SqliteConnectOptions::new()
+ .filename(path)
+ .read_only(true)
+ .busy_timeout(options.busy_timeout())
+ .pragma("query_only", "ON"),
+ )
+ .await
+ .map_err(|source| crate::open::map_database_open_error(&source, plan.database))?;
+ let inspected = async {
+ let metadata = inspect(&mut connection, options.mode(), &plan).await?;
+ if plan.database == RUNTIME_DATABASE && metadata.version > 0 {
+ crate::open::preflight_source_generation(
+ &mut connection,
+ options.mode(),
+ options.source_generation_bootstrap(),
+ )
+ .await?;
+ }
+ Ok::<(), Error>(())
+ }
+ .await;
+ let closed = connection
+ .close()
+ .await
+ .map_err(|_| Error::DatabaseCloseFailed {
+ database: plan.database,
+ });
+ inspected?;
+ closed?;
+ }
+ Ok(())
+}
+
+async fn inspect(
connection: &mut SqliteConnection,
mode: OpenMode,
plan: &MigrationPlan,
-) -> Result<MigrationReport, Error> {
+) -> Result<SchemaMetadata, Error> {
validate_plan(plan)?;
let initial = metadata(connection, plan.database).await?;
validate_metadata(plan, initial)?;
validate_catalog(connection, plan, initial.version).await?;
-
- if initial.version == plan.current_version {
- return Ok(MigrationReport {
- initial_version: initial.version,
- final_version: initial.version,
- applied: 0,
- });
+ if plan.database == RUNTIME_DATABASE && initial.version == 10 {
+ let inspected = authored_v10::inspect(connection).await?;
+ if !inspected.report.is_eligible() {
+ return Err(inspected.report.blocked_error());
+ }
}
- if !mode.is_writable() {
- if plan.database == RUNTIME_DATABASE && initial.version == 10 {
- let inspected = authored_v10::inspect(connection).await?;
- if !inspected.report.is_eligible() {
- return Err(inspected.report.blocked_error());
- }
+ // Private v4's pinned SQL refuses old v2 envelopes without a context
+ // fingerprint. Check the same eligibility before either file is migrated;
+ // keep the SQL guard as the transactional authority.
+ if mode.is_writable() && plan.database == PRIVATE_DATABASE && (1..4).contains(&initial.version)
+ {
+ let unsupported = sqlx::query_scalar::<_, bool>(
+ "SELECT EXISTS(SELECT 1 FROM radroots_private_artifacts WHERE envelope_version = 2)",
+ )
+ .fetch_one(&mut *connection)
+ .await
+ .map_err(|source| schema_metadata_error(&source, PRIVATE_DATABASE))?;
+ if unsupported {
+ return Err(Error::SchemaMigrationFailed {
+ database: PRIVATE_DATABASE,
+ target_version: 4,
+ });
}
+ }
+ if initial.version != plan.current_version && !mode.is_writable() {
return Err(Error::SchemaMigrationRequired {
database: plan.database,
current: plan.current_version,
actual: initial.version,
});
}
+ Ok(initial)
+}
+
+#[cfg_attr(coverage_nightly, coverage(off))]
+async fn migrate(
+ connection: &mut SqliteConnection,
+ mode: OpenMode,
+ plan: &MigrationPlan,
+) -> Result<MigrationReport, Error> {
+ let initial = inspect(connection, mode, plan).await?;
+ if initial.version == plan.current_version {
+ return Ok(MigrationReport {
+ initial_version: initial.version,
+ final_version: initial.version,
+ applied: 0,
+ });
+ }
let mut transaction = connection
.begin_with("BEGIN IMMEDIATE")
@@ -355,11 +428,11 @@ async fn metadata(
let application_id = sqlx::query_scalar::<_, i64>("PRAGMA application_id")
.fetch_one(&mut *connection)
.await
- .map_err(|_| Error::SchemaMetadataUnavailable { database })?;
+ .map_err(|source| schema_metadata_error(&source, database))?;
let version = sqlx::query_scalar::<_, i64>("PRAGMA user_version")
.fetch_one(&mut *connection)
.await
- .map_err(|_| Error::SchemaMetadataUnavailable { database })?;
+ .map_err(|source| schema_metadata_error(&source, database))?;
Ok(SchemaMetadata {
application_id: u32::try_from(application_id)
.map_err(|_| Error::SchemaMetadataUnavailable { database })?,
@@ -368,6 +441,13 @@ async fn metadata(
})
}
+fn schema_metadata_error(source: &sqlx::Error, database: &'static str) -> Error {
+ match crate::open::map_database_open_error(source, database) {
+ error @ Error::DatabaseCorrupt { .. } => error,
+ _ => Error::SchemaMetadataUnavailable { database },
+ }
+}
+
fn validate_plan(plan: &MigrationPlan) -> Result<(), Error> {
let valid = plan.minimum_version > 0
&& plan.minimum_version <= plan.current_version
@@ -544,7 +624,7 @@ mod tests {
.expect("runtime user version");
}
- async fn establish_private_version(connection: &mut SqliteConnection, version: u32) {
+ pub(super) async fn establish_private_version(connection: &mut SqliteConnection, version: u32) {
for migration_version in 1..=version {
sqlx::raw_sql(
private::migration_sql(migration_version).expect("registered private SQL"),
@@ -885,3 +965,7 @@ mod delivery_facts_tests;
#[cfg_attr(coverage_nightly, coverage(off))]
#[path = "migration_delivery_reconciliation_tests.rs"]
mod delivery_reconciliation_tests;
+
+#[cfg(test)]
+#[path = "migration_preflight_tests.rs"]
+mod preflight_tests;
diff --git a/crates/storage_sqlite/src/migration_preflight_tests.rs b/crates/storage_sqlite/src/migration_preflight_tests.rs
@@ -0,0 +1,319 @@
+//! Public-open regressions use only synthetic owned databases.
+
+use super::{tests, *};
+use crate::{OpenOptions, Paths, SqliteStorage};
+use std::path::Path;
+
+async fn file_connection(path: &Path, create: bool) -> SqliteConnection {
+ SqliteConnection::connect_with(
+ &SqliteConnectOptions::new()
+ .filename(path)
+ .create_if_missing(create),
+ )
+ .await
+ .expect("fixture connection")
+}
+
+async fn fixture(runtime_version: u32, private_version: u32) -> (tempfile::TempDir, Paths) {
+ let directory = tempfile::tempdir().expect("temporary directory");
+ let paths = Paths::from_directory(directory.path()).expect("owned paths");
+ let mut runtime = file_connection(paths.runtime(), true).await;
+ tests::establish_runtime_version(&mut runtime, runtime_version).await;
+ sqlx::query("INSERT INTO radroots_runtime_source_generations (generation, state, created_at_unix_ms) VALUES (?, 'active', 10)")
+ .bind([7_u8; 32].as_slice()).execute(&mut runtime).await.expect("retained generation");
+ runtime.close().await.expect("close runtime fixture");
+ let mut private = file_connection(paths.private(), true).await;
+ tests::establish_private_version(&mut private, private_version).await;
+ private.close().await.expect("close private fixture");
+ (directory, paths)
+}
+
+async fn alter(path: &Path, sql: &'static str) {
+ let mut connection = file_connection(path, false).await;
+ sqlx::raw_sql(sql)
+ .execute(&mut connection)
+ .await
+ .expect("fixture alteration");
+ connection.close().await.expect("close fixture alteration");
+}
+
+async fn assert_refused_unchanged(paths: Paths, mode: OpenMode, expected: fn(&Error) -> bool) {
+ let runtime = std::fs::read(paths.runtime()).expect("runtime bytes");
+ let private = std::fs::read(paths.private()).expect("private bytes");
+ for _ in 0..2 {
+ let error = match SqliteStorage::open(OpenOptions::new(paths.clone(), mode)).await {
+ Ok(store) => {
+ store.close().await.expect("close unexpected store");
+ panic!("incompatible pair opened")
+ }
+ Err(error) => error,
+ };
+ assert!(expected(&error), "unexpected error: {error:?}");
+ assert!(
+ std::fs::read(paths.runtime()).expect("retained runtime") == runtime,
+ "runtime bytes changed"
+ );
+ assert!(
+ std::fs::read(paths.private()).expect("retained private") == private,
+ "private bytes changed"
+ );
+ }
+}
+
+#[tokio::test]
+async fn future_private_schema_preserves_pending_runtime_and_both_original_files() {
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ alter(paths.private(), "PRAGMA user_version = 999").await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::SchemaTooNew {
+ database: PRIVATE_DATABASE,
+ actual: 999,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn foreign_private_namespace_preserves_pending_runtime_and_both_original_files() {
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ alter(paths.private(), "PRAGMA application_id = 42").await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::SchemaIdentityMismatch {
+ database: PRIVATE_DATABASE,
+ actual: 42,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn unexpected_private_catalog_preserves_pending_runtime_and_both_original_files() {
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ alter(paths.private(), "CREATE TABLE foreign_state (value TEXT)").await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::SchemaCatalogMismatch {
+ database: PRIVATE_DATABASE,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn future_runtime_schema_preserves_pending_private_and_both_original_files() {
+ let (_directory, paths) = fixture(runtime::CURRENT_VERSION, 3).await;
+ alter(paths.runtime(), "PRAGMA user_version = 999").await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::SchemaTooNew {
+ database: RUNTIME_DATABASE,
+ actual: 999,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn read_only_pending_pair_preserves_both_original_files() {
+ let (_directory, paths) = fixture(16, 3).await;
+ assert_refused_unchanged(paths, OpenMode::ReadOnly, |error| {
+ matches!(
+ error,
+ Error::SchemaMigrationRequired {
+ database: RUNTIME_DATABASE,
+ actual: 16,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn compatible_pending_pair_migrates_and_reopens_with_same_generation() {
+ use radroots_storage::EventStore;
+ let (_directory, paths) = fixture(16, 3).await;
+ for _ in 0..2 {
+ let store =
+ SqliteStorage::open(OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting))
+ .await
+ .expect("supported pair");
+ assert_eq!(
+ store
+ .status()
+ .await
+ .expect("status")
+ .generation()
+ .as_bytes(),
+ &[7; 32]
+ );
+ store.close().await.expect("close supported pair");
+ }
+ for (path, version) in [
+ (paths.runtime(), runtime::CURRENT_VERSION),
+ (paths.private(), private::CURRENT_VERSION),
+ ] {
+ let mut connection = file_connection(path, false).await;
+ assert_eq!(
+ tests::pragma(&mut connection, "user_version").await,
+ i64::from(version)
+ );
+ connection.close().await.expect("close inspected fixture");
+ }
+}
+
+#[tokio::test]
+async fn unversioned_namespace_collision_preserves_both_original_files() {
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ alter(
+ paths.private(),
+ "PRAGMA user_version = 0; PRAGMA application_id = 0",
+ )
+ .await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::UnrecognizedSchema {
+ database: PRIVATE_DATABASE
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn corrupt_private_member_preserves_pending_runtime_and_corrupt_evidence() {
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ std::fs::write(paths.private(), b"synthetic invalid SQLite evidence").expect("corrupt fixture");
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::DatabaseCorrupt {
+ database: PRIVATE_DATABASE
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn incompatible_private_member_prevents_creation_of_missing_runtime() {
+ use radroots_storage::event::SourceGeneration;
+ let (_directory, paths) = fixture(16, private::CURRENT_VERSION).await;
+ std::fs::remove_file(paths.runtime()).expect("remove synthetic runtime");
+ alter(paths.private(), "PRAGMA user_version = 999").await;
+ let before = std::fs::read(paths.private()).expect("private evidence");
+ for _ in 0..2 {
+ let options = OpenOptions::new(paths.clone(), OpenMode::Create)
+ .with_source_generation(SourceGeneration::new([9; 32]).expect("generation"), 10)
+ .expect("options");
+ assert!(matches!(
+ SqliteStorage::open(options).await,
+ Err(Error::SchemaTooNew {
+ database: PRIVATE_DATABASE,
+ actual: 999,
+ ..
+ })
+ ));
+ assert!(!paths.runtime().exists());
+ assert!(std::fs::read(paths.private()).expect("private retained") == before);
+ }
+}
+
+#[tokio::test]
+async fn ineligible_v10_authored_evidence_preserves_both_original_files() {
+ let (_directory, paths) = fixture(10, 3).await;
+ alter(
+ paths.runtime(),
+ "INSERT INTO radroots_runtime_journal_operations (
+ instance_id, operation_id, idempotency_key, input_digest, prepared_at_unix_ms,
+ revision, stage, cancellation_state, updated_at_unix_ms
+ ) VALUES (zeroblob(16), x'ff', 'retained-invalid-journal', zeroblob(32), 10,
+ 1, 'prepared', 'not_requested', 10)",
+ )
+ .await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::AuthoredMigrationBlocked {
+ invalid_or_unsupported: 1,
+ ..
+ }
+ )
+ })
+ .await;
+}
+
+#[tokio::test]
+async fn incompatible_generation_preserves_pending_schemas_and_original_files() {
+ use radroots_storage::event::SourceGeneration;
+ let (_directory, paths) = fixture(16, 3).await;
+ let runtime = std::fs::read(paths.runtime()).expect("runtime evidence");
+ let private = std::fs::read(paths.private()).expect("private evidence");
+ for _ in 0..2 {
+ let options = OpenOptions::new(paths.clone(), OpenMode::ReadWriteExisting)
+ .with_source_generation(SourceGeneration::new([9; 32]).expect("generation"), 10)
+ .expect("options");
+ assert!(matches!(
+ SqliteStorage::open(options).await,
+ Err(Error::SourceGenerationMismatch)
+ ));
+ assert!(
+ std::fs::read(paths.runtime()).expect("retained runtime") == runtime,
+ "runtime bytes changed before rejecting generation"
+ );
+ assert!(
+ std::fs::read(paths.private()).expect("retained private") == private,
+ "private bytes changed before rejecting generation"
+ );
+ }
+}
+
+#[test]
+fn unrelated_metadata_failure_retains_typed_metadata_diagnostic() {
+ assert!(matches!(
+ schema_metadata_error(&sqlx::Error::RowNotFound, PRIVATE_DATABASE),
+ Error::SchemaMetadataUnavailable {
+ database: PRIVATE_DATABASE
+ }
+ ));
+}
+
+#[tokio::test]
+async fn unfingerprinted_private_v2_envelope_preserves_both_original_files() {
+ let (_directory, paths) = fixture(16, 3).await;
+ alter(
+ paths.private(),
+ "INSERT INTO radroots_private_artifacts (
+ artifact_id, artifact_kind, schema_id, commitment, protected_size_bytes,
+ secret_provider, secret_reference, key_version, envelope_version,
+ encrypted_envelope, revision, stage, created_at_unix_ms, updated_at_unix_ms
+ ) VALUES (zeroblob(16), 'trade.private_terms', 'trade.private_terms.v1', zeroblob(32), 1,
+ 'memory', 'synthetic-key', 1, 2, x'03', 1, 'active', 1, 1)",
+ )
+ .await;
+ assert_refused_unchanged(paths, OpenMode::ReadWriteExisting, |error| {
+ matches!(
+ error,
+ Error::SchemaMigrationFailed {
+ database: PRIVATE_DATABASE,
+ target_version: 4
+ }
+ )
+ })
+ .await;
+}
diff --git a/crates/storage_sqlite/src/open.rs b/crates/storage_sqlite/src/open.rs
@@ -644,6 +644,8 @@ impl SqliteStorage {
}
options.validate_filesystem()?;
+ migration::preflight_existing(&options).await?;
+
let runtime_options = connect_options(
options.paths().runtime(),
options.mode(),
@@ -756,6 +758,22 @@ pub(crate) fn map_database_open_error(source: &sqlx::Error, database: &'static s
}
}
+// Inspect the existing generation before migrations without bootstrapping it.
+// The actual writer transaction still validates or installs it after migration.
+pub(crate) async fn preflight_source_generation(
+ connection: &mut SqliteConnection,
+ mode: OpenMode,
+ expected: Option<(SourceGeneration, u64)>,
+) -> Result<(), Error> {
+ let rows = active_generation_rows(connection).await?;
+ if rows.is_empty() && mode.is_writable() {
+ expected.ok_or(Error::SourceGenerationRequired)?;
+ Ok(())
+ } else {
+ existing_source_generation(rows.as_slice(), expected).map(|_| ())
+ }
+}
+
#[cfg_attr(coverage_nightly, coverage(off))]
async fn active_source_generation(
connection: &mut SqliteConnection,