use super::*;
use mj_core::relay::RELAY_EVENT_GENESIS_DIGEST;
pub(super) fn migrate_to_baseline(connection: &Connection) -> Result<()> {
let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
if version < 2 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN resource_allocation TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (2, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 2;
COMMIT;",
)?;
}
if version < 3 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN last_checkpoint_error TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (3, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 3;
COMMIT;",
)?;
}
if version < 4 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN project_directory BLOB;
INSERT INTO schema_migrations(version, applied_at)
VALUES (4, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 4;
COMMIT;",
)?;
}
if version < 5 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE session_targets RENAME TO session_targets_v4;
CREATE TABLE session_targets (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
kind TEXT NOT NULL CHECK(kind IN ('local-bare','local-podman','apple-container','aws-ec2','ssh-bare','ssh-podman')),
host TEXT,
resource_id TEXT,
address TEXT,
workspace BLOB,
worker_id TEXT,
CHECK(
(kind = 'local-bare' AND workspace IS NOT NULL
AND host IS NULL AND resource_id IS NULL AND address IS NULL AND worker_id IS NULL)
OR (kind IN ('local-podman','apple-container') AND resource_id IS NOT NULL
AND host IS NULL AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'aws-ec2' AND resource_id IS NOT NULL
AND host IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'ssh-bare' AND host IS NOT NULL AND workspace IS NOT NULL
AND resource_id IS NULL AND address IS NULL)
OR (kind = 'ssh-podman' AND host IS NOT NULL AND resource_id IS NOT NULL
AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
)
) STRICT;
INSERT INTO session_targets
SELECT * FROM session_targets_v4;
DROP TABLE session_targets_v4;
INSERT INTO schema_migrations(version, applied_at)
VALUES (5, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 5;
COMMIT;",
)?;
}
if version < 6 {
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
ALTER TABLE session_checkpoints
RENAME COLUMN event_sequence TO event_frontier;
ALTER TABLE prompt_history
RENAME COLUMN event_sequence TO event_ordinal;
ALTER TABLE sessions ADD COLUMN detached_after_event_ordinal INTEGER NOT NULL
DEFAULT 0 CHECK(detached_after_event_ordinal >= 0);
ALTER TABLE sessions ADD COLUMN managed_worktree TEXT;
CREATE TABLE materialized_sessions (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
applied_event_ordinal INTEGER NOT NULL DEFAULT 0 CHECK(applied_event_ordinal >= 0),
applied_event_digest TEXT NOT NULL
DEFAULT '{RELAY_EVENT_GENESIS_DIGEST}'
CHECK(length(applied_event_digest) = 64
AND applied_event_digest NOT GLOB '*[^0-9a-f]*'),
last_activity_at_ms INTEGER,
execution_state TEXT NOT NULL DEFAULT 'idle'
CHECK(execution_state IN ('idle','running','closing','closed')),
running_started_at_ms INTEGER,
session_title TEXT CHECK(session_title IS NULL OR length(trim(session_title)) > 0),
configuration_json TEXT NOT NULL DEFAULT '{{}}',
CHECK(
(execution_state = 'running' AND running_started_at_ms IS NOT NULL)
OR (execution_state != 'running' AND running_started_at_ms IS NULL)
)
) STRICT;
CREATE TABLE materialized_transcript_items (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
stable_id TEXT NOT NULL CHECK(length(trim(stable_id)) > 0),
position INTEGER NOT NULL CHECK(position > 0),
latest_content_event_ordinal INTEGER
CHECK(latest_content_event_ordinal IS NULL
OR latest_content_event_ordinal >= position),
created_at_ms INTEGER NOT NULL,
last_changed_at_ms INTEGER NOT NULL CHECK(last_changed_at_ms >= created_at_ms),
body_json TEXT NOT NULL,
PRIMARY KEY(session_id, stable_id)
) STRICT;
CREATE INDEX materialized_transcript_position
ON materialized_transcript_items(session_id, position, stable_id);
CREATE TABLE materialized_queued_prompts (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
ordinal INTEGER NOT NULL CHECK(ordinal >= 0),
command_id TEXT NOT NULL CHECK(length(trim(command_id)) > 0),
content_json TEXT NOT NULL,
queued_at_ms INTEGER NOT NULL,
PRIMARY KEY(session_id, ordinal),
UNIQUE(session_id, command_id)
) STRICT;
INSERT INTO materialized_sessions(session_id)
SELECT session_id FROM sessions;
INSERT INTO schema_migrations(version, applied_at)
VALUES (6, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 6;
COMMIT;",
))?;
}
ensure_managed_worktree_column(connection)?;
if version < 7 {
ensure_relay_projection_schema(connection)?;
migrate_destroying_session_state(connection)?;
}
ensure_projection_digest_column(connection)?;
ensure_session_draft_input_column(connection)?;
if version < 8 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE materialized_queued_prompts
ADD COLUMN kind_json TEXT NOT NULL DEFAULT '\"prompt\"';
INSERT INTO schema_migrations(version, applied_at)
VALUES (8, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 8;
COMMIT;",
)?;
}
if version < 9 {
migrate_grok_harness_kind(connection)?;
}
ensure_session_container_override_columns(connection)?;
ensure_session_mount_read_only_column(connection)?;
ensure_materialized_elicitation_column(connection)?;
connection.execute_batch(
"CREATE TABLE IF NOT EXISTS profile_config_cache (
profile TEXT NOT NULL, model TEXT NOT NULL, fingerprint TEXT NOT NULL,
observed_at INTEGER NOT NULL, body TEXT NOT NULL, PRIMARY KEY(profile, model));
CREATE TABLE IF NOT EXISTS api_config_results (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
command_id TEXT NOT NULL, error TEXT, PRIMARY KEY(session_id, command_id));
CREATE TABLE IF NOT EXISTS session_turn_usage (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
command_id TEXT NOT NULL, completed_ordinal INTEGER NOT NULL, turn_start_position INTEGER,
body TEXT NOT NULL, PRIMARY KEY(session_id, command_id));
CREATE INDEX IF NOT EXISTS session_turn_usage_order ON session_turn_usage(session_id, completed_ordinal);
CREATE TABLE IF NOT EXISTS session_provider_cost (
session_id TEXT PRIMARY KEY REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
body TEXT NOT NULL);",
)?;
if version < 10 {
migrate_stopped_session_state(connection)?;
}
if version < 11 {
migrate_deepseek_harness_kind(connection)?;
}
if version < 12 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions
RENAME COLUMN detached_after_event_ordinal TO viewed_through_event_ordinal;
UPDATE sessions
SET target_template_id = 'localhost'
WHERE target_template_id = 'raw-localhost';
INSERT INTO schema_migrations(version, applied_at)
VALUES (12, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 12;
COMMIT;",
)?;
}
if version < 13 {
connection.execute_batch(
"BEGIN IMMEDIATE;
UPDATE sessions
SET state = 'error'
WHERE state = 'lost'
AND EXISTS(
SELECT 1
FROM session_checkpoints
WHERE session_checkpoints.session_id = sessions.session_id
);
INSERT INTO schema_migrations(version, applied_at)
VALUES (13, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 13;
COMMIT;",
)?;
}
if version < 14 {
ensure_workspace_schema(connection)?;
connection.execute_batch(
"BEGIN IMMEDIATE;
INSERT OR IGNORE INTO schema_migrations(version, applied_at)
VALUES (14, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 14;
COMMIT;",
)?;
}
if version < 15 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE host_container_sizes (
host TEXT PRIMARY KEY CHECK(length(trim(host)) > 0),
cpus INTEGER NOT NULL CHECK(cpus > 0),
memory_bytes INTEGER NOT NULL CHECK(memory_bytes > 0)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (15, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 15;
COMMIT;",
)?;
}
if version < 16 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE second_opinion_defaults (
workspace_id TEXT NOT NULL CHECK(length(trim(workspace_id)) > 0),
profile_id TEXT NOT NULL CHECK(length(trim(profile_id)) > 0),
model TEXT NOT NULL,
effort TEXT NOT NULL,
PRIMARY KEY (workspace_id, profile_id, model)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (16, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 16;
COMMIT;",
)?;
}
if version < 17 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE second_opinion_reviews (
session_id TEXT PRIMARY KEY
REFERENCES sessions(session_id) ON DELETE CASCADE,
workflow TEXT NOT NULL,
generation INTEGER NOT NULL CHECK(generation >= 0),
context_baseline INTEGER NOT NULL CHECK(context_baseline >= 0),
native_lost INTEGER NOT NULL CHECK(native_lost IN (0, 1))
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (17, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 17;
COMMIT;",
)?;
}
if version < 18 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE second_opinion_reviews
ADD COLUMN reviewer_transcript TEXT NOT NULL DEFAULT '[]';
INSERT INTO schema_migrations(version, applied_at)
VALUES (18, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 18;
COMMIT;",
)?;
}
if version < 19 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE turn_review_settings (
workspace_id TEXT PRIMARY KEY
CHECK(length(trim(workspace_id)) > 0),
auto_review INTEGER NOT NULL CHECK(auto_review IN (0, 1)),
tier TEXT NOT NULL CHECK(tier IN ('quick', 'extended'))
) STRICT;
CREATE TABLE turn_review_state (
session_id TEXT PRIMARY KEY
REFERENCES sessions(session_id) ON DELETE CASCADE,
baselines TEXT NOT NULL,
reviewed_through_ordinal INTEGER NOT NULL
CHECK(reviewed_through_ordinal >= 0),
prior_review TEXT,
active TEXT
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (19, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 19;
COMMIT;",
)?;
}
if version < 20 {
let target_table_exists: bool = connection.query_row(
"SELECT EXISTS(
SELECT 1 FROM sqlite_master
WHERE type = 'table' AND name = 'session_targets'
)",
[],
|row| row.get(0),
)?;
let rebuild = if target_table_exists {
"ALTER TABLE session_targets RENAME TO session_targets_v19;"
} else {
""
};
let copy = if target_table_exists {
"INSERT INTO session_targets SELECT * FROM session_targets_v19;
DROP TABLE session_targets_v19;"
} else {
""
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{rebuild}
CREATE TABLE session_targets (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
kind TEXT NOT NULL CHECK(kind IN ('local-bare','local-podman','local-docker','apple-container','aws-ec2','ssh-bare','ssh-podman')),
host TEXT,
resource_id TEXT,
address TEXT,
workspace BLOB,
worker_id TEXT,
CHECK(
(kind = 'local-bare' AND workspace IS NOT NULL
AND host IS NULL AND resource_id IS NULL AND address IS NULL AND worker_id IS NULL)
OR (kind IN ('local-podman','local-docker','apple-container') AND resource_id IS NOT NULL
AND host IS NULL AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'aws-ec2' AND resource_id IS NOT NULL
AND host IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'ssh-bare' AND host IS NOT NULL AND workspace IS NOT NULL
AND resource_id IS NULL AND address IS NULL)
OR (kind = 'ssh-podman' AND host IS NOT NULL AND resource_id IS NOT NULL
AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
)
) STRICT;
{copy}
INSERT INTO schema_migrations(version, applied_at)
VALUES (20, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 20;
COMMIT;"
))?;
}
if version < 21 {
connection.execute_batch(
"BEGIN IMMEDIATE;
DROP TABLE IF EXISTS turn_review_settings;
INSERT INTO schema_migrations(version, applied_at)
VALUES (21, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 21;
COMMIT;",
)?;
}
if version < 22 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE session_targets ADD COLUMN workspace_storage TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (22, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 22;
COMMIT;",
)?;
}
if version < 23 {
let add_pending_forward =
if table_has_column(connection, "turn_review_state", "pending_forward")? {
""
} else {
"ALTER TABLE turn_review_state ADD COLUMN pending_forward TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_pending_forward}
INSERT INTO schema_migrations(version, applied_at)
VALUES (23, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 23;
COMMIT;"
))?;
}
if version < 24 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE session_targets RENAME TO session_targets_v23;
CREATE TABLE session_targets (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
kind TEXT NOT NULL CHECK(kind IN ('local-bare','local-podman','local-docker','apple-container','aws-ec2','ssh-bare','ssh-podman','ssh-docker')),
host TEXT,
resource_id TEXT,
address TEXT,
workspace BLOB,
worker_id TEXT,
workspace_storage TEXT,
CHECK(
(kind = 'local-bare' AND workspace IS NOT NULL
AND host IS NULL AND resource_id IS NULL AND address IS NULL AND worker_id IS NULL)
OR (kind IN ('local-podman','local-docker','apple-container') AND resource_id IS NOT NULL
AND host IS NULL AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'aws-ec2' AND resource_id IS NOT NULL
AND host IS NULL AND workspace IS NULL AND worker_id IS NULL)
OR (kind = 'ssh-bare' AND host IS NOT NULL AND workspace IS NOT NULL
AND resource_id IS NULL AND address IS NULL)
OR (kind IN ('ssh-podman','ssh-docker') AND host IS NOT NULL AND resource_id IS NOT NULL
AND address IS NULL AND workspace IS NULL AND worker_id IS NULL)
)
) STRICT;
INSERT INTO session_targets SELECT * FROM session_targets_v23;
DROP TABLE session_targets_v23;
INSERT INTO schema_migrations(version, applied_at)
VALUES (24, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 24;
COMMIT;",
)?;
}
if version < 25 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE workspace_pane_sizes (
workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
sessions TEXT NOT NULL CHECK(sessions IN ('minimized', 'standard', 'maximized')),
targets TEXT NOT NULL CHECK(targets IN ('minimized', 'standard', 'maximized')),
quota TEXT NOT NULL CHECK(quota IN ('minimized', 'standard', 'maximized')),
CHECK((sessions = 'maximized') + (targets = 'maximized') + (quota = 'maximized') <= 1)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (25, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 25;
COMMIT;",
)?;
}
if version < 26 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE session_moves (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
operation_id TEXT NOT NULL UNIQUE,
operation_json TEXT NOT NULL CHECK(json_valid(operation_json))
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (26, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 26;
COMMIT;",
)?;
}
if version < 27 {
migrate_muse_harness_kind(connection)?;
}
if version < 28 {
migrate_turn_outcome_columns(connection)?;
}
if version < 29 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN create_managed_worktree INTEGER
CHECK(create_managed_worktree IN (0, 1));
INSERT INTO schema_migrations(version, applied_at)
VALUES (29, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 29;
COMMIT;",
)?;
}
if version < 30 {
migrate_compatibility_metadata(connection)?;
}
if version < 31 {
migrate_subagent_sessions(connection)?;
}
if version < 32 {
migrate_zcode_harness_kind(connection)?;
}
if version < 33 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN mjolnir_subagents INTEGER
CHECK(mjolnir_subagents IN (0, 1));
INSERT INTO schema_migrations(version, applied_at)
VALUES (33, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 33;
COMMIT;",
)?;
}
let recorded: Option<i64> =
connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
row.get(0)
})?;
if recorded != Some(33) {
bail!(
"Mjolnir database migration ledger {:?} does not match schema {}",
recorded,
33
);
}
ensure_client_session_state_schema(connection)?;
ensure_api_events_schema(connection)?;
Ok(())
}
fn migrate_compatibility_metadata(connection: &Connection) -> Result<()> {
let transaction = connection.unchecked_transaction()?;
transaction.execute_batch(
"CREATE TABLE schema_compatibility (
singleton INTEGER PRIMARY KEY CHECK(singleton = 1),
minimum_compatible_version INTEGER NOT NULL CHECK(minimum_compatible_version >= 30)
) STRICT;
INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, 30);
INSERT INTO schema_migrations(version, applied_at)
VALUES (30, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 30;",
)?;
transaction.commit()?;
Ok(())
}
fn migrate_subagent_sessions(connection: &Connection) -> Result<()> {
let transaction = connection.unchecked_transaction()?;
transaction.execute_batch(
"CREATE TABLE IF NOT EXISTS subagent_sessions (
child_session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
parent_session_id TEXT NOT NULL REFERENCES sessions(session_id),
request_key TEXT NOT NULL,
record_json TEXT NOT NULL CHECK(json_valid(record_json)),
CHECK(child_session_id <> parent_session_id),
UNIQUE(parent_session_id, request_key)
) STRICT;
CREATE INDEX IF NOT EXISTS subagent_sessions_parent
ON subagent_sessions(parent_session_id, child_session_id);
UPDATE schema_compatibility SET minimum_compatible_version = 31
WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (31, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 31;",
)?;
transaction.commit()?;
Ok(())
}
fn migrate_zcode_harness_kind(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = (|| -> Result<()> {
let transaction = connection.unchecked_transaction()?;
for table in ["sessions", "hidden_native_sessions"] {
let sql: String = transaction.query_row(
"SELECT sql FROM sqlite_schema WHERE type='table' AND name=?1",
[table],
|row| row.get(0),
)?;
let (_, definition) = sql
.split_once('(')
.context("missing harness table definition")?;
if definition.contains("'muse','zcode')") {
continue;
}
ensure!(
definition.contains("'muse')"),
"unexpected {table} harness constraint"
);
let definition = definition.replace("'muse')", "'muse','zcode')");
let objects: Vec<String> = transaction
.prepare(
"SELECT sql FROM sqlite_schema WHERE tbl_name=?1
AND type IN ('index','trigger') AND sql IS NOT NULL",
)?
.query_map([table], |row| row.get(0))?
.collect::<rusqlite::Result<_>>()?;
transaction.execute_batch(&format!(
"CREATE TABLE {table}_zcode_v32 ({definition};
INSERT INTO {table}_zcode_v32 SELECT * FROM {table};
DROP TABLE {table};
ALTER TABLE {table}_zcode_v32 RENAME TO {table};"
))?;
for object in objects {
transaction.execute_batch(&object)?;
}
}
ensure!(
!transaction
.prepare("PRAGMA foreign_key_check")?
.exists([])?,
"foreign key violation in ZCode migration"
);
transaction.execute_batch(
"UPDATE schema_compatibility SET minimum_compatible_version = 32
WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (32, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
PRAGMA user_version = 32;",
)?;
transaction.commit()?;
Ok(())
})();
let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate ZCode harness constraints")?;
restored.context("restore foreign key enforcement after ZCode migration")?;
Ok(())
}
pub(super) fn table_has_column(connection: &Connection, table: &str, column: &str) -> Result<bool> {
connection
.query_row(
"SELECT EXISTS(
SELECT 1 FROM pragma_table_info(?1)
WHERE name = ?2
)",
params![table, column],
|row| row.get(0),
)
.map_err(Into::into)
}
fn ensure_workspace_schema(connection: &Connection) -> Result<()> {
connection.execute_batch(
"CREATE TABLE IF NOT EXISTS workspaces (
workspace_id TEXT PRIMARY KEY CHECK(length(trim(workspace_id)) > 0),
name TEXT NOT NULL CHECK(length(trim(name)) BETWEEN 1 AND 64),
name_key TEXT NOT NULL UNIQUE CHECK(length(trim(name_key)) BETWEEN 1 AND 64),
created_at TEXT NOT NULL,
last_opened_at TEXT NOT NULL
) STRICT;
INSERT OR IGNORE INTO workspaces(
workspace_id, name, name_key, created_at, last_opened_at
) VALUES (
'default', 'default', 'default',
strftime('%Y-%m-%dT%H:%M:%fZ', 'now'),
strftime('%Y-%m-%dT%H:%M:%fZ', 'now')
);",
)?;
if !table_has_column(connection, "session_contexts", "workspace_id")? {
connection.execute_batch(
"ALTER TABLE session_contexts
ADD COLUMN workspace_id TEXT NOT NULL DEFAULT 'default';",
)?;
}
connection.execute_batch(
"CREATE INDEX IF NOT EXISTS session_contexts_workspace
ON session_contexts(workspace_id, session_id);
CREATE TRIGGER IF NOT EXISTS session_contexts_workspace_insert
BEFORE INSERT ON session_contexts
WHEN NOT EXISTS(
SELECT 1 FROM workspaces WHERE workspace_id = NEW.workspace_id
)
BEGIN
SELECT RAISE(ABORT, 'unknown workspace');
END;
CREATE TRIGGER IF NOT EXISTS session_contexts_workspace_update
BEFORE UPDATE OF workspace_id ON session_contexts
WHEN NOT EXISTS(
SELECT 1 FROM workspaces WHERE workspace_id = NEW.workspace_id
)
BEGIN
SELECT RAISE(ABORT, 'unknown workspace');
END;
CREATE TABLE IF NOT EXISTS client_read_frontiers (
client_id TEXT NOT NULL CHECK(length(trim(client_id)) > 0),
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
session_id TEXT NOT NULL REFERENCES session_contexts(session_id) ON DELETE CASCADE,
through_event_ordinal INTEGER NOT NULL DEFAULT 0
CHECK(through_event_ordinal >= 0),
updated_at TEXT NOT NULL,
PRIMARY KEY(client_id, workspace_id, session_id)
) STRICT;
CREATE TABLE IF NOT EXISTS detached_drafts (
draft_id TEXT PRIMARY KEY CHECK(length(trim(draft_id)) > 0),
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id),
session_id TEXT REFERENCES session_contexts(session_id),
source TEXT NOT NULL CHECK(length(trim(source)) > 0),
owner_pid INTEGER CHECK(owner_pid IS NULL OR owner_pid > 0),
saved_at TEXT NOT NULL,
text TEXT NOT NULL CHECK(length(text) > 0),
recovered_at TEXT
) STRICT;
CREATE INDEX IF NOT EXISTS detached_drafts_workspace_recent
ON detached_drafts(workspace_id, saved_at DESC);",
)?;
ensure_client_session_state_schema(connection)?;
Ok(())
}
fn ensure_client_session_state_schema(connection: &Connection) -> Result<()> {
connection.execute_batch(
"CREATE TABLE IF NOT EXISTS client_session_state (
client_id TEXT NOT NULL CHECK(length(trim(client_id)) > 0),
workspace_id TEXT NOT NULL REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
session_id TEXT NOT NULL REFERENCES session_contexts(session_id) ON DELETE CASCADE,
draft TEXT NOT NULL DEFAULT '',
updated_at TEXT NOT NULL,
PRIMARY KEY(client_id, workspace_id, session_id)
) STRICT;
CREATE INDEX IF NOT EXISTS client_session_state_age
ON client_session_state(updated_at);",
)?;
Ok(())
}
fn ensure_session_container_override_columns(connection: &Connection) -> Result<()> {
for column in ["container_cpus", "container_memory"] {
if !table_has_column(connection, "sessions", column)? {
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN {column} TEXT;
COMMIT;"
))?;
}
}
Ok(())
}
fn migrate_turn_outcome_columns(connection: &Connection) -> Result<()> {
let mut statements = String::from("BEGIN IMMEDIATE;\n");
if !table_has_column(connection, "materialized_sessions", "active_turn_json")? {
statements.push_str(
"ALTER TABLE materialized_sessions ADD COLUMN active_turn_json TEXT
CHECK(active_turn_json IS NULL OR json_valid(active_turn_json));\n",
);
}
if !table_has_column(
connection,
"materialized_sessions",
"last_turn_outcome_json",
)? {
statements.push_str(
"ALTER TABLE materialized_sessions ADD COLUMN last_turn_outcome_json TEXT
CHECK(last_turn_outcome_json IS NULL OR json_valid(last_turn_outcome_json));\n",
);
}
if !table_has_column(
connection,
"materialized_queued_prompts",
"accepted_ordinal",
)? {
statements.push_str(
"ALTER TABLE materialized_queued_prompts ADD COLUMN accepted_ordinal INTEGER
CHECK(accepted_ordinal IS NULL OR accepted_ordinal > 0);\n",
);
}
statements.push_str(
"CREATE TABLE IF NOT EXISTS api_idempotency (
key TEXT PRIMARY KEY CHECK(length(trim(key)) BETWEEN 1 AND 128),
session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
created_at_ms INTEGER NOT NULL
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (28, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 28;
COMMIT;",
);
connection.execute_batch(&statements)?;
Ok(())
}
fn ensure_session_mount_read_only_column(connection: &Connection) -> Result<()> {
if !table_has_column(connection, "session_mounts", "read_only")? {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE session_mounts ADD COLUMN read_only INTEGER NOT NULL DEFAULT 0;
COMMIT;",
)?;
}
Ok(())
}
fn ensure_materialized_elicitation_column(connection: &Connection) -> Result<()> {
if !table_has_column(
connection,
"materialized_sessions",
"pending_elicitations_json",
)? {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE materialized_sessions
ADD COLUMN pending_elicitations_json TEXT NOT NULL DEFAULT '[]';
COMMIT;",
)?;
}
Ok(())
}
fn ensure_managed_worktree_column(connection: &Connection) -> Result<()> {
if !table_has_column(connection, "sessions", "managed_worktree")? {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN managed_worktree TEXT;
COMMIT;",
)?;
}
Ok(())
}
fn ensure_relay_projection_schema(connection: &Connection) -> Result<()> {
if table_has_column(connection, "sessions", "detached_after_event_ordinal")? {
return Ok(());
}
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
ALTER TABLE session_checkpoints
RENAME COLUMN event_sequence TO event_frontier;
ALTER TABLE prompt_history
RENAME COLUMN event_sequence TO event_ordinal;
ALTER TABLE sessions ADD COLUMN detached_after_event_ordinal INTEGER NOT NULL
DEFAULT 0 CHECK(detached_after_event_ordinal >= 0);
CREATE TABLE materialized_sessions (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
applied_event_ordinal INTEGER NOT NULL DEFAULT 0 CHECK(applied_event_ordinal >= 0),
applied_event_digest TEXT NOT NULL
DEFAULT '{RELAY_EVENT_GENESIS_DIGEST}'
CHECK(length(applied_event_digest) = 64
AND applied_event_digest NOT GLOB '*[^0-9a-f]*'),
last_activity_at_ms INTEGER,
execution_state TEXT NOT NULL DEFAULT 'idle'
CHECK(execution_state IN ('idle','running','closing','closed')),
running_started_at_ms INTEGER,
session_title TEXT CHECK(session_title IS NULL OR length(trim(session_title)) > 0),
configuration_json TEXT NOT NULL DEFAULT '{{}}',
CHECK(
(execution_state = 'running' AND running_started_at_ms IS NOT NULL)
OR (execution_state != 'running' AND running_started_at_ms IS NULL)
)
) STRICT;
CREATE TABLE materialized_transcript_items (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
stable_id TEXT NOT NULL CHECK(length(trim(stable_id)) > 0),
position INTEGER NOT NULL CHECK(position > 0),
latest_content_event_ordinal INTEGER
CHECK(latest_content_event_ordinal IS NULL
OR latest_content_event_ordinal >= position),
created_at_ms INTEGER NOT NULL,
last_changed_at_ms INTEGER NOT NULL CHECK(last_changed_at_ms >= created_at_ms),
body_json TEXT NOT NULL,
PRIMARY KEY(session_id, stable_id)
) STRICT;
CREATE INDEX materialized_transcript_position
ON materialized_transcript_items(session_id, position, stable_id);
CREATE TABLE materialized_queued_prompts (
session_id TEXT NOT NULL REFERENCES materialized_sessions(session_id) ON DELETE CASCADE,
ordinal INTEGER NOT NULL CHECK(ordinal >= 0),
command_id TEXT NOT NULL CHECK(length(trim(command_id)) > 0),
content_json TEXT NOT NULL,
queued_at_ms INTEGER NOT NULL,
PRIMARY KEY(session_id, ordinal),
UNIQUE(session_id, command_id)
) STRICT;
INSERT INTO materialized_sessions(session_id)
SELECT session_id FROM sessions;
COMMIT;",
))?;
Ok(())
}
fn migrate_destroying_session_state(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE sessions_v7 (
session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
title TEXT NOT NULL CHECK(length(trim(title)) > 0),
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi')),
last_profile TEXT NOT NULL,
target_template_id TEXT NOT NULL,
state TEXT NOT NULL CHECK(state IN (
'provisioning','running','disconnected','checkpointing','closing','destroying',
'archived','lost','error','destroyed-with-data-loss'
)),
native_session_id TEXT,
acp_session_title TEXT CHECK(acp_session_title IS NULL OR length(trim(acp_session_title)) > 0),
session_title_override TEXT CHECK(session_title_override IS NULL OR length(trim(session_title_override)) > 0),
updated_at TEXT NOT NULL,
detached_after_event_ordinal INTEGER NOT NULL DEFAULT 0
CHECK(detached_after_event_ordinal >= 0),
last_error TEXT,
resource_allocation TEXT,
last_checkpoint_error TEXT,
project_directory BLOB,
managed_worktree TEXT
) STRICT;
INSERT INTO sessions_v7(
session_id, title, harness_kind, last_profile, target_template_id, state,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree
)
SELECT
session_id, title, harness_kind, last_profile, target_template_id, state,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree
FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_v7 RENAME TO sessions;
INSERT INTO schema_migrations(version, applied_at)
VALUES (7, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 7;
COMMIT;",
);
if migration.is_err()
&& let Err(error) = connection.execute_batch("ROLLBACK;")
{
tracing::warn!(%error, "could not roll back durable-destroying-session migration");
}
let foreign_keys = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate durable destroying session state")?;
foreign_keys.context("restore foreign key enforcement after schema migration")?;
let mut statement = connection.prepare("PRAGMA foreign_key_check")?;
if statement.exists([])? {
bail!("foreign key violation after migrating durable destroying session state");
}
Ok(())
}
fn migrate_grok_harness_kind(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE sessions_v9 (
session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
title TEXT NOT NULL CHECK(length(trim(title)) > 0),
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi','grok')),
last_profile TEXT NOT NULL,
target_template_id TEXT NOT NULL,
state TEXT NOT NULL CHECK(state IN (
'provisioning','running','disconnected','checkpointing','closing','destroying',
'archived','lost','error','destroyed-with-data-loss'
)),
native_session_id TEXT,
acp_session_title TEXT CHECK(acp_session_title IS NULL OR length(trim(acp_session_title)) > 0),
session_title_override TEXT CHECK(session_title_override IS NULL OR length(trim(session_title_override)) > 0),
updated_at TEXT NOT NULL,
detached_after_event_ordinal INTEGER NOT NULL DEFAULT 0
CHECK(detached_after_event_ordinal >= 0),
last_error TEXT,
resource_allocation TEXT,
last_checkpoint_error TEXT,
project_directory BLOB,
managed_worktree TEXT,
draft_input TEXT NOT NULL DEFAULT ''
) STRICT;
INSERT INTO sessions_v9(
session_id, title, harness_kind, last_profile, target_template_id, state,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree, draft_input
)
SELECT
session_id, title, harness_kind, last_profile, target_template_id, state,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree, draft_input
FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_v9 RENAME TO sessions;
INSERT INTO schema_migrations(version, applied_at)
VALUES (9, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 9;
COMMIT;",
);
if migration.is_err()
&& let Err(error) = connection.execute_batch("ROLLBACK;")
{
tracing::warn!(%error, "could not roll back Grok harness migration");
}
let foreign_keys = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate sessions table for the Grok Build harness")?;
foreign_keys.context("restore foreign key enforcement after schema migration")?;
let mut statement = connection.prepare("PRAGMA foreign_key_check")?;
if statement.exists([])? {
bail!("foreign key violation after migrating the sessions harness list");
}
Ok(())
}
fn migrate_stopped_session_state(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE sessions_v10 (
session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
title TEXT NOT NULL CHECK(length(trim(title)) > 0),
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi','grok')),
last_profile TEXT NOT NULL,
target_template_id TEXT NOT NULL,
state TEXT NOT NULL CHECK(state IN (
'provisioning','running','disconnected','checkpointing','closing','destroying',
'stopped','lost','error','destroyed-with-data-loss'
)),
native_session_id TEXT,
acp_session_title TEXT CHECK(acp_session_title IS NULL OR length(trim(acp_session_title)) > 0),
session_title_override TEXT CHECK(session_title_override IS NULL OR length(trim(session_title_override)) > 0),
updated_at TEXT NOT NULL,
detached_after_event_ordinal INTEGER NOT NULL DEFAULT 0
CHECK(detached_after_event_ordinal >= 0),
last_error TEXT,
resource_allocation TEXT,
last_checkpoint_error TEXT,
project_directory BLOB,
managed_worktree TEXT,
draft_input TEXT NOT NULL DEFAULT '',
container_cpus TEXT,
container_memory TEXT,
archived INTEGER NOT NULL DEFAULT 0 CHECK(archived IN (0, 1))
) STRICT;
INSERT INTO sessions_v10(
session_id, title, harness_kind, last_profile, target_template_id, state,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree, draft_input,
container_cpus, container_memory
)
SELECT
session_id, title, harness_kind, last_profile, target_template_id,
CASE state WHEN 'archived' THEN 'stopped' ELSE state END,
native_session_id, acp_session_title, session_title_override, updated_at,
detached_after_event_ordinal, last_error, resource_allocation,
last_checkpoint_error, project_directory, managed_worktree, draft_input,
container_cpus, container_memory
FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_v10 RENAME TO sessions;
CREATE TABLE hidden_native_sessions (
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi','grok')),
native_session_id TEXT NOT NULL CHECK(length(trim(native_session_id)) > 0),
hidden_at TEXT NOT NULL,
PRIMARY KEY(harness_kind, native_session_id)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (10, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 10;
COMMIT;",
);
if migration.is_err()
&& let Err(error) = connection.execute_batch("ROLLBACK;")
{
tracing::warn!(%error, "could not roll back stopped-session migration");
}
let foreign_keys = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate sessions table for the stopped session state")?;
foreign_keys.context("restore foreign key enforcement after schema migration")?;
let mut statement = connection.prepare("PRAGMA foreign_key_check")?;
if statement.exists([])? {
bail!("foreign key violation after migrating the stopped session state");
}
Ok(())
}
fn migrate_muse_harness_kind(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = (|| -> Result<()> {
let transaction = connection.unchecked_transaction()?;
for table in ["sessions", "hidden_native_sessions"] {
let sql: String = transaction.query_row(
"SELECT sql FROM sqlite_schema WHERE type='table' AND name=?1",
[table],
|row| row.get(0),
)?;
let (_, definition) = sql
.split_once('(')
.context("missing harness table definition")?;
if definition.contains("'deepseek','muse')") {
continue;
}
ensure!(
definition.contains("'deepseek')"),
"unexpected {table} harness constraint"
);
let definition = definition.replace("'deepseek')", "'deepseek','muse')");
let objects: Vec<String> = transaction.prepare("SELECT sql FROM sqlite_schema WHERE tbl_name=?1 AND type IN ('index','trigger') AND sql IS NOT NULL")?
.query_map([table], |row| row.get(0))?.collect::<rusqlite::Result<_>>()?;
transaction.execute_batch(&format!(
"CREATE TABLE {table}_muse_v27 ({definition}; INSERT INTO {table}_muse_v27 SELECT * FROM {table}; DROP TABLE {table}; ALTER TABLE {table}_muse_v27 RENAME TO {table};"
))?;
for object in objects {
transaction.execute_batch(&object)?;
}
}
ensure!(
!transaction
.prepare("PRAGMA foreign_key_check")?
.exists([])?,
"foreign key violation in Muse migration"
);
transaction.execute_batch("INSERT INTO schema_migrations(version, applied_at) VALUES (27, strftime('%Y-%m-%dT%H:%M:%fZ','now')); PRAGMA user_version = 27;")?;
transaction.commit()?;
Ok(())
})();
let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate Muse harness constraints")?;
restored.context("restore foreign key enforcement after Muse migration")?;
Ok(())
}
fn migrate_deepseek_harness_kind(connection: &Connection) -> Result<()> {
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE sessions_v11 (
session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
title TEXT NOT NULL CHECK(length(trim(title)) > 0),
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi','grok','deepseek')),
last_profile TEXT NOT NULL,
target_template_id TEXT NOT NULL,
state TEXT NOT NULL CHECK(state IN (
'provisioning','running','disconnected','checkpointing','closing','destroying',
'stopped','lost','error','destroyed-with-data-loss'
)),
native_session_id TEXT,
acp_session_title TEXT CHECK(acp_session_title IS NULL OR length(trim(acp_session_title)) > 0),
session_title_override TEXT CHECK(session_title_override IS NULL OR length(trim(session_title_override)) > 0),
updated_at TEXT NOT NULL,
detached_after_event_ordinal INTEGER NOT NULL DEFAULT 0
CHECK(detached_after_event_ordinal >= 0),
last_error TEXT,
resource_allocation TEXT,
last_checkpoint_error TEXT,
project_directory BLOB,
managed_worktree TEXT,
draft_input TEXT NOT NULL DEFAULT '',
container_cpus TEXT,
container_memory TEXT,
archived INTEGER NOT NULL DEFAULT 0 CHECK(archived IN (0, 1))
) STRICT;
INSERT INTO sessions_v11 SELECT * FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_v11 RENAME TO sessions;
ALTER TABLE hidden_native_sessions RENAME TO hidden_native_sessions_v10;
CREATE TABLE hidden_native_sessions (
harness_kind TEXT NOT NULL CHECK(harness_kind IN ('codex','claude','kimi','grok','deepseek')),
native_session_id TEXT NOT NULL CHECK(length(trim(native_session_id)) > 0),
hidden_at TEXT NOT NULL,
PRIMARY KEY(harness_kind, native_session_id)
) STRICT;
INSERT INTO hidden_native_sessions SELECT * FROM hidden_native_sessions_v10;
DROP TABLE hidden_native_sessions_v10;
INSERT INTO schema_migrations(version, applied_at)
VALUES (11, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 11;
COMMIT;",
);
if migration.is_err()
&& let Err(error) = connection.execute_batch("ROLLBACK;")
{
tracing::warn!(%error, "could not roll back DeepSeek harness migration");
}
let foreign_keys = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate sessions table for DeepSeek Harness")?;
foreign_keys.context("restore foreign key enforcement after schema migration")?;
let mut statement = connection.prepare("PRAGMA foreign_key_check")?;
if statement.exists([])? {
bail!("foreign key violation after migrating the DeepSeek Harness list");
}
Ok(())
}
fn ensure_session_draft_input_column(connection: &Connection) -> Result<()> {
if !table_has_column(connection, "sessions", "draft_input")? {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN draft_input TEXT NOT NULL DEFAULT '';
COMMIT;",
)?;
}
Ok(())
}
fn ensure_projection_digest_column(connection: &Connection) -> Result<()> {
let present = connection.query_row(
"SELECT EXISTS(
SELECT 1 FROM pragma_table_info('materialized_sessions')
WHERE name = 'applied_event_digest'
)",
[],
|row| row.get::<_, bool>(0),
)?;
if !present {
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
ALTER TABLE materialized_sessions ADD COLUMN applied_event_digest TEXT NOT NULL
DEFAULT '{RELAY_EVENT_GENESIS_DIGEST}'
CHECK(length(applied_event_digest) = 64
AND applied_event_digest NOT GLOB '*[^0-9a-f]*');
COMMIT;",
))?;
}
Ok(())
}
fn ensure_api_events_schema(connection: &Connection) -> Result<()> {
connection.execute_batch(
"CREATE TABLE IF NOT EXISTS api_events (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
recorded_at_ms INTEGER NOT NULL,
body TEXT NOT NULL CHECK(json_valid(body))
) STRICT;
CREATE INDEX IF NOT EXISTS api_events_session ON api_events(session_id, seq);
CREATE TABLE IF NOT EXISTS api_session_activity (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
body TEXT NOT NULL CHECK(json_valid(body))
) STRICT;
CREATE TRIGGER IF NOT EXISTS api_session_error_updated
AFTER UPDATE OF last_error ON sessions
WHEN NEW.last_error IS NOT NULL AND NEW.last_error IS NOT OLD.last_error
BEGIN
INSERT INTO api_events(session_id, recorded_at_ms, body)
VALUES (NEW.session_id, CAST((julianday('now') - 2440587.5) * 86400000 AS INTEGER),
json_object('type', 'error', 'data', json_object('message', NEW.last_error, 'command_id', NULL)));
END;
CREATE TRIGGER IF NOT EXISTS api_session_error_inserted
AFTER INSERT ON sessions WHEN NEW.last_error IS NOT NULL
BEGIN
INSERT INTO api_events(session_id, recorded_at_ms, body)
VALUES (NEW.session_id, CAST((julianday('now') - 2440587.5) * 86400000 AS INTEGER),
json_object('type', 'error', 'data', json_object('message', NEW.last_error, 'command_id', NULL)));
END;"
)?;
Ok(())
}