use std::collections::HashSet;
use std::path::Path;
use anyhow::Context;
use futures_util::future::BoxFuture;
use crate::db::{Connection, params};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum TargetDb {
Core,
Logs,
}
type RustFn = for<'a> fn(&'a Connection, &'a Path) -> 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_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 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 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 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 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 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 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 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 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 DELTA_DROP_CONFIG_ROLE: &str = "DROP TABLE IF EXISTS config_role;";
const DELTA_DROP_TICKET_STAGE_JOBS: &str = "DROP TABLE IF EXISTS ticket_stage_jobs;";
const DELTA_FTS_INDEX: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_title_fts ON tickets \
USING fts (title) WITH (tokenizer = 'ngram');";
const DELTA_BOARD_ACTIVE_INDEX: &str = "CREATE INDEX IF NOT EXISTS idx_tickets_board_active ON tickets \
(is_archived, priority ASC, created_at DESC);";
const 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);";
const 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 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 DELTA_CHAT_HISTORY_TIMESTAMP: &str = "ALTER TABLE chat_history ADD COLUMN timestamp TEXT;";
const 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';";
pub(crate) const CONSOLIDATION_IMPORT_ID: &str = "consolidate_001_import_domain_stores";
pub(crate) const MIGRATIONS: &[Migration] = &[
Migration {
id: "1",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_BOARD_TABLES),
},
Migration {
id: "2",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_SESSION_TABLES),
},
Migration {
id: "3",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_WORKSPACE_TABLES),
},
Migration {
id: "4",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_USERS_TABLES),
},
Migration {
id: "5",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_CONFIG_TABLES),
},
Migration {
id: "6",
target: TargetDb::Core,
body: MigrationBody::Sql(BASELINE_CHAT_HISTORY_TABLES),
},
Migration {
id: "7",
target: TargetDb::Core,
body: MigrationBody::Sql(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(DELTA_DROP_CONFIG_ROLE),
},
Migration {
id: "11",
target: TargetDb::Core,
body: MigrationBody::Sql(DELTA_FTS_INDEX),
},
Migration {
id: "12",
target: TargetDb::Core,
body: MigrationBody::Sql(DELTA_BOARD_ACTIVE_INDEX),
},
Migration {
id: CONSOLIDATION_IMPORT_ID,
target: TargetDb::Core,
body: MigrationBody::Rust(import_domain_stores),
},
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(DELTA_DROP_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(BASELINE_LOGS_TABLES),
},
Migration {
id: "9",
target: TargetDb::Logs,
body: MigrationBody::Sql(BASELINE_LOGS_INDEXES),
},
Migration {
id: "14",
target: TargetDb::Logs,
body: MigrationBody::Sql(DELTA_GREP_TELEMETRY),
},
Migration {
id: "17",
target: TargetDb::Core,
body: MigrationBody::Sql(DELTA_ALARMS),
},
Migration {
id: "18",
target: TargetDb::Core,
body: MigrationBody::Sql(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(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(DELTA_TICKETS_WORKSPACE_PHASE_INDEX),
},
];
pub(crate) async fn run_migrations(
conn: &Connection,
db: TargetDb,
root: &Path,
) -> anyhow::Result<()> {
conn.execute(
"CREATE TABLE IF NOT EXISTS schema_migrations (\
id TEXT PRIMARY KEY,\
applied_at TEXT NOT NULL\
)",
(),
)
.await
.context("Failed to create schema_migrations tracking table")?;
let applied: HashSet<String> = conn
.query("SELECT id FROM schema_migrations", ())
.await
.context("Failed to read applied migrations")?
.into_iter()
.filter_map(|row| row.get::<String>(0).ok())
.collect();
for migration in MIGRATIONS {
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(
"INSERT INTO schema_migrations (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, root)
.await
.with_context(|| format!("Migration '{}' failed", migration.id))?;
conn.execute(
"INSERT INTO schema_migrations (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 import_domain_stores<'a>(
conn: &'a Connection,
root: &'a Path,
) -> BoxFuture<'a, anyhow::Result<()>> {
Box::pin(run_consolidation_import(conn, root))
}
fn cleanup_legacy_ticket_jobs<'a>(
conn: &'a Connection,
_root: &'a Path,
) -> BoxFuture<'a, anyhow::Result<()>> {
Box::pin(run_import_cleanup(conn))
}
fn drop_jobs_paused_frozen<'a>(
conn: &'a Connection,
_root: &'a Path,
) -> BoxFuture<'a, anyhow::Result<()>> {
Box::pin(run_drop_jobs_paused_frozen(conn))
}
async fn run_consolidation_import(conn: &Connection, root: &Path) -> anyhow::Result<()> {
let any_legacy = crate::db::DOMAIN_STORE_NAMES
.iter()
.any(|name| crate::db::legacy_store_db_path(root, name).exists());
if !any_legacy {
return Ok(());
}
crate::db::import_legacy_stores(conn, root).await
}
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(())
}
async fn run_drop_jobs_paused_frozen(conn: &Connection) -> anyhow::Result<()> {
let has = conn
.query("PRAGMA table_info(jobs)", ())
.await
.context("Failed to probe jobs.paused_frozen")?
.into_iter()
.any(|row| row.get::<String>(1).ok().as_deref() == Some("paused_frozen"));
if has {
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<'a>(
conn: &'a Connection,
_root: &'a Path,
) -> BoxFuture<'a, 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<'a>(
conn: &'a Connection,
_root: &'a Path,
) -> BoxFuture<'a, 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(())
}
#[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_%'",
(),
)
.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_%'",
(),
)
.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_%'",
(),
)
.await
.expect("read tables")
.into_iter()
.map(|row| row.get::<String>(0).expect("table name"))
.collect()
}
const LEGACY_0_4_2_BOARD_SCHEMA: &str = r#"
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 INDEX IF NOT EXISTS idx_ticket_comments_ticket_id ON ticket_comments(ticket_id);
CREATE TABLE IF NOT EXISTS ticket_counters (
workspace_name TEXT PRIMARY KEY,
next_id INTEGER NOT NULL DEFAULT 1
);
"#;
const LEGACY_0_4_2_SESSIONS_SCHEMA: &str = r#"
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 INDEX IF NOT EXISTS idx_sessions_agent_id ON sessions(agent_id, id);
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 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 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 UNIQUE INDEX IF NOT EXISTS idx_agents_anchor ON agents(agent_id) WHERE job_id IS NULL;
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 INDEX IF NOT EXISTS idx_pending_jobs_agent_created ON pending_jobs(target_agent_id, created_at);
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 LEGACY_0_4_2_WORKSPACES_SCHEMA: &str = r#"
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 UNIQUE INDEX IF NOT EXISTS workspace_contexts_null_role ON workspace_contexts(workspace_name) WHERE role IS NULL;
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 LEGACY_0_4_2_USERS_SCHEMA: &str = r#"
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 LEGACY_0_4_2_CONFIG_SCHEMA: &str = r#"
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 LEGACY_0_4_2_CHAT_HISTORY_SCHEMA: &str = r#"
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
);
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 LEGACY_0_4_2_LOGS_SCHEMA: &str = r#"
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 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 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 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 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
);
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);
"#;
#[tokio::test]
async fn fresh_install_converges_to_current_shape() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_consolidated_store(tmp.path())
.await
.expect("fresh consolidated store");
for table in [
"tickets",
"ticket_comments",
"ticket_counters",
"sessions",
"session_metadata",
"jobs",
"agents",
"pending_jobs",
"research_jobs",
"workspaces",
"workspace_contexts",
"editor_tabs",
"users",
"user_channels",
"config_kv",
"config_model_routing",
"chat_history",
"ticket_chronicle",
] {
assert!(
crate::db::table_exists(&conn, table).await.unwrap(),
"missing {table}"
);
}
for table in [
"config_role",
"ticket_jobs",
"ticket_stage_jobs",
"user_roles",
] {
assert!(
!crate::db::table_exists(&conn, table).await.unwrap(),
"{table} must be gone"
);
}
let jobs_cols = column_names(&conn, "jobs").await;
assert!(jobs_cols.contains(&"ticket_id".to_string()));
assert!(!jobs_cols.contains(&"paused_frozen".to_string()));
let tickets_cols = column_names(&conn, "tickets").await;
assert!(!tickets_cols.contains(&"pipeline_reservation".to_string()));
assert!(!tickets_cols.contains(&"assigned_to".to_string()));
for id in [
"001_drop_ticket_pipeline_reservation",
"002_drop_ticket_assigned_to",
"003_reset_nonterminal_tickets",
"003_drop_ticket_stage_jobs_round",
"005_drop_ticket_stage_jobs_phase",
"006_rename_ticket_stage_jobs_to_ticket_jobs",
"consolidate_002_drop_ticket_jobs_stage",
"007_jobs_paused_frozen",
"consolidate_003_jobs_ticket_id",
"consolidate_005_drop_ticket_jobs",
"consolidate_006_drop_jobs_paused_frozen",
"consolidate_007_jobs_phase_ticket_index",
"consolidate_008_create_ticket_chronicle",
CONSOLIDATION_IMPORT_ID,
"1",
"2",
"3",
"4",
"5",
"6",
"7",
"10",
"11",
"12",
"13",
"15",
"16",
"19",
"22",
] {
assert!(
applied_ids(&conn).await.contains(&id.to_string()),
"missing applied id {id}"
);
}
}
#[tokio::test]
async fn migration_19_existing_install_seeds_finished_and_normalizes_roles() {
let tmp = tempfile::TempDir::new().unwrap();
let conn = crate::db::open_consolidated_store(tmp.path())
.await
.expect("fresh consolidated store");
conn.execute(
"INSERT INTO users (name, permissions, selected_role) VALUES ('admin', 'full', 'manager')",
(),
)
.await
.expect("seed full admin");
conn.execute(
"INSERT INTO users (name, permissions, selected_role) VALUES ('bob', NULL, 'analyst')",
(),
)
.await
.expect("seed non-full out-of-pool user");
conn.execute(
"INSERT INTO users (name, permissions, selected_role) VALUES ('carol', 'full', NULL)",
(),
)
.await
.expect("seed full admin with NULL selected_role");
run_drop_user_roles_and_seed_onboarding(&conn)
.await
.expect("re-run migration 19 body");
let state: String = conn
.query(
&format!(
"SELECT value FROM config_kv WHERE key = '{}'",
crate::config::CONFIG_KEY_ONBOARDING_STATE,
),
(),
)
.await
.expect("read onboarding_state")
.into_iter()
.next()
.and_then(|row| row.get::<String>(0).ok())
.expect("onboarding_state row");
assert_eq!(state, crate::config::OnboardingState::Finished.as_str());
let admin_sel: String = conn
.query("SELECT selected_role FROM users WHERE name = 'admin'", ())
.await
.expect("read admin role")
.into_iter()
.next()
.and_then(|row| row.get::<String>(0).ok())
.expect("admin row");
assert_eq!(admin_sel, "manager");
let bob_sel: String = conn
.query("SELECT selected_role FROM users WHERE name = 'bob'", ())
.await
.expect("read bob role")
.into_iter()
.next()
.and_then(|row| row.get::<String>(0).ok())
.expect("bob row");
assert_eq!(bob_sel, "assistant");
let carol_sel: String = conn
.query("SELECT selected_role FROM users WHERE name = 'carol'", ())
.await
.expect("read carol role")
.into_iter()
.next()
.and_then(|row| row.get::<String>(0).ok())
.expect("carol row");
assert_eq!(carol_sel, "assistant");
assert!(!crate::db::table_exists(&conn, "user_roles").await.unwrap());
}
#[tokio::test]
async fn fresh_logs_install_runs_catalog() {
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, root)
.await
.expect("run logs catalog");
for table in [
"logs",
"tool_calls",
"llm_requests",
"grep_telemetry",
"schema_migrations",
] {
assert!(
crate::db::table_exists(&conn, table).await.unwrap(),
"missing {table}"
);
}
let applied = applied_ids(&conn).await;
for id in ["8", "9", "14"] {
assert!(applied.contains(&id.to_string()), "missing applied id {id}");
}
}
#[tokio::test]
async fn golden_upgrade_from_0_4_2_per_store() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let now = crate::db::now();
let board = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "board"),
LEGACY_0_4_2_BOARD_SCHEMA,
)
.await
.expect("legacy board");
board
.execute(
"INSERT INTO tickets (id, title, description, phase, assigned_to, workspace_name, \
created_at, updated_at, pipeline_reservation, priority) \
VALUES ('T1', 'Done ticket', 'd', 'done', 'eng', 'ws', ?1, ?1, 0, 1)",
params![now.clone()],
)
.await
.unwrap();
board
.execute(
"INSERT INTO tickets (id, title, description, phase, assigned_to, workspace_name, \
created_at, updated_at, pipeline_reservation, priority) \
VALUES ('T2', 'Wip ticket', 'd', 'in_development', 'eng', 'ws', ?1, ?1, 0, 1)",
params![now.clone()],
)
.await
.unwrap();
board
.execute(
"INSERT INTO ticket_comments (id, ticket_id, role, content, created_at) \
VALUES ('C1', 'T1', 'manager', 'ship it', ?1)",
params![now.clone()],
)
.await
.unwrap();
let sessions = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "sessions"),
LEGACY_0_4_2_SESSIONS_SCHEMA,
)
.await
.expect("legacy sessions");
sessions
.execute(
"INSERT INTO sessions (agent_id, role, content, created_at) \
VALUES ('sess1', 'assistant', 'hello', ?1)",
params![now.clone()],
)
.await
.unwrap();
sessions
.execute(
"INSERT INTO session_metadata (agent_id, last_activity, role, message_count) \
VALUES ('sess1', ?1, 'assistant', 1)",
params![now.clone()],
)
.await
.unwrap();
sessions
.execute(
"INSERT INTO jobs (id, kind, role, workspace_name, task, user_name, channel, \
retry_count, status, created_at, updated_at) \
VALUES ('job_research', 'research', '', 'ws', '', '', '', 0, 'launched', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
sessions
.execute(
"INSERT INTO jobs (id, kind, role, workspace_name, task, user_name, channel, \
retry_count, status, created_at, updated_at) \
VALUES ('job_stage', 'ticket_stage', '', 'ws', '', '', '', 0, 'launched', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
sessions
.execute(
"INSERT INTO ticket_stage_jobs (id, ticket_id, stage, phase, round) \
VALUES ('job_stage', 'T1', 'implementation', 'in_development', 1)",
(),
)
.await
.unwrap();
let workspaces = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "workspaces"),
LEGACY_0_4_2_WORKSPACES_SCHEMA,
)
.await
.expect("legacy workspaces");
workspaces
.execute(
"INSERT INTO workspaces (name, path, created_at, updated_at) \
VALUES ('ws', '/ws', ?1, ?1)",
params![now.clone()],
)
.await
.unwrap();
let users = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "users"),
LEGACY_0_4_2_USERS_SCHEMA,
)
.await
.expect("legacy users");
users
.execute(
"INSERT INTO users (name, selected_workspace) VALUES ('alice', 'ws')",
(),
)
.await
.unwrap();
let config = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "config"),
LEGACY_0_4_2_CONFIG_SCHEMA,
)
.await
.expect("legacy config");
config
.execute("INSERT INTO config_kv (key, value) VALUES ('k', 'v')", ())
.await
.unwrap();
config
.execute(
"INSERT INTO config_role (role, model) VALUES ('manager', 'm')",
(),
)
.await
.unwrap();
let chat_history = crate::db::open_with_schema(
&crate::db::legacy_store_db_path(root, "chat_history"),
LEGACY_0_4_2_CHAT_HISTORY_SCHEMA,
)
.await
.expect("legacy chat_history");
chat_history
.execute(
"INSERT INTO chat_history (message_id, user_name, direction, content, workspace) \
VALUES ('m1', 'alice', 'in', 'hi', 'ws')",
(),
)
.await
.unwrap();
drop(board);
drop(sessions);
drop(workspaces);
drop(users);
drop(config);
drop(chat_history);
let conn = crate::db::open_consolidated_store(root)
.await
.expect("consolidate 0.4.2 install");
assert_eq!(
conn.query("SELECT count(*) FROM tickets WHERE id='T1'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"ticket survives"
);
assert_eq!(
conn.query("SELECT count(*) FROM ticket_comments WHERE id='C1'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"ticket comment survives"
);
assert_eq!(
conn.query("SELECT count(*) FROM sessions WHERE agent_id='sess1'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"session survives"
);
assert_eq!(
conn.query(
"SELECT count(*) FROM session_metadata WHERE agent_id='sess1'",
()
)
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"session_metadata survives"
);
assert_eq!(
conn.query("SELECT count(*) FROM workspaces WHERE name='ws'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"workspace survives"
);
assert_eq!(
conn.query("SELECT count(*) FROM users WHERE name='alice'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"user survives"
);
assert_eq!(
conn.query("SELECT count(*) FROM config_kv WHERE key='k'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"config_kv survives"
);
assert_eq!(
conn.query(
"SELECT count(*) FROM chat_history WHERE message_id='m1'",
()
)
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"chat_history survives"
);
assert_eq!(
conn.query(
"SELECT seq FROM sqlite_sequence WHERE name='chat_history'",
()
)
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"AUTOINCREMENT watermark preserved after consolidation"
);
assert_eq!(
conn.query("SELECT count(*) FROM jobs WHERE kind='research'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
1,
"research job survives"
);
assert_eq!(
conn.query("SELECT phase FROM tickets WHERE id='T2'", ())
.await
.unwrap()[0]
.get::<String>(0)
.unwrap(),
"backlog",
"non-terminal legacy ticket reset to backlog"
);
assert_eq!(
conn.query("SELECT count(*) FROM jobs WHERE kind='ticket_stage'", ())
.await
.unwrap()[0]
.get::<i64>(0)
.unwrap(),
0,
"pre-rework ticket_stage jobs dropped after import"
);
assert!(!crate::db::table_exists(&conn, "config_role").await.unwrap());
assert!(!crate::db::table_exists(&conn, "ticket_jobs").await.unwrap());
assert!(
!crate::db::table_exists(&conn, "ticket_stage_jobs")
.await
.unwrap()
);
assert_eq!(
conn.query("PRAGMA foreign_key_check", ())
.await
.unwrap()
.len(),
0,
"no orphan rows after the import"
);
}
#[tokio::test]
async fn current_consolidated_db_reopen_converges() {
const OLD_DEPLOYED_IDS: &[&str] = &[
"001_session_token_length",
"002_session_message_count",
"003_drop_ticket_stage_jobs_round",
"004_reset_implementation_jobs",
"005_drop_ticket_stage_jobs_phase",
"006_rename_ticket_stage_jobs_to_ticket_jobs",
"consolidate_002_drop_ticket_jobs_stage",
"consolidate_003_jobs_ticket_id",
"consolidate_004_discard_old_ticket_jobs",
"consolidate_005_drop_ticket_jobs",
"consolidate_006_drop_jobs_paused_frozen",
"consolidate_007_jobs_phase_ticket_index",
"001_drop_ticket_pipeline_reservation",
"002_drop_ticket_assigned_to",
"003_reset_nonterminal_tickets",
"004_drop_tickets_review_base_count",
"consolidate_008_create_ticket_chronicle",
CONSOLIDATION_IMPORT_ID,
];
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_migrations(&conn, TargetDb::Core, root)
.await
.expect("fresh catalog");
conn.execute("DELETE FROM schema_migrations", ())
.await
.unwrap();
conn.execute("ALTER TABLE chat_history DROP COLUMN timestamp;", ())
.await
.unwrap();
for id in OLD_DEPLOYED_IDS {
conn.execute(
"INSERT INTO schema_migrations (id, applied_at) VALUES (?1, ?2)",
params![id, crate::db::now()],
)
.await
.unwrap();
}
run_migrations(&conn, TargetDb::Core, root)
.await
.expect("reopen catalog");
assert!(
!crate::db::table_exists(&conn, "ticket_stage_jobs")
.await
.unwrap()
);
assert!(!crate::db::table_exists(&conn, "config_role").await.unwrap());
let jobs_cols = column_names(&conn, "jobs").await;
assert!(jobs_cols.contains(&"ticket_id".to_string()));
assert!(!jobs_cols.contains(&"paused_frozen".to_string()));
let tickets_cols = column_names(&conn, "tickets").await;
assert!(!tickets_cols.contains(&"pipeline_reservation".to_string()));
assert!(!tickets_cols.contains(&"assigned_to".to_string()));
let applied = applied_ids(&conn).await;
for id in ["15", "16"] {
assert!(applied.contains(&id.to_string()), "missing applied id {id}");
}
}
#[tokio::test]
async fn migration_21_rewrites_analysis_verdicts() {
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_migrations(&conn, TargetDb::Core, root)
.await
.expect("fresh catalog");
let now = crate::db::now();
conn.execute(
"INSERT INTO jobs (id, kind, role, workspace_name, created_at, updated_at) \
VALUES ('J1', 'analysis', 'analyst', 'ws', ?1, ?2)",
params![now.clone(), now.clone()],
)
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES ('J1', 'A1', 'analyst', 0, 'done', ?1, 'task')",
params!["{\"verdict\":{\"score\":5,\"issues\":[\"issue a\",\"issue b\"]}}"],
)
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES ('J1', 'A2', 'analyst', 1, 'done', ?1, 'task')",
params!["{\"verdict\":{\"score\":8,\"issues\":[\"note\"]}}"],
)
.await
.unwrap();
conn.execute(
"INSERT INTO agents (job_id, agent_id, kind, idx, status, outcome, task) \
VALUES ('J1', 'A3', 'verifier', 2, 'done', ?1, 'task')",
params!["{\"verdict\":{\"score\":9,\"issues\":[]}}"],
)
.await
.unwrap();
run_rewrite_analysis_verdicts(&conn).await.unwrap();
async fn outcome(conn: &Connection, agent_id: &str) -> String {
conn.query(
"SELECT outcome FROM agents WHERE agent_id = ?1",
params![agent_id],
)
.await
.unwrap()
.into_iter()
.next()
.unwrap()
.get::<String>(0)
.unwrap()
}
assert_eq!(
outcome(&conn, "A1")
.await
.parse::<serde_json::Value>()
.unwrap(),
serde_json::json!({
"verdict": { "issues": [
{ "text": "issue a", "grade": "blocker" },
{ "text": "issue b", "grade": "blocker" },
] }
})
);
assert_eq!(
outcome(&conn, "A2")
.await
.parse::<serde_json::Value>()
.unwrap(),
serde_json::json!({
"verdict": { "issues": [ { "text": "note", "grade": "minor" } ] }
})
);
assert_eq!(
outcome(&conn, "A3").await,
"{\"verdict\":{\"score\":9,\"issues\":[]}}"
);
}
#[tokio::test]
async fn catalog_baseline_matches_verbatim_0_4_2_schema() {
let tmp = tempfile::TempDir::new().unwrap();
let root = tmp.path();
let core = crate::db::open_with_schema(&root.join("baseline-core.db"), "")
.await
.expect("open baseline core");
for sql in [
BASELINE_BOARD_TABLES,
BASELINE_SESSION_TABLES,
BASELINE_WORKSPACE_TABLES,
BASELINE_USERS_TABLES,
BASELINE_CONFIG_TABLES,
BASELINE_CHAT_HISTORY_TABLES,
BASELINE_CORE_INDEXES,
] {
core.execute_batch(sql).await.expect("apply baseline core");
}
let logs = crate::db::open_with_schema(&root.join("baseline-logs.db"), "")
.await
.expect("open baseline logs");
logs.execute_batch(BASELINE_LOGS_TABLES)
.await
.expect("apply baseline logs tables");
logs.execute_batch(BASELINE_LOGS_INDEXES)
.await
.expect("apply baseline logs indexes");
let mut legacy_core_tables: std::collections::HashSet<String> =
std::collections::HashSet::new();
let mut legacy_core_indexes: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
let mut legacy_core_table_defs: std::collections::BTreeMap<String, String> =
std::collections::BTreeMap::new();
for (store, schema) in [
("board", LEGACY_0_4_2_BOARD_SCHEMA),
("sessions", LEGACY_0_4_2_SESSIONS_SCHEMA),
("workspaces", LEGACY_0_4_2_WORKSPACES_SCHEMA),
("users", LEGACY_0_4_2_USERS_SCHEMA),
("config", LEGACY_0_4_2_CONFIG_SCHEMA),
("chat_history", LEGACY_0_4_2_CHAT_HISTORY_SCHEMA),
] {
let legacy =
crate::db::open_with_schema(&crate::db::legacy_store_db_path(root, store), schema)
.await
.expect("open legacy store");
for table in table_names(&legacy).await {
assert!(
legacy_core_tables.insert(table.clone()),
"duplicate table {table} across legacy 0.4.2 stores"
);
}
for (name, sql) in index_defs(&legacy).await {
assert!(
legacy_core_indexes.insert(name.clone(), sql).is_none(),
"duplicate index {name} across legacy 0.4.2 stores"
);
}
for t in table_defs(&legacy).await {
assert!(
legacy_core_table_defs.insert(t.0, t.1).is_none(),
"duplicate table DDL across legacy 0.4.2 stores"
);
}
}
let core_tables: std::collections::HashSet<String> =
table_names(&core).await.into_iter().collect();
assert_eq!(
core_tables, legacy_core_tables,
"baseline core tables diverge from the 0.4.2 per-store schemas",
);
assert_eq!(
index_defs(&core).await,
legacy_core_indexes,
"baseline core indexes diverge from the 0.4.2 per-store schemas",
);
assert_eq!(
table_defs(&core).await,
legacy_core_table_defs,
"baseline core table DDL diverges from the 0.4.2 per-store schemas",
);
let legacy_logs =
crate::db::open_with_schema(&root.join("legacy-logs.db"), LEGACY_0_4_2_LOGS_SCHEMA)
.await
.expect("open legacy logs");
let legacy_logs_tables: std::collections::HashSet<String> =
table_names(&legacy_logs).await.into_iter().collect();
let logs_tables: std::collections::HashSet<String> =
table_names(&logs).await.into_iter().collect();
assert_eq!(
logs_tables, legacy_logs_tables,
"baseline logs tables diverge from the 0.4.2 logs schema",
);
assert_eq!(
index_defs(&logs).await,
index_defs(&legacy_logs).await,
"baseline logs indexes diverge from the 0.4.2 logs schema",
);
assert_eq!(
table_defs(&logs).await,
table_defs(&legacy_logs).await,
"baseline logs table DDL diverges from the 0.4.2 logs schema",
);
}
#[tokio::test]
async fn logs_upgrade_from_0_4_2_adds_grep_telemetry() {
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),
LEGACY_0_4_2_LOGS_SCHEMA,
)
.await
.expect("open legacy logs");
conn.execute(
"INSERT INTO logs (timestamp, level, target, message) \
VALUES (?1, 'INFO', 'test', 'hello')",
params![crate::db::now()],
)
.await
.unwrap();
run_migrations(&conn, TargetDb::Logs, root)
.await
.expect("upgrade logs catalog");
assert!(
crate::db::table_exists(&conn, "grep_telemetry")
.await
.unwrap()
);
assert!(
crate::db::table_exists(&conn, "llm_requests")
.await
.unwrap()
);
let count = conn
.query("SELECT COUNT(*) FROM logs WHERE message = 'hello'", ())
.await
.expect("count logs")
.into_iter()
.next()
.map(|row| row.get::<i64>(0).expect("count"))
.unwrap_or(0);
assert_eq!(count, 1, "existing logs rows must survive the logs upgrade");
let applied = applied_ids(&conn).await;
for id in ["8", "9", "14"] {
assert!(applied.contains(&id.to_string()), "missing applied id {id}");
}
}
}