use super::*;
use rusqlite::OpenFlags;
const COMPATIBILITY_METADATA_VERSION: i64 = 30;
pub(super) struct SchemaState {
pub(super) revision: i64,
minimum_compatible: Option<i64>,
}
impl SchemaState {
pub(super) fn ensure_supported(&self) -> Result<()> {
self.ensure_supported_by(SCHEMA_VERSION)
}
fn ensure_supported_by(&self, supported: i64) -> Result<()> {
let reason = if self.revision < supported {
StoreSchemaMismatchReason::NeedsMigration
} else if let Some(minimum_compatible) = self.minimum_compatible {
if minimum_compatible <= supported {
return Ok(());
}
StoreSchemaMismatchReason::Incompatible { minimum_compatible }
} else {
StoreSchemaMismatchReason::InvalidCompatibilityMetadata
};
Err(StoreSchemaMismatch {
found: self.revision,
supported,
reason,
}
.into())
}
}
pub(super) fn read_schema_state(connection: &Connection) -> Result<SchemaState> {
let snapshot = connection
.unchecked_transaction()
.context("start database compatibility snapshot")?;
let revision: i64 = snapshot
.query_row("PRAGMA user_version", [], |row| row.get(0))
.context("read database migration revision")?;
let minimum_compatible = if revision >= COMPATIBILITY_METADATA_VERSION {
let invalid = || StoreSchemaMismatch {
found: revision,
supported: SCHEMA_VERSION,
reason: StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
};
let (count, singleton, floor, recorded): (i64, Option<i64>, Option<i64>, Option<i64>) =
snapshot
.query_row(
"SELECT count(*), min(singleton), min(minimum_compatible_version),
(SELECT max(version) FROM schema_migrations)
FROM schema_compatibility",
[],
|row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
)
.map_err(|error| {
let structural = match &error {
rusqlite::Error::SqliteFailure(code, _) => {
code.code == rusqlite::ErrorCode::Unknown
}
_ => true,
};
let error = anyhow::Error::new(error);
if structural {
error.context(invalid())
} else {
error.context("read database compatibility metadata")
}
})?;
if count != 1
|| singleton != Some(1)
|| recorded != Some(revision)
|| !floor
.is_some_and(|floor| (COMPATIBILITY_METADATA_VERSION..=revision).contains(&floor))
{
return Err(invalid().into());
}
floor
} else {
None
};
snapshot
.commit()
.context("finish database compatibility snapshot")?;
Ok(SchemaState {
revision,
minimum_compatible,
})
}
pub fn database_path() -> PathBuf {
data_dir().join("mj.sqlite3")
}
pub fn check_read_compatibility() -> Result<()> {
open_reader_strict(&database_path()).map(drop)
}
pub(super) fn open_writer(path: &Path) -> Result<Connection> {
let mut connection = open_writable(path)?;
connection.set_transaction_behavior(rusqlite::TransactionBehavior::Immediate);
Ok(connection)
}
fn open_writable(path: &Path) -> Result<Connection> {
if let Some(store) = path.parent() {
mj_core::config::ensure_may_control_store(store, "open this database for writing")?;
}
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)
.with_context(|| format!("create Mjolnir data directory {}", parent.display()))?;
}
let connection = Connection::open(path)
.with_context(|| format!("open Mjolnir database {}", path.display()))?;
connection.busy_timeout(Duration::from_secs(5))?;
connection.execute_batch(
"PRAGMA foreign_keys = ON;
PRAGMA journal_mode = WAL;
PRAGMA synchronous = FULL;",
)?;
verify_schema_once(path, &connection)?;
committed::observe_connection(&connection, path)?;
Ok(connection)
}
pub(super) fn open(path: &Path) -> Result<Connection> {
open_writer(path)
}
#[cfg(not(test))]
pub(super) fn open_reader(path: &Path) -> Result<Connection> {
open_reader_strict(path)
}
#[cfg(test)]
pub(super) fn open_reader(path: &Path) -> Result<Connection> {
open_writable(path)
}
#[cfg_attr(test, allow(dead_code))]
fn open_reader_strict(path: &Path) -> Result<Connection> {
let connection = Connection::open_with_flags(
path,
OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
)
.with_context(|| format!("open Mjolnir database read-only {}", path.display()))?;
connection.busy_timeout(Duration::from_secs(5))?;
connection.execute_batch(
"PRAGMA foreign_keys = ON;
PRAGMA query_only = ON;",
)?;
read_schema_state(&connection)?.ensure_supported()?;
Ok(connection)
}
fn verified_schemas() -> &'static Mutex<HashSet<PathBuf>> {
static VERIFIED: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
VERIFIED.get_or_init(|| Mutex::new(HashSet::new()))
}
fn schema_cache_key(path: &Path) -> PathBuf {
let Some(parent) = path
.parent()
.filter(|parent| !parent.as_os_str().is_empty())
else {
return path.to_owned();
};
match (fs::canonicalize(parent), path.file_name()) {
(Ok(canonical), Some(name)) => canonical.join(name),
_ => path.to_owned(),
}
}
fn verify_schema_once(path: &Path, connection: &Connection) -> Result<()> {
let key = schema_cache_key(path);
let mut verified = verified_schemas()
.lock()
.unwrap_or_else(PoisonError::into_inner);
let state = read_schema_state(connection)?;
if state.revision > SCHEMA_VERSION
|| (state.revision == SCHEMA_VERSION && verified.contains(&key))
{
return state.ensure_supported();
}
migrate_schema(connection)?;
read_schema_state(connection)?.ensure_supported()?;
verified.insert(key);
Ok(())
}
#[cfg(test)]
pub(super) fn forget_verified_schema(path: &Path) {
verified_schemas()
.lock()
.unwrap_or_else(PoisonError::into_inner)
.remove(&schema_cache_key(path));
}
const BASELINE_SCHEMA_VERSION: i64 = 33;
const BASELINE_MINIMUM_COMPATIBLE_VERSION: i64 = 32;
const ACCOUNTING_MIGRATION_SQL: &str = "
ALTER TABLE session_turn_usage RENAME TO old_session_turn_usage;
CREATE TABLE session_turn_usage (
session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
command_id TEXT NOT NULL,
completed_ordinal INTEGER NOT NULL,
turn_start_position INTEGER,
body TEXT NOT NULL,
PRIMARY KEY(session_id, command_id)
);
INSERT INTO session_turn_usage SELECT * FROM old_session_turn_usage;
DROP TABLE old_session_turn_usage;
CREATE INDEX session_turn_usage_order ON session_turn_usage(session_id, completed_ordinal);
ALTER TABLE session_provider_cost RENAME TO old_session_provider_cost;
CREATE TABLE session_provider_cost (
session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
body TEXT NOT NULL
);
INSERT INTO session_provider_cost SELECT * FROM old_session_provider_cost;
DROP TABLE old_session_provider_cost;
CREATE TABLE subagent_accounting (
child_session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
parent_session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
task_name TEXT NOT NULL,
CHECK(child_session_id <> parent_session_id)
) STRICT;
CREATE INDEX subagent_accounting_parent ON subagent_accounting(parent_session_id);
INSERT INTO subagent_accounting SELECT child_session_id, parent_session_id,
json_extract(record_json, '$.task_name') FROM subagent_sessions;
CREATE TABLE session_turn_selections (
session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
command_id TEXT NOT NULL,
model TEXT,
effort TEXT,
PRIMARY KEY(session_id, command_id)
) STRICT;
";
fn migrate_schema(connection: &Connection) -> Result<()> {
let state = read_schema_state(connection)?;
let version = state.revision;
if version > SCHEMA_VERSION {
return state.ensure_supported();
}
if version == 0 {
create_baseline_schema(connection)?;
} else if version < BASELINE_SCHEMA_VERSION {
super::legacy_schema::migrate_to_baseline(connection)
.context("upgrade historical database schema")?;
}
if version < 34 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS session_mount_access (
session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
source BLOB NOT NULL,
destination BLOB NOT NULL,
access TEXT NOT NULL CHECK(access IN ('rw')),
PRIMARY KEY(session_id, destination)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (34, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 34;
COMMIT;",
)?;
}
if version < 35 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN container_workspace TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (35, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 35;
COMMIT;",
)?;
}
if version < 36 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE sessions ADD COLUMN build_cache_json TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (36, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 36;
COMMIT;",
)?;
}
if version < 37 {
connection.execute_batch(
"BEGIN IMMEDIATE;
ALTER TABLE session_targets ADD COLUMN borrowed_from TEXT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (37, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 37;
COMMIT;",
)?;
}
if version < 38 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS workspace_layouts (
workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
layout TEXT NOT NULL
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (38, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 38;
COMMIT;",
)?;
}
if version < 39 {
connection.execute_batch(
"BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 39 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (39, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 39;
COMMIT;",
)?;
}
if version < 40 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE native_agents (
owner TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
child TEXT NOT NULL,
staging INTEGER NOT NULL CHECK(staging IN (0,1)),
body TEXT NOT NULL CHECK(json_valid(body)),
PRIMARY KEY(owner, child, staging)
) STRICT;
CREATE TABLE native_agent_transcript (
owner TEXT NOT NULL,
child TEXT NOT NULL,
staging INTEGER NOT NULL,
stable_id TEXT NOT NULL,
position INTEGER NOT NULL,
body TEXT NOT NULL CHECK(json_valid(body)),
PRIMARY KEY(owner, child, staging, stable_id),
FOREIGN KEY(owner, child, staging) REFERENCES native_agents(owner, child, staging)
ON DELETE CASCADE ON UPDATE CASCADE
) STRICT;
CREATE INDEX native_agent_transcript_position ON native_agent_transcript(owner, child, staging, position);
CREATE TABLE native_agent_replay (
owner TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE
) STRICT;
UPDATE schema_compatibility SET minimum_compatible_version = 40 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (40, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 40;
COMMIT;",
)?;
}
if version < 41 {
connection.execute_batch(
"BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 41 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (41, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 41;
COMMIT;",
)?;
}
if version < 42 {
connection.execute_batch(
"BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 42 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (42, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 42;
COMMIT;",
)?;
}
if version < 43 {
connection.execute_batch("BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 43 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (43, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 43;
COMMIT;")?;
}
if version < 44 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE quota_reset_cache (identity TEXT PRIMARY KEY, body TEXT NOT NULL);
UPDATE schema_compatibility SET minimum_compatible_version = 44 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (44, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 44;
COMMIT;")?;
}
if version < 45 {
let add_column =
match super::legacy_schema::table_has_column(connection, "sessions", "launch_base")? {
true => "",
false => "ALTER TABLE sessions ADD COLUMN launch_base TEXT;",
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
INSERT INTO schema_migrations(version, applied_at)
VALUES (45, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 45;
COMMIT;"
))?;
}
if version < 46 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"sessions",
"target_runtime_json",
)? {
""
} else {
"ALTER TABLE sessions ADD COLUMN target_runtime_json TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
UPDATE schema_compatibility SET minimum_compatible_version = 46 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (46, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 46;
COMMIT;"
))?;
}
if version < 47 {
let add_branch =
if super::legacy_schema::table_has_column(connection, "sessions", "launch_branch")? {
""
} else {
"ALTER TABLE sessions ADD COLUMN launch_branch TEXT;"
};
let add_publication = if super::legacy_schema::table_has_column(
connection,
"sessions",
"publication_json",
)? {
""
} else {
"ALTER TABLE sessions ADD COLUMN publication_json TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_branch}
{add_publication}
UPDATE schema_compatibility SET minimum_compatible_version = 47 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (47, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 47;
COMMIT;",
))?;
}
if version < 48 {
connection.execute_batch(
"BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 48 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (48, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 48;
COMMIT;",
)?;
}
if version < 49 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS subagent_handbacks (
child_session_id TEXT PRIMARY KEY,
handback_command_id TEXT,
handback_message TEXT,
handback_recorded_at_ms INTEGER,
reminder_command_id TEXT,
reminder_for_command_id TEXT,
reminder_sent_at_ms INTEGER,
reminder_failed_for_command_id TEXT
);
UPDATE schema_compatibility SET minimum_compatible_version = 49 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (49, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 49;
COMMIT;",
)?;
}
if version < 50 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"subagent_handbacks",
"awaited_ordinal",
)? {
""
} else {
"ALTER TABLE subagent_handbacks ADD COLUMN awaited_ordinal INTEGER;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
INSERT INTO schema_migrations(version, applied_at)
VALUES (50, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 50;
COMMIT;"
))?;
}
if version < 51 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"subagent_handbacks",
"report_dir",
)? {
""
} else {
"ALTER TABLE subagent_handbacks ADD COLUMN report_dir TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
INSERT INTO schema_migrations(version, applied_at)
VALUES (51, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 51;
COMMIT;"
))?;
}
if version < 52 {
connection.execute_batch(
"BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS stopped_subagents (
parent_session_id TEXT NOT NULL
REFERENCES sessions(session_id) ON DELETE CASCADE,
child_session_id TEXT NOT NULL,
record_json TEXT NOT NULL CHECK(json_valid(record_json)),
PRIMARY KEY(parent_session_id, child_session_id)
) STRICT;
INSERT INTO schema_migrations(version, applied_at)
VALUES (52, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 52;
COMMIT;",
)?;
}
if version < 53 {
migrate_parked_session_state(connection)?;
}
if version < 54 {
let add_column =
if super::legacy_schema::table_has_column(connection, "sessions", "checkout_json")? {
""
} else {
"ALTER TABLE sessions ADD COLUMN checkout_json TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
UPDATE schema_compatibility SET minimum_compatible_version = 54 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (54, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 54;
COMMIT;"
))?;
}
if version < 55 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"sessions",
"expected_runtime_identity",
)? {
""
} else {
"ALTER TABLE sessions ADD COLUMN expected_runtime_identity TEXT;"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
UPDATE schema_compatibility SET minimum_compatible_version = 55 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (55, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 55;
COMMIT;"
))?;
}
if version < 56 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"sessions",
"subagents",
)? {
""
} else {
"ALTER TABLE sessions ADD COLUMN subagents TEXT CHECK(subagents IS NULL OR json_valid(subagents));"
};
connection.execute_batch(&format!(
"BEGIN IMMEDIATE;
{add_column}
UPDATE sessions SET subagents = CASE WHEN mjolnir_subagents = 1
THEN '{{\"mode\":\"all_models\"}}' ELSE '{{\"mode\":\"native\"}}' END WHERE subagents IS NULL AND mjolnir_subagents IS NOT NULL;
CREATE TABLE IF NOT EXISTS subagent_preference (singleton INTEGER PRIMARY KEY CHECK(singleton = 1), policy TEXT NOT NULL CHECK(json_valid(policy)));
UPDATE schema_compatibility SET minimum_compatible_version = 56 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (56, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 56;
COMMIT;"
))?;
}
if version < 57 {
let tx = connection.unchecked_transaction()?;
tx.execute_batch("DROP TRIGGER IF EXISTS api_session_error_updated;
DROP TRIGGER IF EXISTS api_session_error_inserted;
CREATE TABLE IF NOT EXISTS checkpoint_operations (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
command_id TEXT NOT NULL UNIQUE,
related_command_ids TEXT NOT NULL DEFAULT '[]' CHECK(json_valid(related_command_ids))
) STRICT;")?;
super::events::migrate_event_outcomes(&tx)?;
tx.execute_batch("UPDATE schema_compatibility SET minimum_compatible_version = 57 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (57, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 57;")?;
tx.commit()?;
}
if version < 58 {
connection.execute_batch("BEGIN IMMEDIATE;
UPDATE schema_compatibility SET minimum_compatible_version = 58 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (58, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 58;
COMMIT;")?;
}
if version < 59 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS startup_steps (
sequence INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL,
group_id TEXT,
command_id TEXT NOT NULL UNIQUE,
step_json TEXT NOT NULL,
phase TEXT NOT NULL DEFAULT 'pending'
CHECK (phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting', 'done', 'failed', 'dismissed')),
error TEXT,
accepted_ordinal INTEGER
) STRICT;
CREATE INDEX IF NOT EXISTS startup_steps_session_sequence
ON startup_steps(session_id, sequence);
CREATE INDEX IF NOT EXISTS startup_steps_group_sequence
ON startup_steps(group_id, sequence) WHERE group_id IS NOT NULL;
CREATE INDEX IF NOT EXISTS startup_steps_pending
ON startup_steps(session_id, sequence)
WHERE phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting');
UPDATE schema_compatibility SET minimum_compatible_version = 59 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (59, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 59;
COMMIT;")?;
}
if version < 60 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS delegation_effects (
parent_session_id TEXT NOT NULL,
request_id TEXT NOT NULL,
phase TEXT NOT NULL,
prepared_json TEXT,
result_json TEXT,
receipt_pending INTEGER NOT NULL DEFAULT 0,
PRIMARY KEY(parent_session_id, request_id)
) STRICT;
CREATE INDEX IF NOT EXISTS delegation_effects_receipt_cleanup
ON delegation_effects(parent_session_id, request_id) WHERE receipt_pending = 1;
UPDATE schema_compatibility SET minimum_compatible_version = 60 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (60, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 60;
COMMIT;")?;
}
if version < 61 {
let add_column = if super::legacy_schema::table_has_column(
connection,
"turn_review_state",
"orchestration",
)? {
""
} else {
"ALTER TABLE turn_review_state ADD COLUMN orchestration TEXT;"
};
connection.execute_batch(&format!("BEGIN IMMEDIATE;
{add_column}
UPDATE schema_compatibility SET minimum_compatible_version = 61 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (61, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 61;
COMMIT;"))?;
}
if version < 62 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS worker_restart_intents (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
operation_id TEXT NOT NULL,
target_json TEXT NOT NULL,
desired_build TEXT NOT NULL,
phase TEXT NOT NULL CHECK(phase IN ('prepared', 'swapping', 'awaiting_readiness'))
) STRICT;
UPDATE schema_compatibility SET minimum_compatible_version = 62 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (62, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 62;
COMMIT;")?;
}
if version < 63 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS session_incarnations (
session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
identity TEXT NOT NULL
) STRICT;
INSERT OR IGNORE INTO session_incarnations(session_id, identity)
SELECT session_id, lower(hex(randomblob(16))) FROM sessions;
CREATE TRIGGER IF NOT EXISTS session_incarnation_insert
AFTER INSERT ON sessions BEGIN
INSERT INTO session_incarnations(session_id, identity)
VALUES(NEW.session_id, lower(hex(randomblob(16))));
END;
CREATE TRIGGER IF NOT EXISTS session_incarnation_resume
AFTER UPDATE OF state ON sessions
WHEN (NEW.state = 'provisioning' AND OLD.state <> 'provisioning')
OR (NEW.state = 'running' AND OLD.state IN
('stopped', 'parked', 'error', 'lost', 'destroyed-with-data-loss'))
BEGIN
UPDATE session_incarnations SET identity = lower(hex(randomblob(16)))
WHERE session_id = NEW.session_id;
END;
UPDATE schema_compatibility SET minimum_compatible_version = 63 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (63, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 63;
COMMIT;")?;
}
if version < 64 {
connection.execute_batch("BEGIN IMMEDIATE;
CREATE TABLE IF NOT EXISTS retained_move_sources (
operation_id TEXT PRIMARY KEY,
session_id TEXT NOT NULL,
source_json TEXT NOT NULL,
exclusions_json TEXT NOT NULL,
created_at TEXT NOT NULL
) STRICT;
UPDATE schema_compatibility SET minimum_compatible_version = 64 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (64, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 64;
COMMIT;")?;
}
if version < 65 {
let drop_column = if super::legacy_schema::table_has_column(
connection,
"turn_review_state",
"orchestration",
)? {
"ALTER TABLE turn_review_state DROP COLUMN orchestration;"
} else {
""
};
connection.execute_batch(&format!("BEGIN IMMEDIATE;
DROP TRIGGER IF EXISTS session_incarnation_insert;
DROP TRIGGER IF EXISTS session_incarnation_resume;
DROP TABLE IF EXISTS session_incarnations;
{drop_column}
UPDATE schema_compatibility SET minimum_compatible_version = 65 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (65, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 65;
COMMIT;"))?;
}
if version < 66 {
let drop_column = if super::legacy_schema::table_has_column(
connection,
"sessions",
"expected_runtime_identity",
)? {
"ALTER TABLE sessions DROP COLUMN expected_runtime_identity;"
} else {
""
};
connection.execute_batch(&format!("BEGIN IMMEDIATE;
{drop_column}
UPDATE schema_compatibility SET minimum_compatible_version = 66 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (66, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 66;
COMMIT;"))?;
}
if version < 67 {
connection.execute_batch(&format!("BEGIN IMMEDIATE;
{ACCOUNTING_MIGRATION_SQL}
UPDATE schema_compatibility SET minimum_compatible_version = 67 WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at) VALUES (67, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
PRAGMA user_version = 67;
COMMIT;"))?;
}
if version < 68 {
migrate_startup_cleanup_state(connection)?;
}
if version < 69 {
let has_accounting = connection
.prepare(
"SELECT 1 FROM sqlite_schema WHERE type='table' AND name='subagent_accounting'",
)?
.exists([])?;
let accounting = if has_accounting {
""
} else {
ACCOUNTING_MIGRATION_SQL
};
let add_snapshot = if super::legacy_schema::table_has_column(
connection,
"sessions",
"project_json",
)? {
""
} else {
"ALTER TABLE sessions ADD COLUMN project_json TEXT CHECK(project_json IS NULL OR json_valid(project_json));"
};
connection.execute_batch(&format!("BEGIN IMMEDIATE;
{accounting}
{add_snapshot}
CREATE TABLE IF NOT EXISTS project_catalog (
bundle_id TEXT PRIMARY KEY,
project_key TEXT NOT NULL UNIQUE,
snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
hidden INTEGER NOT NULL DEFAULT 0 CHECK(hidden IN (0,1))
) STRICT;
CREATE TABLE IF NOT EXISTS project_aliases (
bundle_id TEXT PRIMARY KEY,
canonical_id TEXT NOT NULL REFERENCES project_catalog(bundle_id),
snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
config_pending INTEGER NOT NULL DEFAULT 0 CHECK(config_pending IN (0,1))
) STRICT;
CREATE TABLE IF NOT EXISTS project_session_aliases (
session_id TEXT NOT NULL REFERENCES session_contexts(session_id) ON DELETE CASCADE,
bundle_id TEXT NOT NULL, PRIMARY KEY(session_id,bundle_id)
) STRICT;
CREATE TABLE IF NOT EXISTS project_locations (
host TEXT NOT NULL,
directory BLOB NOT NULL,
checkout_root BLOB NOT NULL,
repository_root BLOB NOT NULL,
identity_json TEXT NOT NULL CHECK(json_valid(identity_json)),
seen_at TEXT NOT NULL,
PRIMARY KEY(host, directory)
) STRICT;
CREATE TABLE IF NOT EXISTS project_seed_homes (
harness TEXT NOT NULL,
home BLOB NOT NULL,
PRIMARY KEY(harness, home)
) STRICT;
CREATE TABLE IF NOT EXISTS project_seed_failures (
harness TEXT NOT NULL, home BLOB NOT NULL, directory BLOB NOT NULL, error TEXT NOT NULL, source_file INTEGER NOT NULL CHECK(source_file IN (0,1)),
PRIMARY KEY(harness,home,directory)
) STRICT;
CREATE TABLE IF NOT EXISTS project_discovery_changes (
sequence INTEGER PRIMARY KEY AUTOINCREMENT,
session_id TEXT NOT NULL,
directory BLOB,
managed_worktree TEXT,
target_template_id TEXT NOT NULL
) STRICT;
CREATE TABLE IF NOT EXISTS project_discovery_progress (
singleton INTEGER PRIMARY KEY CHECK(singleton=1),
sequence INTEGER NOT NULL DEFAULT 0
) STRICT;
CREATE TABLE IF NOT EXISTS project_discovery_failures (
sequence INTEGER PRIMARY KEY REFERENCES project_discovery_changes(sequence) ON DELETE CASCADE,
error TEXT NOT NULL
) STRICT;
INSERT OR IGNORE INTO project_discovery_progress(singleton) VALUES(1);
INSERT INTO project_discovery_changes(session_id, directory, managed_worktree, target_template_id)
SELECT session_id, project_directory, managed_worktree, target_template_id FROM sessions
WHERE project_directory IS NOT NULL AND NOT EXISTS (
SELECT 1 FROM project_discovery_changes d WHERE d.session_id=sessions.session_id
);
CREATE TRIGGER IF NOT EXISTS project_discovery_insert AFTER INSERT ON sessions
WHEN NEW.project_directory IS NOT NULL BEGIN
INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
END;
CREATE TRIGGER IF NOT EXISTS project_discovery_update AFTER UPDATE OF project_directory,managed_worktree,target_template_id ON sessions
WHEN NEW.project_directory IS NOT NULL AND
(NEW.project_directory IS NOT OLD.project_directory
OR NEW.managed_worktree IS NOT OLD.managed_worktree
OR NEW.target_template_id IS NOT OLD.target_template_id) BEGIN
INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
END;
UPDATE schema_compatibility SET minimum_compatible_version=69 WHERE singleton=1;
INSERT INTO schema_migrations(version,applied_at) VALUES(69,strftime('%Y-%m-%dT%H:%M:%fZ','now'));
PRAGMA user_version=69;
COMMIT;"))?;
}
if version < 70 {
let transaction = connection.unchecked_transaction()?;
let titles: Vec<(String, String)> = transaction
.prepare("SELECT session_id, acp_session_title FROM sessions WHERE length(acp_session_title) > ?1")?
.query_map([mj_core::state::MAX_SESSION_TITLE_CHARS], |row| {
Ok((row.get(0)?, row.get(1)?))
})?
.collect::<rusqlite::Result<_>>()?;
for (session_id, title) in titles {
transaction.execute(
"UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
params![session_id, mj_core::state::normalize_session_title(&title)],
)?;
}
transaction.execute_batch(
"INSERT INTO schema_migrations(version, applied_at)
VALUES (70, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
PRAGMA user_version = 70;",
)?;
transaction.commit()?;
}
let recorded: Option<i64> =
connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
row.get(0)
})?;
if recorded == Some(SCHEMA_VERSION) && version < SCHEMA_VERSION {
tracing::info!(
from_revision = version,
to_revision = SCHEMA_VERSION,
"database migrations applied"
);
}
if recorded != Some(SCHEMA_VERSION) {
bail!(
"Mjolnir database migration ledger {:?} does not match schema {}",
recorded,
SCHEMA_VERSION
);
}
Ok(())
}
fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
const BEFORE: &str = "'stopped','lost',";
const AFTER: &str = "'stopped','parked','lost',";
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = (|| -> Result<()> {
let transaction = connection.unchecked_transaction()?;
let sql: String = transaction.query_row(
"SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
[],
|row| row.get(0),
)?;
let (_, definition) = sql
.split_once('(')
.context("missing sessions table definition")?;
if !definition.contains(AFTER) {
ensure!(
definition.matches(BEFORE).count() == 1,
"unexpected sessions state constraint"
);
let definition = definition.replace(BEFORE, AFTER);
let objects: Vec<String> = transaction
.prepare(
"SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
AND type IN ('index','trigger') AND sql IS NOT NULL",
)?
.query_map([], |row| row.get(0))?
.collect::<rusqlite::Result<_>>()?;
transaction.execute_batch(&format!(
"CREATE TABLE sessions_parked_v53 ({definition};
INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
))?;
for object in objects {
transaction.execute_batch(&object)?;
}
ensure!(
!transaction
.prepare("PRAGMA foreign_key_check")?
.exists([])?,
"foreign key violation in the parked-state migration"
);
}
transaction.execute_batch(
"UPDATE schema_compatibility SET minimum_compatible_version = 53
WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
PRAGMA user_version = 53;",
)?;
transaction.commit()?;
Ok(())
})();
let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate the sessions state constraint for parked sub-agents")?;
restored.context("restore foreign key enforcement after the parked-state migration")?;
Ok(())
}
fn migrate_startup_cleanup_state(connection: &Connection) -> Result<()> {
const BEFORE: &str = "'stopped','parked','lost',";
const AFTER: &str = "'stopped','parked','startup-cleanup','lost',";
connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
let migration = (|| -> Result<()> {
let transaction = connection.unchecked_transaction()?;
let sql: String = transaction.query_row(
"SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
[],
|row| row.get(0),
)?;
let (_, definition) = sql
.split_once('(')
.context("missing sessions table definition")?;
if !definition.contains(AFTER) {
ensure!(
definition.matches(BEFORE).count() == 1,
"unexpected sessions state constraint"
);
let definition = definition.replace(BEFORE, AFTER);
let objects: Vec<String> = transaction
.prepare(
"SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
AND type IN ('index','trigger') AND sql IS NOT NULL",
)?
.query_map([], |row| row.get(0))?
.collect::<rusqlite::Result<_>>()?;
transaction.execute_batch(&format!(
"CREATE TABLE sessions_startup_cleanup_v68 ({definition};
INSERT INTO sessions_startup_cleanup_v68 SELECT * FROM sessions;
DROP TABLE sessions;
ALTER TABLE sessions_startup_cleanup_v68 RENAME TO sessions;"
))?;
for object in objects {
transaction.execute_batch(&object)?;
}
ensure!(
!transaction
.prepare("PRAGMA foreign_key_check")?
.exists([])?,
"foreign key violation in the startup-cleanup migration"
);
}
transaction.execute_batch(
"UPDATE schema_compatibility SET minimum_compatible_version = 68
WHERE singleton = 1;
INSERT INTO schema_migrations(version, applied_at)
VALUES (68, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
PRAGMA user_version = 68;",
)?;
transaction.commit()?;
Ok(())
})();
let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
migration.context("migrate the sessions state constraint for failed startup cleanup")?;
restored.context("restore foreign key enforcement after the startup-cleanup migration")?;
Ok(())
}
fn create_baseline_schema(connection: &Connection) -> Result<()> {
connection.execute_batch("BEGIN IMMEDIATE;")?;
let created = (|| -> Result<()> {
let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
if version != 0 {
return Ok(());
}
connection.execute_batch(include_str!("baseline.sql"))?;
connection.execute(
"INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
[BASELINE_MINIMUM_COMPATIBLE_VERSION],
)?;
connection.execute(
"INSERT INTO schema_migrations(version, applied_at)
VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
[BASELINE_SCHEMA_VERSION],
)?;
connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
Ok(())
})();
match created {
Ok(()) => connection
.execute_batch("COMMIT;")
.context("commit baseline database schema"),
Err(error) => {
if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
tracing::warn!(%rollback, "could not roll back a failed baseline schema");
}
Err(error.context("create baseline database schema"))
}
}
}
#[cfg(test)]
pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
let connection = Connection::open(path).unwrap();
let transaction = connection.unchecked_transaction().unwrap();
transaction
.execute(
"UPDATE schema_compatibility SET minimum_compatible_version = ?1",
[minimum_compatible],
)
.unwrap();
transaction
.execute(
"INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
[revision],
)
.unwrap();
transaction
.pragma_update(None, "user_version", revision)
.unwrap();
transaction.commit().unwrap();
forget_verified_schema(path);
}
#[cfg(test)]
mod reader_tests {
use super::*;
fn assert_divergent_history_upgrades(revision: i64, project_history: bool, interrupt: bool) {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("divergent-history.sqlite");
let connection = Connection::open(&path).unwrap();
create_baseline_schema(&connection).unwrap();
connection.execute_batch(&format!(
"INSERT INTO session_contexts(session_id,bundle_id,created_at)
VALUES ('kept','project','now');
INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at,project_directory)
VALUES ('kept','Keep my work','codex','codex','local','error','now',X'2F7265706F');
INSERT INTO materialized_sessions(session_id) VALUES ('kept');
INSERT INTO session_turn_usage VALUES ('kept','turn',1,1,'{{\"tokens\":42}}');
INSERT INTO session_provider_cost VALUES ('kept','{{\"amount\":1}}');
CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
WHEN NEW.version > {revision}
BEGIN SELECT RAISE(ABORT,'fixture migration boundary'); END;"
)).unwrap();
assert!(migrate_schema(&connection).is_err());
if !connection.is_autocommit() {
connection.execute_batch("ROLLBACK").unwrap();
}
connection
.execute_batch("DROP TRIGGER stop_at_revision")
.unwrap();
assert_eq!(read_schema_state(&connection).unwrap().revision, revision);
if project_history {
connection
.execute_batch(include_str!("project_catalog_v67.sql"))
.unwrap();
connection
.execute_batch(
"INSERT INTO project_catalog VALUES ('project','key','{}',0);
INSERT INTO project_aliases VALUES ('alias','project','{}',1);
INSERT INTO project_session_aliases VALUES ('kept','alias');
UPDATE sessions SET project_json='{\"kept\":true}';",
)
.unwrap();
}
if interrupt {
connection
.execute_batch(
"CREATE TRIGGER interrupt_reconciliation BEFORE INSERT ON schema_migrations
WHEN NEW.version=69 BEGIN SELECT RAISE(ABORT,'interrupted reconciliation'); END;",
)
.unwrap();
assert!(migrate_schema(&connection).is_err());
drop(connection);
let connection = Connection::open(&path).unwrap();
assert_eq!(read_schema_state(&connection).unwrap().revision, 68);
connection
.execute_batch("DROP TRIGGER interrupt_reconciliation")
.unwrap();
} else {
drop(connection);
}
let writer = open_writer(&path).unwrap();
let state = read_schema_state(&writer).unwrap();
assert_eq!(state.revision, SCHEMA_VERSION);
assert!(state.ensure_supported_by(68).is_err());
assert!(state.ensure_supported_by(67).is_err());
if project_history {
let snapshot: String = writer
.query_row(
"SELECT project_json FROM sessions WHERE session_id='kept'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(snapshot, r#"{"kept":true}"#);
let alias: (String, i64) = writer.query_row(
"SELECT canonical_id,config_pending FROM project_aliases WHERE bundle_id='alias'", [],
|row| Ok((row.get(0)?, row.get(1)?))
).unwrap();
assert_eq!(alias, ("project".to_owned(), 1));
}
writer.execute_batch(
"UPDATE sessions SET state='startup-cleanup',project_directory=X'2F6E6577' WHERE session_id='kept';
INSERT INTO session_turn_selections VALUES ('kept','turn','model','high');"
).unwrap();
let changed: Vec<u8> = writer
.query_row(
"SELECT directory FROM project_discovery_changes ORDER BY sequence DESC LIMIT 1",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(changed, b"/new");
writer
.execute("DELETE FROM sessions WHERE session_id='kept'", [])
.unwrap();
let usage: String = writer
.query_row("SELECT body FROM session_turn_usage", [], |row| row.get(0))
.unwrap();
assert_eq!(usage, r#"{"tokens":42}"#);
let cost: String = writer
.query_row("SELECT body FROM session_provider_cost", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(cost, r#"{"amount":1}"#);
assert!(
!writer
.prepare("PRAGMA foreign_key_check")
.unwrap()
.exists([])
.unwrap()
);
drop(writer);
forget_verified_schema(&path);
assert_eq!(
read_schema_state(&open_writer(&path).unwrap())
.unwrap()
.revision,
SCHEMA_VERSION
);
}
#[test]
fn divergent_accounting_revision_67_preserves_usage_and_adds_projects() {
assert_divergent_history_upgrades(67, false, false);
}
#[test]
fn divergent_cleanup_revision_68_preserves_usage_and_adds_projects() {
assert_divergent_history_upgrades(68, false, false);
}
#[test]
fn divergent_project_revision_67_preserves_aliases_snapshots_and_usage() {
assert_divergent_history_upgrades(66, true, false);
}
#[test]
fn divergent_project_reconciliation_resumes_after_interruption() {
assert_divergent_history_upgrades(66, true, true);
}
#[test]
fn move_ownership_upgrade_retains_sources_and_refuses_previous_daemons() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("move.sqlite");
let connection = open_writer(&path).unwrap();
connection.execute_batch("DROP TABLE retained_move_sources; DELETE FROM schema_migrations WHERE version>=64; UPDATE schema_compatibility SET minimum_compatible_version=63 WHERE singleton=1; DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version=63;").unwrap();
drop(connection);
forget_verified_schema(&path);
let upgraded = open_writer(&path).unwrap();
let schema = read_schema_state(&upgraded).unwrap();
assert!(schema.ensure_supported_by(63).is_err());
upgraded.execute("INSERT INTO retained_move_sources VALUES ('move-one','session-one','{}','[]','now')", []).unwrap();
drop(upgraded);
let reopened = open_writer(&path).unwrap();
let count: i64 = reopened
.query_row("SELECT count(*) FROM retained_move_sources", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(count, 1);
}
#[test]
fn title_migration_caps_old_titles_preserves_short_titles_and_is_compatible() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("title-migration.sqlite3");
let titles = [
("long", Some("word ".repeat(20_000))),
("unicode", Some("界".repeat(257))),
("exact", Some("界".repeat(256))),
("short", Some(" Keep\nthis title ".into())),
("unset", None),
];
for (id, _) in &titles {
let mut session = super::super::tests::session(id, "project");
session.state = SessionState::Stopped;
save_session_to(&path, &session).unwrap();
}
stamp_schema_version(&path, 69);
let old = Connection::open(&path).unwrap();
for (id, title) in &titles {
old.execute(
"UPDATE sessions SET acp_session_title=?2 WHERE session_id=?1",
params![id, title],
)
.unwrap();
}
drop(old);
let upgraded = open_writer(&path).unwrap();
let state = read_schema_state(&upgraded).unwrap();
assert_eq!(state.revision, 70);
assert_eq!(state.minimum_compatible, Some(69));
state.ensure_supported_by(69).unwrap();
for (id, original) in &titles {
let stored: Option<String> = upgraded
.query_row(
"SELECT acp_session_title FROM sessions WHERE session_id=?1",
[id],
|row| row.get(0),
)
.unwrap();
let expected = match *id {
"long" => Some(format!("{}word…", "word ".repeat(50))),
"unicode" => Some(format!("{}…", "界".repeat(255))),
_ => original.clone(),
};
assert_eq!(stored, expected, "session {id}");
assert!(stored.is_none_or(|title| title.chars().count() <= 256));
}
drop(upgraded);
forget_verified_schema(&path);
drop(open_writer(&path).unwrap());
}
#[test]
fn interrupted_title_migration_rolls_back_titles_and_revision() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("title-migration-interrupted.sqlite3");
save_session_to(&path, &super::super::tests::session("old", "project")).unwrap();
stamp_schema_version(&path, 69);
let original = "word ".repeat(20_000);
let old = Connection::open(&path).unwrap();
old.execute("UPDATE sessions SET acp_session_title=?1", [&original])
.unwrap();
old.execute_batch(
"CREATE TRIGGER stop_title_migration BEFORE INSERT ON schema_migrations
WHEN NEW.version=70 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
)
.unwrap();
assert!(migrate_schema(&old).is_err());
assert_eq!(read_schema_state(&old).unwrap().revision, 69);
let stored: String = old
.query_row("SELECT acp_session_title FROM sessions", [], |row| {
row.get(0)
})
.unwrap();
assert_eq!(stored, original);
old.execute_batch("DROP TRIGGER stop_title_migration")
.unwrap();
drop(old);
drop(open_writer(&path).unwrap());
assert!(
load_state_from(&path).unwrap().sessions["old"]
.acp_session_title
.as_ref()
.unwrap()
.chars()
.count()
<= 256
);
}
#[test]
fn recent_revisions_upgrade_directly_and_preserve_user_data() {
const TEST_REVISION_FLOOR: i64 = 47;
for revision in TEST_REVISION_FLOOR..SCHEMA_VERSION {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
let connection = Connection::open(&path).unwrap();
connection
.execute_batch(include_str!("legacy_v1.sql"))
.unwrap();
connection.execute_batch(
"INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
INSERT INTO sessions(session_id, title, harness_kind, last_profile,
target_template_id, state, updated_at, native_session_id)
VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
'2026-01-01T00:00:00Z', 'native-original');
INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
).unwrap();
connection
.execute_batch(&format!(
"CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
WHEN NEW.version > {revision}
BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
))
.unwrap();
super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
assert!(migrate_schema(&connection).is_err());
drop(connection);
let connection = Connection::open(&path).unwrap();
let found: i64 = connection
.query_row("PRAGMA user_version", [], |row| row.get(0))
.unwrap();
assert_eq!(found, revision);
connection
.execute_batch("DROP TRIGGER stop_at_revision")
.unwrap();
drop(connection);
let writer =
open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
assert_eq!(
writer
.query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
.unwrap(),
"ok"
);
assert!(
!writer
.prepare("PRAGMA foreign_key_check")
.unwrap()
.exists([])
.unwrap()
);
drop(writer);
let reader = open_reader_strict(&path).unwrap();
let prompt: String = reader
.query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
.unwrap();
assert_eq!(prompt, "Keep my prompt");
let state = load_state_from(&path).unwrap();
assert_eq!(state.sessions["old-session"].title, "Keep my work");
assert_eq!(
state.sessions["old-session"].native_session_id.as_deref(),
Some("native-original")
);
drop(reader);
forget_verified_schema(&path);
drop(open_writer(&path).unwrap());
}
}
#[test]
fn accounting_migration_preserves_usage_and_changes_its_deletion_owner() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("migration.sqlite");
let connection = Connection::open(&path).unwrap();
connection
.execute_batch(include_str!("legacy_v1.sql"))
.unwrap();
connection.execute_batch("INSERT INTO session_contexts VALUES ('old-session','project','2026-01-01T00:00:00Z');
INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at)
VALUES ('old-session','Retain usage','codex','codex','local','error','2026-01-01T00:00:00Z');
CREATE TRIGGER stop_before_accounting BEFORE INSERT ON schema_migrations WHEN NEW.version=67
BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;").unwrap();
super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
assert!(migrate_schema(&connection).is_err());
connection
.execute_batch(
"ROLLBACK; DROP TRIGGER stop_before_accounting;
INSERT INTO session_turn_usage VALUES ('old-session','turn',7,1,'{}');
INSERT INTO session_provider_cost VALUES ('old-session','{\"amount\":1}');",
)
.unwrap();
drop(connection);
let writer = open_writer(&path).unwrap();
writer
.execute("DELETE FROM sessions WHERE session_id='old-session'", [])
.unwrap();
for table in ["session_turn_usage", "session_provider_cost"] {
assert_eq!(
writer
.query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| row
.get::<_, u64>(0))
.unwrap(),
1
);
}
assert_eq!(
writer
.query_row("SELECT body FROM session_turn_usage", [], |row| row
.get::<_, String>(0))
.unwrap(),
"{}"
);
assert!(
!writer
.prepare("PRAGMA foreign_key_check")
.unwrap()
.exists([])
.unwrap()
);
assert_eq!(
read_schema_state(&writer).unwrap().minimum_compatible,
Some(MINIMUM_COMPATIBLE_VERSION)
);
}
const MINIMUM_COMPATIBLE_VERSION: i64 = 69;
fn stamp_schema_version(path: &Path, version: i64) {
if version > SCHEMA_VERSION {
advance_test_schema(path, version, version);
return;
}
let connection = Connection::open(path).unwrap();
if version < 67 {
connection.execute_batch("DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET;").unwrap();
}
connection
.execute_batch(&format!("PRAGMA user_version = {version};"))
.unwrap();
connection
.execute(
"DELETE FROM schema_migrations WHERE version > ?1",
[version],
)
.unwrap();
if version == 30 {
connection
.execute(
"UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
[],
)
.unwrap();
}
drop(connection);
forget_verified_schema(path);
}
#[test]
fn subagent_policy_migration_preserves_legacy_choices_and_refuses_old_writers() {
use mj_core::subagent::SubagentPolicy;
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("subagent-policy.sqlite3");
for id in ["all", "native", "unset"] {
save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
}
let connection = Connection::open(&path).unwrap();
connection
.execute_batch(
"UPDATE sessions SET mjolnir_subagents = 1 WHERE session_id = 'all';
UPDATE sessions SET mjolnir_subagents = 0 WHERE session_id = 'native';
ALTER TABLE sessions DROP COLUMN subagents;
DROP TABLE subagent_preference;
DELETE FROM schema_migrations WHERE version >= 56;
UPDATE schema_compatibility SET minimum_compatible_version = 55;
DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 55;",
)
.unwrap();
drop(connection);
forget_verified_schema(&path);
let upgraded = open_writer(&path).unwrap();
assert!(
read_schema_state(&upgraded)
.unwrap()
.ensure_supported_by(55)
.is_err()
);
drop(upgraded);
let state = load_state_from(&path).unwrap();
assert_eq!(
state.sessions["all"].subagents,
Some(SubagentPolicy::AllModels)
);
assert_eq!(
state.sessions["native"].subagents,
Some(SubagentPolicy::Native)
);
assert_eq!(state.sessions["unset"].subagents, None);
assert_eq!(state.last_subagent_policy, SubagentPolicy::Native);
}
#[test]
fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("steering-migration.sqlite3");
let record = super::super::tests::session("steered-session", "project");
save_session_to(&path, &record).unwrap();
let connection = Connection::open(&path).unwrap();
connection
.execute_batch(
"DELETE FROM schema_migrations WHERE version >= 48;
UPDATE schema_compatibility SET minimum_compatible_version = 47;
DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 47;",
)
.unwrap();
drop(connection);
forget_verified_schema(&path);
let upgraded = open_writer(&path).unwrap();
let schema = read_schema_state(&upgraded).unwrap();
assert_eq!(schema.revision, SCHEMA_VERSION);
assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
let error = schema.ensure_supported_by(47).unwrap_err();
assert!(matches!(
error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
StoreSchemaMismatchReason::Incompatible {
minimum_compatible: MINIMUM_COMPATIBLE_VERSION
}
));
drop(upgraded);
assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
}
#[test]
fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("target-migration.sqlite3");
let record = super::super::tests::session("preserved-session", "project");
save_session_to(&path, &record).unwrap();
let connection = Connection::open(&path).unwrap();
connection
.execute_batch(
"ALTER TABLE sessions DROP COLUMN target_runtime_json;
DELETE FROM schema_migrations WHERE version >= 46;
UPDATE schema_compatibility SET minimum_compatible_version = 44;
DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 45;",
)
.unwrap();
forget_verified_schema(&path);
let upgraded = open_writer(&path).unwrap();
let schema = read_schema_state(&upgraded).unwrap();
assert_eq!(schema.revision, SCHEMA_VERSION);
assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
let error = schema.ensure_supported_by(45).unwrap_err();
assert!(matches!(
error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
StoreSchemaMismatchReason::Incompatible {
minimum_compatible: MINIMUM_COMPATIBLE_VERSION
}
));
drop(upgraded);
let restored = load_state_from(&path).unwrap();
assert_eq!(restored.sessions[&record.id], record);
}
#[test]
fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
let connection = open_writer(&path).unwrap();
connection
.execute_batch(
"BEGIN IMMEDIATE;
DROP TABLE quota_reset_cache;
DROP TABLE native_agent_transcript;
DROP TABLE native_agents;
DROP TABLE native_agent_replay;
DELETE FROM schema_migrations WHERE version >= 39;
UPDATE schema_compatibility SET minimum_compatible_version = 32;
DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 38;
COMMIT;",
)
.unwrap();
migrate_schema(&connection).unwrap();
let state = read_schema_state(&connection).unwrap();
assert_eq!(state.revision, SCHEMA_VERSION);
assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
let event = ApiEventData::InputRequired {
request: None,
turn_id: Some(1),
};
#[derive(serde::Deserialize)]
struct LegacyInputEvent {
#[serde(rename = "request")]
_request: mj_core::elicitation::ElicitationRequest,
}
let encoded = serde_json::to_value(&event).unwrap();
assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
}
#[test]
fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
let connection = open_writer(&path).unwrap();
connection
.execute_batch(
"CREATE TABLE future_feature(value TEXT NOT NULL);
INSERT INTO future_feature VALUES ('preserve me');",
)
.unwrap();
drop(connection);
advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
let reader = open_reader_strict(&path).unwrap();
assert_eq!(
reader
.query_row("SELECT value FROM future_feature", [], |row| row
.get::<_, String>(0))
.unwrap(),
"preserve me"
);
assert!(reader.execute("DELETE FROM future_feature", []).is_err());
drop(reader);
let raw = Connection::open(&path).unwrap();
raw.execute_batch("DROP TRIGGER session_contexts_workspace_update;")
.unwrap();
drop(raw);
let writer = open_writer(&path).unwrap();
assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'session_contexts_workspace_update')", [], |row| row.get::<_, bool>(0)).unwrap());
assert_eq!(
writer
.query_row("SELECT value FROM future_feature", [], |row| row
.get::<_, String>(0))
.unwrap(),
"preserve me"
);
let state = read_schema_state(&writer).unwrap();
assert_eq!(state.revision, SCHEMA_VERSION + 1);
assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
}
#[test]
fn invalid_compatibility_metadata_refuses_readers_and_writers() {
for alteration in [
"DROP TABLE schema_compatibility",
"DELETE FROM schema_compatibility",
"PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
"UPDATE schema_compatibility SET minimum_compatible_version = 99999",
"PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
"PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
"DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
"DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
] {
for future in [false, true] {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
drop(open_writer(&path).unwrap());
if future {
advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
}
let raw = Connection::open(&path).unwrap();
raw.execute_batch(alteration).unwrap();
let before: i64 = raw
.query_row("PRAGMA schema_version", [], |row| row.get(0))
.unwrap();
for error in [
open_reader_strict(&path).unwrap_err(),
open_writer(&path).unwrap_err(),
] {
let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
assert_eq!(
mismatch.reason,
StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
"{alteration}"
);
}
forget_verified_schema(&path);
assert!(open_writer(&path).is_err(), "{alteration}");
let after: i64 = raw
.query_row("PRAGMA schema_version", [], |row| row.get(0))
.unwrap();
assert_eq!(
before, after,
"a rejected open repaired schema: {alteration}"
);
}
}
}
#[test]
fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
let connection = Connection::open(&path).unwrap();
connection
.execute_batch("CREATE TABLE workspaces(conflict TEXT)")
.unwrap();
let error = migrate_schema(&connection).unwrap_err();
assert!(format!("{error:#}").contains("create baseline database schema"));
assert!(
connection.is_autocommit(),
"the failed baseline left a transaction open"
);
assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
let tables: i64 = connection
.query_row(
"SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
[],
|row| row.get(0),
)
.unwrap();
assert_eq!(tables, 1, "only the conflicting table remains");
connection.execute_batch("DROP TABLE workspaces").unwrap();
drop(connection);
let writer = open_writer(&path).unwrap();
let state = read_schema_state(&writer).unwrap();
assert_eq!(state.revision, SCHEMA_VERSION);
assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
}
#[test]
fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
drop(open_writer(&path).unwrap());
stamp_schema_version(&path, SCHEMA_VERSION + 1);
let error = open_reader_strict(&path).unwrap_err();
let mismatch = error
.chain()
.find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
.expect("the reader reports the mismatch as a typed cause");
assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
assert_eq!(mismatch.supported, SCHEMA_VERSION);
let message = mismatch.to_string();
assert!(message.contains("upgrade Mjolnir"), "got {message}");
assert!(
!message.contains("start the Mjolnir daemon"),
"got {message}"
);
}
#[test]
fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
drop(open_writer(&path).unwrap());
let raw = Connection::open(&path).unwrap();
raw.execute_batch(&format!(
"UPDATE schema_compatibility SET minimum_compatible_version = {0};
DELETE FROM schema_migrations WHERE version > {0};
INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
PRAGMA user_version = {0};",
SCHEMA_VERSION - 1
))
.unwrap();
drop(raw);
let error = open_reader_strict(&path).unwrap_err();
let mismatch = error
.chain()
.find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
.expect("the reader reports the mismatch as a typed cause");
assert_eq!(
mismatch.to_string(),
format!(
"Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
start the Mjolnir daemon to migrate it",
SCHEMA_VERSION - 1
)
);
}
#[test]
fn strict_reader_rejects_mutation() {
let directory = tempfile::tempdir().unwrap();
let path = directory.path().join("mj.sqlite3");
drop(open_writer(&path).unwrap());
let reader = open_reader_strict(&path).unwrap();
let error = reader
.execute("CREATE TABLE forbidden(value TEXT)", [])
.unwrap_err();
assert!(
matches!(
error.sqlite_error_code(),
Some(rusqlite::ErrorCode::ReadOnly)
),
"unexpected mutation error: {error}"
);
}
}