Skip to main content

mj_controller/database/
schema.rs

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