test(archive): harden M4 shape validation and pin the init barrier

Resolve the three review findings on the retention Phase-1 branch.

M4 now validates the COMPLETE expected shape of its objects inside the
BEGIN IMMEDIATE transaction — every column's name, type, nullability and
primary-key position via PRAGMA table_info, and the scope-age index's key
order via pragma_index_info — instead of counting object names. CREATE ...
IF NOT EXISTS preserves a wrong-shaped object that already carries the
expected name, so name presence alone could mark a wrong-shaped named
table valid and hand Phase 2 an unusable table or missing access path. A
wrong-shaped named index is dropped and rebuilt (an index carries no data);
an incompatible named table is rejected and the whole transaction rolls
back with no marker.

The init-barrier and guard-lifetime contracts are now pinned directly
against the ArchiveDb OnceCell/RwLock orchestration through a cfg(test)
path/hook seam, rather than via raw SQLite contention: production-shaped
with_conn callers race a held init and prove none opens a connection until
initialization completes and exactly one initialization runs; a separate
test proves a write-lock contender cannot enter until a with_conn closure
returns and its connection drops.

Adds the full policy lifecycle on one evolving state (active -> kind
disabled -> subscription deleted -> orphan edited -> orphan deleted) and a
deterministic concurrent merge/remove/set interleaving asserting the
subscription stays valid, the explicit choice survives, and no policy row
is deleted by a kinds mutation.

Co-authored-by: Will Pfleger <pfleger.will@gmail.com>
Signed-off-by: Will Pfleger <pfleger.will@gmail.com>
This commit is contained in:
Duncan
2026-08-14 10:38:07 -04:00
co-authored by Will Pfleger
parent 34773f0472
commit 85438f647b
5 changed files with 720 additions and 33 deletions
+90 -7
View File
@@ -37,6 +37,22 @@ use tokio::sync::{OnceCell, RwLock};
use super::store;
use crate::managed_agents::nest_dir;
/// A hook run once on the blocking pool at the start of the single init task.
/// Test-only: lets a test count initializations and hold the winner long
/// enough to prove concurrent callers await it. Never set in production.
#[cfg(test)]
type InitHook = std::sync::Arc<dyn Fn() + Send + Sync>;
/// Test-only overrides so the barrier and guard-lifetime contracts can be
/// exercised without a real nest: a fixed DB path in place of `nest_dir()` and
/// an optional init hook. `Default` leaves this `None`, so production always
/// resolves the path from the nest and runs no hook.
#[cfg(test)]
struct TestSeam {
path: PathBuf,
on_init: Option<InitHook>,
}
/// Gated owner of every production archive DB connection. Lives in
/// [`crate::app_state::AppState`]; commands call [`ArchiveDb::with_conn`].
#[derive(Default)]
@@ -49,26 +65,51 @@ pub struct ArchiveDb {
/// Maintenance lock. Ordinary connections hold the read guard for their
/// whole lifetime; the Phase-4 conversion holds the write guard.
maintenance: RwLock<()>,
/// Test-only path/hook overrides; always `None` in production.
#[cfg(test)]
test_seam: Option<TestSeam>,
}
impl ArchiveDb {
/// Resolve the archive DB path from the nest directory. Errors only when
/// the nest cannot be resolved (fatal for archive access, same as the
/// former `open_db`).
fn db_path() -> Result<PathBuf, String> {
/// Resolve the archive DB path. Production resolves from the nest
/// directory; a test seam (when present) supplies a fixed path so the
/// barrier can be exercised without a real nest. Errors only when the nest
/// cannot be resolved (fatal for archive access, same as the former
/// `open_db`).
fn db_path(&self) -> Result<PathBuf, String> {
#[cfg(test)]
if let Some(seam) = &self.test_seam {
return Ok(seam.path.clone());
}
let nest = nest_dir().ok_or("cannot resolve nest directory for archive")?;
Ok(nest.join("archive").join("archive.db"))
}
/// The init hook, if a test installed one; always `None` in production.
#[cfg(test)]
fn init_hook(&self) -> Option<InitHook> {
self.test_seam.as_ref().and_then(|s| s.on_init.clone())
}
/// Complete the one-time init: open the DB once on the blocking pool,
/// running `SCHEMA` + all migrations (incl. M4), then drop the connection.
/// Concurrent callers await the same single execution. Idempotent and
/// cheap after the first success (the cached `()` short-circuits).
async fn ensure_initialized(&self) -> Result<(), String> {
let path = self.db_path()?;
#[cfg(test)]
let hook = self.init_hook();
self.init
.get_or_try_init(|| async {
tokio::task::spawn_blocking(|| {
let path = Self::db_path()?;
tokio::task::spawn_blocking(move || {
// Test hook runs at the very start of the single init task,
// before the migration opens the DB — this is where a test
// holds the winner past the busy timeout to prove ordinary
// callers await it. No-op in production.
#[cfg(test)]
if let Some(hook) = hook {
hook();
}
// Opening runs every migration; the connection exists only
// to complete them behind the barrier, so drop it here.
let conn = store::open_archive_db(&path)?;
@@ -104,9 +145,9 @@ impl ArchiveDb {
F: FnOnce(&Connection) -> Result<T, String> + Send + 'static,
{
self.ensure_initialized().await?;
let path = self.db_path()?;
let _guard = self.maintenance.read().await;
tokio::task::spawn_blocking(move || {
let path = Self::db_path()?;
let conn = store::open_archive_db(&path)?;
task(&conn)
})
@@ -114,3 +155,45 @@ impl ArchiveDb {
.map_err(|e| format!("archive db task failed: {e}"))?
}
}
#[cfg(test)]
impl ArchiveDb {
/// Build an adapter bound to a fixed DB path (no nest required), so the
/// barrier and guard-lifetime contracts can be exercised in isolation.
fn with_test_path(path: PathBuf) -> Self {
Self {
init: OnceCell::new(),
maintenance: RwLock::new(()),
test_seam: Some(TestSeam {
path,
on_init: None,
}),
}
}
/// Build an adapter bound to a fixed path whose single initialization runs
/// `hook` first — used to count initializations and to hold the init task
/// open across the concurrent-caller window.
fn with_test_hook(path: PathBuf, hook: InitHook) -> Self {
Self {
init: OnceCell::new(),
maintenance: RwLock::new(()),
test_seam: Some(TestSeam {
path,
on_init: Some(hook),
}),
}
}
/// Whether the maintenance WRITE guard can be taken right now. A live
/// `with_conn` connection holds the read guard, so this returns `false`
/// while any ordinary connection is open and `true` once all have dropped —
/// exactly the signal the Phase-4 sole-connection VACUUM will gate on.
fn maintenance_write_available(&self) -> bool {
self.maintenance.try_write().is_ok()
}
}
#[cfg(test)]
#[path = "archive_db_tests.rs"]
mod archive_db_tests;
@@ -0,0 +1,246 @@
//! Behavior tests for the [`ArchiveDb`] init barrier and maintenance-lock
//! guard lifetime — the two contracts Phase 1 introduced and Thufir's pass-1
//! review required to be pinned directly (not via raw SQLite contention).
//!
//! These race PRODUCTION-shaped `with_conn` callers through the real
//! `OnceCell`/`RwLock` orchestration, using the `#[cfg(test)]` path/hook seam
//! on `ArchiveDb` to make timing deterministic instead of relying on a large
//! on-disk fixture or wall-clock sleeps to approach the 5s busy timeout.
use super::*;
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use tempfile::TempDir;
/// A one-shot latch that blocks a blocking-pool thread until the test releases
/// it. Deterministic stand-in for "the M4 winner holds init open past the busy
/// timeout": the init/closure thread parks here, the test observes the frozen
/// state, then releases. `Mutex` + `Condvar` are `Sync`, so an `Arc<Latch>`
/// captured by the `Send + Sync` init hook type-checks.
struct Latch {
open: Mutex<bool>,
cv: Condvar,
}
impl Latch {
fn new() -> Arc<Self> {
Arc::new(Self {
open: Mutex::new(false),
cv: Condvar::new(),
})
}
/// Block until [`release`](Self::release) is called (returns immediately if
/// already released).
fn wait(&self) {
let mut open = self.open.lock().unwrap();
while !*open {
open = self.cv.wait(open).unwrap();
}
}
fn release(&self) {
*self.open.lock().unwrap() = true;
self.cv.notify_all();
}
}
/// Poll `cond` on the async runtime until it holds or `timeout` elapses.
/// Panics on timeout so a broken barrier surfaces as a failure, never a hang.
async fn await_until(what: &str, timeout: Duration, cond: impl Fn() -> bool) {
let deadline = Instant::now() + timeout;
while !cond() {
assert!(Instant::now() < deadline, "timed out waiting for {what}");
tokio::time::sleep(Duration::from_millis(5)).await;
}
}
/// An archive DB path inside a fresh temp dir. The dir is returned so the
/// caller keeps it alive for the whole test (dropping it deletes the file).
fn temp_db() -> (TempDir, std::path::PathBuf) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("archive").join("archive.db");
(dir, path)
}
/// The barrier serializes first-open: exactly one initialization runs while
/// every concurrent `with_conn` caller awaits it, and no ordinary connection
/// opens until that initialization has completed.
///
/// The init hook holds the single init task open on a latch. While it is held
/// we prove no `with_conn` caller has opened its connection (`open_count == 0`)
/// — impossible if callers bypassed the `OnceCell` and opened independently.
/// Releasing the latch lets init finish; all callers then complete, exactly
/// one initialization ran, and every open happened after init.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_first_open_barrier_serializes_init_and_defers_opens() {
let (_dir, path) = temp_db();
let init_count = Arc::new(AtomicUsize::new(0));
let open_count = Arc::new(AtomicUsize::new(0));
let release = Latch::new();
let hook = {
let init_count = Arc::clone(&init_count);
let release = Arc::clone(&release);
Arc::new(move || {
// Runs once, at the head of the single init task. Record the
// initialization, then park so the test can inspect the frozen
// pre-open state.
init_count.fetch_add(1, Ordering::SeqCst);
release.wait();
}) as InitHook
};
let db = Arc::new(ArchiveDb::with_test_hook(path, hook));
// Trigger the single init and wait until it is provably in flight (hook
// ran) and parked on the latch.
let warm = {
let db = Arc::clone(&db);
tokio::spawn(async move { db.warm_init().await })
};
await_until("init to start", Duration::from_secs(10), || {
init_count.load(Ordering::SeqCst) == 1
})
.await;
// Fan out production-shaped callers while init is held. Each records that
// its task was scheduled (`entered`) and that its closure actually opened a
// connection (`open_count`).
let entered = Arc::new(AtomicUsize::new(0));
let callers: Vec<_> = (0..4)
.map(|_| {
let db = Arc::clone(&db);
let open_count = Arc::clone(&open_count);
let entered = Arc::clone(&entered);
tokio::spawn(async move {
entered.fetch_add(1, Ordering::SeqCst);
db.with_conn(move |conn| {
open_count.fetch_add(1, Ordering::SeqCst);
// Touch the migrated schema to prove a usable connection.
conn.query_row("SELECT COUNT(*) FROM retention_policies", [], |r| {
r.get::<_, i64>(0)
})
.map_err(|e| e.to_string())
})
.await
})
})
.collect();
// All four caller tasks are scheduled and running before we judge the
// barrier: they have entered `with_conn` and can only be parked on the
// init `OnceCell`. Without the barrier they would instead open independent
// connections here and bump `open_count` while init is still held.
await_until("callers to be scheduled", Duration::from_secs(10), || {
entered.load(Ordering::SeqCst) == 4
})
.await;
tokio::time::sleep(Duration::from_millis(100)).await;
assert_eq!(
open_count.load(Ordering::SeqCst),
0,
"no ordinary connection may open until initialization completes"
);
// Let init finish; every caller now completes against the migrated DB.
release.release();
assert!(warm.await.unwrap().is_ok(), "warm init must succeed");
for caller in callers {
assert!(
caller.await.unwrap().is_ok(),
"every with_conn must succeed"
);
}
assert_eq!(
init_count.load(Ordering::SeqCst),
1,
"exactly one initialization ran behind the barrier"
);
assert_eq!(
open_count.load(Ordering::SeqCst),
4,
"all callers opened, and only after init"
);
}
/// The maintenance read guard lives for the FULL lifetime of a `with_conn`
/// connection: a write-lock contender cannot enter until the closure returns
/// and its connection has dropped. This is the invariant the Phase-4
/// sole-connection VACUUM depends on.
///
/// A `with_conn` closure parks on a latch while holding its connection; a
/// separate task contends for the maintenance write guard. While the closure
/// is parked the writer must be blocked. Releasing the closure — which returns
/// and drops the connection — lets the writer finally acquire the guard.
#[tokio::test(flavor = "multi_thread", worker_threads = 4)]
async fn test_with_conn_read_guard_blocks_writer_until_connection_drops() {
let (_dir, path) = temp_db();
let db = Arc::new(ArchiveDb::with_test_path(path));
db.warm_init().await.expect("init must succeed");
let in_closure = Arc::new(AtomicUsize::new(0));
let hold = Latch::new();
// A live connection: the closure parks holding it (and thus the read
// guard) until released.
let work = {
let db = Arc::clone(&db);
let in_closure = Arc::clone(&in_closure);
let hold = Arc::clone(&hold);
tokio::spawn(async move {
db.with_conn(move |_conn| {
in_closure.fetch_add(1, Ordering::SeqCst);
hold.wait();
Ok(())
})
.await
})
};
await_until(
"closure to hold the connection",
Duration::from_secs(10),
|| in_closure.load(Ordering::SeqCst) == 1,
)
.await;
// A genuine write-lock contender.
let write_entered = Arc::new(AtomicUsize::new(0));
let writer = {
let db = Arc::clone(&db);
let write_entered = Arc::clone(&write_entered);
tokio::spawn(async move {
let _w = db.maintenance.write().await;
write_entered.fetch_add(1, Ordering::SeqCst);
})
};
// While the connection is held, the writer must not have entered, and the
// write guard must be unavailable.
tokio::time::sleep(Duration::from_millis(200)).await;
assert_eq!(
write_entered.load(Ordering::SeqCst),
0,
"writer must block while a with_conn connection is live"
);
assert!(
!db.maintenance_write_available(),
"write guard is unavailable while the read guard is held"
);
// Release the closure: it returns and its connection drops, releasing the
// read guard so the writer can proceed.
hold.release();
assert!(work.await.unwrap().is_ok(), "held with_conn must succeed");
writer.await.unwrap();
assert_eq!(
write_entered.load(Ordering::SeqCst),
1,
"writer enters once the connection has dropped"
);
assert!(
db.maintenance_write_available(),
"write guard is free again after the connection drops"
);
}
+138 -1
View File
@@ -12,7 +12,7 @@
//! the existing `metric_store.rs` / `pipeline.rs` / `store_migrations.rs`
//! precedent.
use rusqlite::{params, Connection};
use rusqlite::{params, Connection, OptionalExtension};
// ── Constants ────────────────────────────────────────────────────────────────
@@ -63,6 +63,143 @@ CREATE INDEX IF NOT EXISTS idx_archived_event_scopes_age
(identity_pubkey, relay_url, scope_type, scope_value, archived_at, id);
";
/// Name of the scope-age index, shared by the DDL, the shape check, and the
/// M4 drop-and-rebuild repair path.
pub(super) const SCOPE_AGE_INDEX_NAME: &str = "idx_archived_event_scopes_age";
// ── Expected schema shape (the single source of truth M4 validates against) ────
/// One column's expected shape: `(name, declared_type, not_null, pk_position)`.
/// `pk_position` is 0 when the column is not part of the primary key, else its
/// 1-based position in the PK — this is exactly what `PRAGMA table_info` reports
/// in its `pk` field, so PK column ORDER is validated, not just membership.
type ColumnShape = (&'static str, &'static str, bool, i64);
/// Expected `retention_policies` columns in `cid` order. Mirrors
/// [`RETENTION_SCHEMA`]; a drift here or there is caught by the M4 shape check.
const RETENTION_POLICIES_SHAPE: &[ColumnShape] = &[
("identity_pubkey", "TEXT", true, 1),
("relay_url", "TEXT", true, 2),
("scope_type", "TEXT", true, 3),
("scope_value", "TEXT", true, 4),
("kind", "INTEGER", true, 5),
("days", "INTEGER", false, 0),
("updated_at", "INTEGER", true, 0),
];
/// Expected `archive_meta` columns in `cid` order. A `TEXT PRIMARY KEY` is
/// nullable in SQLite (only `INTEGER PRIMARY KEY` implies NOT NULL), so `key`
/// carries `not_null = false` with `pk_position = 1`.
const ARCHIVE_META_SHAPE: &[ColumnShape] = &[("key", "TEXT", false, 1), ("value", "TEXT", true, 0)];
/// Expected key columns of [`SCOPE_AGE_INDEX_NAME`] in seqno order — the
/// age-range access path Phase 2 depends on. Order is load-bearing: a covering
/// seek needs the scope keys before `archived_at`.
const SCOPE_AGE_INDEX_SHAPE: &[&str] = &[
"identity_pubkey",
"relay_url",
"scope_type",
"scope_value",
"archived_at",
"id",
];
// ── Shape validation (used by migration M4) ────────────────────────────────────
/// Whether `table`'s live columns exactly match `expected` (name, declared
/// type, nullability, and PK position, all in `cid` order). `table` is always a
/// compile-time constant from this module, never user input, so interpolating
/// it into the table-valued `pragma_table_info` call carries no injection risk.
fn table_shape_matches(
conn: &Connection,
table: &str,
expected: &[ColumnShape],
) -> Result<bool, String> {
let mut stmt = conn
.prepare(&format!(
"SELECT name, type, \"notnull\", pk FROM pragma_table_info('{table}') ORDER BY cid"
))
.map_err(|e| format!("shape check: prepare table_info({table}): {e}"))?;
let actual: Vec<(String, String, i64, i64)> = stmt
.query_map([], |r| Ok((r.get(0)?, r.get(1)?, r.get(2)?, r.get(3)?)))
.map_err(|e| format!("shape check: query table_info({table}): {e}"))?
.collect::<Result<Vec<_>, _>>()
.map_err(|e| format!("shape check: read table_info({table}): {e}"))?;
Ok(actual.len() == expected.len()
&& actual.iter().zip(expected).all(
|((name, ty, not_null, pk), (exp_name, exp_ty, exp_not_null, exp_pk))| {
name == exp_name
&& ty.eq_ignore_ascii_case(exp_ty)
&& (*not_null != 0) == *exp_not_null
&& pk == exp_pk
},
))
}
/// Whether the scope-age index exists on `archived_event_scopes` with exactly
/// the expected key columns in order. False when the index is absent, sits on
/// the wrong table, or indexes the wrong columns — all of which M4 repairs by
/// dropping and recreating it (an index carries no data, so a rebuild is safe).
pub(super) fn scope_age_index_is_correct(conn: &Connection) -> Result<bool, String> {
let table: Option<String> = conn
.query_row(
"SELECT tbl_name FROM sqlite_master WHERE type = 'index' AND name = ?1",
params![SCOPE_AGE_INDEX_NAME],
|r| r.get(0),
)
.optional()
.map_err(|e| format!("shape check: scope-age index table: {e}"))?;
if table.as_deref() != Some("archived_event_scopes") {
return Ok(false);
}
let mut stmt = conn
.prepare(&format!(
"SELECT name FROM pragma_index_info('{SCOPE_AGE_INDEX_NAME}') ORDER BY seqno"
))
.map_err(|e| format!("shape check: prepare index_info: {e}"))?;
let columns: Vec<String> = stmt
.query_map([], |r| r.get::<_, String>(0))
.map_err(|e| format!("shape check: query index_info: {e}"))?
.collect::<Result<Vec<_>, _>>()
.map_err(|e| format!("shape check: read index_info: {e}"))?;
Ok(columns.len() == SCOPE_AGE_INDEX_SHAPE.len()
&& columns
.iter()
.zip(SCOPE_AGE_INDEX_SHAPE)
.all(|(actual, expected)| actual == expected))
}
/// Validate the COMPLETE expected shape of the M4 objects: both tables' columns
/// (names, types, nullability, PK positions) and the scope-age index's key
/// order. `CREATE ... IF NOT EXISTS` silently preserves a wrong-shaped object
/// that already carries the expected name, so name presence alone cannot
/// certify the schema — a wrong-shaped named table would otherwise let M4 mark
/// itself applied and hand Phase 2 an unusable table or missing access path.
/// Called inside M4's `BEGIN IMMEDIATE`; an `Err` rolls the transaction back
/// with no marker written.
pub(super) fn validate_retention_schema_shape(conn: &Connection) -> Result<(), String> {
if !table_shape_matches(conn, "retention_policies", RETENTION_POLICIES_SHAPE)? {
return Err(
"migration M4: retention_policies has an unexpected column or primary-key shape"
.to_string(),
);
}
if !table_shape_matches(conn, "archive_meta", ARCHIVE_META_SHAPE)? {
return Err(
"migration M4: archive_meta has an unexpected column or primary-key shape".to_string(),
);
}
if !scope_age_index_is_correct(conn)? {
return Err(format!(
"migration M4: {SCOPE_AGE_INDEX_NAME} has an unexpected shape after rebuild"
));
}
Ok(())
}
// ── Types ─────────────────────────────────────────────────────────────────────
/// A retention policy row with its live-subscription status, returned by
@@ -397,6 +397,75 @@ fn test_m4_partial_schema_marker_absent_recovers_on_next_open() {
assert_eq!(objects, 3, "partial schema repaired to the full shape");
}
#[test]
fn test_m4_wrong_shaped_named_table_rejected_no_marker() {
let db = NamedTempFile::new().unwrap();
build_pre_m4_db(db.path());
// Precreate an `archive_meta` table carrying the expected NAME but the
// wrong shape (its required `value` column is missing). `CREATE ... IF NOT
// EXISTS` preserves it, so only the explicit shape check can catch it.
{
let conn = Connection::open(db.path()).unwrap();
conn.execute_batch("CREATE TABLE archive_meta (key TEXT PRIMARY KEY);")
.unwrap();
}
// M4 must reject the incompatible named table and roll back with no marker,
// rather than certify a table Phase 2 could not use.
assert!(
store::open_archive_db(db.path()).is_err(),
"M4 must reject a wrong-shaped named archive_meta"
);
let verify = Connection::open(db.path()).unwrap();
assert_eq!(
m4_marker_count(&verify),
0,
"no marker may certify the incompatible table"
);
// The rollback left the schema untouched: the bad table is still one-column
// and retention_policies was never committed.
assert!(
!scope_age_index_is_correct(&verify).unwrap(),
"the index build rolled back with the rest of the transaction"
);
let value_cols: i64 = verify
.query_row(
"SELECT COUNT(*) FROM pragma_table_info('archive_meta') WHERE name = 'value'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(
value_cols, 0,
"the wrong-shaped table was not silently altered"
);
}
#[test]
fn test_m4_wrong_index_order_dropped_and_rebuilt() {
let db = NamedTempFile::new().unwrap();
build_pre_m4_db(db.path());
// Precreate the scope-age index with the expected NAME on the right table
// but the WRONG key order (`archived_at` first). A covering age-range seek
// needs the scope keys before `archived_at`, so this order is unusable.
{
let conn = Connection::open(db.path()).unwrap();
conn.execute_batch(&format!(
"CREATE INDEX {SCOPE_AGE_INDEX_NAME}
ON archived_event_scopes
(archived_at, identity_pubkey, relay_url, scope_type, scope_value, id);"
))
.unwrap();
}
// M4 rebuilds a wrong-shaped index (safe — an index carries no data) rather
// than rejecting, so the open succeeds and the marker lands.
let conn = fresh(&db);
assert_eq!(m4_marker_count(&conn), 1, "M4 completes after the rebuild");
assert!(
scope_age_index_is_correct(&conn).unwrap(),
"the index was rebuilt into the correct key order"
);
}
// ── Concurrency ───────────────────────────────────────────────────────────────
#[test]
@@ -481,3 +550,156 @@ fn test_concurrent_merge_seeds_both_policies() {
assert_eq!(policy_days(&verify, OBSERVER_FRAME_KIND), Some(30));
assert_eq!(policy_days(&verify, 44200), None);
}
/// Deterministic interleaving of merge / remove / set on ONE `(scope, kind)`
/// across two connections to the same WAL file. Each mutator takes `BEGIN
/// IMMEDIATE`, so a barrier held before both begin forces them to serialize on
/// `busy_timeout` rather than race to a `BUSY_SNAPSHOT`. Whichever order the OS
/// picks, the invariants must hold: the subscription stays valid, the explicit
/// policy choice set by one thread is never silently reset by the other's
/// merge-seed (which is `ON CONFLICT DO NOTHING`), and no policy row is deleted
/// by a `kinds` mutation. Removing the `BEGIN IMMEDIATE` guard from the
/// mutators makes one contender fail with `BUSY_SNAPSHOT`, tripping this test.
#[test]
fn test_concurrent_merge_remove_set_keeps_state_valid_and_choice() {
use std::thread;
let db = NamedTempFile::new().unwrap();
let path = db.path().to_path_buf();
// Start from a subscription that already lists 24200 with an EXPLICIT
// 7-day policy — the choice both racing threads must not clobber.
{
let conn = fresh(&db);
store::merge_owner_p_kinds(&conn, ID, RELAY, ID, 24200, 100).unwrap();
set_policy(&conn, ID, RELAY, OWNER, ID, 24200, Some(7), 100).unwrap();
}
let barrier = Arc::new(Barrier::new(3));
// T1: enable a second kind (seeds its default policy).
// T2: disable 24200 (removes it from `kinds`, must NOT touch its policy).
// T3: re-set 24200's explicit policy to 5 days.
let t1 = {
let (p, b) = (path.clone(), Arc::clone(&barrier));
thread::spawn(move || {
let conn = store::open_archive_db(&p).unwrap();
b.wait();
store::merge_owner_p_kinds(&conn, ID, RELAY, ID, 44200, 200)
})
};
let t2 = {
let (p, b) = (path.clone(), Arc::clone(&barrier));
thread::spawn(move || {
let conn = store::open_archive_db(&p).unwrap();
b.wait();
store::remove_owner_p_kind(&conn, ID, RELAY, ID, 24200)
})
};
let t3 = {
let (p, b) = (path.clone(), Arc::clone(&barrier));
thread::spawn(move || {
let conn = store::open_archive_db(&p).unwrap();
b.wait();
set_policy(&conn, ID, RELAY, OWNER, ID, 24200, Some(5), 300)
})
};
assert!(t1.join().unwrap().is_ok(), "concurrent merge must succeed");
assert!(t2.join().unwrap().is_ok(), "concurrent remove must succeed");
assert!(t3.join().unwrap().is_ok(), "concurrent set must succeed");
let verify = store::open_archive_db(&path).unwrap();
// Subscription is valid JSON and now lists 44200 (added) but not 24200
// (removed) — the two `kinds` mutations composed cleanly.
let subs = store::list_save_subscriptions(&verify, ID, RELAY).unwrap();
assert_eq!(subs.len(), 1, "one owner_p row survives the interleaving");
let sub_kinds: Vec<u32> = serde_json::from_str(&subs[0].kinds).unwrap();
assert!(sub_kinds.contains(&44200), "the merged kind is present");
assert!(!sub_kinds.contains(&24200), "the removed kind is gone");
// Both policy rows still exist — no `kinds` mutation deleted one.
let policies = list_policies(&verify, ID, RELAY).unwrap();
assert_eq!(policies.len(), 2, "both policy rows survive; none deleted");
// T3's explicit 5-day choice for 24200 is the final value: T1's default
// seed for 44200 never touches 24200, and T2's remove leaves policies
// alone, so the observer policy reflects the explicit set, not a reset.
assert_eq!(
policy_days(&verify, 24200),
Some(5),
"explicit choice survives the concurrent merge/remove"
);
assert_eq!(
policy_days(&verify, 44200),
None,
"44200 kept its Forever default"
);
}
/// One evolving state walked through the full policy lifecycle in a single
/// test: an active observer policy → the kind is disabled (policy orphaned but
/// preserved) → the whole subscription is deleted (still orphaned) → the orphan
/// is edited → the orphan is explicitly deleted. This proves the transitions
/// compose on the SAME rows, which the separate per-transition tests cannot.
#[test]
fn test_policy_lifecycle_active_orphaned_edited_deleted_on_one_state() {
let db = NamedTempFile::new().unwrap();
let conn = fresh(&db);
// 1. Active: a subscription lists 24200, its policy is the seeded default.
create_subscription_with_policies(&conn, ID, RELAY, OWNER, ID, &[24200], "[24200]", 100)
.unwrap();
let active = list_policies(&conn, ID, RELAY).unwrap();
assert_eq!(active.len(), 1);
assert!(
active[0].active,
"policy is active while the kind is listed"
);
assert_eq!(active[0].days, Some(30), "seeded observer default");
// 2. Kind disabled: 24200 leaves `kinds`; the policy stays but goes orphaned.
store::remove_owner_p_kind(&conn, ID, RELAY, ID, 24200).unwrap();
// Removing the last kind deleted the subscription row entirely.
assert!(
store::list_save_subscriptions(&conn, ID, RELAY)
.unwrap()
.is_empty(),
"the last kind off removes the subscription row"
);
let disabled = list_policies(&conn, ID, RELAY).unwrap();
assert_eq!(disabled.len(), 1, "the policy is preserved, not deleted");
assert!(
!disabled[0].active,
"policy is orphaned once the kind is gone"
);
assert_eq!(
disabled[0].days,
Some(30),
"it keeps expiring historical data"
);
// 3. Subscription deleted: an explicit delete on an already-absent row is a
// no-op; the orphaned policy is unaffected.
assert!(
!store::delete_save_subscription(&conn, ID, RELAY, OWNER, ID).unwrap(),
"no subscription row remains to delete"
);
let still_orphaned = list_policies(&conn, ID, RELAY).unwrap();
assert_eq!(still_orphaned.len(), 1);
assert!(!still_orphaned[0].active, "policy remains orphaned");
// 4. Edit the orphan: a user can still change retention for historical data.
set_policy(&conn, ID, RELAY, OWNER, ID, 24200, Some(90), 400).unwrap();
let edited = list_policies(&conn, ID, RELAY).unwrap();
assert_eq!(edited.len(), 1);
assert_eq!(edited[0].days, Some(90), "orphaned policy is editable");
assert!(
!edited[0].active,
"editing does not resurrect the subscription"
);
// 5. Delete the orphan: the only path that removes a policy row.
assert!(
delete_policy(&conn, ID, RELAY, OWNER, ID, 24200).unwrap(),
"the orphaned policy is deleted"
);
assert!(
list_policies(&conn, ID, RELAY).unwrap().is_empty(),
"no policy rows remain after the lifecycle completes"
);
}
@@ -379,22 +379,35 @@ fn migrate_add_retention_policies_locked(conn: &Connection) -> Result<(), String
return Ok(());
}
// Idempotent DDL — repairs a partial state left by any interrupted run.
// Idempotent DDL — creates any missing objects, repairing a partial state
// left by an interrupted earlier run. `IF NOT EXISTS` preserves an object
// that already carries the expected name, so it cannot fix a wrong SHAPE —
// that is what the explicit validation below is for.
conn.execute_batch(super::retention::RETENTION_SCHEMA)
.map_err(|e| format!("migration M4: create retention schema: {e}"))?;
conn.execute_batch(super::retention::SCOPE_AGE_INDEX_DDL)
.map_err(|e| format!("migration M4: create scope-age index: {e}"))?;
// Schema-shape recheck: confirm all three objects (two tables + one index)
// actually materialized before seeding or writing the marker, so a marker
// never certifies a half-built schema.
if retention_schema_object_count(conn)? != 3 {
return Err(
"migration M4: retention schema incomplete after DDL (expected 2 tables + 1 index)"
.to_string(),
);
// Repair a missing or wrong-shaped scope-age index by dropping and
// recreating it. An index carries no data, so a rebuild is always safe —
// unlike a table, whose wrong shape we must reject rather than drop.
if !super::retention::scope_age_index_is_correct(conn)? {
conn.execute_batch(&format!(
"DROP INDEX IF EXISTS {}",
super::retention::SCOPE_AGE_INDEX_NAME
))
.map_err(|e| format!("migration M4: drop wrong-shaped scope-age index: {e}"))?;
conn.execute_batch(super::retention::SCOPE_AGE_INDEX_DDL)
.map_err(|e| format!("migration M4: create scope-age index: {e}"))?;
}
// Validate the COMPLETE expected shape (both tables' columns, types,
// nullability, and PK positions, plus the index key order) before seeding
// or writing the marker. A wrong-shaped named TABLE cannot be auto-repaired
// without risking archived data, so it is rejected here and the whole
// transaction rolls back with no marker — the next open re-runs M4 once the
// schema is corrected. This guarantees the marker never certifies a
// half-built or mis-shaped schema that Phase 2 would inherit.
super::retention::validate_retention_schema_shape(conn)?;
let now = std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.unwrap_or_default()
@@ -425,20 +438,6 @@ fn retention_migration_applied(conn: &Connection) -> Result<bool, String> {
Ok(count > 0)
}
/// Count the M4 schema objects present: the `retention_policies` and
/// `archive_meta` tables plus the `idx_archived_event_scopes_age` index. A
/// fully-created schema returns 3.
fn retention_schema_object_count(conn: &Connection) -> Result<i64, String> {
conn.query_row(
"SELECT COUNT(*) FROM sqlite_master
WHERE (type = 'table' AND name IN ('retention_policies', 'archive_meta'))
OR (type = 'index' AND name = 'idx_archived_event_scopes_age')",
[],
|r| r.get(0),
)
.map_err(|e| format!("migration M4: schema-shape check: {e}"))
}
/// Seed a default retention policy for every `(subscription, kind)` pair in
/// `save_subscriptions`: kind 24200 → 30 days, every other kind → an explicit
/// `NULL` (Forever) row. Runs on the caller's transaction so seeding is atomic