Skip to main content

mj_controller/database/
schema.rs

1use super::*;
2use rusqlite::OpenFlags;
3
4const COMPATIBILITY_METADATA_VERSION: i64 = 30;
5
6pub(super) struct SchemaState {
7    pub(super) revision: i64,
8    minimum_compatible: Option<i64>,
9}
10
11impl SchemaState {
12    pub(super) fn ensure_supported(&self) -> Result<()> {
13        self.ensure_supported_by(SCHEMA_VERSION)
14    }
15
16    fn ensure_supported_by(&self, supported: i64) -> Result<()> {
17        let reason = if self.revision < supported {
18            StoreSchemaMismatchReason::NeedsMigration
19        } else if let Some(minimum_compatible) = self.minimum_compatible {
20            if minimum_compatible <= supported {
21                return Ok(());
22            }
23            StoreSchemaMismatchReason::Incompatible { minimum_compatible }
24        } else {
25            StoreSchemaMismatchReason::InvalidCompatibilityMetadata
26        };
27        Err(StoreSchemaMismatch {
28            found: self.revision,
29            supported,
30            reason,
31        }
32        .into())
33    }
34}
35
36/// The revision, ledger, and compatibility floor must describe one snapshot.
37/// A missing floor is only legitimate before compatibility was introduced.
38pub(super) fn read_schema_state(connection: &Connection) -> Result<SchemaState> {
39    let snapshot = connection
40        .unchecked_transaction()
41        .context("start database compatibility snapshot")?;
42    let revision: i64 = snapshot
43        .query_row("PRAGMA user_version", [], |row| row.get(0))
44        .context("read database migration revision")?;
45    let minimum_compatible = if revision >= COMPATIBILITY_METADATA_VERSION {
46        let invalid = || StoreSchemaMismatch {
47            found: revision,
48            supported: SCHEMA_VERSION,
49            reason: StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
50        };
51        let (count, singleton, floor, recorded): (i64, Option<i64>, Option<i64>, Option<i64>) =
52            snapshot
53                .query_row(
54                    "SELECT count(*), min(singleton), min(minimum_compatible_version),
55                    (SELECT max(version) FROM schema_migrations)
56             FROM schema_compatibility",
57                    [],
58                    |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
59                )
60                .map_err(|error| {
61                    // Missing tables/columns and invalid field types are
62                    // structural. Busy, I/O, and interruption errors are not
63                    // evidence of an incompatible migration.
64                    let structural = match &error {
65                        rusqlite::Error::SqliteFailure(code, _) => {
66                            code.code == rusqlite::ErrorCode::Unknown
67                        }
68                        _ => true,
69                    };
70                    let error = anyhow::Error::new(error);
71                    if structural {
72                        error.context(invalid())
73                    } else {
74                        error.context("read database compatibility metadata")
75                    }
76                })?;
77        if count != 1
78            || singleton != Some(1)
79            || recorded != Some(revision)
80            || !floor
81                .is_some_and(|floor| (COMPATIBILITY_METADATA_VERSION..=revision).contains(&floor))
82        {
83            return Err(invalid().into());
84        }
85        floor
86    } else {
87        None
88    };
89    snapshot
90        .commit()
91        .context("finish database compatibility snapshot")?;
92    Ok(SchemaState {
93        revision,
94        minimum_compatible,
95    })
96}
97
98pub fn database_path() -> PathBuf {
99    data_dir().join("mj.sqlite3")
100}
101
102/// Verify that this client can read the daemon's store without creating or
103/// migrating it. Startup must pass this gate before handing out a connection,
104/// even when an older daemon happens to speak the same wire protocol.
105pub fn check_read_compatibility() -> Result<()> {
106    open_reader_strict(&database_path()).map(drop)
107}
108
109/// A writer-capable connection whose transactions take the WAL write lock at
110/// `BEGIN`, where the busy handler applies. A DEFERRED transaction that has
111/// already read cannot wait: SQLite only calls the busy handler when the
112/// connection holds no transaction, so the upgrade to a write returns
113/// `SQLITE_BUSY` at once (issue 1117).
114pub(super) fn open_writer(path: &Path) -> Result<Connection> {
115    let mut connection = open_writable(path)?;
116    connection.set_transaction_behavior(rusqlite::TransactionBehavior::Immediate);
117    Ok(connection)
118}
119
120fn open_writable(path: &Path) -> Result<Connection> {
121    // A writable open migrates the store, so this is where the decision
122    // belongs: test binaries reach it without the controller lock.
123    if let Some(store) = path.parent() {
124        mj_core::config::ensure_may_control_store(store, "open this database for writing")?;
125    }
126    if let Some(parent) = path.parent() {
127        fs::create_dir_all(parent)
128            .with_context(|| format!("create Mjolnir data directory {}", parent.display()))?;
129    }
130    let connection = Connection::open(path)
131        .with_context(|| format!("open Mjolnir database {}", path.display()))?;
132    connection.busy_timeout(Duration::from_secs(5))?;
133    connection.execute_batch(
134        "PRAGMA foreign_keys = ON;
135         PRAGMA journal_mode = WAL;
136         PRAGMA synchronous = FULL;",
137    )?;
138    verify_schema_once(path, &connection)?;
139    committed::observe_connection(&connection, path)?;
140    Ok(connection)
141}
142
143pub(super) fn open(path: &Path) -> Result<Connection> {
144    open_writer(path)
145}
146
147/// Open an existing database without permitting schema or data mutation.
148/// Client processes use this path so an accidental write fails locally
149/// instead of competing with the daemon's writer.
150#[cfg(not(test))]
151pub(super) fn open_reader(path: &Path) -> Result<Connection> {
152    open_reader_strict(path)
153}
154
155#[cfg(test)]
156pub(super) fn open_reader(path: &Path) -> Result<Connection> {
157    // Path-taking database helpers are migration fixtures in unit tests: they
158    // intentionally open old or not-yet-created schemas. Production query
159    // entry points compile against the strict reader above. The connection is
160    // writable but keeps SQLite's DEFERRED default, so a fixture read does not
161    // take the write lock.
162    open_writable(path)
163}
164
165#[cfg_attr(test, allow(dead_code))]
166fn open_reader_strict(path: &Path) -> Result<Connection> {
167    let connection = Connection::open_with_flags(
168        path,
169        OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
170    )
171    .with_context(|| format!("open Mjolnir database read-only {}", path.display()))?;
172    connection.busy_timeout(Duration::from_secs(5))?;
173    connection.execute_batch(
174        "PRAGMA foreign_keys = ON;
175         PRAGMA query_only = ON;",
176    )?;
177    read_schema_state(&connection)?.ensure_supported()?;
178    Ok(connection)
179}
180
181/// Databases this process has already migrated. A controller owns its store
182/// exclusively (`ControllerStoreGuard`), so a schema verified once stays
183/// verified and later connections skip the migration probes entirely.
184fn verified_schemas() -> &'static Mutex<HashSet<PathBuf>> {
185    static VERIFIED: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
186    VERIFIED.get_or_init(|| Mutex::new(HashSet::new()))
187}
188
189/// Stable cache identity for a database. The file itself may not exist yet, so
190/// the canonicalized parent directory carries the identity.
191fn schema_cache_key(path: &Path) -> PathBuf {
192    let Some(parent) = path
193        .parent()
194        .filter(|parent| !parent.as_os_str().is_empty())
195    else {
196        return path.to_owned();
197    };
198    match (fs::canonicalize(parent), path.file_name()) {
199        (Ok(canonical), Some(name)) => canonical.join(name),
200        _ => path.to_owned(),
201    }
202}
203
204/// Run the migration ladder the first time this process opens a database.
205/// Later opens confirm compatibility without repeating schema repairs. A
206/// database behind this build is migrated again, so a recreated file under a
207/// reused path still converges. Compatible future stores are never repaired.
208fn verify_schema_once(path: &Path, connection: &Connection) -> Result<()> {
209    let key = schema_cache_key(path);
210    let mut verified = verified_schemas()
211        .lock()
212        .unwrap_or_else(PoisonError::into_inner);
213    let state = read_schema_state(connection)?;
214    if state.revision > SCHEMA_VERSION
215        || (state.revision == SCHEMA_VERSION && verified.contains(&key))
216    {
217        // An older build must never run its repairs against a newer schema.
218        return state.ensure_supported();
219    }
220    // Holding the lock across the ladder keeps two first opens of the same
221    // database from running the additive migration steps against each other.
222    migrate_schema(connection)?;
223    read_schema_state(connection)?.ensure_supported()?;
224    verified.insert(key);
225    Ok(())
226}
227
228/// Forget that this process verified a database's schema. Only tests need it:
229/// they simulate a store written by an older build by editing the schema of a
230/// database this process has already opened, which no controller can do.
231#[cfg(test)]
232pub(super) fn forget_verified_schema(path: &Path) {
233    verified_schemas()
234        .lock()
235        .unwrap_or_else(PoisonError::into_inner)
236        .remove(&schema_cache_key(path));
237}
238
239/// New stores start at the revision Mjolnir 2.7.2 shipped. Existing stores
240/// retain the complete migration path, even when users skip many releases.
241const BASELINE_SCHEMA_VERSION: i64 = 33;
242
243/// The compatibility floor a baseline store records. Migration 32 (ZCode) was
244/// the last breaking change before the baseline.
245const BASELINE_MINIMUM_COMPATIBLE_VERSION: i64 = 32;
246
247// Data changes from the original accounting revision 67, also needed when
248// upgrading the independently published project-catalog revision 67.
249const ACCOUNTING_MIGRATION_SQL: &str = "
250            ALTER TABLE session_turn_usage RENAME TO old_session_turn_usage;
251            CREATE TABLE session_turn_usage (
252                session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
253                command_id TEXT NOT NULL,
254                completed_ordinal INTEGER NOT NULL,
255                turn_start_position INTEGER,
256                body TEXT NOT NULL,
257                PRIMARY KEY(session_id, command_id)
258            );
259            INSERT INTO session_turn_usage SELECT * FROM old_session_turn_usage;
260            DROP TABLE old_session_turn_usage;
261            CREATE INDEX session_turn_usage_order ON session_turn_usage(session_id, completed_ordinal);
262            ALTER TABLE session_provider_cost RENAME TO old_session_provider_cost;
263            CREATE TABLE session_provider_cost (
264                session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
265                body TEXT NOT NULL
266            );
267            INSERT INTO session_provider_cost SELECT * FROM old_session_provider_cost;
268            DROP TABLE old_session_provider_cost;
269            CREATE TABLE subagent_accounting (
270                child_session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
271                parent_session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
272                task_name TEXT NOT NULL,
273                CHECK(child_session_id <> parent_session_id)
274            ) STRICT;
275            CREATE INDEX subagent_accounting_parent ON subagent_accounting(parent_session_id);
276            INSERT INTO subagent_accounting SELECT child_session_id, parent_session_id,
277                json_extract(record_json, '$.task_name') FROM subagent_sessions;
278            CREATE TABLE session_turn_selections (
279                session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
280                command_id TEXT NOT NULL,
281                model TEXT,
282                effort TEXT,
283                PRIMARY KEY(session_id, command_id)
284            ) STRICT;
285";
286
287fn migrate_schema(connection: &Connection) -> Result<()> {
288    let state = read_schema_state(connection)?;
289    let version = state.revision;
290    if version > SCHEMA_VERSION {
291        return state.ensure_supported();
292    }
293    if version == 0 {
294        create_baseline_schema(connection)?;
295    } else if version < BASELINE_SCHEMA_VERSION {
296        super::legacy_schema::migrate_to_baseline(connection)
297            .context("upgrade historical database schema")?;
298    }
299    // Compatible: adds one table. Older readers ignore it and treat read-write
300    // mounts as copy-on-write, a behaviour difference rather than lost data.
301    // Older writers rewrite `session_mounts` but never touch this table, so its
302    // rows survive their updates; a row only applies while a mount with the same
303    // source and destination is still not read-only, so an older build that
304    // makes the mount read-only or removes it keeps that choice. The
305    // compatibility floor stays where it is.
306    if version < 34 {
307        connection.execute_batch(
308            "BEGIN IMMEDIATE;
309             CREATE TABLE IF NOT EXISTS session_mount_access (
310                 session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
311                 source BLOB NOT NULL,
312                 destination BLOB NOT NULL,
313                 access TEXT NOT NULL CHECK(access IN ('rw')),
314                 PRIMARY KEY(session_id, destination)
315             ) STRICT;
316             INSERT INTO schema_migrations(version, applied_at)
317                 VALUES (34, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
318             PRAGMA user_version = 34;
319             COMMIT;",
320        )?;
321    }
322    // Compatible: adds one nullable column. Older readers ignore it, and the
323    // older writer's session upsert lists columns explicitly, so it preserves
324    // the value. An older executable launching such a session uses the shared
325    // `/workspace` instead of the recorded per-session path, which is a
326    // behaviour difference, not data loss. The compatibility floor stays where
327    // it is.
328    if version < 35 {
329        connection.execute_batch(
330            "BEGIN IMMEDIATE;
331             ALTER TABLE sessions ADD COLUMN container_workspace TEXT;
332             INSERT INTO schema_migrations(version, applied_at)
333                 VALUES (35, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
334             PRAGMA user_version = 35;
335             COMMIT;",
336        )?;
337    }
338    // Compatible: adds one nullable column. Older readers ignore it, and the
339    // older writer's session upsert lists columns explicitly, so it preserves
340    // the value. An older executable launching such a session runs it without
341    // the mbx build cache, which is a behaviour difference, not data loss. The
342    // compatibility floor stays where it is.
343    if version < 36 {
344        connection.execute_batch(
345            "BEGIN IMMEDIATE;
346             ALTER TABLE sessions ADD COLUMN build_cache_json TEXT;
347             INSERT INTO schema_migrations(version, applied_at)
348                 VALUES (36, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
349             PRAGMA user_version = 36;
350             COMMIT;",
351        )?;
352    }
353    // Compatible: adds one nullable column to `session_targets`. Only a
354    // container sub-agent child row ever carries a value, and older builds
355    // could never start such a child, so an older update that rewrites the row
356    // without the column loses nothing usable. Older readers ignore it. The
357    // compatibility floor stays where it is.
358    if version < 37 {
359        connection.execute_batch(
360            "BEGIN IMMEDIATE;
361             ALTER TABLE session_targets ADD COLUMN borrowed_from TEXT;
362             INSERT INTO schema_migrations(version, applied_at)
363                 VALUES (37, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
364             PRAGMA user_version = 37;
365             COMMIT;",
366        )?;
367    }
368    // Compatible: adds one table holding the dashboard's conversation pane
369    // arrangement per workspace. Older readers never select from it and older
370    // writers never touch it, so their updates leave its rows intact; a
371    // workspace deleted by an older build still removes them through the
372    // foreign key. Losing the table only means the conversation area opens as
373    // a single pane. The compatibility floor stays where it is.
374    if version < 38 {
375        connection.execute_batch(
376            "BEGIN IMMEDIATE;
377             CREATE TABLE IF NOT EXISTS workspace_layouts (
378                 workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
379                 layout TEXT NOT NULL
380             ) STRICT;
381             INSERT INTO schema_migrations(version, applied_at)
382                 VALUES (38, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
383             PRAGMA user_version = 38;
384             COMMIT;",
385        )?;
386    }
387    // Breaking: persisted input_required API events can now omit the structured
388    // request. Older readers require it and fail to deserialize the event log;
389    // the shared read/write compatibility floor must advance with the revision.
390    if version < 39 {
391        connection.execute_batch(
392            "BEGIN IMMEDIATE;
393             UPDATE schema_compatibility SET minimum_compatible_version = 39 WHERE singleton = 1;
394             INSERT INTO schema_migrations(version, applied_at)
395                 VALUES (39, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
396             PRAGMA user_version = 39;
397             COMMIT;",
398        )?;
399    }
400    // Breaking: native child projections and relay observations must be preserved
401    // by every reader/writer; older builds cannot interpret their lifecycle.
402    if version < 40 {
403        connection.execute_batch(
404            "BEGIN IMMEDIATE;
405             CREATE TABLE native_agents (
406                 owner TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
407                 child TEXT NOT NULL,
408                 staging INTEGER NOT NULL CHECK(staging IN (0,1)),
409                 body TEXT NOT NULL CHECK(json_valid(body)),
410                 PRIMARY KEY(owner, child, staging)
411             ) STRICT;
412             CREATE TABLE native_agent_transcript (
413                 owner TEXT NOT NULL,
414                 child TEXT NOT NULL,
415                 staging INTEGER NOT NULL,
416                 stable_id TEXT NOT NULL,
417                 position INTEGER NOT NULL,
418                 body TEXT NOT NULL CHECK(json_valid(body)),
419                 PRIMARY KEY(owner, child, staging, stable_id),
420                 FOREIGN KEY(owner, child, staging) REFERENCES native_agents(owner, child, staging)
421                     ON DELETE CASCADE ON UPDATE CASCADE
422             ) STRICT;
423             CREATE INDEX native_agent_transcript_position ON native_agent_transcript(owner, child, staging, position);
424             CREATE TABLE native_agent_replay (
425                 owner TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE
426             ) STRICT;
427             UPDATE schema_compatibility SET minimum_compatible_version = 40 WHERE singleton = 1;
428             INSERT INTO schema_migrations(version, applied_at)
429                 VALUES (40, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
430             PRAGMA user_version = 40;
431             COMMIT;",
432        )?;
433    }
434    // Breaking: clear-context relay commands/outcomes are persisted in event
435    // JSON. Older readers cannot decode them or honor the context boundary.
436    if version < 41 {
437        connection.execute_batch(
438            "BEGIN IMMEDIATE;
439             UPDATE schema_compatibility SET minimum_compatible_version = 41 WHERE singleton = 1;
440             INSERT INTO schema_migrations(version, applied_at)
441                 VALUES (41, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
442             PRAGMA user_version = 41;
443             COMMIT;",
444        )?;
445    }
446
447    // Breaking: older layout writers discard Browse identity and pin badges.
448    if version < 42 {
449        connection.execute_batch(
450            "BEGIN IMMEDIATE;
451             UPDATE schema_compatibility SET minimum_compatible_version = 42 WHERE singleton = 1;
452             INSERT INTO schema_migrations(version, applied_at)
453                 VALUES (42, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
454             PRAGMA user_version = 42;
455             COMMIT;",
456        )?;
457    }
458
459    // Breaking: durable steering commands/observations and native availability
460    // cannot be interpreted or preserved by older readers and writers.
461    if version < 43 {
462        connection.execute_batch("BEGIN IMMEDIATE;
463            UPDATE schema_compatibility SET minimum_compatible_version = 43 WHERE singleton = 1;
464            INSERT INTO schema_migrations(version, applied_at) VALUES (43, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
465            PRAGMA user_version = 43;
466            COMMIT;")?;
467    }
468
469    // Breaking: new quota recovery commands in stored relay JSON cannot be
470    // read or preserved by older binaries, even though the cache is additive.
471    if version < 44 {
472        connection.execute_batch("BEGIN IMMEDIATE;
473            CREATE TABLE quota_reset_cache (identity TEXT PRIMARY KEY, body TEXT NOT NULL);
474            UPDATE schema_compatibility SET minimum_compatible_version = 44 WHERE singleton = 1;
475            INSERT INTO schema_migrations(version, applied_at) VALUES (44, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
476            PRAGMA user_version = 44;
477            COMMIT;")?;
478    }
479
480    // Compatible: adds one nullable column. Older readers ignore it, and the
481    // older writer's session upsert lists columns explicitly, so it preserves
482    // the value. An older executable relaunching such a session starts it at
483    // HEAD or the remote default branch instead of the recorded revision,
484    // which is a behaviour difference, not data loss. The compatibility floor
485    // stays where it is.
486    if version < 45 {
487        // The column is added only when it is absent, the way migration 34
488        // creates its table only when absent: a store rolled back to an older
489        // revision still carries the column, and a second ALTER would refuse.
490        let add_column =
491            match super::legacy_schema::table_has_column(connection, "sessions", "launch_base")? {
492                true => "",
493                false => "ALTER TABLE sessions ADD COLUMN launch_base TEXT;",
494            };
495        connection.execute_batch(&format!(
496            "BEGIN IMMEDIATE;
497             {add_column}
498             INSERT INTO schema_migrations(version, applied_at)
499                 VALUES (45, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
500             PRAGMA user_version = 45;
501             COMMIT;"
502        ))?;
503    }
504
505    // Breaking: older writers can replace a target without updating its saved
506    // connection, leaving access metadata attached to the wrong resource.
507    // Refuse both older readers and writers before they operate that target.
508    if version < 46 {
509        let add_column = if super::legacy_schema::table_has_column(
510            connection,
511            "sessions",
512            "target_runtime_json",
513        )? {
514            ""
515        } else {
516            "ALTER TABLE sessions ADD COLUMN target_runtime_json TEXT;"
517        };
518        connection.execute_batch(&format!(
519            "BEGIN IMMEDIATE;
520             {add_column}
521             UPDATE schema_compatibility SET minimum_compatible_version = 46 WHERE singleton = 1;
522             INSERT INTO schema_migrations(version, applied_at)
523                 VALUES (46, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
524             PRAGMA user_version = 46;
525             COMMIT;"
526        ))?;
527    }
528
529    // Breaking: the managed checkout JSON can now describe an independent
530    // clone. Older readers reject its `kind`, and older writers cannot safely
531    // retain that clone's unpublished commits during lifecycle cleanup.
532    if version < 47 {
533        let add_branch =
534            if super::legacy_schema::table_has_column(connection, "sessions", "launch_branch")? {
535                ""
536            } else {
537                "ALTER TABLE sessions ADD COLUMN launch_branch TEXT;"
538            };
539        let add_publication = if super::legacy_schema::table_has_column(
540            connection,
541            "sessions",
542            "publication_json",
543        )? {
544            ""
545        } else {
546            "ALTER TABLE sessions ADD COLUMN publication_json TEXT;"
547        };
548        connection.execute_batch(&format!(
549            "BEGIN IMMEDIATE;
550             {add_branch}
551             {add_publication}
552             UPDATE schema_compatibility SET minimum_compatible_version = 47 WHERE singleton = 1;
553             INSERT INTO schema_migrations(version, applied_at)
554                 VALUES (47, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
555             PRAGMA user_version = 47;
556             COMMIT;",
557        ))?;
558    }
559
560    // Breaking: stored relay JSON can now hold the steering-returned outcome,
561    // and stored active-turn JSON the prompt a steer continued, which older
562    // readers reject and older writers would drop.
563    if version < 48 {
564        connection.execute_batch(
565            "BEGIN IMMEDIATE;
566             UPDATE schema_compatibility SET minimum_compatible_version = 48 WHERE singleton = 1;
567             INSERT INTO schema_migrations(version, applied_at)
568                 VALUES (48, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
569             PRAGMA user_version = 48;
570             COMMIT;",
571        )?;
572    }
573
574    // Breaking: a sub-agent record can now say its child has the `handback`
575    // tool, and older readers refuse the unknown field. The reports themselves
576    // live in their own table, because a full state save rewrites every
577    // sub-agent record from whatever copy the saving controller holds.
578    if version < 49 {
579        connection.execute_batch(
580            "BEGIN IMMEDIATE;
581             CREATE TABLE IF NOT EXISTS subagent_handbacks (
582                 child_session_id TEXT PRIMARY KEY,
583                 handback_command_id TEXT,
584                 handback_message TEXT,
585                 handback_recorded_at_ms INTEGER,
586                 reminder_command_id TEXT,
587                 reminder_for_command_id TEXT,
588                 reminder_sent_at_ms INTEGER,
589                 reminder_failed_for_command_id TEXT
590             );
591             UPDATE schema_compatibility SET minimum_compatible_version = 49 WHERE singleton = 1;
592             INSERT INTO schema_migrations(version, applied_at)
593                 VALUES (49, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
594             PRAGMA user_version = 49;
595             COMMIT;",
596        )?;
597    }
598
599    // Compatible: a new nullable column that only this build reads. Older
600    // readers ignore it, and older writers name their columns, so an update
601    // from one keeps it.
602    if version < 50 {
603        let add_column = if super::legacy_schema::table_has_column(
604            connection,
605            "subagent_handbacks",
606            "awaited_ordinal",
607        )? {
608            ""
609        } else {
610            "ALTER TABLE subagent_handbacks ADD COLUMN awaited_ordinal INTEGER;"
611        };
612        connection.execute_batch(&format!(
613            "BEGIN IMMEDIATE;
614             {add_column}
615             INSERT INTO schema_migrations(version, applied_at)
616                 VALUES (50, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
617             PRAGMA user_version = 50;
618             COMMIT;"
619        ))?;
620    }
621
622    // Compatible: the directory on the parent's target where a child writes
623    // the details its report points to. A nullable column only this build
624    // reads, like `awaited_ordinal`.
625    if version < 51 {
626        let add_column = if super::legacy_schema::table_has_column(
627            connection,
628            "subagent_handbacks",
629            "report_dir",
630        )? {
631            ""
632        } else {
633            "ALTER TABLE subagent_handbacks ADD COLUMN report_dir TEXT;"
634        };
635        connection.execute_batch(&format!(
636            "BEGIN IMMEDIATE;
637             {add_column}
638             INSERT INTO schema_migrations(version, applied_at)
639                 VALUES (51, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
640             PRAGMA user_version = 51;
641             COMMIT;"
642        ))?;
643    }
644
645    // Compatible: adds one table that only this build reads. It lists the
646    // sub-agents a parent's suspend stopped until the parent's model has been
647    // told. It is not a column of `sessions` because a full state save
648    // rewrites every session row from whatever copy the saving controller
649    // holds, which could drop the list or bring back one already delivered.
650    // Rows go with their parent's session row.
651    if version < 52 {
652        connection.execute_batch(
653            "BEGIN IMMEDIATE;
654             CREATE TABLE IF NOT EXISTS stopped_subagents (
655                 parent_session_id TEXT NOT NULL
656                     REFERENCES sessions(session_id) ON DELETE CASCADE,
657                 child_session_id TEXT NOT NULL,
658                 record_json TEXT NOT NULL CHECK(json_valid(record_json)),
659                 PRIMARY KEY(parent_session_id, child_session_id)
660             ) STRICT;
661             INSERT INTO schema_migrations(version, applied_at)
662                 VALUES (52, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
663             PRAGMA user_version = 52;
664             COMMIT;",
665        )?;
666    }
667
668    // Breaking: a sub-agent whose turn ended can be stored as `parked`. The
669    // `sessions.state` constraint has to be rebuilt to admit it, and an older
670    // reader panics on a state it does not know, so the floor rises with the
671    // revision.
672    if version < 53 {
673        migrate_parked_session_state(connection)?;
674    }
675
676    // Breaking: older daemons cannot enforce exact checkout preparation on
677    // admitted sessions, so they must not recover or provision these records.
678    if version < 54 {
679        let add_column =
680            if super::legacy_schema::table_has_column(connection, "sessions", "checkout_json")? {
681                ""
682            } else {
683                "ALTER TABLE sessions ADD COLUMN checkout_json TEXT;"
684            };
685        connection.execute_batch(&format!(
686            "BEGIN IMMEDIATE;
687             {add_column}
688             UPDATE schema_compatibility SET minimum_compatible_version = 54 WHERE singleton = 1;
689             INSERT INTO schema_migrations(version, applied_at)
690                 VALUES (54, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
691             PRAGMA user_version = 54;
692             COMMIT;"
693        ))?;
694    }
695
696    // Breaking: older daemons cannot enforce the runtime constraint or read
697    // the new runtime receipt stored in relay/API events.
698    if version < 55 {
699        let add_column = if super::legacy_schema::table_has_column(
700            connection,
701            "sessions",
702            "expected_runtime_identity",
703        )? {
704            ""
705        } else {
706            "ALTER TABLE sessions ADD COLUMN expected_runtime_identity TEXT;"
707        };
708        connection.execute_batch(&format!(
709            "BEGIN IMMEDIATE;
710             {add_column}
711             UPDATE schema_compatibility SET minimum_compatible_version = 55 WHERE singleton = 1;
712             INSERT INTO schema_migrations(version, applied_at)
713                 VALUES (55, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
714             PRAGMA user_version = 55;
715             COMMIT;"
716        ))?;
717    }
718
719    // Breaking: older readers/writers cannot honor None or fixed model/effort.
720    if version < 56 {
721        let add_column = if super::legacy_schema::table_has_column(
722            connection,
723            "sessions",
724            "subagents",
725        )? {
726            ""
727        } else {
728            "ALTER TABLE sessions ADD COLUMN subagents TEXT CHECK(subagents IS NULL OR json_valid(subagents));"
729        };
730        connection.execute_batch(&format!(
731            "BEGIN IMMEDIATE;
732             {add_column}
733             UPDATE sessions SET subagents = CASE WHEN mjolnir_subagents = 1
734                 THEN '{{\"mode\":\"all_models\"}}' ELSE '{{\"mode\":\"native\"}}' END WHERE subagents IS NULL AND mjolnir_subagents IS NOT NULL;
735             CREATE TABLE IF NOT EXISTS subagent_preference (singleton INTEGER PRIMARY KEY CHECK(singleton = 1), policy TEXT NOT NULL CHECK(json_valid(policy)));
736             UPDATE schema_compatibility SET minimum_compatible_version = 56 WHERE singleton = 1;
737             INSERT INTO schema_migrations(version, applied_at) VALUES (56, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
738             PRAGMA user_version = 56;
739             COMMIT;"
740        ))?;
741    }
742
743    // Breaking: typed event bodies cannot be decoded by older readers, and
744    // pending checkpoint identities must survive every writer's transitions.
745    if version < 57 {
746        let tx = connection.unchecked_transaction()?;
747        tx.execute_batch("DROP TRIGGER IF EXISTS api_session_error_updated;
748            DROP TRIGGER IF EXISTS api_session_error_inserted;
749            CREATE TABLE IF NOT EXISTS checkpoint_operations (
750                session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
751                command_id TEXT NOT NULL UNIQUE,
752                related_command_ids TEXT NOT NULL DEFAULT '[]' CHECK(json_valid(related_command_ids))
753            ) STRICT;")?;
754        super::events::migrate_event_outcomes(&tx)?;
755        tx.execute_batch("UPDATE schema_compatibility SET minimum_compatible_version = 57 WHERE singleton = 1;
756            INSERT INTO schema_migrations(version, applied_at) VALUES (57, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
757            PRAGMA user_version = 57;")?;
758        tx.commit()?;
759    }
760
761    // Breaking: relay assessment observations and seeded-context commands are
762    // persisted as JSON that older readers cannot decode or preserve.
763    if version < 58 {
764        connection.execute_batch("BEGIN IMMEDIATE;
765            UPDATE schema_compatibility SET minimum_compatible_version = 58 WHERE singleton = 1;
766            INSERT INTO schema_migrations(version, applied_at) VALUES (58, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
767            PRAGMA user_version = 58;
768            COMMIT;")?;
769    }
770
771    // Breaking: an older daemon cannot resume acknowledged startup delivery
772    // and may instead restore or resend its prompt with a different identity.
773    if version < 59 {
774        connection.execute_batch("BEGIN IMMEDIATE;
775            CREATE TABLE IF NOT EXISTS startup_steps (
776                sequence INTEGER PRIMARY KEY AUTOINCREMENT,
777                session_id TEXT NOT NULL,
778                group_id TEXT,
779                command_id TEXT NOT NULL UNIQUE,
780                step_json TEXT NOT NULL,
781                phase TEXT NOT NULL DEFAULT 'pending'
782                    CHECK (phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting', 'done', 'failed', 'dismissed')),
783                error TEXT,
784                accepted_ordinal INTEGER
785            ) STRICT;
786            CREATE INDEX IF NOT EXISTS startup_steps_session_sequence
787                ON startup_steps(session_id, sequence);
788            CREATE INDEX IF NOT EXISTS startup_steps_group_sequence
789                ON startup_steps(group_id, sequence) WHERE group_id IS NOT NULL;
790            CREATE INDEX IF NOT EXISTS startup_steps_pending
791                ON startup_steps(session_id, sequence)
792                WHERE phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting');
793            UPDATE schema_compatibility SET minimum_compatible_version = 59 WHERE singleton = 1;
794            INSERT INTO schema_migrations(version, applied_at) VALUES (59, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
795            PRAGMA user_version = 59;
796            COMMIT;")?;
797    }
798    // Breaking: older dispatchers would resolve a prepared request's target
799    // again instead of honoring its persisted effect identity and outcome.
800    if version < 60 {
801        connection.execute_batch("BEGIN IMMEDIATE;
802            CREATE TABLE IF NOT EXISTS delegation_effects (
803                parent_session_id TEXT NOT NULL,
804                request_id TEXT NOT NULL,
805                phase TEXT NOT NULL,
806                prepared_json TEXT,
807                result_json TEXT,
808                receipt_pending INTEGER NOT NULL DEFAULT 0,
809                PRIMARY KEY(parent_session_id, request_id)
810            ) STRICT;
811            CREATE INDEX IF NOT EXISTS delegation_effects_receipt_cleanup
812                ON delegation_effects(parent_session_id, request_id) WHERE receipt_pending = 1;
813            UPDATE schema_compatibility SET minimum_compatible_version = 60 WHERE singleton = 1;
814            INSERT INTO schema_migrations(version, applied_at) VALUES (60, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
815            PRAGMA user_version = 60;
816            COMMIT;")?;
817    }
818    // Breaking: older review hosts clear active reviews on startup and cannot
819    // reconstruct or preserve the durable orchestration and effect outbox.
820    if version < 61 {
821        let add_column = if super::legacy_schema::table_has_column(
822            connection,
823            "turn_review_state",
824            "orchestration",
825        )? {
826            ""
827        } else {
828            "ALTER TABLE turn_review_state ADD COLUMN orchestration TEXT;"
829        };
830        connection.execute_batch(&format!("BEGIN IMMEDIATE;
831            {add_column}
832            UPDATE schema_compatibility SET minimum_compatible_version = 61 WHERE singleton = 1;
833            INSERT INTO schema_migrations(version, applied_at) VALUES (61, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
834            PRAGMA user_version = 61;
835            COMMIT;"))?;
836    }
837    // Breaking: older recovery cannot distinguish an accepted worker swap
838    // awaiting readiness from a worker that should be recreated or abandoned.
839    if version < 62 {
840        connection.execute_batch("BEGIN IMMEDIATE;
841            CREATE TABLE IF NOT EXISTS worker_restart_intents (
842                session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
843                operation_id TEXT NOT NULL,
844                target_json TEXT NOT NULL,
845                desired_build TEXT NOT NULL,
846                phase TEXT NOT NULL CHECK(phase IN ('prepared', 'swapping', 'awaiting_readiness'))
847            ) STRICT;
848            UPDATE schema_compatibility SET minimum_compatible_version = 62 WHERE singleton = 1;
849            INSERT INTO schema_migrations(version, applied_at) VALUES (62, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
850            PRAGMA user_version = 62;
851            COMMIT;")?;
852    }
853
854    // Breaking: older delegation executors do not fence delayed lifecycle
855    // effects against the durable incarnation they originally selected.
856    if version < 63 {
857        connection.execute_batch("BEGIN IMMEDIATE;
858            CREATE TABLE IF NOT EXISTS session_incarnations (
859                session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
860                identity TEXT NOT NULL
861            ) STRICT;
862            INSERT OR IGNORE INTO session_incarnations(session_id, identity)
863                SELECT session_id, lower(hex(randomblob(16))) FROM sessions;
864            CREATE TRIGGER IF NOT EXISTS session_incarnation_insert
865                AFTER INSERT ON sessions BEGIN
866                    INSERT INTO session_incarnations(session_id, identity)
867                    VALUES(NEW.session_id, lower(hex(randomblob(16))));
868                END;
869            CREATE TRIGGER IF NOT EXISTS session_incarnation_resume
870                AFTER UPDATE OF state ON sessions
871                WHEN (NEW.state = 'provisioning' AND OLD.state <> 'provisioning')
872                  OR (NEW.state = 'running' AND OLD.state IN
873                      ('stopped', 'parked', 'error', 'lost', 'destroyed-with-data-loss'))
874                BEGIN
875                    UPDATE session_incarnations SET identity = lower(hex(randomblob(16)))
876                    WHERE session_id = NEW.session_id;
877                END;
878            UPDATE schema_compatibility SET minimum_compatible_version = 63 WHERE singleton = 1;
879            INSERT INTO schema_migrations(version, applied_at) VALUES (63, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
880            PRAGMA user_version = 63;
881            COMMIT;")?;
882    }
883
884    // Breaking: Move-only handoffs do not contain workspace backups. Older
885    // daemons can mistake a retained source for a target safe to destroy.
886    if version < 64 {
887        connection.execute_batch("BEGIN IMMEDIATE;
888            CREATE TABLE IF NOT EXISTS retained_move_sources (
889                operation_id TEXT PRIMARY KEY,
890                session_id TEXT NOT NULL,
891                source_json TEXT NOT NULL,
892                exclusions_json TEXT NOT NULL,
893                created_at TEXT NOT NULL
894            ) STRICT;
895            UPDATE schema_compatibility SET minimum_compatible_version = 64 WHERE singleton = 1;
896            INSERT INTO schema_migrations(version, applied_at) VALUES (64, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
897            PRAGMA user_version = 64;
898            COMMIT;")?;
899    }
900
901    // Breaking: removes the durable review orchestration and the execution
902    // incarnation fence that migrations 61 and 63 introduced. Reviews running
903    // when the daemon stops are cancelled and re-covered by the next review,
904    // and a sub-agent close is admitted once per request instead of being
905    // fenced against a resumed child. An older daemon would still write both.
906    if version < 65 {
907        let drop_column = if super::legacy_schema::table_has_column(
908            connection,
909            "turn_review_state",
910            "orchestration",
911        )? {
912            "ALTER TABLE turn_review_state DROP COLUMN orchestration;"
913        } else {
914            ""
915        };
916        connection.execute_batch(&format!("BEGIN IMMEDIATE;
917            DROP TRIGGER IF EXISTS session_incarnation_insert;
918            DROP TRIGGER IF EXISTS session_incarnation_resume;
919            DROP TABLE IF EXISTS session_incarnations;
920            {drop_column}
921            UPDATE schema_compatibility SET minimum_compatible_version = 65 WHERE singleton = 1;
922            INSERT INTO schema_migrations(version, applied_at) VALUES (65, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
923            PRAGMA user_version = 65;
924            COMMIT;"))?;
925    }
926
927    // Breaking: removes the saved runtime constraint; older daemons still
928    // read and write it. Historical receipts remain readable in event history.
929    if version < 66 {
930        let drop_column = if super::legacy_schema::table_has_column(
931            connection,
932            "sessions",
933            "expected_runtime_identity",
934        )? {
935            "ALTER TABLE sessions DROP COLUMN expected_runtime_identity;"
936        } else {
937            ""
938        };
939        connection.execute_batch(&format!("BEGIN IMMEDIATE;
940            {drop_column}
941            UPDATE schema_compatibility SET minimum_compatible_version = 66 WHERE singleton = 1;
942            INSERT INTO schema_migrations(version, applied_at) VALUES (66, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
943            PRAGMA user_version = 66;
944            COMMIT;"))?;
945    }
946
947    // Breaking: accounting is owned by durable session identities. Older
948    // projection writers cannot capture selections or preserve tree identity.
949    if version < 67 {
950        connection.execute_batch(&format!("BEGIN IMMEDIATE;
951            {ACCOUNTING_MIGRATION_SQL}
952            UPDATE schema_compatibility SET minimum_compatible_version = 67 WHERE singleton = 1;
953            INSERT INTO schema_migrations(version, applied_at) VALUES (67, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
954            PRAGMA user_version = 67;
955            COMMIT;"))?;
956    }
957
958    // Breaking: older daemons cannot decode or resume startup teardown, and
959    // would release capacity or provision over a surviving failed worker.
960    if version < 68 {
961        migrate_startup_cleanup_state(connection)?;
962    }
963
964    // Breaking: revision 67 existed in two branches (accounting or projects).
965    // Reconcile both histories; neither older writer can preserve the union.
966    if version < 69 {
967        let has_accounting = connection
968            .prepare(
969                "SELECT 1 FROM sqlite_schema WHERE type='table' AND name='subagent_accounting'",
970            )?
971            .exists([])?;
972        let accounting = if has_accounting {
973            ""
974        } else {
975            ACCOUNTING_MIGRATION_SQL
976        };
977        let add_snapshot = if super::legacy_schema::table_has_column(
978            connection,
979            "sessions",
980            "project_json",
981        )? {
982            ""
983        } else {
984            "ALTER TABLE sessions ADD COLUMN project_json TEXT CHECK(project_json IS NULL OR json_valid(project_json));"
985        };
986        connection.execute_batch(&format!("BEGIN IMMEDIATE;
987            {accounting}
988            {add_snapshot}
989            CREATE TABLE IF NOT EXISTS project_catalog (
990                bundle_id TEXT PRIMARY KEY,
991                project_key TEXT NOT NULL UNIQUE,
992                snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
993                hidden INTEGER NOT NULL DEFAULT 0 CHECK(hidden IN (0,1))
994            ) STRICT;
995            CREATE TABLE IF NOT EXISTS project_aliases (
996                bundle_id TEXT PRIMARY KEY,
997                canonical_id TEXT NOT NULL REFERENCES project_catalog(bundle_id),
998                snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
999                config_pending INTEGER NOT NULL DEFAULT 0 CHECK(config_pending IN (0,1))
1000            ) STRICT;
1001            CREATE TABLE IF NOT EXISTS project_session_aliases (
1002                session_id TEXT NOT NULL REFERENCES session_contexts(session_id) ON DELETE CASCADE,
1003                bundle_id TEXT NOT NULL, PRIMARY KEY(session_id,bundle_id)
1004            ) STRICT;
1005            CREATE TABLE IF NOT EXISTS project_locations (
1006                host TEXT NOT NULL,
1007                directory BLOB NOT NULL,
1008                checkout_root BLOB NOT NULL,
1009                repository_root BLOB NOT NULL,
1010                identity_json TEXT NOT NULL CHECK(json_valid(identity_json)),
1011                seen_at TEXT NOT NULL,
1012                PRIMARY KEY(host, directory)
1013            ) STRICT;
1014            CREATE TABLE IF NOT EXISTS project_seed_homes (
1015                harness TEXT NOT NULL,
1016                home BLOB NOT NULL,
1017                PRIMARY KEY(harness, home)
1018            ) STRICT;
1019            CREATE TABLE IF NOT EXISTS project_seed_failures (
1020             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)),
1021             PRIMARY KEY(harness,home,directory)
1022         ) STRICT;
1023         CREATE TABLE IF NOT EXISTS project_discovery_changes (
1024                sequence INTEGER PRIMARY KEY AUTOINCREMENT,
1025                session_id TEXT NOT NULL,
1026                directory BLOB,
1027                managed_worktree TEXT,
1028                target_template_id TEXT NOT NULL
1029            ) STRICT;
1030            CREATE TABLE IF NOT EXISTS project_discovery_progress (
1031                singleton INTEGER PRIMARY KEY CHECK(singleton=1),
1032                sequence INTEGER NOT NULL DEFAULT 0
1033            ) STRICT;
1034            CREATE TABLE IF NOT EXISTS project_discovery_failures (
1035                sequence INTEGER PRIMARY KEY REFERENCES project_discovery_changes(sequence) ON DELETE CASCADE,
1036                error TEXT NOT NULL
1037            ) STRICT;
1038            INSERT OR IGNORE INTO project_discovery_progress(singleton) VALUES(1);
1039            INSERT INTO project_discovery_changes(session_id, directory, managed_worktree, target_template_id)
1040                SELECT session_id, project_directory, managed_worktree, target_template_id FROM sessions
1041                WHERE project_directory IS NOT NULL AND NOT EXISTS (
1042                    SELECT 1 FROM project_discovery_changes d WHERE d.session_id=sessions.session_id
1043                );
1044            CREATE TRIGGER IF NOT EXISTS project_discovery_insert AFTER INSERT ON sessions
1045                WHEN NEW.project_directory IS NOT NULL BEGIN
1046                    INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
1047                    VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
1048                END;
1049            CREATE TRIGGER IF NOT EXISTS project_discovery_update AFTER UPDATE OF project_directory,managed_worktree,target_template_id ON sessions
1050                WHEN NEW.project_directory IS NOT NULL AND
1051                    (NEW.project_directory IS NOT OLD.project_directory
1052                    OR NEW.managed_worktree IS NOT OLD.managed_worktree
1053                    OR NEW.target_template_id IS NOT OLD.target_template_id) BEGIN
1054                    INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
1055                    VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
1056                END;
1057            UPDATE schema_compatibility SET minimum_compatible_version=69 WHERE singleton=1;
1058            INSERT INTO schema_migrations(version,applied_at) VALUES(69,strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1059            PRAGMA user_version=69;
1060            COMMIT;"))?;
1061    }
1062
1063    // Compatible: only shortens existing title text. Older readers and writers
1064    // accept the same column and values; no schema or JSON shape changes.
1065    if version < 70 {
1066        let transaction = connection.unchecked_transaction()?;
1067        let titles: Vec<(String, String)> = transaction
1068            .prepare("SELECT session_id, acp_session_title FROM sessions WHERE length(acp_session_title) > ?1")?
1069            .query_map([mj_core::state::MAX_SESSION_TITLE_CHARS], |row| {
1070                Ok((row.get(0)?, row.get(1)?))
1071            })?
1072            .collect::<rusqlite::Result<_>>()?;
1073        for (session_id, title) in titles {
1074            transaction.execute(
1075                "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
1076                params![session_id, mj_core::state::normalize_session_title(&title)],
1077            )?;
1078        }
1079        transaction.execute_batch(
1080            "INSERT INTO schema_migrations(version, applied_at)
1081                 VALUES (70, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1082             PRAGMA user_version = 70;",
1083        )?;
1084        transaction.commit()?;
1085    }
1086
1087    let recorded: Option<i64> =
1088        connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
1089            row.get(0)
1090        })?;
1091    if recorded == Some(SCHEMA_VERSION) && version < SCHEMA_VERSION {
1092        tracing::info!(
1093            from_revision = version,
1094            to_revision = SCHEMA_VERSION,
1095            "database migrations applied"
1096        );
1097    }
1098    if recorded != Some(SCHEMA_VERSION) {
1099        bail!(
1100            "Mjolnir database migration ledger {:?} does not match schema {}",
1101            recorded,
1102            SCHEMA_VERSION
1103        );
1104    }
1105    Ok(())
1106}
1107
1108/// Migration 53: rebuild `sessions` so its `state` constraint admits
1109/// `'parked'`. SQLite cannot alter a CHECK constraint, so the table is copied
1110/// into one declared with the new constraint, the way migration 32
1111/// (`migrate_zcode_harness_kind`) did. The table's own stored definition is
1112/// edited, so every column later migrations added is kept as it is. A store
1113/// whose constraint already admits the value is only stamped.
1114fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
1115    const BEFORE: &str = "'stopped','lost',";
1116    const AFTER: &str = "'stopped','parked','lost',";
1117    connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1118    let migration = (|| -> Result<()> {
1119        let transaction = connection.unchecked_transaction()?;
1120        let sql: String = transaction.query_row(
1121            "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1122            [],
1123            |row| row.get(0),
1124        )?;
1125        let (_, definition) = sql
1126            .split_once('(')
1127            .context("missing sessions table definition")?;
1128        if !definition.contains(AFTER) {
1129            ensure!(
1130                definition.matches(BEFORE).count() == 1,
1131                "unexpected sessions state constraint"
1132            );
1133            let definition = definition.replace(BEFORE, AFTER);
1134            let objects: Vec<String> = transaction
1135                .prepare(
1136                    "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1137                     AND type IN ('index','trigger') AND sql IS NOT NULL",
1138                )?
1139                .query_map([], |row| row.get(0))?
1140                .collect::<rusqlite::Result<_>>()?;
1141            transaction.execute_batch(&format!(
1142                "CREATE TABLE sessions_parked_v53 ({definition};
1143                 INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
1144                 DROP TABLE sessions;
1145                 ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
1146            ))?;
1147            for object in objects {
1148                transaction.execute_batch(&object)?;
1149            }
1150            ensure!(
1151                !transaction
1152                    .prepare("PRAGMA foreign_key_check")?
1153                    .exists([])?,
1154                "foreign key violation in the parked-state migration"
1155            );
1156        }
1157        transaction.execute_batch(
1158            "UPDATE schema_compatibility SET minimum_compatible_version = 53
1159                 WHERE singleton = 1;
1160             INSERT INTO schema_migrations(version, applied_at)
1161                 VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1162             PRAGMA user_version = 53;",
1163        )?;
1164        transaction.commit()?;
1165        Ok(())
1166    })();
1167    let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1168    migration.context("migrate the sessions state constraint for parked sub-agents")?;
1169    restored.context("restore foreign key enforcement after the parked-state migration")?;
1170    Ok(())
1171}
1172
1173// Rebuild from the stored definition so every shipped column and trigger survives.
1174fn migrate_startup_cleanup_state(connection: &Connection) -> Result<()> {
1175    const BEFORE: &str = "'stopped','parked','lost',";
1176    const AFTER: &str = "'stopped','parked','startup-cleanup','lost',";
1177    connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1178    let migration = (|| -> Result<()> {
1179        let transaction = connection.unchecked_transaction()?;
1180        let sql: String = transaction.query_row(
1181            "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1182            [],
1183            |row| row.get(0),
1184        )?;
1185        let (_, definition) = sql
1186            .split_once('(')
1187            .context("missing sessions table definition")?;
1188        if !definition.contains(AFTER) {
1189            ensure!(
1190                definition.matches(BEFORE).count() == 1,
1191                "unexpected sessions state constraint"
1192            );
1193            let definition = definition.replace(BEFORE, AFTER);
1194            let objects: Vec<String> = transaction
1195                .prepare(
1196                    "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1197                     AND type IN ('index','trigger') AND sql IS NOT NULL",
1198                )?
1199                .query_map([], |row| row.get(0))?
1200                .collect::<rusqlite::Result<_>>()?;
1201            transaction.execute_batch(&format!(
1202                "CREATE TABLE sessions_startup_cleanup_v68 ({definition};
1203                 INSERT INTO sessions_startup_cleanup_v68 SELECT * FROM sessions;
1204                 DROP TABLE sessions;
1205                 ALTER TABLE sessions_startup_cleanup_v68 RENAME TO sessions;"
1206            ))?;
1207            for object in objects {
1208                transaction.execute_batch(&object)?;
1209            }
1210            ensure!(
1211                !transaction
1212                    .prepare("PRAGMA foreign_key_check")?
1213                    .exists([])?,
1214                "foreign key violation in the startup-cleanup migration"
1215            );
1216        }
1217        transaction.execute_batch(
1218            "UPDATE schema_compatibility SET minimum_compatible_version = 68
1219                 WHERE singleton = 1;
1220             INSERT INTO schema_migrations(version, applied_at)
1221                 VALUES (68, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1222             PRAGMA user_version = 68;",
1223        )?;
1224        transaction.commit()?;
1225        Ok(())
1226    })();
1227    let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1228    migration.context("migrate the sessions state constraint for failed startup cleanup")?;
1229    restored.context("restore foreign key enforcement after the startup-cleanup migration")?;
1230    Ok(())
1231}
1232
1233/// Create an empty store at the baseline revision in one immediate transaction.
1234/// The revision is read again under the write lock, so a second process that
1235/// raced to create the same store finds it already created.
1236fn create_baseline_schema(connection: &Connection) -> Result<()> {
1237    connection.execute_batch("BEGIN IMMEDIATE;")?;
1238    let created = (|| -> Result<()> {
1239        let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
1240        if version != 0 {
1241            return Ok(());
1242        }
1243        connection.execute_batch(include_str!("baseline.sql"))?;
1244        connection.execute(
1245            "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
1246            [BASELINE_MINIMUM_COMPATIBLE_VERSION],
1247        )?;
1248        connection.execute(
1249            "INSERT INTO schema_migrations(version, applied_at)
1250             VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
1251            [BASELINE_SCHEMA_VERSION],
1252        )?;
1253        connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
1254        Ok(())
1255    })();
1256    match created {
1257        Ok(()) => connection
1258            .execute_batch("COMMIT;")
1259            .context("commit baseline database schema"),
1260        Err(error) => {
1261            if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
1262                tracing::warn!(%rollback, "could not roll back a failed baseline schema");
1263            }
1264            Err(error.context("create baseline database schema"))
1265        }
1266    }
1267}
1268
1269#[cfg(test)]
1270pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
1271    let connection = Connection::open(path).unwrap();
1272    let transaction = connection.unchecked_transaction().unwrap();
1273    transaction
1274        .execute(
1275            "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
1276            [minimum_compatible],
1277        )
1278        .unwrap();
1279    transaction
1280        .execute(
1281            "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
1282            [revision],
1283        )
1284        .unwrap();
1285    transaction
1286        .pragma_update(None, "user_version", revision)
1287        .unwrap();
1288    transaction.commit().unwrap();
1289    forget_verified_schema(path);
1290}
1291
1292#[cfg(test)]
1293mod reader_tests {
1294    use super::*;
1295
1296    fn assert_divergent_history_upgrades(revision: i64, project_history: bool, interrupt: bool) {
1297        let directory = tempfile::tempdir().unwrap();
1298        let path = directory.path().join("divergent-history.sqlite");
1299        let connection = Connection::open(&path).unwrap();
1300        create_baseline_schema(&connection).unwrap();
1301        connection.execute_batch(&format!(
1302            "INSERT INTO session_contexts(session_id,bundle_id,created_at)
1303                 VALUES ('kept','project','now');
1304             INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at,project_directory)
1305                 VALUES ('kept','Keep my work','codex','codex','local','error','now',X'2F7265706F');
1306             INSERT INTO materialized_sessions(session_id) VALUES ('kept');
1307             INSERT INTO session_turn_usage VALUES ('kept','turn',1,1,'{{\"tokens\":42}}');
1308             INSERT INTO session_provider_cost VALUES ('kept','{{\"amount\":1}}');
1309             CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1310                 WHEN NEW.version > {revision}
1311                 BEGIN SELECT RAISE(ABORT,'fixture migration boundary'); END;"
1312        )).unwrap();
1313        assert!(migrate_schema(&connection).is_err());
1314        if !connection.is_autocommit() {
1315            connection.execute_batch("ROLLBACK").unwrap();
1316        }
1317        connection
1318            .execute_batch("DROP TRIGGER stop_at_revision")
1319            .unwrap();
1320        assert_eq!(read_schema_state(&connection).unwrap().revision, revision);
1321        if project_history {
1322            connection
1323                .execute_batch(include_str!("project_catalog_v67.sql"))
1324                .unwrap();
1325            connection
1326                .execute_batch(
1327                    "INSERT INTO project_catalog VALUES ('project','key','{}',0);
1328                 INSERT INTO project_aliases VALUES ('alias','project','{}',1);
1329                 INSERT INTO project_session_aliases VALUES ('kept','alias');
1330                 UPDATE sessions SET project_json='{\"kept\":true}';",
1331                )
1332                .unwrap();
1333        }
1334        if interrupt {
1335            connection
1336                .execute_batch(
1337                    "CREATE TRIGGER interrupt_reconciliation BEFORE INSERT ON schema_migrations
1338                 WHEN NEW.version=69 BEGIN SELECT RAISE(ABORT,'interrupted reconciliation'); END;",
1339                )
1340                .unwrap();
1341            assert!(migrate_schema(&connection).is_err());
1342            drop(connection);
1343            let connection = Connection::open(&path).unwrap();
1344            assert_eq!(read_schema_state(&connection).unwrap().revision, 68);
1345            connection
1346                .execute_batch("DROP TRIGGER interrupt_reconciliation")
1347                .unwrap();
1348        } else {
1349            drop(connection);
1350        }
1351
1352        let writer = open_writer(&path).unwrap();
1353        let state = read_schema_state(&writer).unwrap();
1354        assert_eq!(state.revision, SCHEMA_VERSION);
1355        assert!(state.ensure_supported_by(68).is_err());
1356        assert!(state.ensure_supported_by(67).is_err());
1357        if project_history {
1358            let snapshot: String = writer
1359                .query_row(
1360                    "SELECT project_json FROM sessions WHERE session_id='kept'",
1361                    [],
1362                    |row| row.get(0),
1363                )
1364                .unwrap();
1365            assert_eq!(snapshot, r#"{"kept":true}"#);
1366            let alias: (String, i64) = writer.query_row(
1367                "SELECT canonical_id,config_pending FROM project_aliases WHERE bundle_id='alias'", [],
1368                |row| Ok((row.get(0)?, row.get(1)?))
1369            ).unwrap();
1370            assert_eq!(alias, ("project".to_owned(), 1));
1371        }
1372        // Rebuilding sessions must retain discovery triggers and admit the new state.
1373        writer.execute_batch(
1374            "UPDATE sessions SET state='startup-cleanup',project_directory=X'2F6E6577' WHERE session_id='kept';
1375             INSERT INTO session_turn_selections VALUES ('kept','turn','model','high');"
1376        ).unwrap();
1377        let changed: Vec<u8> = writer
1378            .query_row(
1379                "SELECT directory FROM project_discovery_changes ORDER BY sequence DESC LIMIT 1",
1380                [],
1381                |row| row.get(0),
1382            )
1383            .unwrap();
1384        assert_eq!(changed, b"/new");
1385        writer
1386            .execute("DELETE FROM sessions WHERE session_id='kept'", [])
1387            .unwrap();
1388        let usage: String = writer
1389            .query_row("SELECT body FROM session_turn_usage", [], |row| row.get(0))
1390            .unwrap();
1391        assert_eq!(usage, r#"{"tokens":42}"#);
1392        let cost: String = writer
1393            .query_row("SELECT body FROM session_provider_cost", [], |row| {
1394                row.get(0)
1395            })
1396            .unwrap();
1397        assert_eq!(cost, r#"{"amount":1}"#);
1398        assert!(
1399            !writer
1400                .prepare("PRAGMA foreign_key_check")
1401                .unwrap()
1402                .exists([])
1403                .unwrap()
1404        );
1405        drop(writer);
1406        forget_verified_schema(&path);
1407        assert_eq!(
1408            read_schema_state(&open_writer(&path).unwrap())
1409                .unwrap()
1410                .revision,
1411            SCHEMA_VERSION
1412        );
1413    }
1414
1415    #[test]
1416    fn divergent_accounting_revision_67_preserves_usage_and_adds_projects() {
1417        assert_divergent_history_upgrades(67, false, false);
1418    }
1419
1420    #[test]
1421    fn divergent_cleanup_revision_68_preserves_usage_and_adds_projects() {
1422        assert_divergent_history_upgrades(68, false, false);
1423    }
1424
1425    #[test]
1426    fn divergent_project_revision_67_preserves_aliases_snapshots_and_usage() {
1427        assert_divergent_history_upgrades(66, true, false);
1428    }
1429
1430    #[test]
1431    fn divergent_project_reconciliation_resumes_after_interruption() {
1432        assert_divergent_history_upgrades(66, true, true);
1433    }
1434
1435    #[test]
1436    fn move_ownership_upgrade_retains_sources_and_refuses_previous_daemons() {
1437        let directory = tempfile::tempdir().unwrap();
1438        let path = directory.path().join("move.sqlite");
1439        let connection = open_writer(&path).unwrap();
1440        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();
1441        drop(connection);
1442        forget_verified_schema(&path);
1443        let upgraded = open_writer(&path).unwrap();
1444        let schema = read_schema_state(&upgraded).unwrap();
1445        assert!(schema.ensure_supported_by(63).is_err());
1446        upgraded.execute("INSERT INTO retained_move_sources VALUES ('move-one','session-one','{}','[]','now')", []).unwrap();
1447        drop(upgraded);
1448        let reopened = open_writer(&path).unwrap();
1449        let count: i64 = reopened
1450            .query_row("SELECT count(*) FROM retained_move_sources", [], |row| {
1451                row.get(0)
1452            })
1453            .unwrap();
1454        assert_eq!(count, 1);
1455    }
1456
1457    #[test]
1458    fn title_migration_caps_old_titles_preserves_short_titles_and_is_compatible() {
1459        let directory = tempfile::tempdir().unwrap();
1460        let path = directory.path().join("title-migration.sqlite3");
1461        let titles = [
1462            ("long", Some("word ".repeat(20_000))),
1463            ("unicode", Some("界".repeat(257))),
1464            ("exact", Some("界".repeat(256))),
1465            ("short", Some("  Keep\nthis title  ".into())),
1466            ("unset", None),
1467        ];
1468        for (id, _) in &titles {
1469            let mut session = super::super::tests::session(id, "project");
1470            session.state = SessionState::Stopped;
1471            save_session_to(&path, &session).unwrap();
1472        }
1473        stamp_schema_version(&path, 69);
1474        let old = Connection::open(&path).unwrap();
1475        for (id, title) in &titles {
1476            old.execute(
1477                "UPDATE sessions SET acp_session_title=?2 WHERE session_id=?1",
1478                params![id, title],
1479            )
1480            .unwrap();
1481        }
1482        drop(old);
1483
1484        let upgraded = open_writer(&path).unwrap();
1485        let state = read_schema_state(&upgraded).unwrap();
1486        assert_eq!(state.revision, 70);
1487        assert_eq!(state.minimum_compatible, Some(69));
1488        state.ensure_supported_by(69).unwrap();
1489        for (id, original) in &titles {
1490            let stored: Option<String> = upgraded
1491                .query_row(
1492                    "SELECT acp_session_title FROM sessions WHERE session_id=?1",
1493                    [id],
1494                    |row| row.get(0),
1495                )
1496                .unwrap();
1497            let expected = match *id {
1498                "long" => Some(format!("{}word…", "word ".repeat(50))),
1499                "unicode" => Some(format!("{}…", "界".repeat(255))),
1500                _ => original.clone(),
1501            };
1502            assert_eq!(stored, expected, "session {id}");
1503            assert!(stored.is_none_or(|title| title.chars().count() <= 256));
1504        }
1505        drop(upgraded);
1506        forget_verified_schema(&path);
1507        drop(open_writer(&path).unwrap());
1508    }
1509
1510    #[test]
1511    fn interrupted_title_migration_rolls_back_titles_and_revision() {
1512        let directory = tempfile::tempdir().unwrap();
1513        let path = directory.path().join("title-migration-interrupted.sqlite3");
1514        save_session_to(&path, &super::super::tests::session("old", "project")).unwrap();
1515        stamp_schema_version(&path, 69);
1516        let original = "word ".repeat(20_000);
1517        let old = Connection::open(&path).unwrap();
1518        old.execute("UPDATE sessions SET acp_session_title=?1", [&original])
1519            .unwrap();
1520        old.execute_batch(
1521            "CREATE TRIGGER stop_title_migration BEFORE INSERT ON schema_migrations
1522             WHEN NEW.version=70 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1523        )
1524        .unwrap();
1525        assert!(migrate_schema(&old).is_err());
1526        assert_eq!(read_schema_state(&old).unwrap().revision, 69);
1527        let stored: String = old
1528            .query_row("SELECT acp_session_title FROM sessions", [], |row| {
1529                row.get(0)
1530            })
1531            .unwrap();
1532        assert_eq!(stored, original);
1533        old.execute_batch("DROP TRIGGER stop_title_migration")
1534            .unwrap();
1535        drop(old);
1536        drop(open_writer(&path).unwrap());
1537        assert!(
1538            load_state_from(&path).unwrap().sessions["old"]
1539                .acp_session_title
1540                .as_ref()
1541                .unwrap()
1542                .chars()
1543                .count()
1544                <= 256
1545        );
1546    }
1547
1548    #[test]
1549    fn recent_revisions_upgrade_directly_and_preserve_user_data() {
1550        // Revision 47 was current on 2026-09-23. Keep the exhaustive
1551        // interruption matrix focused on the last week's migration history.
1552        const TEST_REVISION_FLOOR: i64 = 47;
1553        for revision in TEST_REVISION_FLOOR..SCHEMA_VERSION {
1554            let directory = tempfile::tempdir().unwrap();
1555            let path = directory.path().join("mj.sqlite3");
1556            let connection = Connection::open(&path).unwrap();
1557            connection
1558                .execute_batch(include_str!("legacy_v1.sql"))
1559                .unwrap();
1560            connection.execute_batch(
1561                "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
1562                 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
1563                     target_template_id, state, updated_at, native_session_id)
1564                     VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
1565                         '2026-01-01T00:00:00Z', 'native-original');
1566                 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
1567                     VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
1568            ).unwrap();
1569            // Interrupt at each historical transaction boundary. Closing the
1570            // connection rolls back the interrupted step, exactly as a killed
1571            // updater would. Reopening must resume without manual cleanup.
1572            connection
1573                .execute_batch(&format!(
1574                    "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1575                 WHEN NEW.version > {revision}
1576                 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
1577                ))
1578                .unwrap();
1579            super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
1580            assert!(migrate_schema(&connection).is_err());
1581            drop(connection);
1582            let connection = Connection::open(&path).unwrap();
1583            let found: i64 = connection
1584                .query_row("PRAGMA user_version", [], |row| row.get(0))
1585                .unwrap();
1586            assert_eq!(found, revision);
1587            connection
1588                .execute_batch("DROP TRIGGER stop_at_revision")
1589                .unwrap();
1590            drop(connection);
1591
1592            let writer =
1593                open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
1594            assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
1595            assert_eq!(
1596                writer
1597                    .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
1598                    .unwrap(),
1599                "ok"
1600            );
1601            assert!(
1602                !writer
1603                    .prepare("PRAGMA foreign_key_check")
1604                    .unwrap()
1605                    .exists([])
1606                    .unwrap()
1607            );
1608            drop(writer);
1609            let reader = open_reader_strict(&path).unwrap();
1610            let prompt: String = reader
1611                .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
1612                .unwrap();
1613            assert_eq!(prompt, "Keep my prompt");
1614            let state = load_state_from(&path).unwrap();
1615            assert_eq!(state.sessions["old-session"].title, "Keep my work");
1616            assert_eq!(
1617                state.sessions["old-session"].native_session_id.as_deref(),
1618                Some("native-original")
1619            );
1620            drop(reader);
1621            // A fresh writer process must not apply any migration twice.
1622            forget_verified_schema(&path);
1623            drop(open_writer(&path).unwrap());
1624        }
1625    }
1626
1627    #[test]
1628    fn accounting_migration_preserves_usage_and_changes_its_deletion_owner() {
1629        let dir = tempfile::tempdir().unwrap();
1630        let path = dir.path().join("migration.sqlite");
1631        let connection = Connection::open(&path).unwrap();
1632        connection
1633            .execute_batch(include_str!("legacy_v1.sql"))
1634            .unwrap();
1635        connection.execute_batch("INSERT INTO session_contexts VALUES ('old-session','project','2026-01-01T00:00:00Z');
1636            INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at)
1637            VALUES ('old-session','Retain usage','codex','codex','local','error','2026-01-01T00:00:00Z');
1638            CREATE TRIGGER stop_before_accounting BEFORE INSERT ON schema_migrations WHEN NEW.version=67
1639            BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;").unwrap();
1640        super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
1641        assert!(migrate_schema(&connection).is_err());
1642        connection
1643            .execute_batch(
1644                "ROLLBACK; DROP TRIGGER stop_before_accounting;
1645            INSERT INTO session_turn_usage VALUES ('old-session','turn',7,1,'{}');
1646            INSERT INTO session_provider_cost VALUES ('old-session','{\"amount\":1}');",
1647            )
1648            .unwrap();
1649        drop(connection);
1650        let writer = open_writer(&path).unwrap();
1651        writer
1652            .execute("DELETE FROM sessions WHERE session_id='old-session'", [])
1653            .unwrap();
1654        for table in ["session_turn_usage", "session_provider_cost"] {
1655            assert_eq!(
1656                writer
1657                    .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| row
1658                        .get::<_, u64>(0))
1659                    .unwrap(),
1660                1
1661            );
1662        }
1663        assert_eq!(
1664            writer
1665                .query_row("SELECT body FROM session_turn_usage", [], |row| row
1666                    .get::<_, String>(0))
1667                .unwrap(),
1668            "{}"
1669        );
1670        assert!(
1671            !writer
1672                .prepare("PRAGMA foreign_key_check")
1673                .unwrap()
1674                .exists([])
1675                .unwrap()
1676        );
1677        assert_eq!(
1678            read_schema_state(&writer).unwrap().minimum_compatible,
1679            Some(MINIMUM_COMPATIBLE_VERSION)
1680        );
1681    }
1682
1683    /// The oldest executable revision that can still read and write a store at
1684    /// `SCHEMA_VERSION`. Migration 69 reconciles accounting and project histories.
1685    const MINIMUM_COMPATIBLE_VERSION: i64 = 69;
1686
1687    /// Rewrites a store's recorded schema version the way another build's
1688    /// migration ladder would, and forgets that this process verified it.
1689    fn stamp_schema_version(path: &Path, version: i64) {
1690        if version > SCHEMA_VERSION {
1691            advance_test_schema(path, version, version);
1692            return;
1693        }
1694        let connection = Connection::open(path).unwrap();
1695        if version < 67 {
1696            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();
1697        }
1698        connection
1699            .execute_batch(&format!("PRAGMA user_version = {version};"))
1700            .unwrap();
1701        connection
1702            .execute(
1703                "DELETE FROM schema_migrations WHERE version > ?1",
1704                [version],
1705            )
1706            .unwrap();
1707        if version == 30 {
1708            connection
1709                .execute(
1710                    "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
1711                    [],
1712                )
1713                .unwrap();
1714        }
1715        drop(connection);
1716        forget_verified_schema(path);
1717    }
1718
1719    #[test]
1720    fn subagent_policy_migration_preserves_legacy_choices_and_refuses_old_writers() {
1721        use mj_core::subagent::SubagentPolicy;
1722        let directory = tempfile::tempdir().unwrap();
1723        let path = directory.path().join("subagent-policy.sqlite3");
1724        for id in ["all", "native", "unset"] {
1725            save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
1726        }
1727        let connection = Connection::open(&path).unwrap();
1728        connection
1729            .execute_batch(
1730                "UPDATE sessions SET mjolnir_subagents = 1 WHERE session_id = 'all';
1731             UPDATE sessions SET mjolnir_subagents = 0 WHERE session_id = 'native';
1732             ALTER TABLE sessions DROP COLUMN subagents;
1733             DROP TABLE subagent_preference;
1734             DELETE FROM schema_migrations WHERE version >= 56;
1735             UPDATE schema_compatibility SET minimum_compatible_version = 55;
1736             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;",
1737            )
1738            .unwrap();
1739        drop(connection);
1740        forget_verified_schema(&path);
1741        let upgraded = open_writer(&path).unwrap();
1742        assert!(
1743            read_schema_state(&upgraded)
1744                .unwrap()
1745                .ensure_supported_by(55)
1746                .is_err()
1747        );
1748        drop(upgraded);
1749        let state = load_state_from(&path).unwrap();
1750        assert_eq!(
1751            state.sessions["all"].subagents,
1752            Some(SubagentPolicy::AllModels)
1753        );
1754        assert_eq!(
1755            state.sessions["native"].subagents,
1756            Some(SubagentPolicy::Native)
1757        );
1758        assert_eq!(state.sessions["unset"].subagents, None);
1759        assert_eq!(state.last_subagent_policy, SubagentPolicy::Native);
1760    }
1761
1762    #[test]
1763    fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
1764        let directory = tempfile::tempdir().unwrap();
1765        let path = directory.path().join("steering-migration.sqlite3");
1766        let record = super::super::tests::session("steered-session", "project");
1767        save_session_to(&path, &record).unwrap();
1768        let connection = Connection::open(&path).unwrap();
1769        connection
1770            .execute_batch(
1771                "DELETE FROM schema_migrations WHERE version >= 48;
1772             UPDATE schema_compatibility SET minimum_compatible_version = 47;
1773             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;",
1774            )
1775            .unwrap();
1776        drop(connection);
1777        forget_verified_schema(&path);
1778        let upgraded = open_writer(&path).unwrap();
1779        let schema = read_schema_state(&upgraded).unwrap();
1780        assert_eq!(schema.revision, SCHEMA_VERSION);
1781        assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1782        let error = schema.ensure_supported_by(47).unwrap_err();
1783        assert!(matches!(
1784            error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
1785            StoreSchemaMismatchReason::Incompatible {
1786                minimum_compatible: MINIMUM_COMPATIBLE_VERSION
1787            }
1788        ));
1789        drop(upgraded);
1790        assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
1791    }
1792
1793    #[test]
1794    fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
1795        let directory = tempfile::tempdir().unwrap();
1796        let path = directory.path().join("target-migration.sqlite3");
1797        let record = super::super::tests::session("preserved-session", "project");
1798        save_session_to(&path, &record).unwrap();
1799        let connection = Connection::open(&path).unwrap();
1800        connection
1801            .execute_batch(
1802                "ALTER TABLE sessions DROP COLUMN target_runtime_json;
1803             DELETE FROM schema_migrations WHERE version >= 46;
1804             UPDATE schema_compatibility SET minimum_compatible_version = 44;
1805             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;",
1806            )
1807            .unwrap();
1808        forget_verified_schema(&path);
1809        let upgraded = open_writer(&path).unwrap();
1810        let schema = read_schema_state(&upgraded).unwrap();
1811        assert_eq!(schema.revision, SCHEMA_VERSION);
1812        assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1813        let error = schema.ensure_supported_by(45).unwrap_err();
1814        assert!(matches!(
1815            error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
1816            StoreSchemaMismatchReason::Incompatible {
1817                minimum_compatible: MINIMUM_COMPATIBLE_VERSION
1818            }
1819        ));
1820        drop(upgraded);
1821        let restored = load_state_from(&path).unwrap();
1822        assert_eq!(restored.sessions[&record.id], record);
1823    }
1824
1825    #[test]
1826    fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
1827        let directory = tempfile::tempdir().unwrap();
1828        let path = directory.path().join("mj.sqlite3");
1829        let connection = open_writer(&path).unwrap();
1830        connection
1831            .execute_batch(
1832                "BEGIN IMMEDIATE;
1833             DROP TABLE quota_reset_cache;
1834             DROP TABLE native_agent_transcript;
1835             DROP TABLE native_agents;
1836             DROP TABLE native_agent_replay;
1837             DELETE FROM schema_migrations WHERE version >= 39;
1838             UPDATE schema_compatibility SET minimum_compatible_version = 32;
1839             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;
1840             COMMIT;",
1841            )
1842            .unwrap();
1843        migrate_schema(&connection).unwrap();
1844        let state = read_schema_state(&connection).unwrap();
1845        assert_eq!(state.revision, SCHEMA_VERSION);
1846        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1847        let event = ApiEventData::InputRequired {
1848            request: None,
1849            turn_id: Some(1),
1850        };
1851        #[derive(serde::Deserialize)]
1852        struct LegacyInputEvent {
1853            #[serde(rename = "request")]
1854            _request: mj_core::elicitation::ElicitationRequest,
1855        }
1856        let encoded = serde_json::to_value(&event).unwrap();
1857        assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
1858    }
1859
1860    #[test]
1861    fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
1862        let directory = tempfile::tempdir().unwrap();
1863        let path = directory.path().join("mj.sqlite3");
1864        let connection = open_writer(&path).unwrap();
1865        connection
1866            .execute_batch(
1867                "CREATE TABLE future_feature(value TEXT NOT NULL);
1868                 INSERT INTO future_feature VALUES ('preserve me');",
1869            )
1870            .unwrap();
1871        drop(connection);
1872        advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1873
1874        let reader = open_reader_strict(&path).unwrap();
1875        assert_eq!(
1876            reader
1877                .query_row("SELECT value FROM future_feature", [], |row| row
1878                    .get::<_, String>(0))
1879                .unwrap(),
1880            "preserve me"
1881        );
1882        assert!(reader.execute("DELETE FROM future_feature", []).is_err());
1883        drop(reader);
1884
1885        // A repair would recreate this deliberately removed trigger. A future
1886        // schema is authoritative even when it differs from our own repairs.
1887        let raw = Connection::open(&path).unwrap();
1888        raw.execute_batch("DROP TRIGGER session_contexts_workspace_update;")
1889            .unwrap();
1890        drop(raw);
1891        let writer = open_writer(&path).unwrap();
1892        assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'session_contexts_workspace_update')", [], |row| row.get::<_, bool>(0)).unwrap());
1893        assert_eq!(
1894            writer
1895                .query_row("SELECT value FROM future_feature", [], |row| row
1896                    .get::<_, String>(0))
1897                .unwrap(),
1898            "preserve me"
1899        );
1900        let state = read_schema_state(&writer).unwrap();
1901        assert_eq!(state.revision, SCHEMA_VERSION + 1);
1902        assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
1903    }
1904
1905    #[test]
1906    fn invalid_compatibility_metadata_refuses_readers_and_writers() {
1907        for alteration in [
1908            "DROP TABLE schema_compatibility",
1909            "DELETE FROM schema_compatibility",
1910            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
1911            "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
1912            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
1913            "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
1914            "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
1915            "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
1916        ] {
1917            for future in [false, true] {
1918                let directory = tempfile::tempdir().unwrap();
1919                let path = directory.path().join("mj.sqlite3");
1920                drop(open_writer(&path).unwrap());
1921                if future {
1922                    advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1923                }
1924                let raw = Connection::open(&path).unwrap();
1925                raw.execute_batch(alteration).unwrap();
1926                let before: i64 = raw
1927                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
1928                    .unwrap();
1929                // Exercise the cached path as well as a fresh writer open.
1930                for error in [
1931                    open_reader_strict(&path).unwrap_err(),
1932                    open_writer(&path).unwrap_err(),
1933                ] {
1934                    let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
1935                    assert_eq!(
1936                        mismatch.reason,
1937                        StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
1938                        "{alteration}"
1939                    );
1940                }
1941                forget_verified_schema(&path);
1942                assert!(open_writer(&path).is_err(), "{alteration}");
1943                let after: i64 = raw
1944                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
1945                    .unwrap();
1946                assert_eq!(
1947                    before, after,
1948                    "a rejected open repaired schema: {alteration}"
1949                );
1950            }
1951        }
1952    }
1953
1954    #[test]
1955    fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
1956        let directory = tempfile::tempdir().unwrap();
1957        let path = directory.path().join("mj.sqlite3");
1958        let connection = Connection::open(&path).unwrap();
1959        // A table the baseline also creates makes its batch fail part way.
1960        connection
1961            .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
1962            .unwrap();
1963
1964        let error = migrate_schema(&connection).unwrap_err();
1965
1966        assert!(format!("{error:#}").contains("create baseline database schema"));
1967        assert!(
1968            connection.is_autocommit(),
1969            "the failed baseline left a transaction open"
1970        );
1971        assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
1972        let tables: i64 = connection
1973            .query_row(
1974                "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
1975                [],
1976                |row| row.get(0),
1977            )
1978            .unwrap();
1979        assert_eq!(tables, 1, "only the conflicting table remains");
1980
1981        connection.execute_batch("DROP TABLE workspaces").unwrap();
1982        drop(connection);
1983        let writer = open_writer(&path).unwrap();
1984        let state = read_schema_state(&writer).unwrap();
1985        assert_eq!(state.revision, SCHEMA_VERSION);
1986        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1987    }
1988
1989    /// A store ahead of this build cannot be fixed by starting a daemon of
1990    /// this build, so the reader must not say so. This is the message the
1991    /// incident in #24 printed twice a second for an hour.
1992    #[test]
1993    fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
1994        let directory = tempfile::tempdir().unwrap();
1995        let path = directory.path().join("mj.sqlite3");
1996        drop(open_writer(&path).unwrap());
1997        stamp_schema_version(&path, SCHEMA_VERSION + 1);
1998
1999        let error = open_reader_strict(&path).unwrap_err();
2000
2001        let mismatch = error
2002            .chain()
2003            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2004            .expect("the reader reports the mismatch as a typed cause");
2005        assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
2006        assert_eq!(mismatch.supported, SCHEMA_VERSION);
2007        let message = mismatch.to_string();
2008        assert!(message.contains("upgrade Mjolnir"), "got {message}");
2009        assert!(
2010            !message.contains("start the Mjolnir daemon"),
2011            "got {message}"
2012        );
2013    }
2014
2015    /// A store behind this build keeps the advice that works, verbatim, so
2016    /// existing log greps and runbooks keep matching.
2017    #[test]
2018    fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
2019        let directory = tempfile::tempdir().unwrap();
2020        let path = directory.path().join("mj.sqlite3");
2021        drop(open_writer(&path).unwrap());
2022        let raw = Connection::open(&path).unwrap();
2023        raw.execute_batch(&format!(
2024            "UPDATE schema_compatibility SET minimum_compatible_version = {0};
2025             DELETE FROM schema_migrations WHERE version > {0};
2026             INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
2027             PRAGMA user_version = {0};",
2028            SCHEMA_VERSION - 1
2029        ))
2030        .unwrap();
2031        drop(raw);
2032
2033        let error = open_reader_strict(&path).unwrap_err();
2034
2035        let mismatch = error
2036            .chain()
2037            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2038            .expect("the reader reports the mismatch as a typed cause");
2039        assert_eq!(
2040            mismatch.to_string(),
2041            format!(
2042                "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
2043                 start the Mjolnir daemon to migrate it",
2044                SCHEMA_VERSION - 1
2045            )
2046        );
2047    }
2048
2049    #[test]
2050    fn strict_reader_rejects_mutation() {
2051        let directory = tempfile::tempdir().unwrap();
2052        let path = directory.path().join("mj.sqlite3");
2053        drop(open_writer(&path).unwrap());
2054
2055        let reader = open_reader_strict(&path).unwrap();
2056        let error = reader
2057            .execute("CREATE TABLE forbidden(value TEXT)", [])
2058            .unwrap_err();
2059        assert!(
2060            matches!(
2061                error.sqlite_error_code(),
2062                Some(rusqlite::ErrorCode::ReadOnly)
2063            ),
2064            "unexpected mutation error: {error}"
2065        );
2066    }
2067}