Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
dfbf198ef3 | ||
|
|
c47faf0af9 | ||
|
|
7d36ec5e1f |
5 changed files with 467 additions and 128 deletions
|
|
@ -10,6 +10,8 @@ use hive_sh4re::wire_time::now_unix;
|
|||
use hive_sh4re::{Approval, ApprovalKind, ApprovalStatus};
|
||||
use rusqlite::{Connection, OptionalExtension, params};
|
||||
|
||||
use crate::db::Migration;
|
||||
|
||||
const SCHEMA: &str = r"
|
||||
CREATE TABLE IF NOT EXISTS approvals (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
|
|
@ -24,25 +26,33 @@ 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).
|
||||
const MIGRATIONS: &[&str] = &[
|
||||
// `kind` (pre-Phase-8 dbs): legacy rows default to `apply_commit`,
|
||||
// which matches their actual semantics.
|
||||
"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).
|
||||
"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).
|
||||
"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.
|
||||
"ALTER TABLE approvals ADD COLUMN submitter TEXT",
|
||||
/// Ordered schema migrations tracked in `schema_versions` (key `"approvals"`).
|
||||
/// Each migration declares the column it adds, so a legacy DB (fully or
|
||||
/// partially migrated before versioning) converges by skipping the migrations
|
||||
/// whose column already exists. New columns go here as v5, v6, …
|
||||
const MIGRATIONS: &[Migration] = &[
|
||||
// v1: `kind` (pre-Phase-8 dbs): legacy rows default to `apply_commit`.
|
||||
Migration {
|
||||
sql: "ALTER TABLE approvals ADD COLUMN \
|
||||
kind TEXT NOT NULL DEFAULT 'apply_commit'",
|
||||
adds_column: Some(("approvals", "kind")),
|
||||
},
|
||||
// v2: `description`: manager-supplied note on the dashboard card.
|
||||
Migration {
|
||||
sql: "ALTER TABLE approvals ADD COLUMN description TEXT",
|
||||
adds_column: Some(("approvals", "description")),
|
||||
},
|
||||
// v3: `fetched_sha`: canonical sha hive-c0re resolved at submit time.
|
||||
Migration {
|
||||
sql: "ALTER TABLE approvals ADD COLUMN fetched_sha TEXT",
|
||||
adds_column: Some(("approvals", "fetched_sha")),
|
||||
},
|
||||
// v4: `submitter`: authenticated agent that submitted the approval.
|
||||
// Legacy rows are NULL → callers fall back to the root agent.
|
||||
Migration {
|
||||
sql: "ALTER TABLE approvals ADD COLUMN submitter TEXT",
|
||||
adds_column: Some(("approvals", "submitter")),
|
||||
},
|
||||
];
|
||||
|
||||
pub struct Approvals {
|
||||
|
|
@ -54,7 +64,7 @@ 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", MIGRATIONS)?;
|
||||
Ok(Self {
|
||||
conn: Mutex::new(conn),
|
||||
})
|
||||
|
|
|
|||
|
|
@ -10,6 +10,8 @@ use anyhow::{Context, Result};
|
|||
use chrono::{DateTime, Utc};
|
||||
use hive_sh4re::wire_time::now_unix;
|
||||
use hive_sh4re::{InboxRow, Message};
|
||||
|
||||
use crate::db::Migration;
|
||||
use rusqlite::{Connection, OptionalExtension, params};
|
||||
use serde::Serialize;
|
||||
use tokio::sync::broadcast;
|
||||
|
|
@ -156,12 +158,61 @@ pub struct Broker {
|
|||
inflight: Mutex<HashMap<String, RecipientInflight>>,
|
||||
}
|
||||
|
||||
/// Ordered schema migrations for the broker. Tracked in `schema_versions`
|
||||
/// under key `"broker"`. Each migration declares the column it adds, so a
|
||||
/// legacy DB (fully or partially migrated via the old per-column approach)
|
||||
/// is converged by skipping the migrations whose column already exists.
|
||||
const BROKER_MIGRATIONS: &[Migration] = &[
|
||||
// v1: acked_at on messages, with backfill so existing delivered rows
|
||||
// are not phantom-requeued on the next open. The BEGIN/COMMIT block
|
||||
// makes ALTER + UPDATE atomic — either both land or neither, so the
|
||||
// guard column (acked_at) faithfully marks the whole step as done.
|
||||
Migration {
|
||||
sql: "BEGIN;\
|
||||
ALTER TABLE messages ADD COLUMN acked_at INTEGER;\
|
||||
UPDATE messages SET acked_at = delivered_at \
|
||||
WHERE delivered_at IS NOT NULL AND acked_at IS NULL;\
|
||||
COMMIT;",
|
||||
adds_column: Some(("messages", "acked_at")),
|
||||
},
|
||||
// v2: in_reply_to for thread-parent tracking. NULL = root of a thread.
|
||||
Migration {
|
||||
sql: "ALTER TABLE messages ADD COLUMN in_reply_to INTEGER",
|
||||
adds_column: Some(("messages", "in_reply_to")),
|
||||
},
|
||||
// v3: priority for operator-message fast-path. Rebuild the delivery
|
||||
// index to include priority as a secondary sort key. Atomic BEGIN/COMMIT
|
||||
// so priority existing implies the index rebuild also committed.
|
||||
Migration {
|
||||
sql: "BEGIN;\
|
||||
ALTER TABLE messages ADD COLUMN \
|
||||
priority INTEGER NOT NULL DEFAULT 0;\
|
||||
DROP INDEX IF EXISTS idx_messages_undelivered;\
|
||||
CREATE INDEX IF NOT EXISTS idx_messages_undelivered \
|
||||
ON messages (recipient, priority DESC, id) WHERE delivered_at IS NULL;\
|
||||
COMMIT;",
|
||||
adds_column: Some(("messages", "priority")),
|
||||
},
|
||||
// v4: attempt_count on reminders for the MAX_REMINDER_ATTEMPTS cap.
|
||||
Migration {
|
||||
sql: "ALTER TABLE reminders ADD COLUMN \
|
||||
attempt_count INTEGER NOT NULL DEFAULT 0",
|
||||
adds_column: Some(("reminders", "attempt_count")),
|
||||
},
|
||||
// v5: last_error on reminders — last delivery failure surfaced on the
|
||||
// dashboard so a stuck reminder is visible without digging in logs.
|
||||
Migration {
|
||||
sql: "ALTER TABLE reminders ADD COLUMN last_error TEXT",
|
||||
adds_column: Some(("reminders", "last_error")),
|
||||
},
|
||||
];
|
||||
|
||||
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", BROKER_MIGRATIONS)
|
||||
.context("broker migrations")?;
|
||||
let (events, _) = broadcast::channel(EVENT_CHANNEL);
|
||||
Ok(Self {
|
||||
conn: Mutex::new(conn),
|
||||
|
|
@ -1104,70 +1155,6 @@ impl Broker {
|
|||
}
|
||||
}
|
||||
|
||||
/// Idempotent messages-table migrations. Adds `acked_at` and
|
||||
/// 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 {
|
||||
use super::*;
|
||||
|
|
|
|||
|
|
@ -1,47 +1,146 @@
|
|||
//! 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) need no special marker:
|
||||
//! each column-adding migration carries the `(table, column)` it
|
||||
//! introduces, so re-walking a pre-versioning DB from version 0 simply
|
||||
//! skips the migrations whose column already exists — converging a fully-
|
||||
//! or partially-migrated legacy DB without ever re-running an `ADD COLUMN`
|
||||
//! against an existing column. See each store's `MIGRATIONS` constant.
|
||||
|
||||
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
|
||||
)";
|
||||
|
||||
/// One numbered schema migration.
|
||||
///
|
||||
/// `sql` runs once, at the index whose value is the current stored version.
|
||||
/// `adds_column`, when set, is the `(table, column)` this migration
|
||||
/// introduces: if that column already exists — a pre-versioning DB migrated
|
||||
/// by the old try-and-ignore path, or a partially-migrated one — `sql` is
|
||||
/// skipped and only the version counter advances. This makes column-adding
|
||||
/// migrations idempotent *without* `SQLite`'s unsupported `ALTER TABLE … ADD
|
||||
/// COLUMN IF NOT EXISTS`, so re-walking a legacy DB from version 0 never
|
||||
/// re-runs an `ADD COLUMN` against a column that is already there (which
|
||||
/// would fail with "duplicate column name" and wedge startup permanently).
|
||||
///
|
||||
/// Multi-statement migrations wrap their `ALTER` + backfill / index rebuild
|
||||
/// in an explicit `BEGIN; … COMMIT;` batch, so "the added column exists" is a
|
||||
/// faithful proxy for "this whole migration committed".
|
||||
pub struct Migration {
|
||||
/// SQL for this version — a single statement or a `BEGIN; … COMMIT;` batch.
|
||||
pub sql: &'static str,
|
||||
/// `(table, column)` this migration adds, if any. Present → the migration
|
||||
/// is skipped when the column already exists (idempotent replay).
|
||||
pub adds_column: Option<(&'static str, &'static str)>,
|
||||
}
|
||||
|
||||
/// 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"`).
|
||||
/// - `migrations` — ordered list of [`Migration`]s; migration `i` runs when
|
||||
/// the stored version is `i`. A pre-versioning DB (no `schema_versions`
|
||||
/// row) starts at version 0 and re-walks every migration, but each
|
||||
/// column-adding migration is skipped when its `adds_column` is already
|
||||
/// present — so a fully- or partially-migrated legacy DB converges without
|
||||
/// re-running an `ADD COLUMN` against an existing column.
|
||||
///
|
||||
/// # Errors
|
||||
///
|
||||
/// Returns an error if any migration statement fails or the version update
|
||||
/// fails. Each migration + its version bump is a single `execute_batch`
|
||||
/// (multi-statement migrations wrap themselves in `BEGIN; … COMMIT;`), so a
|
||||
/// failure leaves the DB at the last successfully committed version.
|
||||
pub fn apply_versioned_migrations(
|
||||
conn: &Connection,
|
||||
subsystem: &str,
|
||||
migrations: &[Migration],
|
||||
) -> Result<()> {
|
||||
// Ensure the version-tracking table exists (idempotent).
|
||||
conn.execute_batch(SCHEMA_VERSIONS_DDL)
|
||||
.context("create schema_versions table")?;
|
||||
|
||||
// Current version for this subsystem — 0 (and a fresh row) when absent.
|
||||
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 {
|
||||
conn.execute(
|
||||
"INSERT INTO schema_versions (store, version) VALUES (?1, 0)",
|
||||
[subsystem],
|
||||
)
|
||||
.context("insert schema_versions row")?;
|
||||
0
|
||||
};
|
||||
|
||||
if version >= migrations.len() {
|
||||
return Ok(());
|
||||
}
|
||||
|
||||
for (i, mig) in migrations.iter().enumerate().skip(version) {
|
||||
// Skip a column-adding migration whose column is already present:
|
||||
// a legacy DB migrated by the old path, or a resumed partial one.
|
||||
let already_applied = match mig.adds_column {
|
||||
Some((table, column)) => column_exists(conn, table, column)?,
|
||||
None => false,
|
||||
};
|
||||
if !already_applied {
|
||||
conn.execute_batch(mig.sql)
|
||||
.with_context(|| format!("{subsystem} migration v{} failed: {}", i + 1, mig.sql))?;
|
||||
}
|
||||
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(())
|
||||
}
|
||||
|
||||
/// Whether `column` exists on `table` in this connection's schema. Both
|
||||
/// arguments are bound parameters — `pragma_table_info(?1)` accepts the table
|
||||
/// name as a bound value — so this is not a SQL-injection surface even if a
|
||||
/// caller ever passes a non-static name.
|
||||
fn column_exists(conn: &Connection, table: &str, column: &str) -> Result<bool> {
|
||||
conn.prepare("SELECT 1 FROM pragma_table_info(?1) WHERE name = ?2")
|
||||
.context("prepare column probe")?
|
||||
.exists(params![table, column])
|
||||
.context("check column existence")
|
||||
}
|
||||
|
||||
/// Open a connection to the sqlite file at `path`, creating the parent
|
||||
/// directory if needed. `subsystem` labels error contexts (`"broker"`,
|
||||
/// `"approvals"`, …).
|
||||
|
|
@ -56,3 +155,222 @@ 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: &[Migration] = &[
|
||||
Migration {
|
||||
sql: "ALTER TABLE things ADD COLUMN alpha TEXT",
|
||||
adds_column: Some(("things", "alpha")),
|
||||
},
|
||||
Migration {
|
||||
sql: "ALTER TABLE things ADD COLUMN beta INTEGER NOT NULL DEFAULT 0",
|
||||
adds_column: Some(("things", "beta")),
|
||||
},
|
||||
Migration {
|
||||
sql: "ALTER TABLE things ADD COLUMN gamma TEXT",
|
||||
adds_column: Some(("things", "gamma")),
|
||||
},
|
||||
];
|
||||
|
||||
#[test]
|
||||
fn fresh_install_runs_all_migrations() {
|
||||
let conn = open_tmp();
|
||||
conn.execute_batch(TEST_SCHEMA).unwrap();
|
||||
apply_versioned_migrations(&conn, "things", 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", 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", 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", TEST_MIGRATIONS).unwrap();
|
||||
// Second call is a no-op.
|
||||
apply_versioned_migrations(&conn, "things", 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",
|
||||
&[Migration {
|
||||
sql: "ALTER TABLE a ADD COLUMN col1 TEXT",
|
||||
adds_column: Some(("a", "col1")),
|
||||
}],
|
||||
)
|
||||
.unwrap();
|
||||
apply_versioned_migrations(
|
||||
&conn,
|
||||
"store_b",
|
||||
&[Migration {
|
||||
sql: "ALTER TABLE b ADD COLUMN colx INTEGER NOT NULL DEFAULT 0",
|
||||
adds_column: Some(("b", "colx")),
|
||||
}],
|
||||
)
|
||||
.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);
|
||||
}
|
||||
|
||||
/// A partially-migrated legacy database: some columns from the old
|
||||
/// try-and-ignore path exist, with no `schema_versions` row yet. The
|
||||
/// per-migration `adds_column` guard must skip the already-present
|
||||
/// columns and add only the missing ones — no "duplicate column name".
|
||||
#[test]
|
||||
fn partial_legacy_completes_without_error() {
|
||||
let conn = open_tmp();
|
||||
// Simulate a DB that has `alpha` (v1) but not `beta` (v2) or
|
||||
// `gamma` (v3). No schema_versions row exists → start from v0 and
|
||||
// guard-skip alpha while adding beta + gamma.
|
||||
conn.execute_batch(
|
||||
"CREATE TABLE IF NOT EXISTS things (id INTEGER PRIMARY KEY, alpha TEXT)",
|
||||
)
|
||||
.unwrap();
|
||||
apply_versioned_migrations(&conn, "things", TEST_MIGRATIONS).unwrap();
|
||||
// All three columns must now 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 partial-legacy migration"
|
||||
);
|
||||
}
|
||||
// Version must be at the latest.
|
||||
let v: i64 = conn
|
||||
.query_row(
|
||||
"SELECT version FROM schema_versions WHERE store = 'things'",
|
||||
[],
|
||||
|r| r.get(0),
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(
|
||||
v, 3,
|
||||
"must reach latest version after completing partial legacy"
|
||||
);
|
||||
}
|
||||
}
|
||||
|
|
|
|||
|
|
@ -19,6 +19,8 @@ use hive_sh4re::wire_time::now_unix;
|
|||
use rusqlite::{Connection, OptionalExtension, params};
|
||||
use serde::Serialize;
|
||||
|
||||
use crate::db::Migration;
|
||||
|
||||
const SCHEMA: &str = r"
|
||||
CREATE TABLE IF NOT EXISTS operator_questions (
|
||||
id INTEGER PRIMARY KEY AUTOINCREMENT,
|
||||
|
|
@ -33,17 +35,29 @@ 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).
|
||||
const MIGRATIONS: &[&str] = &[
|
||||
"ALTER TABLE operator_questions ADD COLUMN multi INTEGER NOT NULL DEFAULT 0",
|
||||
"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.
|
||||
"ALTER TABLE operator_questions ADD COLUMN target TEXT",
|
||||
/// 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: &[Migration] = &[
|
||||
// v1: `multi` — checkbox-style multi-option questions.
|
||||
Migration {
|
||||
sql: "ALTER TABLE operator_questions ADD COLUMN \
|
||||
multi INTEGER NOT NULL DEFAULT 0",
|
||||
adds_column: Some(("operator_questions", "multi")),
|
||||
},
|
||||
// v2: `deadline_at` — optional TTL after which the watchdog auto-resolves.
|
||||
Migration {
|
||||
sql: "ALTER TABLE operator_questions ADD COLUMN deadline_at INTEGER",
|
||||
adds_column: Some(("operator_questions", "deadline_at")),
|
||||
},
|
||||
// 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.
|
||||
Migration {
|
||||
sql: "ALTER TABLE operator_questions ADD COLUMN target TEXT",
|
||||
adds_column: Some(("operator_questions", "target")),
|
||||
},
|
||||
];
|
||||
|
||||
#[derive(Debug, Clone, Serialize)]
|
||||
|
|
@ -77,7 +91,7 @@ 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", MIGRATIONS)?;
|
||||
Ok(Self {
|
||||
conn: Mutex::new(conn),
|
||||
})
|
||||
|
|
|
|||
|
|
@ -20,6 +20,8 @@ use hive_sh4re::wire_time::now_unix;
|
|||
use rusqlite::{Connection, OptionalExtension, params};
|
||||
use serde::{Deserialize, Serialize};
|
||||
|
||||
use crate::db::Migration;
|
||||
|
||||
/// Typed error returned by [`ScheduledPrompts::pause`] and
|
||||
/// [`ScheduledPrompts::resume`] when the target row does not exist or
|
||||
/// is already cancelled. Handlers downcast on this type to emit 404
|
||||
|
|
@ -186,11 +188,19 @@ 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. The per-migration
|
||||
// `adds_column` guard skips the ALTER on a legacy DB that already
|
||||
// has the column (it was the only migration before versioning).
|
||||
crate::db::apply_versioned_migrations(
|
||||
&conn,
|
||||
"scheduled_prompts",
|
||||
&["ALTER TABLE scheduled_prompts ADD COLUMN paused_at_unix INTEGER"],
|
||||
&[
|
||||
// v1: add paused_at_unix for per-schedule pause support.
|
||||
Migration {
|
||||
sql: "ALTER TABLE scheduled_prompts ADD COLUMN paused_at_unix INTEGER",
|
||||
adds_column: Some(("scheduled_prompts", "paused_at_unix")),
|
||||
},
|
||||
],
|
||||
)?;
|
||||
// Migration: recreate the due-rows index to also exclude paused
|
||||
// rows. `CREATE INDEX IF NOT EXISTS` won't update an existing
|
||||
|
|
|
|||
Loading…
Reference in a new issue