use anyhow::Context;
use futures_util::future::BoxFuture;
use std::collections::HashSet;
use crate::db::{Connection, params};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum TargetDb {
Core,
Logs,
}
type RustFn = for<'a> fn(&'a Connection) -> BoxFuture<'a, anyhow::Result<()>>;
#[derive(Debug, Clone, Copy)]
pub(crate) enum MigrationBody {
Sql(&'static str),
Rust(RustFn),
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct Migration {
pub(crate) id: &'static str,
pub(crate) target: TargetDb,
pub(crate) body: MigrationBody,
}
const BASELINE_CORE_SCHEMA: &str = "\
-- ── Board (tickets, comments, counters) ────────────────────────────────
CREATE TABLE IF NOT EXISTS tickets (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
description TEXT NOT NULL,
phase TEXT NOT NULL DEFAULT 'backlog',
workspace_name TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
prerequisites TEXT NOT NULL DEFAULT '[]',
supersedes TEXT,
superseded_by TEXT,
commit_hash TEXT,
lines_added INTEGER,
lines_removed INTEGER,
reporter TEXT NOT NULL DEFAULT '',
is_archived INTEGER NOT NULL DEFAULT 0,
embedding BLOB,
priority INTEGER NOT NULL DEFAULT 1,
reviewed_head TEXT,
reviewed_tree TEXT,
done_at TEXT,
bounce_count INTEGER NOT NULL DEFAULT 0,
last_transition_actor TEXT
);
CREATE TABLE IF NOT EXISTS ticket_comments (
id TEXT PRIMARY KEY,
ticket_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL,
FOREIGN KEY (ticket_id) REFERENCES tickets(id)
);
CREATE TABLE IF NOT EXISTS ticket_counters (
workspace_name TEXT PRIMARY KEY,
next_id INTEGER NOT NULL DEFAULT 1
);
-- ── Sessions / jobs / agents ───────────────────────────────────────────
CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS session_metadata (
agent_id TEXT PRIMARY KEY,
last_activity TEXT NOT NULL,
channel TEXT,
user_name TEXT,
workspace_name TEXT,
role TEXT,
active_models TEXT,
token_length INTEGER,
message_count INTEGER NOT NULL DEFAULT 0,
created_at TEXT,
sleep_ended TEXT
);
CREATE TABLE IF NOT EXISTS jobs (
id TEXT PRIMARY KEY,
kind TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'launched',
task TEXT NOT NULL DEFAULT '',
workspace_name TEXT NOT NULL,
user_name TEXT NOT NULL DEFAULT '',
channel TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL,
retry_count INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
ticket_id TEXT REFERENCES tickets(id),
caller_agent_id TEXT,
mode TEXT
);
CREATE TABLE IF NOT EXISTS agents (
job_id TEXT REFERENCES jobs(id) ON DELETE CASCADE,
agent_id TEXT NOT NULL,
kind TEXT NOT NULL,
idx INTEGER,
status TEXT NOT NULL DEFAULT 'launched',
outcome TEXT,
task TEXT NOT NULL,
PRIMARY KEY (job_id, agent_id)
);
CREATE TABLE IF NOT EXISTS pending_jobs (
id TEXT PRIMARY KEY,
target_agent_id TEXT NOT NULL,
envelope TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS research_jobs (
id TEXT PRIMARY KEY REFERENCES jobs(id) ON DELETE CASCADE,
state TEXT NOT NULL
);
-- ── Workspaces / editor ───────────────────────────────────────────────
CREATE TABLE IF NOT EXISTS workspaces (
name TEXT PRIMARY KEY,
path TEXT NOT NULL UNIQUE,
status TEXT NOT NULL DEFAULT 'pending',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
maintenance INTEGER NOT NULL DEFAULT 0,
paused INTEGER NOT NULL DEFAULT 1,
maintainer_debounce_mins INTEGER NOT NULL DEFAULT 5,
maintainer_last_run_at TEXT,
maintainer_recommendations TEXT,
diagnostics TEXT,
diagnostics_generation INTEGER NOT NULL DEFAULT 0,
notes TEXT NOT NULL DEFAULT '',
last_analyzed_commit TEXT,
discovery_generation INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS workspace_contexts (
workspace_name TEXT NOT NULL REFERENCES workspaces(name) ON DELETE CASCADE,
role TEXT,
content TEXT NOT NULL,
created_at TEXT NOT NULL,
UNIQUE(workspace_name, role)
);
CREATE TABLE IF NOT EXISTS editor_tabs (
workspace_name TEXT NOT NULL REFERENCES workspaces(name) ON DELETE CASCADE,
file_path TEXT NOT NULL,
tab_order INTEGER NOT NULL DEFAULT 0,
is_active INTEGER NOT NULL DEFAULT 0,
is_dirty INTEGER NOT NULL DEFAULT 0,
dirty_content TEXT,
PRIMARY KEY (workspace_name, file_path)
);
-- ── Users / channels ───────────────────────────────────────────────────
CREATE TABLE IF NOT EXISTS users (
name TEXT PRIMARY KEY,
selected_workspace TEXT,
image_gen_model TEXT,
video_model TEXT,
granted_tools TEXT
);
CREATE TABLE IF NOT EXISTS user_channels (
user_name TEXT NOT NULL REFERENCES users(name),
channel TEXT NOT NULL,
identifier TEXT NOT NULL,
reply_target TEXT,
UNIQUE(channel, identifier)
);
-- ── Config ─────────────────────────────────────────────────────────────
CREATE TABLE IF NOT EXISTS config_kv (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS config_model_routing (
model TEXT PRIMARY KEY,
provider_order TEXT,
allow_fallbacks INTEGER
);
-- ── Chat history ───────────────────────────────────────────────────────
CREATE TABLE IF NOT EXISTS chat_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
message_id TEXT NOT NULL UNIQUE,
user_name TEXT NOT NULL,
direction TEXT NOT NULL,
content TEXT NOT NULL,
agent_role TEXT,
broadcast_id TEXT,
workspace TEXT NOT NULL,
timestamp TEXT,
reply_author TEXT,
reply_snippet TEXT
);
-- ── Ticket chronicle / alarms ──────────────────────────────────────────
CREATE TABLE IF NOT EXISTS ticket_chronicle (
id INTEGER PRIMARY KEY AUTOINCREMENT,
ticket_id TEXT NOT NULL,
workspace_name TEXT NOT NULL,
source_phase TEXT NOT NULL,
target_phase TEXT NOT NULL,
at TEXT NOT NULL,
actor TEXT
);
CREATE TABLE IF NOT EXISTS alarms (
id TEXT PRIMARY KEY,
session_id TEXT NOT NULL,
user_name TEXT NOT NULL,
kind TEXT NOT NULL,
text TEXT NOT NULL,
fire_at TEXT,
interval_seconds INTEGER,
next_fire_at TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TEXT NOT NULL,
trigger TEXT
);
-- ── Indexes ────────────────────────────────────────────────────────────
CREATE INDEX IF NOT EXISTS idx_ticket_comments_ticket_id ON ticket_comments(ticket_id);
CREATE INDEX IF NOT EXISTS idx_sessions_agent_id ON sessions(agent_id, id);
CREATE INDEX IF NOT EXISTS idx_jobs_kind_status ON jobs(kind, status);
CREATE INDEX IF NOT EXISTS idx_jobs_updated_at ON jobs(updated_at);
CREATE UNIQUE INDEX IF NOT EXISTS idx_agents_anchor ON agents(agent_id) WHERE job_id IS NULL;
CREATE INDEX IF NOT EXISTS idx_pending_jobs_agent_created ON pending_jobs(target_agent_id, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS workspace_contexts_null_role ON workspace_contexts(workspace_name) WHERE role IS NULL;
CREATE INDEX IF NOT EXISTS idx_chat_history_user ON chat_history(user_name);
CREATE INDEX IF NOT EXISTS idx_chat_history_workspace ON chat_history(workspace);
CREATE INDEX IF NOT EXISTS idx_chat_history_user_ws_id ON chat_history(user_name, workspace, id);
CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets USING fts (title) WITH (tokenizer = 'ngram');
CREATE INDEX IF NOT EXISTS idx_tickets_board_active ON tickets (is_archived, priority ASC, created_at DESC);
CREATE UNIQUE INDEX IF NOT EXISTS idx_jobs_phase_ticket ON jobs(kind, ticket_id) WHERE ticket_id IS NOT NULL;
CREATE INDEX IF NOT EXISTS idx_ticket_chronicle_ws_id ON ticket_chronicle(workspace_name, id);
CREATE UNIQUE INDEX IF NOT EXISTS idx_ticket_chronicle_dedup ON ticket_chronicle(ticket_id, workspace_name, source_phase, target_phase, at);
CREATE INDEX IF NOT EXISTS idx_alarms_due ON alarms(status, next_fire_at);
CREATE INDEX IF NOT EXISTS idx_tickets_workspace_phase ON tickets (workspace_name, phase, is_archived, priority ASC, created_at DESC);";
const BASELINE_LOGS_SCHEMA: &str = "\
-- ── Logs / tool calls / LLM requests ───────────────────────────────────
CREATE TABLE IF NOT EXISTS logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp TEXT NOT NULL,
level TEXT NOT NULL,
target TEXT NOT NULL,
message TEXT NOT NULL,
fields TEXT NOT NULL DEFAULT '{}',
agent_id TEXT NOT NULL DEFAULT '',
agent_role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT ''
);
CREATE TABLE IF NOT EXISTS tool_calls (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
tool_name TEXT NOT NULL,
arguments TEXT NOT NULL DEFAULT '{}',
duration_ms INTEGER NOT NULL DEFAULT 0,
success INTEGER NOT NULL DEFAULT 1,
error_message TEXT,
workspace TEXT NOT NULL DEFAULT '',
recorded_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS llm_requests (
id INTEGER PRIMARY KEY AUTOINCREMENT,
recorded_at TEXT NOT NULL,
purpose TEXT NOT NULL,
agent_id TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
ticket_id TEXT,
model TEXT NOT NULL,
routing TEXT NOT NULL DEFAULT '',
input_tokens INTEGER,
output_tokens INTEGER,
cached_input_tokens INTEGER,
cache_miss_tokens INTEGER,
duration_ms INTEGER NOT NULL,
retry_attempts INTEGER NOT NULL,
finish_reason TEXT,
failure_class TEXT,
success INTEGER NOT NULL DEFAULT 1,
cost REAL,
cost_details TEXT,
upstream_provider TEXT,
system_fingerprint TEXT
);
-- ── Grep telemetry ─────────────────────────────────────────────────────
CREATE TABLE IF NOT EXISTS grep_telemetry (
id INTEGER PRIMARY KEY AUTOINCREMENT,
recorded_at TEXT NOT NULL,
command TEXT NOT NULL,
served INTEGER NOT NULL DEFAULT 0,
reason TEXT NOT NULL DEFAULT '',
recursive INTEGER NOT NULL DEFAULT 0,
piped INTEGER NOT NULL DEFAULT 0,
operand_count INTEGER NOT NULL DEFAULT 0,
flags TEXT NOT NULL DEFAULT '',
mode TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
grep_count INTEGER NOT NULL DEFAULT 0,
served_count INTEGER NOT NULL DEFAULT 0,
skipped_count INTEGER NOT NULL DEFAULT 0,
duration_ms INTEGER,
exit_code INTEGER
);
-- ── Indexes ────────────────────────────────────────────────────────────
CREATE INDEX IF NOT EXISTS idx_logs_timestamp ON logs(timestamp);
CREATE INDEX IF NOT EXISTS idx_logs_level ON logs(level);
CREATE INDEX IF NOT EXISTS idx_logs_target ON logs(target);
CREATE INDEX IF NOT EXISTS idx_logs_agent_role ON logs(agent_role);
CREATE INDEX IF NOT EXISTS idx_logs_agent_id ON logs(agent_id);
CREATE INDEX IF NOT EXISTS idx_logs_workspace ON logs(workspace);
CREATE INDEX IF NOT EXISTS idx_tool_calls_agent_id ON tool_calls(agent_id);
CREATE INDEX IF NOT EXISTS idx_tool_calls_role ON tool_calls(role);
CREATE INDEX IF NOT EXISTS idx_tool_calls_tool_name ON tool_calls(tool_name);
CREATE INDEX IF NOT EXISTS idx_tool_calls_recorded_at ON tool_calls(recorded_at);
CREATE INDEX IF NOT EXISTS idx_tool_calls_workspace ON tool_calls(workspace);
CREATE INDEX IF NOT EXISTS idx_tool_calls_error_message ON tool_calls(error_message);
CREATE INDEX IF NOT EXISTS idx_llm_requests_recorded_at ON llm_requests(recorded_at);
CREATE INDEX IF NOT EXISTS idx_llm_requests_agent_id ON llm_requests(agent_id);
CREATE INDEX IF NOT EXISTS idx_llm_requests_model ON llm_requests(model);
CREATE INDEX IF NOT EXISTS idx_llm_requests_purpose ON llm_requests(purpose);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_recorded_at ON grep_telemetry(recorded_at);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_served ON grep_telemetry(served);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_reason ON grep_telemetry(reason);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_command ON grep_telemetry(command);
-- ── Per-attempt LLM failure trail ──────────────────────────────────────
CREATE TABLE IF NOT EXISTS llm_failures (
id INTEGER PRIMARY KEY AUTOINCREMENT,
operation_id TEXT NOT NULL,
recorded_at TEXT NOT NULL,
purpose TEXT NOT NULL,
agent_id TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
ticket_id TEXT,
model TEXT NOT NULL,
routing TEXT NOT NULL DEFAULT '',
attempt INTEGER NOT NULL,
attempt_duration_ms INTEGER,
failure_class TEXT NOT NULL,
finish_reason TEXT,
retry_after_ms INTEGER,
error_chain TEXT NOT NULL DEFAULT ''
);
CREATE INDEX IF NOT EXISTS idx_llm_failures_recorded_at ON llm_failures(recorded_at);
CREATE INDEX IF NOT EXISTS idx_llm_failures_operation_id ON llm_failures(operation_id);
CREATE INDEX IF NOT EXISTS idx_llm_failures_failure_class ON llm_failures(failure_class);";
const ADD_LLM_FAILURES: &str = "CREATE TABLE IF NOT EXISTS llm_failures (
id INTEGER PRIMARY KEY AUTOINCREMENT,
operation_id TEXT NOT NULL,
recorded_at TEXT NOT NULL,
purpose TEXT NOT NULL,
agent_id TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
ticket_id TEXT,
model TEXT NOT NULL,
routing TEXT NOT NULL DEFAULT '',
attempt INTEGER NOT NULL,
attempt_duration_ms INTEGER,
failure_class TEXT NOT NULL,
finish_reason TEXT,
retry_after_ms INTEGER,
error_chain TEXT NOT NULL DEFAULT ''
);
CREATE INDEX IF NOT EXISTS idx_llm_failures_recorded_at ON llm_failures(recorded_at);
CREATE INDEX IF NOT EXISTS idx_llm_failures_operation_id ON llm_failures(operation_id);
CREATE INDEX IF NOT EXISTS idx_llm_failures_failure_class ON llm_failures(failure_class);";
const REWRITE_LEGACY_DEEPSEEK_ROUTING_SLUG: &str = "UPDATE config_model_routing SET provider_order = 'deepseek' WHERE provider_order = 'DeepSeek';";
const REWRITE_RETIRED_STAGES_TO_VERIFICATION: &str = "\
DELETE FROM jobs WHERE kind IN ('in_review', 'in_qa');\
DELETE FROM ticket_chronicle WHERE source_phase IN ('in_review', 'in_qa') \
AND target_phase IN ('in_review', 'in_qa');\
UPDATE tickets SET phase = 'verification' WHERE phase IN ('in_review', 'in_qa');\
UPDATE ticket_chronicle SET source_phase = 'verification' WHERE source_phase IN ('in_review', 'in_qa');\
UPDATE ticket_chronicle SET target_phase = 'verification' WHERE target_phase IN ('in_review', 'in_qa');";
pub(crate) const MIGRATIONS: &[Migration] = &[
Migration {
id: "24",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_CORE_SCHEMA),
},
Migration {
id: "25",
target: TargetDb::Core,
body: MigrationBody::Rust(add_chat_history_reply_columns),
},
Migration {
id: "26",
target: TargetDb::Logs,
body: MigrationBody::Sql(BASELINE_LOGS_SCHEMA),
},
Migration {
id: "27",
target: TargetDb::Core,
body: MigrationBody::Rust(add_workspaces_maintainer_recommendations),
},
Migration {
id: "28",
target: TargetDb::Core,
body: MigrationBody::Rust(add_jobs_caller_agent_and_session_created_at),
},
Migration {
id: "29",
target: TargetDb::Core,
body: MigrationBody::Rust(add_users_image_and_video_models),
},
Migration {
id: "30",
target: TargetDb::Core,
body: MigrationBody::Rust(add_jobs_mode),
},
Migration {
id: "31",
target: TargetDb::Core,
body: MigrationBody::Rust(add_session_metadata_sleep_ended),
},
Migration {
id: "32",
target: TargetDb::Core,
body: MigrationBody::Rust(add_alarms_command),
},
Migration {
id: "33",
target: TargetDb::Core,
body: MigrationBody::Rust(add_chat_history_broadcast_id),
},
Migration {
id: "34",
target: TargetDb::Core,
body: MigrationBody::Rust(add_ticket_transition_actor),
},
Migration {
id: "35",
target: TargetDb::Logs,
body: MigrationBody::Sql(ADD_LLM_FAILURES),
},
Migration {
id: "36",
target: TargetDb::Core,
body: MigrationBody::Rust(detach_guest_workspaces),
},
Migration {
id: "37",
target: TargetDb::Core,
body: MigrationBody::Sql(REWRITE_LEGACY_DEEPSEEK_ROUTING_SLUG),
},
Migration {
id: "38",
target: TargetDb::Core,
body: MigrationBody::Rust(add_users_granted_tools),
},
Migration {
id: "39",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_users_permissions),
},
Migration {
id: "40",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_users_selected_role),
},
Migration {
id: "41",
target: TargetDb::Core,
body: MigrationBody::Rust(remove_command_alarms),
},
Migration {
id: "42",
target: TargetDb::Core,
body: MigrationBody::Rust(add_alarms_trigger),
},
Migration {
id: "43",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_alarms_command),
},
Migration {
id: "44",
target: TargetDb::Core,
body: MigrationBody::Sql(REWRITE_RETIRED_STAGES_TO_VERIFICATION),
},
Migration {
id: "45",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_retired_lifecycle_descriptions),
},
];
pub(crate) async fn run_migrations(conn: &Connection, db: TargetDb) -> anyhow::Result<()> {
run_catalog(conn, db, MIGRATIONS).await
}
pub(crate) const MIGRATIONS_TABLE: &str = "schema_migrations";
async fn run_catalog(conn: &Connection, db: TargetDb, catalog: &[Migration]) -> anyhow::Result<()> {
conn.execute(
&format!(
"CREATE TABLE IF NOT EXISTS {MIGRATIONS_TABLE} (\
id TEXT PRIMARY KEY,\
applied_at TEXT NOT NULL\
)"
),
(),
)
.await
.context("Failed to create schema_migrations tracking table")?;
let applied: HashSet<String> = conn
.query(&format!("SELECT id FROM {MIGRATIONS_TABLE}"), ())
.await
.context("Failed to read applied migrations")?
.into_iter()
.filter_map(|row| row.get::<String>(0).ok())
.collect();
for migration in catalog {
if migration.target != db {
continue;
}
if applied.contains(migration.id) {
continue;
}
match migration.body {
MigrationBody::Sql(sql) => {
let tx = conn.begin_tx().await.with_context(|| {
format!("Migration '{}': failed to begin transaction", migration.id)
})?;
tx.execute_batch(sql)
.await
.with_context(|| format!("Migration '{}' failed", migration.id))?;
tx.execute(
&format!("INSERT INTO {MIGRATIONS_TABLE} (id, applied_at) VALUES (?1, ?2)"),
params![migration.id, crate::db::now()],
)
.await
.with_context(|| {
format!("Migration '{}': failed to record as applied", migration.id)
})?;
tx.commit()
.await
.with_context(|| format!("Migration '{}': failed to commit", migration.id))?;
}
MigrationBody::Rust(run) => {
run(conn)
.await
.with_context(|| format!("Migration '{}' failed", migration.id))?;
conn.execute(
&format!("INSERT INTO {MIGRATIONS_TABLE} (id, applied_at) VALUES (?1, ?2)"),
params![migration.id, crate::db::now()],
)
.await
.with_context(|| {
format!("Migration '{}': failed to record as applied", migration.id)
})?;
}
}
}
Ok(())
}
fn add_chat_history_reply_columns(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_chat_history_reply_columns(conn))
}
async fn run_add_chat_history_reply_columns(conn: &Connection) -> anyhow::Result<()> {
for column in ["reply_author", "reply_snippet"] {
add_column_if_missing(conn, "chat_history", column).await?;
}
Ok(())
}
fn add_chat_history_broadcast_id(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_chat_history_broadcast_id(conn))
}
async fn run_add_chat_history_broadcast_id(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "chat_history", "broadcast_id").await
}
fn add_ticket_transition_actor(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_ticket_transition_actor(conn))
}
async fn run_add_ticket_transition_actor(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "tickets", "last_transition_actor").await?;
add_column_if_missing(conn, "ticket_chronicle", "actor").await
}
fn add_workspaces_maintainer_recommendations(
conn: &Connection,
) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_workspaces_maintainer_recommendations(conn))
}
async fn run_add_workspaces_maintainer_recommendations(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "workspaces", "maintainer_recommendations").await
}
async fn column_exists(conn: &Connection, table: &str, column: &str) -> anyhow::Result<bool> {
let found = conn
.query(&format!("PRAGMA table_info({table})"), ())
.await
.with_context(|| format!("Failed to probe {table}.{column}"))?
.into_iter()
.any(|row| row.get::<String>(1).ok().as_deref() == Some(column));
Ok(found)
}
async fn add_column_if_missing(conn: &Connection, table: &str, column: &str) -> anyhow::Result<()> {
if !column_exists(conn, table, column).await? {
conn.execute(&format!("ALTER TABLE {table} ADD COLUMN {column} TEXT"), ())
.await
.with_context(|| format!("Failed to add {table}.{column}"))?;
}
Ok(())
}
async fn drop_column_if_missing(
conn: &Connection,
table: &str,
column: &str,
) -> anyhow::Result<()> {
if column_exists(conn, table, column).await? {
conn.execute(&format!("ALTER TABLE {table} DROP COLUMN {column}"), ())
.await
.with_context(|| format!("Failed to drop {table}.{column}"))?;
}
Ok(())
}
fn add_jobs_caller_agent_and_session_created_at(
conn: &Connection,
) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_jobs_caller_agent_and_session_created_at(conn))
}
async fn run_add_jobs_caller_agent_and_session_created_at(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "jobs", "caller_agent_id").await?;
add_column_if_missing(conn, "session_metadata", "created_at").await?;
conn.execute_batch(
"CREATE INDEX IF NOT EXISTS idx_jobs_caller_agent \
ON jobs(caller_agent_id) WHERE caller_agent_id IS NOT NULL;",
)
.await
.with_context(|| "Failed to create idx_jobs_caller_agent")?;
Ok(())
}
fn add_users_image_and_video_models(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_users_image_and_video_models(conn))
}
async fn run_add_users_image_and_video_models(conn: &Connection) -> anyhow::Result<()> {
for column in ["image_gen_model", "video_model"] {
add_column_if_missing(conn, "users", column).await?;
}
Ok(())
}
fn add_jobs_mode(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_jobs_mode(conn))
}
async fn run_add_jobs_mode(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "jobs", "mode").await?;
conn.execute(
"UPDATE jobs \
SET mode = CASE WHEN caller_agent_id IS NOT NULL THEN 'sync' ELSE 'async' END \
WHERE mode IS NULL",
(),
)
.await
.with_context(|| "Failed to backfill jobs.mode")?;
Ok(())
}
fn add_session_metadata_sleep_ended(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_session_metadata_sleep_ended(conn))
}
async fn run_add_session_metadata_sleep_ended(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "session_metadata", "sleep_ended").await
}
fn add_alarms_command(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_alarms_command(conn))
}
async fn run_add_alarms_command(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "alarms", "command").await
}
fn detach_guest_workspaces(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_detach_guest_workspaces(conn))
}
async fn run_detach_guest_workspaces(conn: &Connection) -> anyhow::Result<()> {
conn.execute(
"UPDATE users SET selected_workspace = NULL \
WHERE selected_workspace IS NOT NULL AND name <> ?1",
params![crate::users::ADMIN_USER_NAME],
)
.await
.with_context(|| "Failed to detach guests from shared workspaces")?;
Ok(())
}
fn add_users_granted_tools(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_users_granted_tools(conn))
}
async fn run_add_users_granted_tools(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "users", "granted_tools").await
}
fn drop_users_permissions(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_users_permissions(conn))
}
async fn run_drop_users_permissions(conn: &Connection) -> anyhow::Result<()> {
drop_column_if_missing(conn, "users", "permissions").await
}
fn drop_users_selected_role(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_users_selected_role(conn))
}
async fn run_drop_users_selected_role(conn: &Connection) -> anyhow::Result<()> {
drop_column_if_missing(conn, "users", "selected_role").await
}
fn remove_command_alarms(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_remove_command_alarms(conn))
}
async fn run_remove_command_alarms(conn: &Connection) -> anyhow::Result<()> {
if !column_exists(conn, "alarms", "command").await? {
return Ok(());
}
conn.execute("DELETE FROM alarms WHERE command IS NOT NULL", ())
.await
.with_context(|| "Failed to delete command-armed alarms")?;
Ok(())
}
fn add_alarms_trigger(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_add_alarms_trigger(conn))
}
async fn run_add_alarms_trigger(conn: &Connection) -> anyhow::Result<()> {
add_column_if_missing(conn, "alarms", "trigger").await
}
fn drop_alarms_command(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_alarms_command(conn))
}
async fn run_drop_alarms_command(conn: &Connection) -> anyhow::Result<()> {
drop_column_if_missing(conn, "alarms", "command").await
}
const RETIRED_LIFECYCLE_MARKERS: &[&str] = &[
"in_review",
"in review",
"inreview",
"in_qa",
"in qa",
"inqa",
"→ review",
"review →",
"→ qa",
"qa →",
"review/qa",
"qa/review",
];
fn describes_retired_lifecycle(content: &str) -> bool {
let mut folded = content.to_ascii_lowercase();
for separator in [" and ", ", ", " then "] {
folded = folded.replace(separator, "/");
}
folded = folded
.replace("reviewers", "role_name")
.replace("reviewer", "role_name");
RETIRED_LIFECYCLE_MARKERS
.iter()
.any(|marker| folded.contains(marker))
}
fn drop_retired_lifecycle_descriptions(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_retired_lifecycle_descriptions(conn))
}
async fn run_drop_retired_lifecycle_descriptions(conn: &Connection) -> anyhow::Result<()> {
let rows = conn
.query("SELECT rowid, content FROM workspace_contexts", ())
.await
.context("Failed to read stored workspace contexts")?;
for row in rows {
let content = row
.get::<String>(1)
.context("Failed to read workspace_contexts.content")?;
if !describes_retired_lifecycle(&content) {
continue;
}
let rowid = row
.get::<i64>(0)
.context("Failed to read workspace_contexts.rowid")?;
conn.execute(
"DELETE FROM workspace_contexts WHERE rowid = ?1",
params![rowid],
)
.await
.context("Failed to drop a stale stored workspace context")?;
}
Ok(())
}
#[cfg(test)]
mod tests {
use super::*;
use crate::db::Connection;
async fn column_names(conn: &Connection, table: &str) -> Vec<String> {
conn.query(&format!("PRAGMA table_info({table})"), ())
.await
.expect("read table_info")
.into_iter()
.map(|row| row.get::<String>(1).expect("column name"))
.collect()
}
async fn applied_ids(conn: &Connection) -> Vec<String> {
conn.query("SELECT id FROM schema_migrations", ())
.await
.expect("read schema_migrations")
.into_iter()
.map(|row| row.get::<String>(0).expect("id"))
.collect()
}
fn normalize_ddl(sql: &str) -> String {
let mut out = String::new();
let mut in_string = false;
let mut chars = sql.chars().peekable();
while let Some(c) = chars.next() {
if in_string {
if c == '\'' {
in_string = false;
}
out.push(c);
} else if c == '\'' {
in_string = true;
out.push(c);
} else if c == '-' && chars.peek() == Some(&'-') {
for c2 in chars.by_ref() {
if c2 == '\n' {
break;
}
}
} else if !c.is_whitespace() {
out.push(c);
}
}
out
}
async fn index_defs(conn: &Connection) -> std::collections::BTreeMap<String, String> {
conn.query(
"SELECT name, sql FROM sqlite_master \
WHERE type = 'index' AND name NOT LIKE 'sqlite_autoindex_%' \
AND name NOT LIKE '__turso_internal_%' AND name NOT LIKE 'turso_cdc%'",
(),
)
.await
.expect("read index definitions")
.into_iter()
.map(|row| {
let name = row.get::<String>(0).expect("index name");
let sql = row
.get::<Option<String>>(1)
.expect("index sql")
.unwrap_or_default();
(name, normalize_ddl(&sql))
})
.collect()
}
async fn table_defs(conn: &Connection) -> std::collections::BTreeMap<String, String> {
conn.query(
"SELECT name, sql FROM sqlite_master \
WHERE type = 'table' AND name NOT LIKE 'sqlite_%' \
AND name NOT LIKE '__turso_internal_%' AND name NOT LIKE 'turso_cdc%'",
(),
)
.await
.expect("read table definitions")
.into_iter()
.map(|row| {
let name = row.get::<String>(0).expect("table name");
let sql = row
.get::<Option<String>>(1)
.expect("table sql")
.unwrap_or_default();
(name, normalize_ddl(&sql))
})
.collect()
}
async fn table_names(conn: &Connection) -> Vec<String> {
conn.query(
"SELECT name FROM sqlite_master \
WHERE type = 'table' AND name NOT LIKE 'sqlite_%' \
AND name NOT LIKE '__turso_internal_%' AND name NOT LIKE 'turso_cdc%'",
(),
)
.await
.expect("read tables")
.into_iter()
.map(|row| row.get::<String>(0).expect("table name"))
.collect()
}
const EXPECTED_CORE_TABLE_COLUMNS: &[(&str, &[&str])] = &[
(
"agents",
&[
"job_id", "agent_id", "kind", "idx", "status", "outcome", "task",
],
),
(
"alarms",
&[
"id",
"session_id",
"user_name",
"kind",
"text",
"fire_at",
"interval_seconds",
"next_fire_at",
"status",
"created_at",
"trigger",
],
),
(
"chat_history",
&[
"id",
"message_id",
"user_name",
"direction",
"content",
"agent_role",
"broadcast_id",
"workspace",
"timestamp",
"reply_author",
"reply_snippet",
],
),
("config_kv", &["key", "value"]),
(
"config_model_routing",
&["model", "provider_order", "allow_fallbacks"],
),
(
"editor_tabs",
&[
"workspace_name",
"file_path",
"tab_order",
"is_active",
"is_dirty",
"dirty_content",
],
),
(
"jobs",
&[
"id",
"kind",
"status",
"task",
"workspace_name",
"user_name",
"channel",
"role",
"retry_count",
"created_at",
"updated_at",
"ticket_id",
"caller_agent_id",
"mode",
],
),
(
"pending_jobs",
&["id", "target_agent_id", "envelope", "created_at"],
),
("research_jobs", &["id", "state"]),
("schema_migrations", &["id", "applied_at"]),
(
"session_metadata",
&[
"agent_id",
"last_activity",
"channel",
"user_name",
"workspace_name",
"role",
"active_models",
"token_length",
"message_count",
"created_at",
"sleep_ended",
],
),
(
"sessions",
&["id", "agent_id", "role", "content", "created_at"],
),
(
"ticket_chronicle",
&[
"id",
"ticket_id",
"workspace_name",
"source_phase",
"target_phase",
"at",
"actor",
],
),
(
"ticket_comments",
&["id", "ticket_id", "role", "content", "created_at"],
),
("ticket_counters", &["workspace_name", "next_id"]),
(
"tickets",
&[
"id",
"title",
"description",
"phase",
"workspace_name",
"created_at",
"updated_at",
"prerequisites",
"supersedes",
"superseded_by",
"commit_hash",
"lines_added",
"lines_removed",
"reporter",
"is_archived",
"embedding",
"priority",
"reviewed_head",
"reviewed_tree",
"done_at",
"bounce_count",
"last_transition_actor",
],
),
(
"user_channels",
&["user_name", "channel", "identifier", "reply_target"],
),
(
"users",
&[
"name",
"selected_workspace",
"image_gen_model",
"video_model",
"granted_tools",
],
),
(
"workspace_contexts",
&["workspace_name", "role", "content", "created_at"],
),
(
"workspaces",
&[
"name",
"path",
"status",
"created_at",
"updated_at",
"maintenance",
"paused",
"maintainer_debounce_mins",
"maintainer_last_run_at",
"maintainer_recommendations",
"diagnostics",
"diagnostics_generation",
"notes",
"last_analyzed_commit",
"discovery_generation",
],
),
];
const EXPECTED_CORE_INDEXES: &[&str] = &[
"idx_ticket_comments_ticket_id",
"idx_sessions_agent_id",
"idx_jobs_kind_status",
"idx_jobs_updated_at",
"idx_agents_anchor",
"idx_pending_jobs_agent_created",
"workspace_contexts_null_role",
"idx_chat_history_user",
"idx_chat_history_workspace",
"idx_chat_history_user_ws_id",
"idx_tickets_title_fts",
"idx_tickets_board_active",
"idx_jobs_phase_ticket",
"idx_jobs_caller_agent",
"idx_ticket_chronicle_ws_id",
"idx_ticket_chronicle_dedup",
"idx_alarms_due",
"idx_tickets_workspace_phase",
];
const EXPECTED_LOGS_TABLE_COLUMNS: &[(&str, &[&str])] = &[
(
"grep_telemetry",
&[
"id",
"recorded_at",
"command",
"served",
"reason",
"recursive",
"piped",
"operand_count",
"flags",
"mode",
"workspace",
"grep_count",
"served_count",
"skipped_count",
"duration_ms",
"exit_code",
],
),
(
"llm_failures",
&[
"id",
"operation_id",
"recorded_at",
"purpose",
"agent_id",
"role",
"workspace",
"ticket_id",
"model",
"routing",
"attempt",
"attempt_duration_ms",
"failure_class",
"finish_reason",
"retry_after_ms",
"error_chain",
],
),
(
"llm_requests",
&[
"id",
"recorded_at",
"purpose",
"agent_id",
"role",
"workspace",
"ticket_id",
"model",
"routing",
"input_tokens",
"output_tokens",
"cached_input_tokens",
"cache_miss_tokens",
"duration_ms",
"retry_attempts",
"finish_reason",
"failure_class",
"success",
"cost",
"cost_details",
"upstream_provider",
"system_fingerprint",
],
),
(
"logs",
&[
"id",
"timestamp",
"level",
"target",
"message",
"fields",
"agent_id",
"agent_role",
"workspace",
],
),
("schema_migrations", &["id", "applied_at"]),
(
"tool_calls",
&[
"id",
"agent_id",
"role",
"tool_name",
"arguments",
"duration_ms",
"success",
"error_message",
"workspace",
"recorded_at",
],
),
];
const EXPECTED_LOGS_INDEXES: &[&str] = &[
"idx_logs_timestamp",
"idx_logs_level",
"idx_logs_target",
"idx_logs_agent_role",
"idx_logs_agent_id",
"idx_logs_workspace",
"idx_tool_calls_agent_id",
"idx_tool_calls_role",
"idx_tool_calls_tool_name",
"idx_tool_calls_recorded_at",
"idx_tool_calls_workspace",
"idx_tool_calls_error_message",
"idx_llm_requests_recorded_at",
"idx_llm_requests_agent_id",
"idx_llm_requests_model",
"idx_llm_requests_purpose",
"idx_llm_failures_recorded_at",
"idx_llm_failures_operation_id",
"idx_llm_failures_failure_class",
"idx_grep_telemetry_recorded_at",
"idx_grep_telemetry_served",
"idx_grep_telemetry_reason",
"idx_grep_telemetry_command",
];
fn expected_core_table_names() -> Vec<&'static str> {
EXPECTED_CORE_TABLE_COLUMNS
.iter()
.map(|(n, _)| *n)
.collect()
}
fn expected_core_domain_tables() -> Vec<&'static str> {
expected_core_table_names()
.into_iter()
.filter(|n| *n != "schema_migrations")
.collect()
}
fn expected_logs_domain_tables() -> Vec<&'static str> {
EXPECTED_LOGS_TABLE_COLUMNS
.iter()
.map(|(n, _)| *n)
.filter(|n| *n != "schema_migrations")
.collect()
}
async fn assert_schema_matches(
conn: &Connection,
expected_tables: &[(&str, &[&str])],
expected_indexes: &[&str],
context: &str,
) {
let actual_names: Vec<String> = table_names(conn).await;
let actual_set: std::collections::HashSet<&str> =
actual_names.iter().map(String::as_str).collect();
let expected_set: std::collections::HashSet<&str> =
expected_tables.iter().map(|(n, _)| *n).collect();
let missing: Vec<&str> = expected_set.difference(&actual_set).copied().collect();
let extra: Vec<&str> = actual_set.difference(&expected_set).copied().collect();
assert!(missing.is_empty(), "{context}: missing tables {missing:?}");
assert!(extra.is_empty(), "{context}: unexpected tables {extra:?}");
for (table, expected_cols) in expected_tables {
let actual = column_names(conn, table).await;
let actual_set: std::collections::HashSet<&str> =
actual.iter().map(String::as_str).collect();
let expected_set: std::collections::HashSet<&str> =
expected_cols.iter().copied().collect();
let miss: Vec<&str> = expected_set.difference(&actual_set).copied().collect();
let ex: Vec<&str> = actual_set.difference(&expected_set).copied().collect();
assert!(
miss.is_empty(),
"{context}: table '{table}' missing columns {miss:?}"
);
assert!(
ex.is_empty(),
"{context}: table '{table}' has extra columns {ex:?}"
);
}
let idx_names: Vec<String> = index_defs(conn).await.into_keys().collect();
let actual_idx: std::collections::HashSet<&str> =
idx_names.iter().map(String::as_str).collect();
let expected_idx: std::collections::HashSet<&str> =
expected_indexes.iter().copied().collect();
let missing_idx: Vec<&str> = expected_idx.difference(&actual_idx).copied().collect();
let extra_idx: Vec<&str> = actual_idx.difference(&expected_idx).copied().collect();
assert!(
missing_idx.is_empty(),
"{context}: missing indexes {missing_idx:?}"
);
assert!(
extra_idx.is_empty(),
"{context}: unexpected indexes {extra_idx:?}"
);
}
async fn column_sets(
conn: &Connection,
tables: &[&str],
) -> std::collections::BTreeMap<String, Vec<String>> {
let mut m = std::collections::BTreeMap::new();
for t in tables {
m.insert(t.to_string(), column_names(conn, t).await);
}
m
}
async fn table_row_counts(
conn: &Connection,
tables: &[&str],
) -> std::collections::BTreeMap<String, i64> {
let mut m = std::collections::BTreeMap::new();
for t in tables {
let count: i64 = conn
.query(&format!("SELECT COUNT(*) FROM {t}"), ())
.await
.expect("count rows")[0]
.get::<i64>(0)
.expect("count");
m.insert(t.to_string(), count);
}
m
}
const OLD_BASELINE_BOARD_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS tickets (
id TEXT PRIMARY KEY,
title TEXT NOT NULL,
description TEXT NOT NULL,
phase TEXT NOT NULL DEFAULT 'backlog',
assigned_to TEXT,
workspace_name TEXT NOT NULL,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
prerequisites TEXT NOT NULL DEFAULT '[]',
supersedes TEXT,
superseded_by TEXT,
commit_hash TEXT,
lines_added INTEGER,
lines_removed INTEGER,
reporter TEXT NOT NULL DEFAULT '',
is_archived INTEGER NOT NULL DEFAULT 0,
embedding BLOB,
pipeline_reservation INTEGER NOT NULL DEFAULT 0,
priority INTEGER NOT NULL DEFAULT 1,
reviewed_head TEXT,
reviewed_tree TEXT,
done_at TEXT,
bounce_count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS ticket_comments (
id TEXT PRIMARY KEY,
ticket_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL,
FOREIGN KEY (ticket_id) REFERENCES tickets(id)
);
CREATE TABLE IF NOT EXISTS ticket_counters (
workspace_name TEXT PRIMARY KEY,
next_id INTEGER NOT NULL DEFAULT 1
);";
const OLD_BASELINE_SESSION_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS sessions (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
content TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS session_metadata (
agent_id TEXT PRIMARY KEY,
last_activity TEXT NOT NULL,
channel TEXT,
user_name TEXT,
workspace_name TEXT,
role TEXT,
active_models TEXT,
token_length INTEGER,
message_count INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS jobs (
id TEXT PRIMARY KEY,
kind TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'launched',
task TEXT NOT NULL DEFAULT '',
workspace_name TEXT NOT NULL,
user_name TEXT NOT NULL DEFAULT '',
channel TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL,
retry_count INTEGER NOT NULL DEFAULT 0,
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS agents (
job_id TEXT REFERENCES jobs(id) ON DELETE CASCADE,
agent_id TEXT NOT NULL,
kind TEXT NOT NULL,
idx INTEGER,
status TEXT NOT NULL DEFAULT 'launched',
outcome TEXT,
task TEXT NOT NULL,
PRIMARY KEY (job_id, agent_id)
);
CREATE TABLE IF NOT EXISTS pending_jobs (
id TEXT PRIMARY KEY,
target_agent_id TEXT NOT NULL,
envelope TEXT NOT NULL,
created_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS ticket_stage_jobs (
id TEXT PRIMARY KEY REFERENCES jobs(id) ON DELETE CASCADE,
ticket_id TEXT NOT NULL,
stage TEXT NOT NULL,
phase TEXT NOT NULL,
round INTEGER NOT NULL
);
CREATE TABLE IF NOT EXISTS research_jobs (
id TEXT PRIMARY KEY REFERENCES jobs(id) ON DELETE CASCADE,
state TEXT NOT NULL
);";
const OLD_BASELINE_WORKSPACE_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS workspaces (
name TEXT PRIMARY KEY,
path TEXT NOT NULL UNIQUE,
status TEXT NOT NULL DEFAULT 'pending',
created_at TEXT NOT NULL,
updated_at TEXT NOT NULL,
maintenance INTEGER NOT NULL DEFAULT 0,
paused INTEGER NOT NULL DEFAULT 1,
maintainer_debounce_mins INTEGER NOT NULL DEFAULT 5,
maintainer_last_run_at TEXT,
diagnostics TEXT,
diagnostics_generation INTEGER NOT NULL DEFAULT 0,
notes TEXT NOT NULL DEFAULT '',
last_analyzed_commit TEXT,
discovery_generation INTEGER NOT NULL DEFAULT 0
);
CREATE TABLE IF NOT EXISTS workspace_contexts (
workspace_name TEXT NOT NULL REFERENCES workspaces(name) ON DELETE CASCADE,
role TEXT,
content TEXT NOT NULL,
created_at TEXT NOT NULL,
UNIQUE(workspace_name, role)
);
CREATE TABLE IF NOT EXISTS editor_tabs (
workspace_name TEXT NOT NULL REFERENCES workspaces(name) ON DELETE CASCADE,
file_path TEXT NOT NULL,
tab_order INTEGER NOT NULL DEFAULT 0,
is_active INTEGER NOT NULL DEFAULT 0,
is_dirty INTEGER NOT NULL DEFAULT 0,
dirty_content TEXT,
PRIMARY KEY (workspace_name, file_path)
);";
const OLD_BASELINE_USERS_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS users (
name TEXT PRIMARY KEY,
permissions TEXT,
selected_workspace TEXT,
selected_role TEXT
);
CREATE TABLE IF NOT EXISTS user_channels (
user_name TEXT NOT NULL REFERENCES users(name),
channel TEXT NOT NULL,
identifier TEXT NOT NULL,
reply_target TEXT,
UNIQUE(channel, identifier)
);
CREATE TABLE IF NOT EXISTS user_roles (
user_name TEXT NOT NULL REFERENCES users(name),
role TEXT NOT NULL,
PRIMARY KEY (user_name, role)
);";
const OLD_BASELINE_CONFIG_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS config_kv (
key TEXT PRIMARY KEY,
value TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS config_role (
role TEXT PRIMARY KEY,
model TEXT,
reasoning_effort TEXT
);
CREATE TABLE IF NOT EXISTS config_model_routing (
model TEXT PRIMARY KEY,
provider_order TEXT,
allow_fallbacks INTEGER
);";
const OLD_BASELINE_CHAT_HISTORY_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS chat_history (
id INTEGER PRIMARY KEY AUTOINCREMENT,
message_id TEXT NOT NULL UNIQUE,
user_name TEXT NOT NULL,
direction TEXT NOT NULL,
content TEXT NOT NULL,
agent_role TEXT,
workspace TEXT NOT NULL
);";
const OLD_BASELINE_CORE_INDEXES: &str = "\
CREATE INDEX IF NOT EXISTS idx_ticket_comments_ticket_id ON ticket_comments(ticket_id);
CREATE INDEX IF NOT EXISTS idx_sessions_agent_id ON sessions(agent_id, id);
CREATE INDEX IF NOT EXISTS idx_jobs_kind_status ON jobs(kind, status);
CREATE INDEX IF NOT EXISTS idx_jobs_updated_at ON jobs(updated_at);
CREATE UNIQUE INDEX IF NOT EXISTS idx_agents_anchor ON agents(agent_id) WHERE job_id IS NULL;
CREATE INDEX IF NOT EXISTS idx_pending_jobs_agent_created ON pending_jobs(target_agent_id, created_at);
CREATE UNIQUE INDEX IF NOT EXISTS workspace_contexts_null_role ON workspace_contexts(workspace_name) WHERE role IS NULL;
CREATE INDEX IF NOT EXISTS idx_chat_history_user ON chat_history(user_name);
CREATE INDEX IF NOT EXISTS idx_chat_history_workspace ON chat_history(workspace);
CREATE INDEX IF NOT EXISTS idx_chat_history_user_ws_id ON chat_history(user_name, workspace, id);";
const OLD_BASELINE_LOGS_TABLES: &str = "\
CREATE TABLE IF NOT EXISTS logs (
id INTEGER PRIMARY KEY AUTOINCREMENT,
timestamp TEXT NOT NULL,
level TEXT NOT NULL,
target TEXT NOT NULL,
message TEXT NOT NULL,
fields TEXT NOT NULL DEFAULT '{}',
agent_id TEXT NOT NULL DEFAULT '',
agent_role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT ''
);
CREATE TABLE IF NOT EXISTS tool_calls (
id INTEGER PRIMARY KEY AUTOINCREMENT,
agent_id TEXT NOT NULL,
role TEXT NOT NULL,
tool_name TEXT NOT NULL,
arguments TEXT NOT NULL DEFAULT '{}',
duration_ms INTEGER NOT NULL DEFAULT 0,
success INTEGER NOT NULL DEFAULT 1,
error_message TEXT,
workspace TEXT NOT NULL DEFAULT '',
recorded_at TEXT NOT NULL
);
CREATE TABLE IF NOT EXISTS llm_requests (
id INTEGER PRIMARY KEY AUTOINCREMENT,
recorded_at TEXT NOT NULL,
purpose TEXT NOT NULL,
agent_id TEXT NOT NULL DEFAULT '',
role TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
ticket_id TEXT,
model TEXT NOT NULL,
routing TEXT NOT NULL DEFAULT '',
input_tokens INTEGER,
output_tokens INTEGER,
cached_input_tokens INTEGER,
cache_miss_tokens INTEGER,
duration_ms INTEGER NOT NULL,
retry_attempts INTEGER NOT NULL,
finish_reason TEXT,
failure_class TEXT,
success INTEGER NOT NULL DEFAULT 1,
cost REAL,
cost_details TEXT,
upstream_provider TEXT,
system_fingerprint TEXT
);";
const OLD_BASELINE_LOGS_INDEXES: &str = "\
CREATE INDEX IF NOT EXISTS idx_logs_timestamp ON logs(timestamp);
CREATE INDEX IF NOT EXISTS idx_logs_level ON logs(level);
CREATE INDEX IF NOT EXISTS idx_logs_target ON logs(target);
CREATE INDEX IF NOT EXISTS idx_logs_agent_role ON logs(agent_role);
CREATE INDEX IF NOT EXISTS idx_logs_agent_id ON logs(agent_id);
CREATE INDEX IF NOT EXISTS idx_logs_workspace ON logs(workspace);
CREATE INDEX IF NOT EXISTS idx_tool_calls_agent_id ON tool_calls(agent_id);
CREATE INDEX IF NOT EXISTS idx_tool_calls_role ON tool_calls(role);
CREATE INDEX IF NOT EXISTS idx_tool_calls_tool_name ON tool_calls(tool_name);
CREATE INDEX IF NOT EXISTS idx_tool_calls_recorded_at ON tool_calls(recorded_at);
CREATE INDEX IF NOT EXISTS idx_tool_calls_workspace ON tool_calls(workspace);
CREATE INDEX IF NOT EXISTS idx_tool_calls_error_message ON tool_calls(error_message);
CREATE INDEX IF NOT EXISTS idx_llm_requests_recorded_at ON llm_requests(recorded_at);
CREATE INDEX IF NOT EXISTS idx_llm_requests_agent_id ON llm_requests(agent_id);
CREATE INDEX IF NOT EXISTS idx_llm_requests_model ON llm_requests(model);
CREATE INDEX IF NOT EXISTS idx_llm_requests_purpose ON llm_requests(purpose);";
const OLD_DELTA_DROP_CONFIG_ROLE: &str = "DROP TABLE IF EXISTS config_role;";
const OLD_DELTA_FTS_INDEX: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram');";
const OLD_DELTA_BOARD_ACTIVE_INDEX: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_board_active ON tickets \
(is_archived, priority ASC, created_at DESC);";
const OLD_DELTA_GREP_TELEMETRY: &str = "\
CREATE TABLE IF NOT EXISTS grep_telemetry (
id INTEGER PRIMARY KEY AUTOINCREMENT,
recorded_at TEXT NOT NULL,
command TEXT NOT NULL,
served INTEGER NOT NULL DEFAULT 0,
reason TEXT NOT NULL DEFAULT '',
recursive INTEGER NOT NULL DEFAULT 0,
piped INTEGER NOT NULL DEFAULT 0,
operand_count INTEGER NOT NULL DEFAULT 0,
flags TEXT NOT NULL DEFAULT '',
mode TEXT NOT NULL DEFAULT '',
workspace TEXT NOT NULL DEFAULT '',
grep_count INTEGER NOT NULL DEFAULT 0,
served_count INTEGER NOT NULL DEFAULT 0,
skipped_count INTEGER NOT NULL DEFAULT 0,
duration_ms INTEGER,
exit_code INTEGER
);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_recorded_at ON grep_telemetry(recorded_at);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_served ON grep_telemetry(served);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_reason ON grep_telemetry(reason);
CREATE INDEX IF NOT EXISTS idx_grep_telemetry_command ON grep_telemetry(command);";
const OLD_DELTA_ALARMS: &str = "\
CREATE TABLE IF NOT EXISTS alarms (
id TEXT PRIMARY KEY,
session_id TEXT NOT NULL,
user_name TEXT NOT NULL,
kind TEXT NOT NULL,
text TEXT NOT NULL,
fire_at TEXT,
interval_seconds INTEGER,
next_fire_at TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'active',
created_at TEXT NOT NULL
);
CREATE INDEX IF NOT EXISTS idx_alarms_due ON alarms(status, next_fire_at);";
const OLD_DELTA_CHAT_HISTORY_TIMESTAMP: &str =
"ALTER TABLE chat_history ADD COLUMN timestamp TEXT;";
const OLD_DELTA_RENAME_READY_FOR_DEVELOPMENT_TO_QUEUED: &str = "\
UPDATE tickets SET phase = 'queued' WHERE phase = 'ready_for_development';\
UPDATE ticket_chronicle SET source_phase = 'queued' WHERE source_phase = 'ready_for_development';\
UPDATE ticket_chronicle SET target_phase = 'queued' WHERE target_phase = 'ready_for_development';\
UPDATE jobs SET kind = 'queued' WHERE kind = 'ready_for_development';";
const OLD_DELTA_TICKETS_WORKSPACE_PHASE_INDEX: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_workspace_phase \
ON tickets (workspace_name, phase, is_archived, priority ASC, created_at DESC);";
fn noop_import(_conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(async { Ok(()) })
}
fn cleanup_legacy_ticket_jobs(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_import_cleanup(conn))
}
async fn run_import_cleanup(conn: &Connection) -> anyhow::Result<()> {
conn.execute(
"DELETE FROM jobs WHERE kind IN \
('ticket_stage', 'ticket_analysis', 'ticket_implementation')",
(),
)
.await
.context("Failed to discard legacy ticket jobs after consolidation import")?;
let violations = conn
.query("PRAGMA foreign_key_check", ())
.await
.context("Failed to run PRAGMA foreign_key_check after consolidation import")?;
if !violations.is_empty() {
tracing::warn!(
count = violations.len(),
"consolidation import found orphan rows violating the new FKs; \
preserving them as soft references (future writes still enforce the FK)",
);
}
Ok(())
}
fn drop_jobs_paused_frozen(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_jobs_paused_frozen(conn))
}
async fn run_drop_jobs_paused_frozen(conn: &Connection) -> anyhow::Result<()> {
if column_exists(conn, "jobs", "paused_frozen").await? {
conn.execute("ALTER TABLE jobs DROP COLUMN paused_frozen", ())
.await
.context("Failed to drop jobs.paused_frozen")?;
}
Ok(())
}
fn drop_user_roles_and_seed_onboarding(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_drop_user_roles_and_seed_onboarding(conn))
}
async fn run_drop_user_roles_and_seed_onboarding(conn: &Connection) -> anyhow::Result<()> {
let user_count: i64 = conn
.query("SELECT COUNT(*) FROM users", ())
.await
.context("Failed to probe users count")?
.into_iter()
.next()
.and_then(|row| row.get::<i64>(0).ok())
.unwrap_or(0);
if user_count > 0 {
conn.execute(
"INSERT OR REPLACE INTO config_kv (key, value) VALUES (?1, ?2)",
params![
crate::config::CONFIG_KEY_ONBOARDING_STATE,
crate::config::OnboardingState::Finished.as_str(),
],
)
.await
.context("Failed to set onboarding_state=finished")?;
let rows = conn
.query("SELECT name, permissions, selected_role FROM users", ())
.await?;
for row in rows {
let name: String = row.get(0)?;
let permissions: Option<String> = row.get(1)?;
let selected_role: Option<String> = row.get(2)?;
let sel = selected_role.unwrap_or_default();
let needs_default = sel.is_empty() || {
let in_pool = if permissions.as_deref() == Some("full") {
matches!(sel.as_str(), "support" | "assistant" | "manager" | "artist")
} else {
matches!(sel.as_str(), "assistant" | "artist")
};
!in_pool
};
if needs_default {
conn.execute(
"UPDATE users SET selected_role = 'assistant' WHERE name = ?1",
params![name],
)
.await
.context("Failed to normalize selected_role")?;
}
}
}
conn.execute("DROP TABLE IF EXISTS user_roles", ())
.await
.context("Failed to drop user_roles")?;
Ok(())
}
fn rewrite_analysis_verdicts(conn: &Connection) -> BoxFuture<'_, anyhow::Result<()>> {
Box::pin(run_rewrite_analysis_verdicts(conn))
}
async fn run_rewrite_analysis_verdicts(conn: &Connection) -> anyhow::Result<()> {
let rows = conn
.query(
"SELECT job_id, agent_id, outcome FROM agents \
WHERE kind = 'analyst' AND outcome LIKE '{\"verdict\":%'",
(),
)
.await
.context("Failed to read analysis verdict rows for migration")?;
for row in rows {
let job_id: String = row.get(0)?;
let agent_id: String = row.get(1)?;
let outcome: String = row.get(2)?;
let Ok(value) = serde_json::from_str::<serde_json::Value>(&outcome) else {
continue;
};
let Some(verdict) = value.get("verdict") else {
continue;
};
let Some(score) = verdict.get("score").and_then(serde_json::Value::as_u64) else {
continue;
};
let Some(issues) = verdict.get("issues").and_then(serde_json::Value::as_array) else {
continue;
};
let grade = if score < 7 { "blocker" } else { "minor" };
let graded: Vec<serde_json::Value> = issues
.iter()
.filter_map(serde_json::Value::as_str)
.map(|text| serde_json::json!({ "text": text, "grade": grade }))
.collect();
let rewritten = serde_json::json!({ "verdict": { "issues": graded } }).to_string();
conn.execute(
"UPDATE agents SET outcome = ?1 WHERE job_id = ?2 AND agent_id = ?3",
params![rewritten, job_id, agent_id],
)
.await
.context("Failed to rewrite analysis verdict row")?;
}
Ok(())
}
const OLD_CATALOG: &[Migration] = &[
Migration {
id: "1",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_BOARD_TABLES),
},
Migration {
id: "2",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_SESSION_TABLES),
},
Migration {
id: "3",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_WORKSPACE_TABLES),
},
Migration {
id: "4",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_USERS_TABLES),
},
Migration {
id: "5",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_CONFIG_TABLES),
},
Migration {
id: "6",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_CHAT_HISTORY_TABLES),
},
Migration {
id: "7",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_BASELINE_CORE_INDEXES),
},
Migration {
id: "001_drop_ticket_pipeline_reservation",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE tickets DROP COLUMN pipeline_reservation;"),
},
Migration {
id: "002_drop_ticket_assigned_to",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE tickets DROP COLUMN assigned_to;"),
},
Migration {
id: "003_drop_ticket_stage_jobs_round",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE ticket_stage_jobs DROP COLUMN round;"),
},
Migration {
id: "005_drop_ticket_stage_jobs_phase",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE ticket_stage_jobs DROP COLUMN phase;"),
},
Migration {
id: "006_rename_ticket_stage_jobs_to_ticket_jobs",
target: TargetDb::Core,
body: MigrationBody::Sql(
"DROP TABLE IF EXISTS ticket_jobs; ALTER TABLE ticket_stage_jobs RENAME TO ticket_jobs;",
),
},
Migration {
id: "consolidate_002_drop_ticket_jobs_stage",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE ticket_jobs DROP COLUMN stage;"),
},
Migration {
id: "007_jobs_paused_frozen",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE jobs ADD COLUMN paused_frozen INTEGER;"),
},
Migration {
id: "consolidate_003_jobs_ticket_id",
target: TargetDb::Core,
body: MigrationBody::Sql(
"ALTER TABLE jobs ADD COLUMN ticket_id TEXT REFERENCES tickets(id);",
),
},
Migration {
id: "consolidate_005_drop_ticket_jobs",
target: TargetDb::Core,
body: MigrationBody::Sql("DROP TABLE IF EXISTS ticket_jobs;"),
},
Migration {
id: "consolidate_006_drop_jobs_paused_frozen",
target: TargetDb::Core,
body: MigrationBody::Sql("ALTER TABLE jobs DROP COLUMN paused_frozen;"),
},
Migration {
id: "consolidate_007_jobs_phase_ticket_index",
target: TargetDb::Core,
body: MigrationBody::Sql(
"CREATE UNIQUE INDEX IF NOT EXISTS idx_jobs_phase_ticket ON jobs(kind, ticket_id) WHERE ticket_id IS NOT NULL;",
),
},
Migration {
id: "consolidate_008_create_ticket_chronicle",
target: TargetDb::Core,
body: MigrationBody::Sql(
"CREATE TABLE IF NOT EXISTS ticket_chronicle (\
id INTEGER PRIMARY KEY AUTOINCREMENT,\
ticket_id TEXT NOT NULL,\
workspace_name TEXT NOT NULL,\
source_phase TEXT NOT NULL,\
target_phase TEXT NOT NULL,\
at TEXT NOT NULL\
);\
CREATE INDEX IF NOT EXISTS idx_ticket_chronicle_ws_id \
ON ticket_chronicle(workspace_name, id);\
CREATE UNIQUE INDEX IF NOT EXISTS idx_ticket_chronicle_dedup \
ON ticket_chronicle(ticket_id, workspace_name, source_phase, target_phase, at);",
),
},
Migration {
id: "10",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_DROP_CONFIG_ROLE),
},
Migration {
id: "11",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_FTS_INDEX),
},
Migration {
id: "12",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_BOARD_ACTIVE_INDEX),
},
Migration {
id: "consolidate_001_import_domain_stores",
target: TargetDb::Core,
body: MigrationBody::Rust(noop_import),
},
Migration {
id: "003_reset_nonterminal_tickets",
target: TargetDb::Core,
body: MigrationBody::Sql(
"UPDATE tickets SET phase = 'backlog' WHERE phase NOT IN ('done','cancelled','failed') AND is_archived = 0;",
),
},
Migration {
id: "13",
target: TargetDb::Core,
body: MigrationBody::Rust(cleanup_legacy_ticket_jobs),
},
Migration {
id: "15",
target: TargetDb::Core,
body: MigrationBody::Sql("DROP TABLE IF EXISTS ticket_stage_jobs;"),
},
Migration {
id: "16",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_jobs_paused_frozen),
},
Migration {
id: "8",
target: TargetDb::Logs,
body: MigrationBody::Sql(OLD_BASELINE_LOGS_TABLES),
},
Migration {
id: "9",
target: TargetDb::Logs,
body: MigrationBody::Sql(OLD_BASELINE_LOGS_INDEXES),
},
Migration {
id: "14",
target: TargetDb::Logs,
body: MigrationBody::Sql(OLD_DELTA_GREP_TELEMETRY),
},
Migration {
id: "17",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_ALARMS),
},
Migration {
id: "18",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_CHAT_HISTORY_TIMESTAMP),
},
Migration {
id: "19",
target: TargetDb::Core,
body: MigrationBody::Rust(drop_user_roles_and_seed_onboarding),
},
Migration {
id: "20",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_RENAME_READY_FOR_DEVELOPMENT_TO_QUEUED),
},
Migration {
id: "21",
target: TargetDb::Core,
body: MigrationBody::Rust(rewrite_analysis_verdicts),
},
Migration {
id: "22",
target: TargetDb::Core,
body: MigrationBody::Sql(OLD_DELTA_TICKETS_WORKSPACE_PHASE_INDEX),
},
Migration {
id: "23",
target: TargetDb::Core,
body: MigrationBody::Sql(
"ALTER TABLE chat_history ADD COLUMN reply_author TEXT; ALTER TABLE chat_history ADD COLUMN reply_snippet TEXT;",
),
},
];
fn old_catalog_without_reply_delta() -> Vec<Migration> {
OLD_CATALOG
.iter()
.filter(|m| m.id != "23")
.copied()
.collect()
}
#[tokio::test]
async fn fresh_install_converges_to_expected_shape() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_consolidated_store(tmp.path())
.await
.expect("fresh consolidated store");
assert_schema_matches(
&conn,
EXPECTED_CORE_TABLE_COLUMNS,
EXPECTED_CORE_INDEXES,
"fresh core",
)
.await;
let mut applied = applied_ids(&conn).await;
applied.sort();
assert_eq!(
applied,
[
"24", "25", "27", "28", "29", "30", "31", "32", "33", "34", "36", "37", "38", "39",
"40", "41", "42", "43", "44", "45"
]
.map(String::from),
"fresh core applies the 24–34 baseline + the 36–45 tail exactly"
);
}
#[tokio::test]
async fn fresh_logs_install_converges_to_expected_shape() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::LOG_DB_NAME),
"",
)
.await
.expect("open logs");
run_migrations(&conn, TargetDb::Logs)
.await
.expect("run logs catalog");
assert_schema_matches(
&conn,
EXPECTED_LOGS_TABLE_COLUMNS,
EXPECTED_LOGS_INDEXES,
"fresh logs",
)
.await;
let mut applied = applied_ids(&conn).await;
applied.sort();
assert_eq!(
applied,
vec!["26".to_string(), "35".to_string()],
"fresh logs applies baseline 26 + delta 35 exactly"
);
}
#[tokio::test]
async fn baseline_logs_db_gains_llm_failures_via_delta_35() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::LOG_DB_NAME),
"",
)
.await
.expect("open logs");
run_catalog(&conn, TargetDb::Logs, &MIGRATIONS[..3])
.await
.expect("baseline only");
conn.execute("DROP TABLE llm_failures", ()).await.unwrap();
assert!(!table_defs(&conn).await.contains_key("llm_failures"));
run_migrations(&conn, TargetDb::Logs)
.await
.expect("delta 35");
let mut applied = applied_ids(&conn).await;
applied.sort();
assert_eq!(
applied,
vec!["26".to_string(), "35".to_string()],
"reopen must record exactly 26 + 35"
);
let (_, cols) = EXPECTED_LOGS_TABLE_COLUMNS
.iter()
.find(|(n, _)| *n == "llm_failures")
.expect("pinned llm_failures shape");
let expected_cols: Vec<String> = cols.iter().map(|s| (*s).to_string()).collect();
assert_eq!(
column_names(&conn, "llm_failures").await,
expected_cols,
"delta 35 must create llm_failures with the pinned shape"
);
for idx in [
"idx_llm_failures_recorded_at",
"idx_llm_failures_operation_id",
"idx_llm_failures_failure_class",
] {
assert!(
index_defs(&conn).await.contains_key(idx),
"delta 35 must create index {idx}"
);
}
}
async fn seed_current_core_rows(conn: &Connection) {
let now = crate::db::now();
conn.execute(
"INSERT INTO tickets (id, title, description, phase, workspace_name, created_at, updated_at) \
VALUES ('T1', 't', 'd', 'backlog', 'ws', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_comments (id, ticket_id, role, content, created_at) \
VALUES ('C1', 'T1', 'manager', 'ship it', ?1)",
params![now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO jobs (id, kind, role, workspace_name, task, user_name, channel, retry_count, \
status, created_at, updated_at, ticket_id) \
VALUES ('J1', 'analysis', 'analyst', 'ws', '', 'bob', '', 0, 'launched', ?1, ?1, 'T1')",
params![now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES ('J1', 'A1', 'analyst', 0, 'done', NULL, 'task')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO users (name, selected_workspace) VALUES ('bob', 'ws')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO chat_history (message_id, user_name, direction, content, agent_role, workspace) \
VALUES ('m1', 'bob', 'in', 'hello', NULL, 'ws')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO ticket_chronicle (ticket_id, workspace_name, source_phase, target_phase, at) \
VALUES ('T1', 'ws', 'backlog', 'queued', ?1)",
params![now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO alarms (id, session_id, user_name, kind, text, next_fire_at, created_at) \
VALUES ('a1', 's1', 'bob', 'reminder', 'ping', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
conn.execute("INSERT INTO config_kv (key, value) VALUES ('k', 'v')", ())
.await
.unwrap();
}
struct CoreSnapshot {
tables: std::collections::BTreeMap<String, String>,
indexes: std::collections::BTreeMap<String, String>,
cols: std::collections::BTreeMap<String, Vec<String>>,
counts: std::collections::BTreeMap<String, i64>,
chat: Vec<String>,
tickets: Vec<(String, String)>,
}
async fn snapshot_core_state(conn: &Connection) -> CoreSnapshot {
CoreSnapshot {
tables: table_defs(conn).await,
indexes: index_defs(conn).await,
cols: column_sets(conn, &expected_core_table_names()).await,
counts: table_row_counts(conn, &expected_core_domain_tables()).await,
chat: conn
.query("SELECT content FROM chat_history", ())
.await
.unwrap()
.into_iter()
.map(|r| r.get::<String>(0).unwrap())
.collect(),
tickets: conn
.query("SELECT title, phase FROM tickets", ())
.await
.unwrap()
.into_iter()
.map(|r| (r.get::<String>(0).unwrap(), r.get::<String>(1).unwrap()))
.collect(),
}
}
fn without<V: Clone>(
map: &std::collections::BTreeMap<String, V>,
exclude: &[&str],
) -> std::collections::BTreeMap<String, V> {
map.iter()
.filter(|(name, _)| !exclude.contains(&name.as_str()))
.map(|(name, v)| (name.clone(), v.clone()))
.collect()
}
async fn assert_core_catalog_unchanged(
conn: &Connection,
before: &CoreSnapshot,
except_tables: &[&str],
except_indexes: &[&str],
) {
assert_eq!(
without(&table_defs(conn).await, except_tables),
without(&before.tables, except_tables),
"core table DDL (except {except_tables:?}) must be unchanged on reopen"
);
assert_eq!(
without(&index_defs(conn).await, except_indexes),
without(&before.indexes, except_indexes),
"core index definitions (except {except_indexes:?}) must be unchanged on reopen"
);
assert_eq!(
without(
&column_sets(conn, &expected_core_table_names()).await,
except_tables
),
without(&before.cols, except_tables),
"core column sets (except {except_tables:?}) must be unchanged on reopen"
);
assert_eq!(
table_row_counts(conn, &expected_core_domain_tables()).await,
before.counts,
"core row counts must be unchanged on reopen"
);
let after_chat: Vec<String> = conn
.query("SELECT content FROM chat_history", ())
.await
.unwrap()
.into_iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert_eq!(
after_chat, before.chat,
"chat_history content must be unchanged"
);
let after_tickets: Vec<(String, String)> = conn
.query("SELECT title, phase FROM tickets", ())
.await
.unwrap()
.into_iter()
.map(|r| (r.get::<String>(0).unwrap(), r.get::<String>(1).unwrap()))
.collect();
assert_eq!(
after_tickets, before.tickets,
"tickets content must be unchanged"
);
}
#[tokio::test]
#[expect(clippy::too_many_lines)]
async fn old_catalog_current_db_reopens_as_noop() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
run_catalog(&conn, TargetDb::Core, OLD_CATALOG)
.await
.expect("old catalog");
seed_current_core_rows(&conn).await;
let before_ids = applied_ids(&conn).await;
let before = snapshot_core_state(&conn).await;
run_migrations(&conn, TargetDb::Core)
.await
.expect("new catalog");
let mut expected_ids = before_ids.clone();
for id in [
"24", "25", "27", "28", "29", "30", "31", "32", "33", "34", "36", "37", "38", "39",
"40", "41", "42", "43", "44", "45",
] {
expected_ids.push(id.to_string());
}
expected_ids.sort();
let mut after_ids = applied_ids(&conn).await;
after_ids.sort();
assert_eq!(
after_ids, expected_ids,
"reopen must record exactly old ids ∪ 24/25/27..34/36..45"
);
assert_core_catalog_unchanged(
&conn,
&before,
&[
"workspaces",
"jobs",
"session_metadata",
"users",
"alarms",
"chat_history",
"tickets",
"ticket_chronicle",
],
&["idx_jobs_caller_agent"],
)
.await;
let after_ws_cols = column_sets(&conn, &["workspaces"]).await;
let mut expected_ws_cols = before.cols["workspaces"].clone();
expected_ws_cols.push("maintainer_recommendations".to_string());
assert_eq!(
after_ws_cols["workspaces"], expected_ws_cols,
"reopen must append exactly maintainer_recommendations to workspaces columns"
);
let after_jobs_cols = column_sets(&conn, &["jobs"]).await;
let mut expected_jobs_cols = before.cols["jobs"].clone();
expected_jobs_cols.push("caller_agent_id".to_string());
expected_jobs_cols.push("mode".to_string());
assert_eq!(
after_jobs_cols["jobs"], expected_jobs_cols,
"reopen must append exactly caller_agent_id/mode to jobs columns"
);
let after_sm_cols = column_sets(&conn, &["session_metadata"]).await;
let mut expected_sm_cols = before.cols["session_metadata"].clone();
expected_sm_cols.push("created_at".to_string());
expected_sm_cols.push("sleep_ended".to_string());
assert_eq!(
after_sm_cols["session_metadata"], expected_sm_cols,
"reopen must append exactly created_at/sleep_ended to session_metadata columns"
);
let after_users_cols = column_sets(&conn, &["users"]).await;
let mut expected_users_cols = before.cols["users"].clone();
expected_users_cols.retain(|c| c != "permissions" && c != "selected_role");
expected_users_cols.push("image_gen_model".to_string());
expected_users_cols.push("video_model".to_string());
expected_users_cols.push("granted_tools".to_string());
assert_eq!(
after_users_cols["users"], expected_users_cols,
"reopen must append exactly image_gen_model/video_model/granted_tools \
and drop permissions/selected_role from users columns"
);
let after_alarms_cols = column_sets(&conn, &["alarms"]).await;
let mut expected_alarms_cols = before.cols["alarms"].clone();
expected_alarms_cols.push("trigger".to_string());
assert_eq!(
after_alarms_cols["alarms"], expected_alarms_cols,
"reopen must trade command (delta 32) for trigger (deltas 42/43) on alarms columns"
);
let after_chat_cols = column_sets(&conn, &["chat_history"]).await;
let mut expected_chat_cols = before.cols["chat_history"].clone();
expected_chat_cols.push("broadcast_id".to_string());
assert_eq!(
after_chat_cols["chat_history"], expected_chat_cols,
"reopen must append exactly broadcast_id to chat_history columns"
);
let after_tickets_cols = column_sets(&conn, &["tickets"]).await;
let mut expected_tickets_cols = before.cols["tickets"].clone();
expected_tickets_cols.push("last_transition_actor".to_string());
assert_eq!(
after_tickets_cols["tickets"], expected_tickets_cols,
"reopen must append exactly last_transition_actor to tickets columns"
);
let after_chronicle_cols = column_sets(&conn, &["ticket_chronicle"]).await;
let mut expected_chronicle_cols = before.cols["ticket_chronicle"].clone();
expected_chronicle_cols.push("actor".to_string());
assert_eq!(
after_chronicle_cols["ticket_chronicle"], expected_chronicle_cols,
"reopen must append exactly actor to ticket_chronicle columns"
);
}
#[tokio::test]
async fn old_catalog_logs_db_reopens_as_noop() {
const NEW_TABLE: &str = "llm_failures";
const NEW_INDEXES: [&str; 3] = [
"idx_llm_failures_recorded_at",
"idx_llm_failures_operation_id",
"idx_llm_failures_failure_class",
];
let counts_tables: Vec<&str> = expected_logs_domain_tables()
.into_iter()
.filter(|n| *n != NEW_TABLE)
.collect();
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::LOG_DB_NAME),
"",
)
.await
.expect("open logs");
run_catalog(&conn, TargetDb::Logs, OLD_CATALOG)
.await
.expect("old catalog");
conn.execute(
"INSERT INTO grep_telemetry (recorded_at, command) VALUES (?1, 'grep foo')",
params![crate::db::now()],
)
.await
.unwrap();
let before_ids = applied_ids(&conn).await;
let before_tables = table_defs(&conn).await;
let before_indexes = index_defs(&conn).await;
let before_counts = table_row_counts(&conn, &counts_tables).await;
run_migrations(&conn, TargetDb::Logs)
.await
.expect("new catalog");
let mut expected_ids = before_ids.clone();
expected_ids.push("26".to_string());
expected_ids.push("35".to_string());
expected_ids.sort();
let mut after_ids = applied_ids(&conn).await;
after_ids.sort();
assert_eq!(
after_ids, expected_ids,
"logs reopen must record exactly old ids ∪ 26 ∪ 35"
);
assert_eq!(
without(&table_defs(&conn).await, &[NEW_TABLE]),
without(&before_tables, &[NEW_TABLE]),
"logs table DDL (except {NEW_TABLE}) must be unchanged on reopen"
);
assert_eq!(
without(&index_defs(&conn).await, &NEW_INDEXES),
without(&before_indexes, &NEW_INDEXES),
"logs index definitions (except {NEW_INDEXES:?}) must be unchanged on reopen"
);
assert_eq!(
table_row_counts(&conn, &counts_tables).await,
before_counts,
"logs row counts must be unchanged on reopen"
);
let after_tables = table_defs(&conn).await;
assert!(
after_tables.contains_key(NEW_TABLE),
"delta 35 must create {NEW_TABLE} on existing installs"
);
let after_indexes = index_defs(&conn).await;
for idx in NEW_INDEXES {
assert!(
after_indexes.contains_key(idx),
"delta 35 must create index {idx}"
);
}
}
#[expect(clippy::too_many_lines)] #[tokio::test]
async fn one_delta_behind_db_upgrades_reply_columns() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
run_catalog(&conn, TargetDb::Core, &old_catalog_without_reply_delta())
.await
.expect("old catalog minus 23");
let now = crate::db::now();
conn.execute(
"INSERT INTO chat_history (message_id, user_name, direction, content, agent_role, workspace, timestamp) \
VALUES ('m1', 'alice', 'in', 'hello', NULL, 'ws', ?1)",
params![now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO chat_history (message_id, user_name, direction, content, agent_role, workspace, timestamp) \
VALUES ('m2', 'assistant', 'out', 'hi alice', 'assistant', 'ws', ?1)",
params![now.clone()],
)
.await
.unwrap();
let before_chat_cols = column_names(&conn, "chat_history").await;
assert!(
!before_chat_cols.contains(&"reply_author".to_string()),
"0.5.0 DB must not have reply_author"
);
assert!(
!before_chat_cols.contains(&"reply_snippet".to_string()),
"0.5.0 DB must not have reply_snippet"
);
let untouched_tables: Vec<&str> = expected_core_table_names()
.into_iter()
.filter(|n| {
*n != "chat_history"
&& *n != "workspaces"
&& *n != "jobs"
&& *n != "session_metadata"
&& *n != "users"
&& *n != "alarms"
&& *n != "tickets"
&& *n != "ticket_chronicle"
&& *n != "schema_migrations"
})
.collect();
let before_other = column_sets(&conn, &untouched_tables).await;
let mut before_ids = applied_ids(&conn).await;
before_ids.sort();
run_migrations(&conn, TargetDb::Core)
.await
.expect("new catalog");
let after_chat_cols = column_names(&conn, "chat_history").await;
assert!(
after_chat_cols.contains(&"reply_author".to_string()),
"reply_author must be added"
);
assert!(
after_chat_cols.contains(&"reply_snippet".to_string()),
"reply_snippet must be added"
);
let after_ws_cols = column_names(&conn, "workspaces").await;
assert!(
after_ws_cols.contains(&"maintainer_recommendations".to_string()),
"maintainer_recommendations must be added to workspaces"
);
let reply_nulls: i64 = conn
.query(
"SELECT COUNT(*) FROM chat_history \
WHERE reply_author IS NULL AND reply_snippet IS NULL",
(),
)
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap();
assert_eq!(reply_nulls, 2, "both rows must keep NULL reply columns");
let contents: Vec<String> = conn
.query("SELECT content FROM chat_history ORDER BY id", ())
.await
.unwrap()
.into_iter()
.map(|r| r.get::<String>(0).unwrap())
.collect();
assert_eq!(contents, vec!["hello".to_string(), "hi alice".to_string()]);
let mut expected_ids = before_ids.clone();
for id in [
"24", "25", "27", "28", "29", "30", "31", "32", "33", "34", "36", "37", "38", "39",
"40", "41", "42", "43", "44", "45",
] {
expected_ids.push(id.to_string());
}
expected_ids.sort();
let mut after_ids = applied_ids(&conn).await;
after_ids.sort();
assert_eq!(
after_ids, expected_ids,
"upgrade must record exactly old ids ∪ 24/25/27..34/36..45"
);
let after_users_cols = column_names(&conn, "users").await;
assert!(
after_users_cols.contains(&"image_gen_model".to_string()),
"image_gen_model must be added to users"
);
assert!(
after_users_cols.contains(&"video_model".to_string()),
"video_model must be added to users"
);
assert!(
after_users_cols.contains(&"granted_tools".to_string()),
"granted_tools must be added to users"
);
assert!(
!after_users_cols.contains(&"permissions".to_string())
&& !after_users_cols.contains(&"selected_role".to_string()),
"the account-kind columns must be dropped from users"
);
let after_alarms_cols = column_names(&conn, "alarms").await;
assert!(
after_alarms_cols.contains(&"trigger".to_string())
&& !after_alarms_cols.contains(&"command".to_string()),
"the retired command column must be replaced by trigger on alarms"
);
let after_chat_cols = column_names(&conn, "chat_history").await;
assert!(
after_chat_cols.contains(&"broadcast_id".to_string()),
"broadcast_id must be added to chat_history"
);
let after_tickets_cols = column_names(&conn, "tickets").await;
assert!(
after_tickets_cols.contains(&"last_transition_actor".to_string()),
"last_transition_actor must be added to tickets"
);
let after_chronicle_cols = column_names(&conn, "ticket_chronicle").await;
assert!(
after_chronicle_cols.contains(&"actor".to_string()),
"actor must be added to ticket_chronicle"
);
assert_eq!(
column_sets(&conn, &untouched_tables).await,
before_other,
"only chat_history (delta 25/33), workspaces (delta 27), \
jobs/session_metadata (delta 28), users (deltas 29/38), jobs.mode \
(delta 30), session_metadata.sleep_ended (delta 31), \
alarms.command → trigger (deltas 32/41/42/43) and \
tickets/ticket_chronicle (delta 34) \
columns may change on the 0.5.0 upgrade"
);
}
async fn users_selected_workspaces(
conn: &Connection,
) -> std::collections::BTreeMap<String, Option<String>> {
conn.query("SELECT name, selected_workspace FROM users", ())
.await
.expect("read users")
.into_iter()
.map(|row| {
(
row.get::<String>(0).expect("user name"),
row.get::<Option<String>>(1).expect("selected_workspace"),
)
})
.collect()
}
async fn model_routing_orders(conn: &Connection) -> Vec<(String, String)> {
conn.query(
"SELECT model, provider_order FROM config_model_routing ORDER BY model",
(),
)
.await
.expect("read config_model_routing")
.into_iter()
.map(|row| {
(
row.get::<String>(0).expect("model"),
row.get::<String>(1).expect("provider_order"),
)
})
.collect()
}
#[tokio::test]
async fn detach_guest_workspaces_clears_only_guests() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(root, crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
run_catalog(&conn, TargetDb::Core, &MIGRATIONS[..2])
.await
.expect("baseline catalog");
conn.execute(
"INSERT INTO users (name, selected_workspace) VALUES ('admin', 'shared_ws')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO users (name, selected_workspace) VALUES ('guest_shared', 'shared_ws')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO users (name, selected_workspace) VALUES ('guest_null', NULL)",
(),
)
.await
.unwrap();
run_detach_guest_workspaces(&conn).await.unwrap();
let after = users_selected_workspaces(&conn).await;
assert_eq!(
after["admin"],
Some("shared_ws".to_string()),
"the admin keeps their shared workspace"
);
assert_eq!(
after["guest_shared"], None,
"guests are detached from shared workspaces"
);
assert_eq!(after["guest_null"], None, "already-personal stays NULL");
run_detach_guest_workspaces(&conn).await.unwrap();
assert_eq!(users_selected_workspaces(&conn).await, after);
}
#[tokio::test]
async fn drop_account_kind_columns_preserves_row_data() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
conn.execute(
"CREATE TABLE users (\
name TEXT PRIMARY KEY, permissions TEXT, selected_workspace TEXT, \
selected_role TEXT, image_gen_model TEXT, video_model TEXT, granted_tools TEXT\
)",
(),
)
.await
.expect("create pre-39 users table");
conn.execute(
"INSERT INTO users \
(name, permissions, selected_workspace, selected_role, image_gen_model, \
video_model, granted_tools) \
VALUES ('admin', 'full', 'proj', 'assistant', 'img-model', 'vid-model', '[\"t\"]')",
(),
)
.await
.unwrap();
conn.execute(
"INSERT INTO users (name, permissions, selected_workspace, selected_role) \
VALUES ('guest', NULL, NULL, 'engineer')",
(),
)
.await
.unwrap();
run_drop_users_permissions(&conn).await.unwrap();
run_drop_users_selected_role(&conn).await.unwrap();
let expected_columns: Vec<String> = [
"name",
"selected_workspace",
"image_gen_model",
"video_model",
"granted_tools",
]
.iter()
.map(|c| (*c).to_string())
.collect();
assert_eq!(column_names(&conn, "users").await, expected_columns);
let rows = conn
.query(
"SELECT name, selected_workspace, image_gen_model, video_model, granted_tools \
FROM users ORDER BY name",
(),
)
.await
.unwrap();
assert_eq!(rows.len(), 2);
assert_eq!(rows[0].get::<String>(0).unwrap(), "admin");
assert_eq!(
rows[0].get::<Option<String>>(1).unwrap(),
Some("proj".to_string()),
"the admin's selected workspace must survive the table rewrite"
);
assert_eq!(
rows[0].get::<Option<String>>(2).unwrap(),
Some("img-model".to_string())
);
assert_eq!(
rows[0].get::<Option<String>>(3).unwrap(),
Some("vid-model".to_string())
);
assert_eq!(
rows[0].get::<Option<String>>(4).unwrap(),
Some("[\"t\"]".to_string())
);
assert_eq!(rows[1].get::<String>(0).unwrap(), "guest");
assert_eq!(rows[1].get::<Option<String>>(1).unwrap(), None);
assert_eq!(rows[1].get::<Option<String>>(4).unwrap(), None);
run_drop_users_permissions(&conn).await.unwrap();
run_drop_users_selected_role(&conn).await.unwrap();
assert_eq!(column_names(&conn, "users").await, expected_columns);
}
#[tokio::test]
async fn command_armed_alarms_are_deleted_and_the_command_column_is_dropped() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
let through_command = MIGRATIONS
.iter()
.position(|m| m.id == "32")
.expect("entry 32 in the catalog");
run_catalog(&conn, TargetDb::Core, &MIGRATIONS[..=through_command])
.await
.expect("catalog through delta 32");
let now = crate::db::now();
for (id, status, command) in [
("armed_active", "active", Some("rm -rf /tmp/x")),
("armed_settled", "fired", Some("say done")),
("reminder", "active", None),
] {
conn.execute(
"INSERT INTO alarms \
(id, session_id, user_name, kind, text, next_fire_at, status, created_at, command) \
VALUES (?1, 's1', 'bob', 'reminder', 'ping', ?2, ?3, ?2, ?4)",
params![id, now.clone(), status, command],
)
.await
.unwrap();
}
run_migrations(&conn, TargetDb::Core)
.await
.expect("new catalog");
let remaining: Vec<String> = conn
.query("SELECT id FROM alarms ORDER BY id", ())
.await
.unwrap()
.into_iter()
.map(|row| row.get::<String>(0).unwrap())
.collect();
assert_eq!(
remaining,
vec!["reminder".to_string()],
"delta 41 must delete every command-armed alarm row and keep the rest"
);
let cols = column_names(&conn, "alarms").await;
assert!(!cols.contains(&"command".to_string()), "got: {cols:?}");
}
#[tokio::test]
async fn rewrite_legacy_deepseek_routing_slug_rewrites_only_the_display_name() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
run_catalog(&conn, TargetDb::Core, &MIGRATIONS[..1])
.await
.expect("baseline catalog");
for (model, order) in [
("deepseek/deepseek-v4-flash", "DeepSeek"),
("deepseek/deepseek-v4.1-flash", "deepseek"),
("qwen/qwen3.8-flash", "GMICloud"),
] {
conn.execute(
"INSERT INTO config_model_routing (model, provider_order) VALUES (?1, ?2)",
params![model, order],
)
.await
.unwrap();
}
let entry_37 = MIGRATIONS
.iter()
.find(|m| m.id == "37")
.expect("entry 37 in the catalog");
run_catalog(&conn, TargetDb::Core, std::slice::from_ref(entry_37))
.await
.expect("apply entry 37");
let rewritten = model_routing_orders(&conn).await;
assert_eq!(
rewritten,
vec![
(
"deepseek/deepseek-v4-flash".to_string(),
"deepseek".to_string()
),
(
"deepseek/deepseek-v4.1-flash".to_string(),
"deepseek".to_string()
),
("qwen/qwen3.8-flash".to_string(), "GMICloud".to_string()),
],
"only the legacy display-name value may be rewritten"
);
conn.execute(REWRITE_LEGACY_DEEPSEEK_ROUTING_SLUG, ())
.await
.expect("re-run entry 37 body");
assert_eq!(model_routing_orders(&conn).await, rewritten);
}
#[test]
fn retired_lifecycle_description_probe_handles_the_stored_shapes() {
for retired in [
"a fixed lifecycle — `Backlog → Analysis → InReview → InQa → InSanitation → Done`",
"Each phase (`in_development`, `in_diagnostics`, `in_review`, `in_qa`) owns a job",
"the lifecycle is in development → in diagnostics → in review → in QA → in sanitation",
"(Analysis → Planning → Queued → Development → Diagnostics → Review → QA → Sanitation)",
"the dev/review/QA pipeline",
"the board keeps IN_REVIEW and IN_QA columns",
"comments come from the analysis, engineer, diagnostics, review and QA",
"two separate checks: review, qa",
"review then QA",
"the docs are updated on each review → merge cycle",
] {
assert!(
describes_retired_lifecycle(retired),
"must be invalidated: {retired:?}"
);
}
for kept in [
"the code reviewers check the change and one functional tester runs the product",
"every phase owns a short-lived `jobs` row",
"reviewing the diff is what the reviewer does",
"the reviewer and the analyst read the same tree",
"the model slots split inspector roles (analyst/coder/qa/reviewer/sanitation) from the rest",
"`review.rs` and `qa.rs` were merged",
] {
assert!(
!describes_retired_lifecycle(kept),
"must survive verbatim: {kept:?}"
);
}
}
async fn workspace_context_rows(
conn: &Connection,
) -> std::collections::BTreeMap<String, String> {
conn.query("SELECT role, content FROM workspace_contexts", ())
.await
.expect("read workspace_contexts")
.into_iter()
.map(|row| {
(
row.get::<Option<String>>(0)
.expect("role")
.unwrap_or_default(),
row.get::<String>(1).expect("content"),
)
})
.collect()
}
#[tokio::test]
#[expect(clippy::too_many_lines)] async fn retired_stages_are_merged_into_verification() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_with_schema(
&crate::db::store_db_path(tmp.path(), crate::db::CONSOLIDATED_DB_NAME),
"",
)
.await
.expect("open core");
run_catalog(&conn, TargetDb::Core, &MIGRATIONS[..1])
.await
.expect("baseline catalog");
let now = crate::db::now();
conn.execute(
"INSERT INTO workspaces (name, path, created_at, updated_at) VALUES ('ws', '/ws', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
for (id, phase) in [("T1", "in_review"), ("T2", "in_qa"), ("T3", "backlog")] {
conn.execute(
"INSERT INTO tickets \
(id, title, description, phase, workspace_name, created_at, updated_at) \
VALUES (?1, ?1, '', ?2, 'ws', ?3, ?3)",
params![id, phase, now.clone()],
)
.await
.unwrap();
}
for (job, kind, ticket) in [
("J1", "in_review", Some("T1")),
("J2", "in_qa", Some("T1")),
("J3", "in_qa", Some("T2")),
("J4", "analysis", Some("T3")),
] {
conn.execute(
"INSERT INTO jobs \
(id, kind, role, workspace_name, task, user_name, channel, retry_count, status, \
created_at, updated_at, ticket_id) \
VALUES (?1, ?2, 'qa', 'ws', '', 'bob', '', 0, 'launched', ?3, ?3, ?4)",
params![job, kind, now.clone(), ticket],
)
.await
.unwrap();
}
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, task) \
VALUES ('J1', 'A1', 'verifier', 0, 'done', 'task')",
(),
)
.await
.unwrap();
for (ticket, source, target) in [
("T1", "in_diagnostics", "in_review"),
("T2", "in_qa", "in_sanitation"),
("T3", "backlog", "queued"),
("T1", "in_review", "in_qa"),
] {
conn.execute(
"INSERT INTO ticket_chronicle \
(ticket_id, workspace_name, source_phase, target_phase, at) VALUES (?1, 'ws', ?2, ?3, ?4)",
params![ticket, source, target, now.clone()],
)
.await
.unwrap();
}
let arrow_lifecycle =
"lifecycle — `Backlog → Analysis → InReview → InQa → InSanitation → Done`";
let phase_list =
"Each phase (`in_development`, `in_diagnostics`, `in_review`, `in_qa`) owns a job";
let spaced_lifecycle =
"the pipeline runs development → diagnostics → review → QA → sanitation";
for (role, content) in [
(None, arrow_lifecycle),
(Some("engineer"), phase_list),
(Some("qa"), "the board keeps IN_REVIEW and IN_QA columns"),
(Some("manager"), spaced_lifecycle),
(Some("coder"), "no retired stage is named here"),
] {
conn.execute(
"INSERT INTO workspace_contexts (workspace_name, role, content, created_at) \
VALUES ('ws', ?1, ?2, ?3)",
params![role, content, now.clone()],
)
.await
.unwrap();
}
let merge_entries: Vec<Migration> = MIGRATIONS
.iter()
.filter(|m| m.id == "44" || m.id == "45")
.copied()
.collect();
run_catalog(&conn, TargetDb::Core, &merge_entries)
.await
.expect("stage-merge tail");
let phases: Vec<(String, String)> = conn
.query("SELECT id, phase FROM tickets ORDER BY id", ())
.await
.unwrap()
.into_iter()
.map(|row| (row.get::<String>(0).unwrap(), row.get::<String>(1).unwrap()))
.collect();
assert_eq!(
phases,
vec![
("T1".to_string(), "verification".to_string()),
("T2".to_string(), "verification".to_string()),
("T3".to_string(), "backlog".to_string()),
],
"the retired phases become verification; every other phase survives"
);
let remaining_jobs: Vec<String> = conn
.query("SELECT id FROM jobs ORDER BY id", ())
.await
.unwrap()
.into_iter()
.map(|row| row.get::<String>(0).unwrap())
.collect();
assert_eq!(
remaining_jobs,
vec!["J4".to_string()],
"the retired job rows must be deleted, never renamed"
);
let orphan_rosters: i64 = conn.query("SELECT COUNT(*) FROM agents", ()).await.unwrap()[0]
.get::<i64>(0)
.unwrap();
assert_eq!(
orphan_rosters, 0,
"the deleted jobs must cascade their roster away"
);
let chronicle: Vec<(String, String, String)> = conn
.query(
"SELECT ticket_id, source_phase, target_phase FROM ticket_chronicle ORDER BY ticket_id",
(),
)
.await
.unwrap()
.into_iter()
.map(|row| {
(
row.get::<String>(0).unwrap(),
row.get::<String>(1).unwrap(),
row.get::<String>(2).unwrap(),
)
})
.collect();
assert_eq!(
chronicle,
vec![
(
"T1".to_string(),
"in_diagnostics".to_string(),
"verification".to_string()
),
(
"T2".to_string(),
"verification".to_string(),
"in_sanitation".to_string()
),
(
"T3".to_string(),
"backlog".to_string(),
"queued".to_string()
),
],
"both chronicle sides must be rewritten, the retired→retired hop \
dropped, and unretired edges kept"
);
let contexts = workspace_context_rows(&conn).await;
assert!(
!contexts.contains_key(""),
"the arrow lifecycle list names a retired stage — the row must go, not be patched"
);
assert!(
!contexts.contains_key("engineer"),
"the code-span slug list names a retired stage"
);
assert!(
!contexts.contains_key("qa"),
"an upper-cased retired stage name still instructs the old arrangement"
);
assert!(
!contexts.contains_key("manager"),
"the spaced lifecycle prose names a retired stage"
);
assert_eq!(
contexts.get("coder"),
Some(&"no retired stage is named here".to_string()),
"a description naming no retired stage survives verbatim"
);
conn.execute_batch(REWRITE_RETIRED_STAGES_TO_VERIFICATION)
.await
.expect("re-run entry 44 body");
run_drop_retired_lifecycle_descriptions(&conn)
.await
.expect("re-run entry 45 body");
assert_eq!(workspace_context_rows(&conn).await, contexts);
}
}