Drop the run_events table and narrow the runs row

Every run is a Petri run whose history is `petri_records` and
`platform_records`, so the legacy run event log has no reader left.
The new migration drops `run_events` (with its indexes), the two
one-time activation tables, and rebuilds `runs` without the columns
only that log wrote or read: `source_last_seq` and the six token and
file-count columns nothing read, as VIEWS.md records. Pre-cutover
development runs are discarded, as decided; the surviving columns of
existing rows are copied across.

fabro-db loses the three migration consts of the dropped schema and the
session-owner preflight that inspected `run_events`, and gains
`DROP_RUN_EVENTS_MIGRATION_SQL` so fixtures that install the runs
schema reach the production shape. The run summary upsert binds only
the surviving columns.

Tests: the `run_events` schema, query-plan and preflight tests are
deleted; `runs_schema_has_its_final_shape_without_the_legacy_event_log`
pins the final columns and indexes, and
`dropping_the_event_log_keeps_the_run_rows` migrates a database left by
an older binary and checks the run row survives.

Co-Authored-By: Claude Fable 5.1 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 14:37:18 -04:00
parent 0d74fdf01d
commit 4942788297
No known key found for this signature in database
6 changed files with 199 additions and 450 deletions

View file

@ -124,6 +124,7 @@ fn pool() -> DbPool {
test_support::in_memory_pool_with(&[
fabro_db::BLOBS_MIGRATION_SQL,
fabro_db::RUNS_MIGRATION_SQL,
fabro_db::DROP_RUN_EVENTS_MIGRATION_SQL,
fabro_db::PETRI_RECORDS_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
])

View file

@ -17,17 +17,15 @@ use crate::run_summary::{build_summary, projected_usage};
use crate::{Error, Result, RunProjection};
/// The `runs` row of a Petri run, written by its projector: every column the
/// list views and the scheduler read, and never `source_last_seq`, which the
/// legacy event path owns while it still writes the row.
/// list views and the scheduler read.
const UPSERT_PETRI_RUN_SQL: &str = r"
INSERT INTO runs (
id, source_last_seq, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms,
id, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms,
status, archived_at_ms, parent_id, title, workflow_slug, workflow_name,
repository_name, automation_id, diff_files_changed, diff_additions, diff_deletions,
input_tokens, output_tokens, reasoning_tokens, cache_read_tokens, cache_write_tokens,
repository_name, automation_id, diff_additions, diff_deletions,
total_usd_micros, summary_json
) VALUES (
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?
)
ON CONFLICT(id) DO UPDATE SET
created_at_ms = excluded.created_at_ms,
@ -403,15 +401,10 @@ pub struct RunSummaryIdentity {
#[derive(Debug)]
struct PreparedRunSummary {
run: Run,
workflow_name: Option<String>,
repository_name: Option<String>,
input_tokens: i64,
output_tokens: i64,
reasoning_tokens: i64,
cache_read_tokens: i64,
cache_write_tokens: i64,
total_usd_micros: Option<i64>,
run: Run,
workflow_name: Option<String>,
repository_name: Option<String>,
total_usd_micros: Option<i64>,
}
impl PreparedRunSummary {
@ -435,11 +428,6 @@ impl PreparedRunSummary {
run,
workflow_name,
repository_name,
input_tokens: column_count(usage.tokens.input),
output_tokens: column_count(usage.tokens.output),
reasoning_tokens: column_count(usage.tokens.reasoning),
cache_read_tokens: column_count(usage.tokens.cache_read),
cache_write_tokens: column_count(usage.tokens.cache_write),
total_usd_micros: usage.cost.map(|cost| column_count(cost.usd_micros)),
}
}
@ -451,9 +439,8 @@ fn column_count(count: u64) -> i64 {
i64::try_from(count).unwrap_or(i64::MAX)
}
/// Binds the `runs` columns shared by the insert, upsert, and update
/// statements, in the positional order those statements declare them
/// (`source_last_seq` through `summary_json`).
/// Binds the `runs` columns after `id`, in the positional order the upsert
/// declares them (`created_at_ms` through `summary_json`).
fn bind_run_columns<'q>(
query: Query<'q, Sqlite, SqliteArguments>,
record: &'q PreparedRunSummary,
@ -462,9 +449,6 @@ fn bind_run_columns<'q>(
let diff = run.diff.unwrap_or_default();
let summary_json = serde_json::to_string(run)?;
Ok(query
// `source_last_seq`: a column the legacy event log owned; `1` until
// the migration that drops it.
.bind(1_i64)
.bind(run.timestamps.created_at.timestamp_millis())
.bind(
run.timestamps
@ -494,14 +478,8 @@ fn bind_run_columns<'q>(
.bind(&record.workflow_name)
.bind(&record.repository_name)
.bind(run.automation.as_ref().map(|automation| &automation.id))
.bind(diff.files_changed)
.bind(diff.additions)
.bind(diff.deletions)
.bind(record.input_tokens)
.bind(record.output_tokens)
.bind(record.reasoning_tokens)
.bind(record.cache_read_tokens)
.bind(record.cache_write_tokens)
.bind(record.total_usd_micros)
.bind(summary_json))
}
@ -1121,8 +1099,8 @@ mod tests {
let row = sqlx::query(
"SELECT created_at_ms, last_event_at_ms, status, title, workflow_slug, \
automation_id, input_tokens, reasoning_tokens, cache_read_tokens, total_usd_micros, \
diff_files_changed, diff_additions, diff_deletions FROM runs WHERE id = ?",
automation_id, total_usd_micros, diff_additions, diff_deletions \
FROM runs WHERE id = ?",
)
.bind(run_id.to_string())
.fetch_one(&store.pool)
@ -1146,14 +1124,10 @@ mod tests {
sqlx::Row::get::<String, _>(&row, "automation_id"),
"nightly"
);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "input_tokens"), 100);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "reasoning_tokens"), 5);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "cache_read_tokens"), 10);
assert_eq!(
sqlx::Row::get::<i64, _>(&row, "total_usd_micros"),
21_000_000
);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "diff_files_changed"), 2);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "diff_additions"), 10);
assert_eq!(sqlx::Row::get::<i64, _>(&row, "diff_deletions"), 3);

View file

@ -25,6 +25,7 @@ pub fn test_blob_store() -> Arc<BlobStore> {
/// platform records and the projection tables.
const RUN_SUMMARY_MIGRATIONS: &[&str] = &[
fabro_db::RUNS_MIGRATION_SQL,
fabro_db::DROP_RUN_EVENTS_MIGRATION_SQL,
fabro_db::PETRI_PROJECTION_MIGRATION_SQL,
];

View file

@ -0,0 +1,69 @@
-- The legacy executor's run event log is gone: every run is a Petri run
-- whose history is `petri_records` and `platform_records`. Fabro is
-- greenfield here, so the rows are dropped, not converted, and the
-- one-time activation bookkeeping of that log goes with them.
DROP TABLE IF EXISTS run_events;
DROP TABLE IF EXISTS legacy_run_history_activation;
DROP TABLE IF EXISTS legacy_run_history_deletions;
-- The `runs` row loses the columns only that log wrote or read:
-- `source_last_seq`, its write-path concurrency guard (the projection's
-- per-log positions replace it), and the six token and file-count columns
-- nothing read (`summary_json` keeps the values). SQLite cannot drop a
-- column a table CHECK names, so the table is rebuilt.
CREATE TABLE runs_next (
id TEXT PRIMARY KEY NOT NULL,
created_at_ms INTEGER NOT NULL,
started_at_ms INTEGER,
last_event_at_ms INTEGER NOT NULL,
completed_at_ms INTEGER,
status TEXT NOT NULL,
archived_at_ms INTEGER,
parent_id TEXT,
title TEXT NOT NULL,
workflow_slug TEXT,
workflow_name TEXT,
repository_name TEXT,
automation_id TEXT,
diff_additions INTEGER NOT NULL DEFAULT 0,
diff_deletions INTEGER NOT NULL DEFAULT 0,
total_usd_micros INTEGER,
summary_json TEXT NOT NULL,
CHECK (status IN (
'submitted',
'pending',
'runnable',
'starting',
'running',
'blocked',
'paused',
'removing',
'succeeded',
'failed',
'dead'
)),
CHECK (diff_additions >= 0),
CHECK (diff_deletions >= 0),
CHECK (total_usd_micros IS NULL OR total_usd_micros >= 0),
CHECK (json_valid(summary_json))
);
INSERT INTO runs_next (
id, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms, status,
archived_at_ms, parent_id, title, workflow_slug, workflow_name, repository_name,
automation_id, diff_additions, diff_deletions, total_usd_micros, summary_json
)
SELECT
id, created_at_ms, started_at_ms, last_event_at_ms, completed_at_ms, status,
archived_at_ms, parent_id, title, workflow_slug, workflow_name, repository_name,
automation_id, diff_additions, diff_deletions, total_usd_micros, summary_json
FROM runs;
DROP TABLE runs;
ALTER TABLE runs_next RENAME TO runs;
CREATE INDEX runs_by_created_at ON runs(created_at_ms DESC, id DESC);
CREATE INDEX runs_by_updated_at ON runs(last_event_at_ms DESC, id DESC);
CREATE INDEX runs_by_status ON runs(archived_at_ms, status, last_event_at_ms DESC, id DESC);
CREATE INDEX runs_by_parent ON runs(parent_id, created_at_ms DESC, id DESC);
CREATE INDEX runs_by_automation ON runs(automation_id, created_at_ms DESC, id DESC);

View file

@ -16,8 +16,6 @@ pub type DbPool = sqlx::SqlitePool;
static MIGRATOR: Migrator = sqlx::migrate!("./migrations");
const SESSION_OWNER_INDEX_MIGRATION_VERSION: i64 = 2_026_083_101;
/// The blob-table migration, exposed so fixtures in other crates can install
/// the production blob schema without a filesystem path into this crate.
pub const BLOBS_MIGRATION_SQL: &str = include_str!("../migrations/2026081301_blobs.sql");
@ -26,14 +24,11 @@ pub const BLOBS_MIGRATION_SQL: &str = include_str!("../migrations/2026081301_blo
/// the production schema without a filesystem path into this crate.
pub const RUNS_MIGRATION_SQL: &str = include_str!("../migrations/2026071104_runs.sql");
/// The run-event migration, exposed so fixtures in other crates can install
/// the production schema without a filesystem path into this crate.
pub const RUN_EVENTS_MIGRATION_SQL: &str = include_str!("../migrations/2026082701_run_events.sql");
/// The run-session owner index migration, exposed so fixtures in other crates
/// can install the production run-history indexes.
pub const RUN_EVENT_SESSION_OWNER_MIGRATION_SQL: &str =
include_str!("../migrations/2026083101_run_event_session_owner.sql");
/// The migration that drops the legacy run event log and narrows the `runs`
/// row to its final shape, exposed so fixtures that install
/// [`RUNS_MIGRATION_SQL`] can apply it next and get the production schema.
pub const DROP_RUN_EVENTS_MIGRATION_SQL: &str =
include_str!("../migrations/2026091803_drop_run_events.sql");
/// The Ask Fabro session record migration, exposed so fixtures in other
/// crates can install the production schema without a filesystem path into
@ -59,11 +54,6 @@ pub const PETRI_RECORDS_MIGRATION_SQL: &str =
pub const PETRI_PROJECTION_MIGRATION_SQL: &str =
include_str!("../migrations/2026091801_petri_projection.sql");
/// The temporary run-history activation migration, exposed so fixtures in
/// other crates can install the production compatibility schema.
pub const RUN_HISTORY_ACTIVATION_MIGRATION_SQL: &str =
include_str!("../migrations/2026082802_run_history_activation.sql");
#[derive(Clone)]
pub struct Database {
pool: DbPool,
@ -98,9 +88,6 @@ impl Database {
pub async fn migrate(&self) -> anyhow::Result<()> {
let applied = applied_migration_versions(&self.pool).await?;
self.preflight_session_owner_index(&applied)
.await
.context("checking session ownership before SQLite migrations")?;
self.snapshot_before_new_migrations(&applied)
.await
.context("snapshotting SQLite database before migrations")?;
@ -110,53 +97,6 @@ impl Database {
.context("running SQLite migrations")
}
/// Refuse the unique owner index when old event history contains
/// collisions. The diagnostic is deliberately count-only because session
/// identifiers and event contents are not safe startup-log fields.
///
/// Temporary compatibility guard: once every supported database has
/// applied the session-owner index migration the version check below
/// always short-circuits, and this preflight can be deleted along with
/// the run-history compatibility window.
async fn preflight_session_owner_index(&self, applied: &HashSet<i64>) -> anyhow::Result<()> {
if applied.contains(&SESSION_OWNER_INDEX_MIGRATION_VERSION) {
return Ok(());
}
let run_events_exists: bool = sqlx::query_scalar(
"SELECT EXISTS(SELECT 1 FROM sqlite_master WHERE type = 'table' AND name = 'run_events')",
)
.fetch_one(&self.pool)
.await
.context("checking for the run event table")?;
if !run_events_exists {
return Ok(());
}
let collision_groups: i64 = sqlx::query_scalar(
r"
SELECT COUNT(*)
FROM (
SELECT session_id
FROM run_events
WHERE session_id IS NOT NULL
AND event_name = 'run.session.created'
GROUP BY session_id
HAVING COUNT(*) > 1
)
",
)
.fetch_one(&self.pool)
.await
.context("counting duplicate session ownership groups")?;
if collision_groups > 0 {
anyhow::bail!(
"cannot create the unique session owner index: found {collision_groups} duplicate session ownership groups"
);
}
Ok(())
}
/// Copy the database aside before applying migrations it has not seen.
///
/// A binary downgrade after new migrations have been applied fails sqlx's

View file

@ -769,13 +769,13 @@ async fn runs_schema_creates_indexes_and_rejects_invalid_rows() -> anyhow::Resul
assert_eq!(index_count, 5);
insert_minimal_run(database.pool(), "submitted", 0, r#"{"id":"run"}"#).await?;
for (status, input_tokens, summary_json) in [
for (status, diff_additions, summary_json) in [
("unknown", 0, r#"{"id":"run-2"}"#),
("submitted", -1, r#"{"id":"run-3"}"#),
("submitted", 0, "not-json"),
] {
assert!(
insert_minimal_run(database.pool(), status, input_tokens, summary_json)
insert_minimal_run(database.pool(), status, diff_additions, summary_json)
.await
.is_err()
);
@ -784,384 +784,148 @@ async fn runs_schema_creates_indexes_and_rejects_invalid_rows() -> anyhow::Resul
Ok(())
}
/// The `runs` row is the projection's summary of a Petri run: the columns
/// the list views filter and sort by, and the JSON the API serves. The
/// legacy event log's tables and the columns only it wrote are gone.
#[tokio::test]
async fn session_owner_schema_has_final_shape_constraints_and_indexes() -> anyhow::Result<()> {
async fn runs_schema_has_its_final_shape_without_the_legacy_event_log() -> anyhow::Result<()> {
let dir = tempfile::tempdir()?;
let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?;
database.migrate().await?;
for table in [
"run_events",
"legacy_run_history_activation",
"legacy_run_history_deletions",
"runs_next",
] {
assert!(
!table_exists(database.pool(), table).await?,
"{table} must not exist"
);
}
let run_columns = sqlx::query("PRAGMA table_info(runs)")
.fetch_all(database.pool())
.await?;
.await?
.iter()
.map(|column| column.get::<String, _>("name"))
.collect::<Vec<_>>();
assert_eq!(run_columns, [
"id",
"created_at_ms",
"started_at_ms",
"last_event_at_ms",
"completed_at_ms",
"status",
"archived_at_ms",
"parent_id",
"title",
"workflow_slug",
"workflow_name",
"repository_name",
"automation_id",
"diff_additions",
"diff_deletions",
"total_usd_micros",
"summary_json",
]);
let index_names = sqlx::query("PRAGMA index_list(runs)")
.fetch_all(database.pool())
.await?
.iter()
.map(|index| index.get::<String, _>("name"))
.filter(|name| name.starts_with("runs_by_"))
.collect::<std::collections::BTreeSet<_>>();
assert_eq!(
run_columns.len(),
24,
"the existing runs row must stay unchanged"
index_names.iter().map(String::as_str).collect::<Vec<_>>(),
[
"runs_by_automation",
"runs_by_created_at",
"runs_by_parent",
"runs_by_status",
"runs_by_updated_at",
]
);
let event_columns = sqlx::query("PRAGMA table_info(run_events)")
.fetch_all(database.pool())
Ok(())
}
/// A database written before the event log was dropped keeps its run rows:
/// the migration rebuilds the table and copies every surviving column.
#[tokio::test]
async fn dropping_the_event_log_keeps_the_run_rows() -> anyhow::Result<()> {
let dir = tempfile::tempdir()?;
let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?;
database.migrate().await?;
// Rewind to the schema an older binary left: the `runs` table with its
// legacy columns, the event log referencing it and the activation
// bookkeeping, with only the drop migration pending again. Those
// migrations' own rows stay applied, so sqlx's checksum validation
// still passes.
sqlx::raw_sql("DROP TABLE runs; DELETE FROM _sqlx_migrations WHERE version = 2026091803;")
.execute(database.pool())
.await?;
let event_column_contract = event_columns
.iter()
.map(|column| {
(
column.get::<String, _>("name"),
column.get::<String, _>("type"),
column.get::<i64, _>("notnull"),
column.get::<i64, _>("pk"),
)
})
.collect::<Vec<_>>();
assert_eq!(event_column_contract, vec![
("run_id".to_string(), "TEXT".to_string(), 1, 1),
("seq".to_string(), "INTEGER".to_string(), 1, 2),
("event_name".to_string(), "TEXT".to_string(), 1, 0),
("node_id".to_string(), "TEXT".to_string(), 0, 0),
("stage_id".to_string(), "TEXT".to_string(), 0, 0),
("session_id".to_string(), "TEXT".to_string(), 0, 0),
("event_json".to_string(), "TEXT".to_string(), 1, 0),
]);
let foreign_keys = sqlx::query("PRAGMA foreign_key_list(run_events)")
.fetch_all(database.pool())
.await?;
assert_eq!(foreign_keys.len(), 1);
assert_eq!(foreign_keys[0].get::<String, _>("table"), "runs");
assert_eq!(foreign_keys[0].get::<String, _>("from"), "run_id");
assert_eq!(foreign_keys[0].get::<String, _>("to"), "id");
assert_eq!(foreign_keys[0].get::<String, _>("on_delete"), "CASCADE");
let indexes = sqlx::query("PRAGMA index_list(run_events)")
.fetch_all(database.pool())
.await?;
let named_indexes = indexes
.iter()
.filter_map(|index| {
let name = index.get::<String, _>("name");
name.starts_with("run_events_by_").then_some((
name,
index.get::<i64, _>("unique"),
index.get::<i64, _>("partial"),
))
})
.collect::<Vec<_>>();
assert_eq!(named_indexes, vec![
("run_events_by_session_owner".to_string(), 1, 1),
(
"run_events_by_pull_request_creation_request".to_string(),
0,
1,
),
("run_events_by_session".to_string(), 0, 1),
("run_events_by_legacy_node".to_string(), 0, 1),
("run_events_by_stage".to_string(), 0, 1),
]);
assert!(indexes.iter().all(|index| {
index.get::<i64, _>("unique") == 0
|| index.get::<String, _>("name") == "run_events_by_session_owner"
|| index.get::<String, _>("name") == "sqlite_autoindex_run_events_1"
}));
insert_run_with_id(database.pool(), "parent", None).await?;
insert_run_with_id(database.pool(), "child", Some("parent")).await?;
insert_run_event(database.pool(), "parent", 1, "run.created").await?;
for invalid in [
insert_run_event(database.pool(), "parent", 1, "run.created").await,
insert_run_event(database.pool(), "missing", 1, "run.created").await,
insert_run_event(database.pool(), "parent", 0, "run.created").await,
insert_run_event(database.pool(), "parent", 1_000_000, "run.created").await,
for migration in [
fabro_db::RUNS_MIGRATION_SQL,
include_str!("../migrations/2026082701_run_events.sql"),
include_str!("../migrations/2026082802_run_history_activation.sql"),
include_str!("../migrations/2026083101_run_event_session_owner.sql"),
] {
assert!(invalid.is_err());
sqlx::raw_sql(migration).execute(database.pool()).await?;
}
let invalid_json = sqlx::query(
"INSERT INTO run_events (run_id, seq, event_name, event_json) VALUES (?, ?, ?, ?)",
sqlx::query(
r#"
INSERT INTO runs (
id, source_last_seq, created_at_ms, last_event_at_ms, status, title, input_tokens,
diff_additions, total_usd_micros, summary_json
) VALUES ('kept', 7, 1, 2, 'succeeded', 'Kept run', 99, 3, 4, '{"id":"kept"}')
"#,
)
.bind("parent")
.bind(2_i64)
.bind("run.started")
.bind("not-json")
.execute(database.pool())
.await;
assert!(invalid_json.is_err());
sqlx::query("INSERT INTO blobs (hash, data) VALUES (?, ?)")
.bind("a".repeat(64))
.bind(vec![1_u8])
.execute(database.pool())
.await?;
sqlx::query("DELETE FROM runs WHERE id = ?")
.bind("parent")
.execute(database.pool())
.await?;
let event_count: i64 =
sqlx::query_scalar("SELECT COUNT(*) FROM run_events WHERE run_id = 'parent'")
.fetch_one(database.pool())
.await?;
let child_parent: Option<String> =
sqlx::query_scalar("SELECT parent_id FROM runs WHERE id = 'child'")
.fetch_one(database.pool())
.await?;
let blob_count: i64 = sqlx::query_scalar("SELECT COUNT(*) FROM blobs")
.fetch_one(database.pool())
.await?;
assert_eq!(event_count, 0);
assert_eq!(child_parent.as_deref(), Some("parent"));
assert_eq!(blob_count, 1);
Ok(())
}
#[tokio::test]
async fn run_events_schema_query_plans_use_candidate_indexes_including_session_owner()
-> anyhow::Result<()> {
let dir = tempfile::tempdir()?;
let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?;
database.migrate().await?;
for (sql, expected_index) in [
(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND seq > ? ORDER BY seq ASC LIMIT ?",
"sqlite_autoindex_run_events_1",
),
(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND seq = ?",
"sqlite_autoindex_run_events_1",
),
(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND stage_id = ? ORDER BY seq ASC LIMIT ?",
"run_events_by_stage",
),
(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND stage_id IS NULL AND node_id = ? ORDER BY seq ASC LIMIT ?",
"run_events_by_legacy_node",
),
(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND session_id = ? AND event_name GLOB 'run.session.*' ORDER BY seq ASC LIMIT ?",
"run_events_by_session",
),
(
"EXPLAIN QUERY PLAN SELECT run_id, seq, event_name, node_id, stage_id, session_id, event_json FROM run_events WHERE session_id = ? AND event_name = 'run.session.created'",
"run_events_by_session_owner",
),
(
"EXPLAIN QUERY PLAN SELECT DISTINCT run_id FROM run_events WHERE event_name = 'pull_request.creation_requested'",
"run_events_by_pull_request_creation_request",
),
] {
let details = sqlx::query(sql)
.bind("run")
.bind("value")
.bind(10_i64)
.fetch_all(database.pool())
.await?
.into_iter()
.map(|row| row.get::<String, _>("detail"))
.collect::<Vec<_>>()
.join("; ");
assert!(
details.contains(expected_index),
"expected {expected_index} in query plan: {details}"
);
}
// The first-visit stage listing unions both shapes so each arm keeps its
// own partial index instead of scanning the run's primary key range.
let details = sqlx::query(
"EXPLAIN QUERY PLAN SELECT * FROM run_events WHERE run_id = ? AND seq >= ? AND stage_id = ? \
UNION ALL SELECT * FROM run_events WHERE run_id = ? AND seq >= ? AND stage_id IS NULL AND node_id = ? \
ORDER BY seq ASC LIMIT ?",
)
.bind("run")
.bind(1_i64)
.bind("stage")
.bind("run")
.bind(1_i64)
.bind("node")
.bind(10_i64)
.fetch_all(database.pool())
.await?
.into_iter()
.map(|row| row.get::<String, _>("detail"))
.collect::<Vec<_>>()
.join("; ");
for expected_index in ["run_events_by_stage", "run_events_by_legacy_node"] {
assert!(
details.contains(expected_index),
"expected {expected_index} in query plan: {details}"
);
}
Ok(())
}
#[tokio::test]
async fn session_owner_migration_preflight_is_count_only_retriable_and_idempotent()
-> anyhow::Result<()> {
let dir = tempfile::tempdir()?;
let database = fabro_db::Database::connect(dir.path().join("fabro.sqlite3")).await?;
database.migrate().await?;
sqlx::query("DROP INDEX IF EXISTS run_events_by_session_owner")
.execute(database.pool())
.await?;
sqlx::query("DELETE FROM _sqlx_migrations WHERE version = 2026083101")
.execute(database.pool())
.await?;
for run_id in ["first", "second", "third", "fourth"] {
insert_run_with_id(database.pool(), run_id, None).await?;
}
for (run_id, session_id) in [
("first", "collision-alpha"),
("second", "collision-alpha"),
("third", "collision-beta"),
("fourth", "collision-beta"),
] {
insert_session_creation_claim(database.pool(), run_id, session_id).await?;
}
let error = database
.migrate()
.await
.expect_err("duplicate session owners must abort migration");
let rendered = format!("{error:#}");
assert!(rendered.contains("2 duplicate session ownership groups"));
assert!(!rendered.contains("collision-alpha"));
assert!(!rendered.contains("collision-beta"));
assert!(!rendered.contains("sensitive event contents"));
assert_eq!(
sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM run_events WHERE event_name = 'run.session.created'"
)
.fetch_one(database.pool())
.await?,
4
);
assert_eq!(
sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = 'run_events_by_session_owner'"
)
.fetch_one(database.pool())
.await?,
0
);
sqlx::query("DELETE FROM run_events WHERE run_id IN ('second', 'fourth')")
.execute(database.pool())
.await?;
database.migrate().await?;
database.migrate().await?;
assert_eq!(
sqlx::query_scalar::<_, i64>(
"SELECT COUNT(*) FROM sqlite_master WHERE type = 'index' AND name = 'run_events_by_session_owner'"
)
.fetch_one(database.pool())
.await?,
1
);
Ok(())
}
async fn insert_session_creation_claim(
pool: &fabro_db::DbPool,
run_id: &str,
session_id: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r"
INSERT INTO run_events (run_id, seq, event_name, session_id, event_json)
VALUES (?, 1, 'run.session.created', ?, json_object(
'run_id', ?,
'event', 'run.session.created',
'session_id', ?,
'properties', json_object('note', 'sensitive event contents')
))
",
)
.bind(run_id)
.bind(session_id)
.bind(run_id)
.bind(session_id)
.execute(pool)
.await?;
Ok(())
}
async fn insert_run_event(
pool: &fabro_db::DbPool,
run_id: &str,
seq: i64,
event_name: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r"
INSERT INTO run_events (run_id, seq, event_name, event_json)
VALUES (?, ?, ?, '{}')
",
"INSERT INTO run_events (run_id, seq, event_name, event_json) VALUES ('kept', 1, 'run.created', '{}')",
)
.bind(run_id)
.bind(seq)
.bind(event_name)
.execute(pool)
.execute(database.pool())
.await?;
database.migrate().await?;
assert!(!table_exists(database.pool(), "run_events").await?);
let row = sqlx::query(
"SELECT created_at_ms, last_event_at_ms, status, title, diff_additions, total_usd_micros, summary_json FROM runs WHERE id = 'kept'",
)
.fetch_one(database.pool())
.await?;
assert_eq!(row.get::<i64, _>("created_at_ms"), 1);
assert_eq!(row.get::<i64, _>("last_event_at_ms"), 2);
assert_eq!(row.get::<String, _>("status"), "succeeded");
assert_eq!(row.get::<String, _>("title"), "Kept run");
assert_eq!(row.get::<i64, _>("diff_additions"), 3);
assert_eq!(row.get::<i64, _>("total_usd_micros"), 4);
assert_eq!(row.get::<String, _>("summary_json"), r#"{"id":"kept"}"#);
Ok(())
}
async fn insert_minimal_run(
pool: &fabro_db::DbPool,
status: &str,
input_tokens: i64,
summary_json: &str,
) -> Result<(), sqlx::Error> {
insert_run_row(
pool,
&format!("run-{status}-{input_tokens}"),
None,
status,
input_tokens,
summary_json,
)
.await
}
async fn insert_run_with_id(
pool: &fabro_db::DbPool,
id: &str,
parent_id: Option<&str>,
) -> Result<(), sqlx::Error> {
insert_run_row(
pool,
id,
parent_id,
"submitted",
0,
&format!(r#"{{"id":"{id}"}}"#),
)
.await
}
async fn insert_run_row(
pool: &fabro_db::DbPool,
id: &str,
parent_id: Option<&str>,
status: &str,
input_tokens: i64,
diff_additions: i64,
summary_json: &str,
) -> Result<(), sqlx::Error> {
sqlx::query(
r"
INSERT INTO runs (
id, source_last_seq, created_at_ms, last_event_at_ms, status, parent_id, title,
input_tokens, summary_json
) VALUES (?, 1, 0, 0, ?, ?, 'title', ?, ?)
id, created_at_ms, last_event_at_ms, status, title, diff_additions, summary_json
) VALUES (?, 0, 0, ?, 'title', ?, ?)
",
)
.bind(id)
.bind(format!("run-{status}-{diff_additions}"))
.bind(status)
.bind(parent_id)
.bind(input_tokens)
.bind(diff_additions)
.bind(summary_json)
.execute(pool)
.await?;