fix(#2279): versioned DB migrations via schema_versions table

Replace the try-and-ignore-duplicate-column approach in apply_migrations
with proper schema versioning using a shared schema_versions table.

## mechanism

New function: db::apply_versioned_migrations(conn, subsystem,
legacy_column, migrations). Tracks the applied-migration count in a
schema_versions table (one row per subsystem key). Only migrations past
the stored version run.

Legacy detection: pre-versioning databases have no schema_versions row.
The legacy_column tuple (table, column) identifies a column that exists
only in a fully-migrated legacy database. If present, all known
migrations are skipped. If absent, migrations start from 0.

## stores migrated

- broker: removes bespoke ensure_message_columns / ensure_reminder_columns.
  Unified into BROKER_MIGRATIONS (v1-v5). Legacy detector: messages.priority
  (added in the last pre-versioning migration).
- approvals: 4 historical migrations (v1-v4). Legacy detector:
  approvals.submitter.
- operator_questions: 3 historical migrations (v1-v3). Legacy detector:
  operator_questions.target.
- scheduled_prompts: 1 historical migration (v1). Legacy detector:
  scheduled_prompts.paused_at_unix.

apply_migrations removed (no callers).

## tests (db.rs)

- fresh_install_runs_all_migrations
- legacy_install_skips_all_migrations
- partial_migration_resumes_from_version
- already_at_latest_is_noop
- multiple_stores_in_same_db
This commit is contained in:
atlas 2026-07-10 14:50:39 +02:00 committed by mara
commit 7d36ec5e1f
5 changed files with 333 additions and 115 deletions

View file

@ -24,24 +24,19 @@ CREATE INDEX IF NOT EXISTS idx_approvals_pending
ON approvals (id) WHERE status = 'pending';
";
/// Additive column migrations for pre-existing databases, applied via
/// `db::apply_migrations` (try-and-ignore-duplicate-column).
/// Ordered schema migrations tracked in `schema_versions` (key `"approvals"`).
/// Legacy databases are detected via the `submitter` column — the last
/// column added before versioning was introduced — and fast-forwarded past
/// all known migrations. New columns go here as v5, v6, …
const MIGRATIONS: &[&str] = &[
// `kind` (pre-Phase-8 dbs): legacy rows default to `apply_commit`,
// which matches their actual semantics.
// v1: `kind` (pre-Phase-8 dbs): legacy rows default to `apply_commit`.
"ALTER TABLE approvals ADD COLUMN kind TEXT NOT NULL DEFAULT 'apply_commit'",
// `description`: manager-supplied note shown on the dashboard
// approval card at submission time (distinct from `note`, set on
// denial/failure).
// v2: `description`: manager-supplied note on the dashboard card.
"ALTER TABLE approvals ADD COLUMN description TEXT",
// `fetched_sha`: the canonical sha hive-c0re vouched for at
// `request_apply_commit` time. Distinct from `commit_ref`
// (manager-supplied, may not even resolve by approve time).
// v3: `fetched_sha`: canonical sha hive-c0re resolved at submit time.
"ALTER TABLE approvals ADD COLUMN fetched_sha TEXT",
// `submitter`: the agent that submitted the approval (the
// authenticated socket caller); approval-scoped helper events
// route to it. Legacy rows are NULL → callers fall back to the
// root agent.
// v4: `submitter`: authenticated agent that submitted the approval.
// Legacy rows are NULL → callers fall back to the root agent.
"ALTER TABLE approvals ADD COLUMN submitter TEXT",
];
@ -54,7 +49,12 @@ impl Approvals {
let conn = crate::db::open(path, "approvals")?;
conn.execute_batch(SCHEMA)
.context("apply approvals schema")?;
crate::db::apply_migrations(&conn, "approvals", MIGRATIONS)?;
crate::db::apply_versioned_migrations(
&conn,
"approvals",
("approvals", "submitter"),
MIGRATIONS,
)?;
Ok(Self {
conn: Mutex::new(conn),
})

View file

@ -156,12 +156,42 @@ pub struct Broker {
inflight: Mutex<HashMap<String, RecipientInflight>>,
}
/// Ordered schema migrations for the broker. Tracked in `schema_versions`
/// under key `"broker"`. Legacy databases (fully migrated via the old
/// per-column check approach) are detected via the `priority` column on
/// `messages` — the last column added before versioning.
const BROKER_MIGRATIONS: &[&str] = &[
// v1: acked_at on messages, with backfill so existing delivered rows
// are not phantom-requeued on the next open. The multi-statement
// batch runs atomically (execute_batch wraps in a transaction).
"ALTER TABLE messages ADD COLUMN acked_at INTEGER;\
UPDATE messages SET acked_at = delivered_at WHERE delivered_at IS NOT NULL;",
// v2: in_reply_to for thread-parent tracking. NULL = root of a thread.
"ALTER TABLE messages ADD COLUMN in_reply_to INTEGER",
// v3: priority for operator-message fast-path. Rebuild the delivery
// index to include priority as a secondary sort key.
"ALTER TABLE messages ADD COLUMN priority INTEGER NOT NULL DEFAULT 0;\
DROP INDEX IF EXISTS idx_messages_undelivered;\
CREATE INDEX idx_messages_undelivered \
ON messages (recipient, priority DESC, id) WHERE delivered_at IS NULL",
// v4: attempt_count on reminders for the MAX_REMINDER_ATTEMPTS cap.
"ALTER TABLE reminders ADD COLUMN attempt_count INTEGER NOT NULL DEFAULT 0",
// v5: last_error on reminders — last delivery failure surfaced on the
// dashboard so a stuck reminder is visible without digging in logs.
"ALTER TABLE reminders ADD COLUMN last_error TEXT",
];
impl Broker {
pub fn open(path: &Path) -> Result<Self> {
let conn = crate::db::open(path, "broker")?;
conn.execute_batch(SCHEMA).context("apply broker schema")?;
ensure_message_columns(&conn).context("migrate messages columns")?;
ensure_reminder_columns(&conn).context("migrate reminders columns")?;
crate::db::apply_versioned_migrations(
&conn,
"broker",
("messages", "priority"),
BROKER_MIGRATIONS,
)
.context("broker migrations")?;
let (events, _) = broadcast::channel(EVENT_CHANNEL);
Ok(Self {
conn: Mutex::new(conn),
@ -1108,65 +1138,6 @@ impl Broker {
/// back-fills it for every already-delivered row, so the
/// pre-migration sessions count as "fully handled" and won't be
/// resurfaced by the first `requeue_inflight` after upgrade.
fn ensure_message_columns(conn: &Connection) -> Result<()> {
let has_acked: bool = conn
.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = 'acked_at'")?
.exists([])?;
if !has_acked {
conn.execute_batch("ALTER TABLE messages ADD COLUMN acked_at INTEGER;")
.context("add messages.acked_at column")?;
// Backfill: treat every existing delivered row as acked. The
// session it was delivered to is gone, so requeue would just
// surface phantom traffic to whatever harness reads next.
conn.execute(
"UPDATE messages SET acked_at = delivered_at \
WHERE delivered_at IS NOT NULL AND acked_at IS NULL",
[],
)
.context("backfill messages.acked_at from delivered_at")?;
}
let has_reply: bool = conn
.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = 'in_reply_to'")?
.exists([])?;
if !has_reply {
conn.execute_batch("ALTER TABLE messages ADD COLUMN in_reply_to INTEGER;")
.context("add messages.in_reply_to column")?;
// No backfill needed — existing messages simply have NULL here,
// meaning "root of a new thread", which is correct.
}
let has_priority: bool = conn
.prepare("SELECT 1 FROM pragma_table_info('messages') WHERE name = 'priority'")?
.exists([])?;
if !has_priority {
// Add column + rebuild the index with priority as the second key so
// recv_batch can serve operator messages ahead of queued wakes without
// a separate sort pass.
conn.execute_batch(
"ALTER TABLE messages ADD COLUMN priority INTEGER NOT NULL DEFAULT 0;
DROP INDEX IF EXISTS idx_messages_undelivered;
CREATE INDEX idx_messages_undelivered
ON messages (recipient, priority DESC, id) WHERE delivered_at IS NULL;",
)
.context("add messages.priority column and update delivery index")?;
// No backfill needed — all existing rows are correctly at priority 0.
}
Ok(())
}
/// Idempotent reminder-table migrations — plain additive columns, via
/// `db::apply_migrations`. (The messages-table migration above stays
/// bespoke: its backfill must run only when `acked_at` was just
/// created, which try-and-ignore can't express.)
fn ensure_reminder_columns(conn: &Connection) -> Result<()> {
crate::db::apply_migrations(
conn,
"broker reminders",
&[
"ALTER TABLE reminders ADD COLUMN attempt_count INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE reminders ADD COLUMN last_error TEXT",
],
)
}
#[cfg(test)]
mod tests {

View file

@ -1,43 +1,114 @@
//! Shared sqlite connection setup for hive-c0re's host-side stores.
//! Shared sqlite connection setup and versioned schema migrations for
//! hive-c0re's host-side stores.
//!
//! Several modules keep their own tables — and their own
//! `Mutex<Connection>` — in the coordinator DB
//! (`db/broker.sqlite`: broker, approvals, operator questions,
//! scheduled prompts, agent power) or in a sibling file under the same
//! `db/` dir (`build_logs.sqlite`, `audit_log.sqlite`). The open dance
//! is identical everywhere: ensure the parent dir exists, open the
//! connection, set a busy timeout so concurrent same-process writers
//! wait each other out instead of surfacing `SQLITE_BUSY`. This helper
//! owns that dance; schema creation + column migrations stay with each
//! store (they're per-table concerns).
//! All stores share one sqlite file (`db/broker.sqlite`) but each owns
//! its own `Mutex<Connection>`. [`open`] handles the connection open
//! dance; [`apply_versioned_migrations`] handles schema evolution.
//!
//! ## Migration strategy
//!
//! Each store calls [`apply_versioned_migrations`] from its `open`. It
//! tracks the applied count in a `schema_versions` table (one row per
//! `subsystem` key). Only migrations past the stored version run.
//!
//! Legacy databases (written before versioning) are detected via
//! `legacy_column = (table, column)` — a column that exists only in
//! a fully-migrated legacy DB. Present → skip all migrations. Absent →
//! start from 0. See each store's `MIGRATIONS` constant for the list.
use std::path::Path;
use std::time::Duration;
use anyhow::{Context, Result};
use rusqlite::Connection;
use rusqlite::{Connection, OptionalExtension, params};
/// How long a write waits on another connection's lock before erroring.
/// Generous relative to the stores' tiny transactions — a timeout here
/// means something is genuinely wedged, not ordinary contention.
const BUSY_TIMEOUT: Duration = Duration::from_secs(5);
/// Apply additive, idempotent migrations: run each statement and
/// ignore `duplicate column name` errors (sqlite has no
/// `ADD COLUMN IF NOT EXISTS`; try-and-ignore is the portable path —
/// same pattern as hive-ag3nt's `turn_stats`). Any other error
/// propagates. Fits plain `ALTER TABLE ADD COLUMN` (new columns must
/// carry a default or tolerate NULL) and `IF NOT EXISTS` index DDL;
/// migrations with creation-conditional backfills (broker's
/// `acked_at`) stay bespoke in their store.
pub fn apply_migrations(conn: &Connection, subsystem: &str, statements: &[&str]) -> Result<()> {
for stmt in statements {
if let Err(e) = conn.execute_batch(stmt) {
if e.to_string().contains("duplicate column name") {
continue;
}
return Err(e).with_context(|| format!("{subsystem} migration failed: {stmt}"));
}
const SCHEMA_VERSIONS_DDL: &str = "
CREATE TABLE IF NOT EXISTS schema_versions (
store TEXT PRIMARY KEY,
version INTEGER NOT NULL DEFAULT 0
)";
/// Apply numbered schema migrations for `subsystem`, tracking progress in
/// the shared `schema_versions` table.
///
/// - `subsystem` — unique key for this store in `schema_versions`
/// (e.g. `"approvals"`, `"broker"`).
/// - `legacy_column` — `(table_name, column_name)` of a column that exists
/// in a fully-migrated pre-versioning database. When there is no
/// `schema_versions` row yet and this column is present, all `migrations`
/// are skipped (they were previously applied via the old try-and-ignore
/// approach). When the column is absent, we start at version 0.
/// - `migrations` — ordered list of SQL batches; each batch runs once, at
/// the index whose value is the current version.
///
/// # Errors
///
/// Returns an error if any migration statement fails or the version update
/// fails. A failed migration leaves the DB at the last successfully
/// committed version (each step is its own implicit transaction via
/// `execute_batch`).
pub fn apply_versioned_migrations(
conn: &Connection,
subsystem: &str,
legacy_column: (&str, &str),
migrations: &[&str],
) -> Result<()> {
// Ensure the version-tracking table exists (idempotent).
conn.execute_batch(SCHEMA_VERSIONS_DDL)
.context("create schema_versions table")?;
// Look up the stored version for this subsystem.
let stored: Option<i64> = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = ?1",
[subsystem],
|row| row.get(0),
)
.optional()
.context("read schema_versions")?;
let version = if let Some(v) = stored {
usize::try_from(v.max(0)).unwrap_or(0)
} else {
// No entry yet. Check the legacy-column to decide where to start.
// If it exists, all known migrations were already applied by the
// old try-and-ignore path — skip them. If it doesn't exist, start
// at 0 (fresh install or partial state).
let (legacy_table, legacy_col) = legacy_column;
let is_legacy: bool = conn
.prepare(&format!(
"SELECT 1 FROM pragma_table_info('{legacy_table}') WHERE name = ?1"
))
.context("prepare legacy-column probe")?
.exists([legacy_col])
.context("check legacy column")?;
let initial = if is_legacy { migrations.len() } else { 0 };
conn.execute(
"INSERT INTO schema_versions (store, version) VALUES (?1, ?2)",
params![subsystem, i64::try_from(initial).unwrap_or(0)],
)
.context("insert schema_versions row")?;
initial
};
if version >= migrations.len() {
return Ok(());
}
for (i, stmt) in migrations.iter().enumerate().skip(version) {
conn.execute_batch(stmt)
.with_context(|| format!("{subsystem} migration v{} failed: {stmt}", i + 1))?;
let next = i64::try_from(i + 1).unwrap_or(i64::MAX);
conn.execute(
"UPDATE schema_versions SET version = ?1 WHERE store = ?2",
params![next, subsystem],
)
.with_context(|| format!("{subsystem} update schema_versions to v{}", i + 1))?;
}
Ok(())
}
@ -56,3 +127,166 @@ pub fn open(path: &Path, subsystem: &str) -> Result<Connection> {
.with_context(|| format!("set {subsystem} busy_timeout"))?;
Ok(conn)
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::atomic::{AtomicU64, Ordering};
static TEST_COUNTER: AtomicU64 = AtomicU64::new(0);
fn open_tmp() -> Connection {
let n = TEST_COUNTER.fetch_add(1, Ordering::Relaxed);
let path = format!("/tmp/hive-db-test-{}-{}.sqlite", std::process::id(), n);
Connection::open(&path).expect("open tmp db")
}
const TEST_SCHEMA: &str = "CREATE TABLE IF NOT EXISTS things (id INTEGER PRIMARY KEY)";
const TEST_MIGRATIONS: &[&str] = &[
"ALTER TABLE things ADD COLUMN alpha TEXT",
"ALTER TABLE things ADD COLUMN beta INTEGER NOT NULL DEFAULT 0",
"ALTER TABLE things ADD COLUMN gamma TEXT",
];
#[test]
fn fresh_install_runs_all_migrations() {
let conn = open_tmp();
conn.execute_batch(TEST_SCHEMA).unwrap();
apply_versioned_migrations(&conn, "things", ("things", "gamma"), TEST_MIGRATIONS).unwrap();
// All three columns should exist.
for col in ["alpha", "beta", "gamma"] {
let exists: bool = conn
.prepare(&format!(
"SELECT 1 FROM pragma_table_info('things') WHERE name = '{col}'"
))
.unwrap()
.exists([])
.unwrap();
assert!(exists, "column {col} missing after fresh migration");
}
// Version should be 3.
let v: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'things'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(v, 3);
}
#[test]
fn legacy_install_skips_all_migrations() {
let conn = open_tmp();
// Simulate a fully-migrated legacy DB: all columns exist,
// but schema_versions does not yet.
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS things (id INTEGER PRIMARY KEY,
alpha TEXT, beta INTEGER NOT NULL DEFAULT 0, gamma TEXT)",
)
.unwrap();
apply_versioned_migrations(&conn, "things", ("things", "gamma"), TEST_MIGRATIONS).unwrap();
// Version should be set to migrations.len() (3), not 0.
let v: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'things'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(v, 3, "legacy install should skip to latest version");
}
#[test]
fn partial_migration_resumes_from_version() {
let conn = open_tmp();
// Simulate a DB that has only the first migration applied.
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS things (id INTEGER PRIMARY KEY, alpha TEXT);
CREATE TABLE IF NOT EXISTS schema_versions (store TEXT PRIMARY KEY, version INTEGER NOT NULL DEFAULT 0);
INSERT INTO schema_versions (store, version) VALUES ('things', 1)",
)
.unwrap();
apply_versioned_migrations(&conn, "things", ("things", "gamma"), TEST_MIGRATIONS).unwrap();
// Only beta and gamma should have been added (alpha already existed).
for col in ["alpha", "beta", "gamma"] {
let exists: bool = conn
.prepare(&format!(
"SELECT 1 FROM pragma_table_info('things') WHERE name = '{col}'"
))
.unwrap()
.exists([])
.unwrap();
assert!(
exists,
"column {col} missing after partial migration resume"
);
}
let v: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'things'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(v, 3);
}
#[test]
fn already_at_latest_is_noop() {
let conn = open_tmp();
conn.execute_batch(TEST_SCHEMA).unwrap();
apply_versioned_migrations(&conn, "things", ("things", "gamma"), TEST_MIGRATIONS).unwrap();
// Second call is a no-op.
apply_versioned_migrations(&conn, "things", ("things", "gamma"), TEST_MIGRATIONS).unwrap();
let v: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'things'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(v, 3);
}
#[test]
fn multiple_stores_in_same_db() {
let conn = open_tmp();
conn.execute_batch(
"CREATE TABLE IF NOT EXISTS a (id INTEGER PRIMARY KEY);
CREATE TABLE IF NOT EXISTS b (id INTEGER PRIMARY KEY)",
)
.unwrap();
apply_versioned_migrations(
&conn,
"store_a",
("a", "col1"),
&["ALTER TABLE a ADD COLUMN col1 TEXT"],
)
.unwrap();
apply_versioned_migrations(
&conn,
"store_b",
("b", "colx"),
&["ALTER TABLE b ADD COLUMN colx INTEGER NOT NULL DEFAULT 0"],
)
.unwrap();
// Each store tracks independently.
let va: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'store_a'",
[],
|r| r.get(0),
)
.unwrap();
let vb: i64 = conn
.query_row(
"SELECT version FROM schema_versions WHERE store = 'store_b'",
[],
|r| r.get(0),
)
.unwrap();
assert_eq!(va, 1);
assert_eq!(vb, 1);
}
}

View file

@ -33,16 +33,18 @@ CREATE INDEX IF NOT EXISTS idx_operator_questions_pending
ON operator_questions (id) WHERE answered_at IS NULL;
";
/// Additive column migrations for pre-existing databases, applied via
/// `db::apply_migrations` (try-and-ignore-duplicate-column).
/// Ordered schema migrations tracked in `schema_versions` (key
/// `"operator_questions"`). Legacy databases are detected via the `target`
/// column — the last column added before versioning — and fast-forwarded
/// past all known migrations.
const MIGRATIONS: &[&str] = &[
// v1: `multi` — checkbox-style multi-option questions.
"ALTER TABLE operator_questions ADD COLUMN multi INTEGER NOT NULL DEFAULT 0",
// v2: `deadline_at` — optional TTL after which the watchdog auto-resolves.
"ALTER TABLE operator_questions ADD COLUMN deadline_at INTEGER",
// `target` = recipient of the question. NULL = operator
// (back-compat default for rows written before agent-to-agent
// questions existed); a non-null agent name = peer-to-peer
// question. Dashboard's `pending()` filters on `target IS NULL`
// so peer questions never leak into the operator's queue.
// v3: `target` — recipient of the question. NULL = operator (back-compat
// default); non-null = peer-to-peer question. Dashboard's `pending()`
// filters on `target IS NULL` so peer questions never leak to the operator.
"ALTER TABLE operator_questions ADD COLUMN target TEXT",
];
@ -77,7 +79,12 @@ impl OperatorQuestions {
let conn = crate::db::open(path, "operator_questions")?;
conn.execute_batch(SCHEMA)
.context("apply operator_questions schema")?;
crate::db::apply_migrations(&conn, "operator_questions", MIGRATIONS)?;
crate::db::apply_versioned_migrations(
&conn,
"operator_questions",
("operator_questions", "target"),
MIGRATIONS,
)?;
Ok(Self {
conn: Mutex::new(conn),
})

View file

@ -186,11 +186,17 @@ impl ScheduledPrompts {
.context("enable foreign keys")?;
conn.execute_batch(SCHEMA)
.context("apply scheduled_prompts schema")?;
// Migration: add paused_at_unix to existing databases.
crate::db::apply_migrations(
// Versioned migration: add paused_at_unix. Legacy databases are
// detected via the presence of this column (it was the only
// migration before versioning was introduced).
crate::db::apply_versioned_migrations(
&conn,
"scheduled_prompts",
&["ALTER TABLE scheduled_prompts ADD COLUMN paused_at_unix INTEGER"],
("scheduled_prompts", "paused_at_unix"),
&[
// v1: add paused_at_unix for per-schedule pause support.
"ALTER TABLE scheduled_prompts ADD COLUMN paused_at_unix INTEGER",
],
)?;
// Migration: recreate the due-rows index to also exclude paused
// rows. `CREATE INDEX IF NOT EXISTS` won't update an existing