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    let recorded: Option<i64> =
1289        connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
1290            row.get(0)
1291        })?;
1292    if recorded == Some(SCHEMA_VERSION) && version < SCHEMA_VERSION {
1293        tracing::info!(
1294            from_revision = version,
1295            to_revision = SCHEMA_VERSION,
1296            "database migrations applied"
1297        );
1298    }
1299    if recorded != Some(SCHEMA_VERSION) {
1300        bail!(
1301            "Mjolnir database migration ledger {:?} does not match schema {}",
1302            recorded,
1303            SCHEMA_VERSION
1304        );
1305    }
1306    Ok(())
1307}
1308
1309/// Migration 53: rebuild `sessions` so its `state` constraint admits
1310/// `'parked'`. SQLite cannot alter a CHECK constraint, so the table is copied
1311/// into one declared with the new constraint, the way migration 32
1312/// (`migrate_zcode_harness_kind`) did. The table's own stored definition is
1313/// edited, so every column later migrations added is kept as it is. A store
1314/// whose constraint already admits the value is only stamped.
1315fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
1316    const BEFORE: &str = "'stopped','lost',";
1317    const AFTER: &str = "'stopped','parked','lost',";
1318    connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1319    let migration = (|| -> Result<()> {
1320        let transaction = connection.unchecked_transaction()?;
1321        let sql: String = transaction.query_row(
1322            "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1323            [],
1324            |row| row.get(0),
1325        )?;
1326        let (_, definition) = sql
1327            .split_once('(')
1328            .context("missing sessions table definition")?;
1329        if !definition.contains(AFTER) {
1330            ensure!(
1331                definition.matches(BEFORE).count() == 1,
1332                "unexpected sessions state constraint"
1333            );
1334            let definition = definition.replace(BEFORE, AFTER);
1335            let objects: Vec<String> = transaction
1336                .prepare(
1337                    "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1338                     AND type IN ('index','trigger') AND sql IS NOT NULL",
1339                )?
1340                .query_map([], |row| row.get(0))?
1341                .collect::<rusqlite::Result<_>>()?;
1342            transaction.execute_batch(&format!(
1343                "CREATE TABLE sessions_parked_v53 ({definition};
1344                 INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
1345                 DROP TABLE sessions;
1346                 ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
1347            ))?;
1348            for object in objects {
1349                transaction.execute_batch(&object)?;
1350            }
1351            ensure!(
1352                !transaction
1353                    .prepare("PRAGMA foreign_key_check")?
1354                    .exists([])?,
1355                "foreign key violation in the parked-state migration"
1356            );
1357        }
1358        transaction.execute_batch(
1359            "UPDATE schema_compatibility SET minimum_compatible_version = 53
1360                 WHERE singleton = 1;
1361             INSERT INTO schema_migrations(version, applied_at)
1362                 VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1363             PRAGMA user_version = 53;",
1364        )?;
1365        transaction.commit()?;
1366        Ok(())
1367    })();
1368    let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1369    migration.context("migrate the sessions state constraint for parked sub-agents")?;
1370    restored.context("restore foreign key enforcement after the parked-state migration")?;
1371    Ok(())
1372}
1373
1374/// Migration 74: admit `'opencode'` in the harness-kind constraints of
1375/// `sessions` and `hidden_native_sessions`, the way migrations 27 (Muse) and
1376/// 32 (ZCode) admitted their kinds. SQLite cannot widen a CHECK constraint in
1377/// place, so each table is rebuilt from its stored definition, which keeps
1378/// every column and trigger later migrations added.
1379///
1380/// Breaking: older daemons skip a session whose harness they do not know when
1381/// loading state, and their next full state save deletes it.
1382fn migrate_opencode_harness_kind(connection: &Connection) -> Result<()> {
1383    connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1384    let migration = (|| -> Result<()> {
1385        let transaction = connection.unchecked_transaction()?;
1386        for table in ["sessions", "hidden_native_sessions"] {
1387            let sql: String = transaction.query_row(
1388                "SELECT sql FROM sqlite_schema WHERE type='table' AND name=?1",
1389                [table],
1390                |row| row.get(0),
1391            )?;
1392            let (_, definition) = sql
1393                .split_once('(')
1394                .context("missing harness table definition")?;
1395            if definition.contains("'opencode')") {
1396                continue;
1397            }
1398            ensure!(
1399                definition.contains("'zcode')"),
1400                "unexpected {table} harness constraint"
1401            );
1402            let definition = definition.replace("'zcode')", "'zcode','opencode')");
1403            let objects: Vec<String> = transaction
1404                .prepare(
1405                    "SELECT sql FROM sqlite_schema WHERE tbl_name=?1
1406                     AND type IN ('index','trigger') AND sql IS NOT NULL",
1407                )?
1408                .query_map([table], |row| row.get(0))?
1409                .collect::<rusqlite::Result<_>>()?;
1410            transaction.execute_batch(&format!(
1411                "CREATE TABLE {table}_opencode_v74 ({definition};
1412                 INSERT INTO {table}_opencode_v74 SELECT * FROM {table};
1413                 DROP TABLE {table};
1414                 ALTER TABLE {table}_opencode_v74 RENAME TO {table};"
1415            ))?;
1416            for object in objects {
1417                transaction.execute_batch(&object)?;
1418            }
1419        }
1420        ensure!(
1421            !transaction
1422                .prepare("PRAGMA foreign_key_check")?
1423                .exists([])?,
1424            "foreign key violation in OpenCode migration"
1425        );
1426        transaction.execute_batch(
1427            "UPDATE schema_compatibility SET minimum_compatible_version = 74
1428                 WHERE singleton = 1;
1429             INSERT INTO schema_migrations(version, applied_at)
1430                 VALUES (74, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1431             PRAGMA user_version = 74;",
1432        )?;
1433        transaction.commit()?;
1434        Ok(())
1435    })();
1436    let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1437    migration.context("migrate OpenCode harness constraints")?;
1438    restored.context("restore foreign key enforcement after OpenCode migration")?;
1439    Ok(())
1440}
1441
1442// Rebuild from the stored definition so every shipped column and trigger survives.
1443fn migrate_startup_cleanup_state(connection: &Connection) -> Result<()> {
1444    const BEFORE: &str = "'stopped','parked','lost',";
1445    const AFTER: &str = "'stopped','parked','startup-cleanup','lost',";
1446    connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1447    let migration = (|| -> Result<()> {
1448        let transaction = connection.unchecked_transaction()?;
1449        let sql: String = transaction.query_row(
1450            "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1451            [],
1452            |row| row.get(0),
1453        )?;
1454        let (_, definition) = sql
1455            .split_once('(')
1456            .context("missing sessions table definition")?;
1457        if !definition.contains(AFTER) {
1458            ensure!(
1459                definition.matches(BEFORE).count() == 1,
1460                "unexpected sessions state constraint"
1461            );
1462            let definition = definition.replace(BEFORE, AFTER);
1463            let objects: Vec<String> = transaction
1464                .prepare(
1465                    "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1466                     AND type IN ('index','trigger') AND sql IS NOT NULL",
1467                )?
1468                .query_map([], |row| row.get(0))?
1469                .collect::<rusqlite::Result<_>>()?;
1470            transaction.execute_batch(&format!(
1471                "CREATE TABLE sessions_startup_cleanup_v68 ({definition};
1472                 INSERT INTO sessions_startup_cleanup_v68 SELECT * FROM sessions;
1473                 DROP TABLE sessions;
1474                 ALTER TABLE sessions_startup_cleanup_v68 RENAME TO sessions;"
1475            ))?;
1476            for object in objects {
1477                transaction.execute_batch(&object)?;
1478            }
1479            ensure!(
1480                !transaction
1481                    .prepare("PRAGMA foreign_key_check")?
1482                    .exists([])?,
1483                "foreign key violation in the startup-cleanup migration"
1484            );
1485        }
1486        transaction.execute_batch(
1487            "UPDATE schema_compatibility SET minimum_compatible_version = 68
1488                 WHERE singleton = 1;
1489             INSERT INTO schema_migrations(version, applied_at)
1490                 VALUES (68, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1491             PRAGMA user_version = 68;",
1492        )?;
1493        transaction.commit()?;
1494        Ok(())
1495    })();
1496    let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1497    migration.context("migrate the sessions state constraint for failed startup cleanup")?;
1498    restored.context("restore foreign key enforcement after the startup-cleanup migration")?;
1499    Ok(())
1500}
1501
1502/// Create an empty store at the baseline revision in one immediate transaction.
1503/// The revision is read again under the write lock, so a second process that
1504/// raced to create the same store finds it already created.
1505fn create_baseline_schema(connection: &Connection) -> Result<()> {
1506    connection.execute_batch("BEGIN IMMEDIATE;")?;
1507    let created = (|| -> Result<()> {
1508        let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
1509        if version != 0 {
1510            return Ok(());
1511        }
1512        connection.execute_batch(include_str!("baseline.sql"))?;
1513        connection.execute(
1514            "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
1515            [BASELINE_MINIMUM_COMPATIBLE_VERSION],
1516        )?;
1517        connection.execute(
1518            "INSERT INTO schema_migrations(version, applied_at)
1519             VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
1520            [BASELINE_SCHEMA_VERSION],
1521        )?;
1522        connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
1523        Ok(())
1524    })();
1525    match created {
1526        Ok(()) => connection
1527            .execute_batch("COMMIT;")
1528            .context("commit baseline database schema"),
1529        Err(error) => {
1530            if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
1531                tracing::warn!(%rollback, "could not roll back a failed baseline schema");
1532            }
1533            Err(error.context("create baseline database schema"))
1534        }
1535    }
1536}
1537
1538#[cfg(test)]
1539pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
1540    let connection = Connection::open(path).unwrap();
1541    let transaction = connection.unchecked_transaction().unwrap();
1542    transaction
1543        .execute(
1544            "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
1545            [minimum_compatible],
1546        )
1547        .unwrap();
1548    transaction
1549        .execute(
1550            "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
1551            [revision],
1552        )
1553        .unwrap();
1554    transaction
1555        .pragma_update(None, "user_version", revision)
1556        .unwrap();
1557    transaction.commit().unwrap();
1558    forget_verified_schema(path);
1559}
1560
1561#[cfg(test)]
1562mod reader_tests {
1563    use super::*;
1564
1565    #[test]
1566    fn restart_intent_schema_matches_for_fresh_and_upgraded_stores() {
1567        let directory = tempfile::tempdir().unwrap();
1568        let fresh = Connection::open(directory.path().join("fresh.sqlite3")).unwrap();
1569        migrate_schema(&fresh).unwrap();
1570
1571        let upgraded = Connection::open(directory.path().join("upgraded.sqlite3")).unwrap();
1572        create_baseline_schema(&upgraded).unwrap();
1573        upgraded
1574            .execute_batch(
1575                "CREATE TRIGGER stop_before_restart_migration BEFORE INSERT ON schema_migrations
1576                 WHEN NEW.version > 72
1577                 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;",
1578            )
1579            .unwrap();
1580        assert!(migrate_schema(&upgraded).is_err());
1581        if !upgraded.is_autocommit() {
1582            upgraded.execute_batch("ROLLBACK").unwrap();
1583        }
1584        upgraded
1585            .execute_batch("DROP TRIGGER stop_before_restart_migration")
1586            .unwrap();
1587        assert_eq!(read_schema_state(&upgraded).unwrap().revision, 72);
1588        assert!(migrate_schema(&upgraded).is_ok());
1589
1590        let table_sql = |connection: &Connection| {
1591            connection
1592                .query_row(
1593                    "SELECT sql FROM sqlite_schema WHERE type='table' AND name='session_restart_intents'",
1594                    [],
1595                    |row| row.get::<_, String>(0),
1596                )
1597                .unwrap()
1598        };
1599        assert_eq!(table_sql(&fresh), table_sql(&upgraded));
1600        assert_eq!(read_schema_state(&fresh).unwrap().revision, SCHEMA_VERSION);
1601        assert_eq!(
1602            read_schema_state(&upgraded).unwrap().revision,
1603            SCHEMA_VERSION
1604        );
1605    }
1606
1607    fn assert_divergent_history_upgrades(revision: i64, project_history: bool, interrupt: bool) {
1608        let directory = tempfile::tempdir().unwrap();
1609        let path = directory.path().join("divergent-history.sqlite");
1610        let connection = Connection::open(&path).unwrap();
1611        create_baseline_schema(&connection).unwrap();
1612        connection.execute_batch(&format!(
1613            "INSERT INTO session_contexts(session_id,bundle_id,created_at)
1614                 VALUES ('kept','project','now');
1615             INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at,project_directory)
1616                 VALUES ('kept','Keep my work','codex','codex','local','error','now',X'2F7265706F');
1617             INSERT INTO materialized_sessions(session_id) VALUES ('kept');
1618             INSERT INTO session_turn_usage VALUES ('kept','turn',1,1,'{{\"tokens\":42}}');
1619             INSERT INTO session_provider_cost VALUES ('kept','{{\"amount\":1}}');
1620             CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1621                 WHEN NEW.version > {revision}
1622                 BEGIN SELECT RAISE(ABORT,'fixture migration boundary'); END;"
1623        )).unwrap();
1624        assert!(migrate_schema(&connection).is_err());
1625        if !connection.is_autocommit() {
1626            connection.execute_batch("ROLLBACK").unwrap();
1627        }
1628        connection
1629            .execute_batch("DROP TRIGGER stop_at_revision")
1630            .unwrap();
1631        assert_eq!(read_schema_state(&connection).unwrap().revision, revision);
1632        if project_history {
1633            connection
1634                .execute_batch(include_str!("project_catalog_v67.sql"))
1635                .unwrap();
1636            connection
1637                .execute_batch(
1638                    "INSERT INTO project_catalog VALUES ('project','key','{}',0);
1639                 INSERT INTO project_aliases VALUES ('alias','project','{}',1);
1640                 INSERT INTO project_session_aliases VALUES ('kept','alias');
1641                 UPDATE sessions SET project_json='{\"kept\":true}';",
1642                )
1643                .unwrap();
1644        }
1645        if interrupt {
1646            connection
1647                .execute_batch(
1648                    "CREATE TRIGGER interrupt_reconciliation BEFORE INSERT ON schema_migrations
1649                 WHEN NEW.version=69 BEGIN SELECT RAISE(ABORT,'interrupted reconciliation'); END;",
1650                )
1651                .unwrap();
1652            assert!(migrate_schema(&connection).is_err());
1653            drop(connection);
1654            let connection = Connection::open(&path).unwrap();
1655            assert_eq!(read_schema_state(&connection).unwrap().revision, 68);
1656            connection
1657                .execute_batch("DROP TRIGGER interrupt_reconciliation")
1658                .unwrap();
1659        } else {
1660            drop(connection);
1661        }
1662
1663        let writer = open_writer(&path).unwrap();
1664        let state = read_schema_state(&writer).unwrap();
1665        assert_eq!(state.revision, SCHEMA_VERSION);
1666        assert!(state.ensure_supported_by(68).is_err());
1667        assert!(state.ensure_supported_by(67).is_err());
1668        if project_history {
1669            let snapshot: String = writer
1670                .query_row(
1671                    "SELECT project_json FROM sessions WHERE session_id='kept'",
1672                    [],
1673                    |row| row.get(0),
1674                )
1675                .unwrap();
1676            assert_eq!(snapshot, r#"{"kept":true}"#);
1677            let alias: (String, i64) = writer.query_row(
1678                "SELECT canonical_id,config_pending FROM project_aliases WHERE bundle_id='alias'", [],
1679                |row| Ok((row.get(0)?, row.get(1)?))
1680            ).unwrap();
1681            assert_eq!(alias, ("project".to_owned(), 1));
1682        }
1683        // Rebuilding sessions must retain discovery triggers and admit the new state.
1684        writer.execute_batch(
1685            "UPDATE sessions SET state='startup-cleanup',project_directory=X'2F6E6577' WHERE session_id='kept';
1686             INSERT INTO session_turn_selections VALUES ('kept','turn','model','high');"
1687        ).unwrap();
1688        let changed: Vec<u8> = writer
1689            .query_row(
1690                "SELECT directory FROM project_discovery_changes ORDER BY sequence DESC LIMIT 1",
1691                [],
1692                |row| row.get(0),
1693            )
1694            .unwrap();
1695        assert_eq!(changed, b"/new");
1696        writer
1697            .execute("DELETE FROM sessions WHERE session_id='kept'", [])
1698            .unwrap();
1699        let usage: String = writer
1700            .query_row("SELECT body FROM session_turn_usage", [], |row| row.get(0))
1701            .unwrap();
1702        assert_eq!(usage, r#"{"tokens":42}"#);
1703        let cost: String = writer
1704            .query_row("SELECT body FROM session_provider_cost", [], |row| {
1705                row.get(0)
1706            })
1707            .unwrap();
1708        assert_eq!(cost, r#"{"amount":1}"#);
1709        assert!(
1710            !writer
1711                .prepare("PRAGMA foreign_key_check")
1712                .unwrap()
1713                .exists([])
1714                .unwrap()
1715        );
1716        drop(writer);
1717        forget_verified_schema(&path);
1718        assert_eq!(
1719            read_schema_state(&open_writer(&path).unwrap())
1720                .unwrap()
1721                .revision,
1722            SCHEMA_VERSION
1723        );
1724    }
1725
1726    #[test]
1727    fn divergent_accounting_revision_67_preserves_usage_and_adds_projects() {
1728        assert_divergent_history_upgrades(67, false, false);
1729    }
1730
1731    #[test]
1732    fn divergent_cleanup_revision_68_preserves_usage_and_adds_projects() {
1733        assert_divergent_history_upgrades(68, false, false);
1734    }
1735
1736    #[test]
1737    fn divergent_project_revision_67_preserves_aliases_snapshots_and_usage() {
1738        assert_divergent_history_upgrades(66, true, false);
1739    }
1740
1741    #[test]
1742    fn divergent_project_reconciliation_resumes_after_interruption() {
1743        assert_divergent_history_upgrades(66, true, true);
1744    }
1745
1746    #[test]
1747    fn move_ownership_upgrade_retains_sources_and_refuses_previous_daemons() {
1748        let directory = tempfile::tempdir().unwrap();
1749        let path = directory.path().join("move.sqlite");
1750        let connection = open_writer(&path).unwrap();
1751        connection.execute_batch("DROP TABLE retained_move_sources; DELETE FROM schema_migrations WHERE version>=64; UPDATE schema_compatibility SET minimum_compatible_version=63 WHERE singleton=1; DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version=63;").unwrap();
1752        drop(connection);
1753        forget_verified_schema(&path);
1754        let upgraded = open_writer(&path).unwrap();
1755        let schema = read_schema_state(&upgraded).unwrap();
1756        assert!(schema.ensure_supported_by(63).is_err());
1757        upgraded.execute("INSERT INTO retained_move_sources VALUES ('move-one','session-one','{}','[]','now')", []).unwrap();
1758        drop(upgraded);
1759        let reopened = open_writer(&path).unwrap();
1760        let count: i64 = reopened
1761            .query_row("SELECT count(*) FROM retained_move_sources", [], |row| {
1762                row.get(0)
1763            })
1764            .unwrap();
1765        assert_eq!(count, 1);
1766    }
1767
1768    #[test]
1769    fn ec2_move_ownership_upgrade_is_atomic_and_preserves_existing_moves() {
1770        let directory = tempfile::tempdir().unwrap();
1771        let path = directory.path().join("ec2-move-migration.sqlite3");
1772        save_session_to(&path, &super::super::tests::session("source", "project")).unwrap();
1773        stamp_schema_version(&path, 70);
1774        let old = Connection::open(&path).unwrap();
1775        old.execute_batch(
1776            "UPDATE schema_compatibility SET minimum_compatible_version=69;
1777             INSERT INTO session_moves(session_id,operation_id,operation_json)
1778                 VALUES('source','existing-move','{\"operation_id\":\"existing-move\"}');
1779             CREATE TRIGGER stop_ec2_move_migration BEFORE INSERT ON schema_migrations
1780             WHEN NEW.version=71 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1781        )
1782        .unwrap();
1783        assert!(migrate_schema(&old).is_err());
1784        let state = read_schema_state(&old).unwrap();
1785        assert_eq!(state.revision, 70);
1786        assert_eq!(state.minimum_compatible, Some(69));
1787        old.execute_batch("DROP TRIGGER stop_ec2_move_migration")
1788            .unwrap();
1789        drop(old);
1790        let upgraded = open_writer(&path).unwrap();
1791        let state = read_schema_state(&upgraded).unwrap();
1792        assert_eq!(state.revision, SCHEMA_VERSION);
1793        assert_eq!(state.minimum_compatible, Some(74));
1794        assert!(state.ensure_supported_by(70).is_err());
1795        let preserved: String = upgraded
1796            .query_row(
1797                "SELECT operation_json FROM session_moves WHERE session_id='source'",
1798                [],
1799                |row| row.get(0),
1800            )
1801            .unwrap();
1802        assert_eq!(preserved, r#"{"operation_id":"existing-move"}"#);
1803        assert_eq!(
1804            upgraded
1805                .query_row(
1806                    "SELECT count(*) FROM sessions WHERE session_id='source'",
1807                    [],
1808                    |row| row.get::<_, i64>(0)
1809                )
1810                .unwrap(),
1811            1
1812        );
1813    }
1814
1815    #[test]
1816    fn title_migration_caps_old_titles_preserves_short_titles_and_is_compatible() {
1817        let directory = tempfile::tempdir().unwrap();
1818        let path = directory.path().join("title-migration.sqlite3");
1819        let titles = [
1820            ("long", Some("word ".repeat(20_000))),
1821            ("unicode", Some("界".repeat(257))),
1822            ("exact", Some("界".repeat(256))),
1823            ("short", Some("  Keep\nthis title  ".into())),
1824            ("unset", None),
1825        ];
1826        for (id, _) in &titles {
1827            let mut session = super::super::tests::session(id, "project");
1828            session.state = SessionState::Stopped;
1829            save_session_to(&path, &session).unwrap();
1830        }
1831        stamp_schema_version(&path, 69);
1832        let old = Connection::open(&path).unwrap();
1833        old.execute_batch(
1834            "UPDATE schema_compatibility SET minimum_compatible_version=69;
1835             CREATE TRIGGER stop_after_title_migration BEFORE INSERT ON schema_migrations
1836             WHEN NEW.version=71 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1837        )
1838        .unwrap();
1839        for (id, title) in &titles {
1840            old.execute(
1841                "UPDATE sessions SET acp_session_title=?2 WHERE session_id=?1",
1842                params![id, title],
1843            )
1844            .unwrap();
1845        }
1846        assert!(migrate_schema(&old).is_err());
1847        let upgraded = old;
1848        let state = read_schema_state(&upgraded).unwrap();
1849        assert_eq!(state.revision, 70);
1850        assert_eq!(state.minimum_compatible, Some(69));
1851        state.ensure_supported_by(69).unwrap();
1852        for (id, original) in &titles {
1853            let stored: Option<String> = upgraded
1854                .query_row(
1855                    "SELECT acp_session_title FROM sessions WHERE session_id=?1",
1856                    [id],
1857                    |row| row.get(0),
1858                )
1859                .unwrap();
1860            let expected = match *id {
1861                "long" => Some(format!("{}word…", "word ".repeat(50))),
1862                "unicode" => Some(format!("{}…", "界".repeat(255))),
1863                _ => original.clone(),
1864            };
1865            assert_eq!(stored, expected, "session {id}");
1866            assert!(stored.is_none_or(|title| title.chars().count() <= 256));
1867        }
1868        upgraded
1869            .execute_batch("DROP TRIGGER stop_after_title_migration")
1870            .unwrap();
1871        drop(upgraded);
1872        forget_verified_schema(&path);
1873        drop(open_writer(&path).unwrap());
1874    }
1875
1876    #[test]
1877    fn interrupted_title_migration_rolls_back_titles_and_revision() {
1878        let directory = tempfile::tempdir().unwrap();
1879        let path = directory.path().join("title-migration-interrupted.sqlite3");
1880        save_session_to(&path, &super::super::tests::session("old", "project")).unwrap();
1881        stamp_schema_version(&path, 69);
1882        let original = "word ".repeat(20_000);
1883        let old = Connection::open(&path).unwrap();
1884        old.execute("UPDATE sessions SET acp_session_title=?1", [&original])
1885            .unwrap();
1886        old.execute_batch(
1887            "CREATE TRIGGER stop_title_migration BEFORE INSERT ON schema_migrations
1888             WHEN NEW.version=70 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1889        )
1890        .unwrap();
1891        assert!(migrate_schema(&old).is_err());
1892        assert_eq!(read_schema_state(&old).unwrap().revision, 69);
1893        let stored: String = old
1894            .query_row("SELECT acp_session_title FROM sessions", [], |row| {
1895                row.get(0)
1896            })
1897            .unwrap();
1898        assert_eq!(stored, original);
1899        old.execute_batch("DROP TRIGGER stop_title_migration")
1900            .unwrap();
1901        drop(old);
1902        drop(open_writer(&path).unwrap());
1903        assert!(
1904            load_state_from(&path).unwrap().sessions["old"]
1905                .acp_session_title
1906                .as_ref()
1907                .unwrap()
1908                .chars()
1909                .count()
1910                <= 256
1911        );
1912    }
1913
1914    #[test]
1915    fn recent_revisions_upgrade_directly_and_preserve_user_data() {
1916        // Revision 47 was current on 2026-09-23. Keep the exhaustive
1917        // interruption matrix focused on the last week's migration history.
1918        const TEST_REVISION_FLOOR: i64 = 47;
1919        for revision in TEST_REVISION_FLOOR..SCHEMA_VERSION {
1920            let directory = tempfile::tempdir().unwrap();
1921            let path = directory.path().join("mj.sqlite3");
1922            let connection = Connection::open(&path).unwrap();
1923            connection
1924                .execute_batch(include_str!("legacy_v1.sql"))
1925                .unwrap();
1926            connection.execute_batch(
1927                "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
1928                 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
1929                     target_template_id, state, updated_at, native_session_id)
1930                     VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
1931                         '2026-01-01T00:00:00Z', 'native-original');
1932                 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
1933                     VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
1934            ).unwrap();
1935            // Interrupt at each historical transaction boundary. Closing the
1936            // connection rolls back the interrupted step, exactly as a killed
1937            // updater would. Reopening must resume without manual cleanup.
1938            connection
1939                .execute_batch(&format!(
1940                    "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1941                 WHEN NEW.version > {revision}
1942                 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
1943                ))
1944                .unwrap();
1945            super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
1946            assert!(migrate_schema(&connection).is_err());
1947            drop(connection);
1948            let connection = Connection::open(&path).unwrap();
1949            let found: i64 = connection
1950                .query_row("PRAGMA user_version", [], |row| row.get(0))
1951                .unwrap();
1952            assert_eq!(found, revision);
1953            connection
1954                .execute_batch("DROP TRIGGER stop_at_revision")
1955                .unwrap();
1956            drop(connection);
1957
1958            let writer =
1959                open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
1960            assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
1961            assert_eq!(
1962                writer
1963                    .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
1964                    .unwrap(),
1965                "ok"
1966            );
1967            assert!(
1968                !writer
1969                    .prepare("PRAGMA foreign_key_check")
1970                    .unwrap()
1971                    .exists([])
1972                    .unwrap()
1973            );
1974            drop(writer);
1975            let reader = open_reader_strict(&path).unwrap();
1976            let prompt: String = reader
1977                .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
1978                .unwrap();
1979            assert_eq!(prompt, "Keep my prompt");
1980            let state = load_state_from(&path).unwrap();
1981            assert_eq!(state.sessions["old-session"].title, "Keep my work");
1982            assert_eq!(
1983                state.sessions["old-session"].native_session_id.as_deref(),
1984                Some("native-original")
1985            );
1986            drop(reader);
1987            // A fresh writer process must not apply any migration twice.
1988            forget_verified_schema(&path);
1989            drop(open_writer(&path).unwrap());
1990        }
1991    }
1992
1993    #[test]
1994    fn accounting_migration_preserves_usage_and_changes_its_deletion_owner() {
1995        let dir = tempfile::tempdir().unwrap();
1996        let path = dir.path().join("migration.sqlite");
1997        let connection = Connection::open(&path).unwrap();
1998        connection
1999            .execute_batch(include_str!("legacy_v1.sql"))
2000            .unwrap();
2001        connection.execute_batch("INSERT INTO session_contexts VALUES ('old-session','project','2026-01-01T00:00:00Z');
2002            INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at)
2003            VALUES ('old-session','Retain usage','codex','codex','local','error','2026-01-01T00:00:00Z');
2004            CREATE TRIGGER stop_before_accounting BEFORE INSERT ON schema_migrations WHEN NEW.version=67
2005            BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;").unwrap();
2006        super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
2007        assert!(migrate_schema(&connection).is_err());
2008        connection
2009            .execute_batch(
2010                "ROLLBACK; DROP TRIGGER stop_before_accounting;
2011            INSERT INTO session_turn_usage VALUES ('old-session','turn',7,1,'{}');
2012            INSERT INTO session_provider_cost VALUES ('old-session','{\"amount\":1}');",
2013            )
2014            .unwrap();
2015        drop(connection);
2016        let writer = open_writer(&path).unwrap();
2017        writer
2018            .execute("DELETE FROM sessions WHERE session_id='old-session'", [])
2019            .unwrap();
2020        for table in ["session_turn_usage", "session_provider_cost"] {
2021            assert_eq!(
2022                writer
2023                    .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| row
2024                        .get::<_, u64>(0))
2025                    .unwrap(),
2026                1
2027            );
2028        }
2029        assert_eq!(
2030            writer
2031                .query_row("SELECT body FROM session_turn_usage", [], |row| row
2032                    .get::<_, String>(0))
2033                .unwrap(),
2034            "{}"
2035        );
2036        assert!(
2037            !writer
2038                .prepare("PRAGMA foreign_key_check")
2039                .unwrap()
2040                .exists([])
2041                .unwrap()
2042        );
2043        assert_eq!(
2044            read_schema_state(&writer).unwrap().minimum_compatible,
2045            Some(MINIMUM_COMPATIBLE_VERSION)
2046        );
2047    }
2048
2049    /// The oldest executable revision that can still read and write a store at
2050    /// `SCHEMA_VERSION`. Migration 74 adds the OpenCode harness kind.
2051    const MINIMUM_COMPATIBLE_VERSION: i64 = 74;
2052
2053    /// Rewrites a store's recorded schema version the way another build's
2054    /// migration ladder would, and forgets that this process verified it.
2055    fn stamp_schema_version(path: &Path, version: i64) {
2056        if version > SCHEMA_VERSION {
2057            advance_test_schema(path, version, version);
2058            return;
2059        }
2060        let connection = Connection::open(path).unwrap();
2061        if version < 67 {
2062            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();
2063        }
2064        connection
2065            .execute_batch(&format!("PRAGMA user_version = {version};"))
2066            .unwrap();
2067        connection
2068            .execute(
2069                "DELETE FROM schema_migrations WHERE version > ?1",
2070                [version],
2071            )
2072            .unwrap();
2073        if version == 30 {
2074            connection
2075                .execute(
2076                    "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
2077                    [],
2078                )
2079                .unwrap();
2080        }
2081        if matches!(version, 69 | 70) {
2082            connection.execute(
2083                "UPDATE schema_compatibility SET minimum_compatible_version=69 WHERE singleton=1", [],
2084            ).unwrap();
2085        }
2086        drop(connection);
2087        forget_verified_schema(path);
2088    }
2089
2090    #[test]
2091    fn subagent_policy_migration_preserves_legacy_choices_and_refuses_old_writers() {
2092        use mj_core::subagent::SubagentPolicy;
2093        let directory = tempfile::tempdir().unwrap();
2094        let path = directory.path().join("subagent-policy.sqlite3");
2095        for id in ["all", "native", "unset"] {
2096            save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
2097        }
2098        let connection = Connection::open(&path).unwrap();
2099        connection
2100            .execute_batch(
2101                "UPDATE sessions SET mjolnir_subagents = 1 WHERE session_id = 'all';
2102             UPDATE sessions SET mjolnir_subagents = 0 WHERE session_id = 'native';
2103             ALTER TABLE sessions DROP COLUMN subagents;
2104             DROP TABLE subagent_preference;
2105             DELETE FROM schema_migrations WHERE version >= 56;
2106             UPDATE schema_compatibility SET minimum_compatible_version = 55;
2107             DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 55;",
2108            )
2109            .unwrap();
2110        drop(connection);
2111        forget_verified_schema(&path);
2112        let upgraded = open_writer(&path).unwrap();
2113        assert!(
2114            read_schema_state(&upgraded)
2115                .unwrap()
2116                .ensure_supported_by(55)
2117                .is_err()
2118        );
2119        drop(upgraded);
2120        let state = load_state_from(&path).unwrap();
2121        assert_eq!(
2122            state.sessions["all"].subagents,
2123            Some(SubagentPolicy::AllModels)
2124        );
2125        assert_eq!(
2126            state.sessions["native"].subagents,
2127            Some(SubagentPolicy::Native)
2128        );
2129        assert_eq!(state.sessions["unset"].subagents, None);
2130        assert_eq!(state.last_subagent_policy, SubagentPolicy::Native);
2131    }
2132
2133    #[test]
2134    fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
2135        let directory = tempfile::tempdir().unwrap();
2136        let path = directory.path().join("steering-migration.sqlite3");
2137        let record = super::super::tests::session("steered-session", "project");
2138        save_session_to(&path, &record).unwrap();
2139        let connection = Connection::open(&path).unwrap();
2140        connection
2141            .execute_batch(
2142                "DELETE FROM schema_migrations WHERE version >= 48;
2143             UPDATE schema_compatibility SET minimum_compatible_version = 47;
2144             DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 47;",
2145            )
2146            .unwrap();
2147        drop(connection);
2148        forget_verified_schema(&path);
2149        let upgraded = open_writer(&path).unwrap();
2150        let schema = read_schema_state(&upgraded).unwrap();
2151        assert_eq!(schema.revision, SCHEMA_VERSION);
2152        assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2153        let error = schema.ensure_supported_by(47).unwrap_err();
2154        assert!(matches!(
2155            error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
2156            StoreSchemaMismatchReason::Incompatible {
2157                minimum_compatible: MINIMUM_COMPATIBLE_VERSION
2158            }
2159        ));
2160        drop(upgraded);
2161        assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
2162    }
2163
2164    #[test]
2165    fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
2166        let directory = tempfile::tempdir().unwrap();
2167        let path = directory.path().join("target-migration.sqlite3");
2168        let record = super::super::tests::session("preserved-session", "project");
2169        save_session_to(&path, &record).unwrap();
2170        let connection = Connection::open(&path).unwrap();
2171        connection
2172            .execute_batch(
2173                "ALTER TABLE sessions DROP COLUMN target_runtime_json;
2174             DELETE FROM schema_migrations WHERE version >= 46;
2175             UPDATE schema_compatibility SET minimum_compatible_version = 44;
2176             DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 45;",
2177            )
2178            .unwrap();
2179        forget_verified_schema(&path);
2180        let upgraded = open_writer(&path).unwrap();
2181        let schema = read_schema_state(&upgraded).unwrap();
2182        assert_eq!(schema.revision, SCHEMA_VERSION);
2183        assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2184        let error = schema.ensure_supported_by(45).unwrap_err();
2185        assert!(matches!(
2186            error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
2187            StoreSchemaMismatchReason::Incompatible {
2188                minimum_compatible: MINIMUM_COMPATIBLE_VERSION
2189            }
2190        ));
2191        drop(upgraded);
2192        let restored = load_state_from(&path).unwrap();
2193        assert_eq!(restored.sessions[&record.id], record);
2194    }
2195
2196    #[test]
2197    fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
2198        let directory = tempfile::tempdir().unwrap();
2199        let path = directory.path().join("mj.sqlite3");
2200        let connection = open_writer(&path).unwrap();
2201        connection
2202            .execute_batch(
2203                "BEGIN IMMEDIATE;
2204             DROP TABLE quota_reset_cache;
2205             DROP TABLE native_agent_transcript;
2206             DROP TABLE native_agents;
2207             DROP TABLE native_agent_replay;
2208             DELETE FROM schema_migrations WHERE version >= 39;
2209             UPDATE schema_compatibility SET minimum_compatible_version = 32;
2210             DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; PRAGMA writable_schema=ON; UPDATE sqlite_schema SET sql=replace(sql, '''startup-cleanup'',', '') WHERE type='table' AND name='sessions'; PRAGMA writable_schema=RESET; PRAGMA user_version = 38;
2211             COMMIT;",
2212            )
2213            .unwrap();
2214        migrate_schema(&connection).unwrap();
2215        let state = read_schema_state(&connection).unwrap();
2216        assert_eq!(state.revision, SCHEMA_VERSION);
2217        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2218        let event = ApiEventData::InputRequired {
2219            request: None,
2220            turn_id: Some(1),
2221        };
2222        #[derive(serde::Deserialize)]
2223        struct LegacyInputEvent {
2224            #[serde(rename = "request")]
2225            _request: mj_core::elicitation::ElicitationRequest,
2226        }
2227        let encoded = serde_json::to_value(&event).unwrap();
2228        assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
2229    }
2230
2231    #[test]
2232    fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
2233        let directory = tempfile::tempdir().unwrap();
2234        let path = directory.path().join("mj.sqlite3");
2235        let connection = open_writer(&path).unwrap();
2236        connection
2237            .execute_batch(
2238                "CREATE TABLE future_feature(value TEXT NOT NULL);
2239                 INSERT INTO future_feature VALUES ('preserve me');",
2240            )
2241            .unwrap();
2242        drop(connection);
2243        advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
2244
2245        let reader = open_reader_strict(&path).unwrap();
2246        assert_eq!(
2247            reader
2248                .query_row("SELECT value FROM future_feature", [], |row| row
2249                    .get::<_, String>(0))
2250                .unwrap(),
2251            "preserve me"
2252        );
2253        assert!(reader.execute("DELETE FROM future_feature", []).is_err());
2254        drop(reader);
2255
2256        // A repair would recreate this deliberately removed trigger. A future
2257        // schema is authoritative even when it differs from our own repairs.
2258        let raw = Connection::open(&path).unwrap();
2259        raw.execute_batch("DROP TRIGGER session_contexts_workspace_update;")
2260            .unwrap();
2261        drop(raw);
2262        let writer = open_writer(&path).unwrap();
2263        assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'session_contexts_workspace_update')", [], |row| row.get::<_, bool>(0)).unwrap());
2264        assert_eq!(
2265            writer
2266                .query_row("SELECT value FROM future_feature", [], |row| row
2267                    .get::<_, String>(0))
2268                .unwrap(),
2269            "preserve me"
2270        );
2271        let state = read_schema_state(&writer).unwrap();
2272        assert_eq!(state.revision, SCHEMA_VERSION + 1);
2273        assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
2274    }
2275
2276    #[test]
2277    fn invalid_compatibility_metadata_refuses_readers_and_writers() {
2278        for alteration in [
2279            "DROP TABLE schema_compatibility",
2280            "DELETE FROM schema_compatibility",
2281            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
2282            "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
2283            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
2284            "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
2285            "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
2286            "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
2287        ] {
2288            for future in [false, true] {
2289                let directory = tempfile::tempdir().unwrap();
2290                let path = directory.path().join("mj.sqlite3");
2291                drop(open_writer(&path).unwrap());
2292                if future {
2293                    advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
2294                }
2295                let raw = Connection::open(&path).unwrap();
2296                raw.execute_batch(alteration).unwrap();
2297                let before: i64 = raw
2298                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
2299                    .unwrap();
2300                // Exercise the cached path as well as a fresh writer open.
2301                for error in [
2302                    open_reader_strict(&path).unwrap_err(),
2303                    open_writer(&path).unwrap_err(),
2304                ] {
2305                    let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
2306                    assert_eq!(
2307                        mismatch.reason,
2308                        StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
2309                        "{alteration}"
2310                    );
2311                }
2312                forget_verified_schema(&path);
2313                assert!(open_writer(&path).is_err(), "{alteration}");
2314                let after: i64 = raw
2315                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
2316                    .unwrap();
2317                assert_eq!(
2318                    before, after,
2319                    "a rejected open repaired schema: {alteration}"
2320                );
2321            }
2322        }
2323    }
2324
2325    #[test]
2326    fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
2327        let directory = tempfile::tempdir().unwrap();
2328        let path = directory.path().join("mj.sqlite3");
2329        let connection = Connection::open(&path).unwrap();
2330        // A table the baseline also creates makes its batch fail part way.
2331        connection
2332            .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
2333            .unwrap();
2334
2335        let error = migrate_schema(&connection).unwrap_err();
2336
2337        assert!(format!("{error:#}").contains("create baseline database schema"));
2338        assert!(
2339            connection.is_autocommit(),
2340            "the failed baseline left a transaction open"
2341        );
2342        assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
2343        let tables: i64 = connection
2344            .query_row(
2345                "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
2346                [],
2347                |row| row.get(0),
2348            )
2349            .unwrap();
2350        assert_eq!(tables, 1, "only the conflicting table remains");
2351
2352        connection.execute_batch("DROP TABLE workspaces").unwrap();
2353        drop(connection);
2354        let writer = open_writer(&path).unwrap();
2355        let state = read_schema_state(&writer).unwrap();
2356        assert_eq!(state.revision, SCHEMA_VERSION);
2357        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2358    }
2359
2360    /// A store ahead of this build cannot be fixed by starting a daemon of
2361    /// this build, so the reader must not say so. This is the message the
2362    /// incident in #24 printed twice a second for an hour.
2363    #[test]
2364    fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
2365        let directory = tempfile::tempdir().unwrap();
2366        let path = directory.path().join("mj.sqlite3");
2367        drop(open_writer(&path).unwrap());
2368        stamp_schema_version(&path, SCHEMA_VERSION + 1);
2369
2370        let error = open_reader_strict(&path).unwrap_err();
2371
2372        let mismatch = error
2373            .chain()
2374            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2375            .expect("the reader reports the mismatch as a typed cause");
2376        assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
2377        assert_eq!(mismatch.supported, SCHEMA_VERSION);
2378        let message = mismatch.to_string();
2379        assert!(message.contains("upgrade Mjolnir"), "got {message}");
2380        assert!(
2381            !message.contains("start the Mjolnir daemon"),
2382            "got {message}"
2383        );
2384    }
2385
2386    /// A store behind this build keeps the advice that works, verbatim, so
2387    /// existing log greps and runbooks keep matching.
2388    #[test]
2389    fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
2390        let directory = tempfile::tempdir().unwrap();
2391        let path = directory.path().join("mj.sqlite3");
2392        drop(open_writer(&path).unwrap());
2393        let raw = Connection::open(&path).unwrap();
2394        raw.execute_batch(&format!(
2395            "UPDATE schema_compatibility SET minimum_compatible_version = {0};
2396             DELETE FROM schema_migrations WHERE version > {0};
2397             INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
2398             PRAGMA user_version = {0};",
2399            SCHEMA_VERSION - 1
2400        ))
2401        .unwrap();
2402        drop(raw);
2403
2404        let error = open_reader_strict(&path).unwrap_err();
2405
2406        let mismatch = error
2407            .chain()
2408            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2409            .expect("the reader reports the mismatch as a typed cause");
2410        assert_eq!(
2411            mismatch.to_string(),
2412            format!(
2413                "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
2414                 start the Mjolnir daemon to migrate it",
2415                SCHEMA_VERSION - 1
2416            )
2417        );
2418    }
2419
2420    #[test]
2421    fn strict_reader_rejects_mutation() {
2422        let directory = tempfile::tempdir().unwrap();
2423        let path = directory.path().join("mj.sqlite3");
2424        drop(open_writer(&path).unwrap());
2425
2426        let reader = open_reader_strict(&path).unwrap();
2427        let error = reader
2428            .execute("CREATE TABLE forbidden(value TEXT)", [])
2429            .unwrap_err();
2430        assert!(
2431            matches!(
2432                error.sqlite_error_code(),
2433                Some(rusqlite::ErrorCode::ReadOnly)
2434            ),
2435            "unexpected mutation error: {error}"
2436        );
2437    }
2438
2439    fn opened_readers() -> usize {
2440        OPENED_READERS.with(std::cell::Cell::get)
2441    }
2442
2443    fn has_workspace(reader: &Connection, name: &str) -> bool {
2444        reader
2445            .query_row(
2446                "SELECT EXISTS(SELECT 1 FROM workspaces WHERE name = ?1)",
2447                [name],
2448                |row| row.get(0),
2449            )
2450            .unwrap()
2451    }
2452
2453    fn add_workspace(path: &Path, name: &str) {
2454        open_writer(path)
2455            .unwrap()
2456            .execute(
2457                "INSERT INTO workspaces(workspace_id, name, name_key, created_at, last_opened_at)
2458                 VALUES (?1, ?1, ?1, 'now', 'now')",
2459                [name],
2460            )
2461            .unwrap();
2462    }
2463
2464    /// Reads reuse an idle connection instead of opening one and parsing the
2465    /// schema again, and still see every write committed since.
2466    #[test]
2467    fn strict_readers_are_reused_and_see_later_commits() {
2468        let directory = tempfile::tempdir().unwrap();
2469        let path = directory.path().join("mj.sqlite3");
2470        drop(open_writer(&path).unwrap());
2471        let before = opened_readers();
2472        assert!(!has_workspace(&open_reader_strict(&path).unwrap(), "later"));
2473        add_workspace(&path, "later");
2474        for _ in 0..20 {
2475            assert!(has_workspace(&open_reader_strict(&path).unwrap(), "later"));
2476        }
2477        assert_eq!(
2478            opened_readers() - before,
2479            1,
2480            "one connection served every read"
2481        );
2482
2483        // Two readers at once need two connections; both are kept afterwards.
2484        let first = open_reader_strict(&path).unwrap();
2485        let second = open_reader_strict(&path).unwrap();
2486        drop((first, second));
2487        assert_eq!(opened_readers() - before, 2);
2488        drop(open_reader_strict(&path).unwrap());
2489        drop(open_reader_strict(&path).unwrap());
2490        assert_eq!(opened_readers() - before, 2);
2491
2492        // A connection given back inside a transaction is not reused.
2493        let reader = open_reader_strict(&path).unwrap();
2494        reader.execute_batch("BEGIN").unwrap();
2495        drop(reader);
2496        let (first, second) = (
2497            open_reader_strict(&path).unwrap(),
2498            open_reader_strict(&path).unwrap(),
2499        );
2500        assert!(first.is_autocommit() && second.is_autocommit());
2501        assert_eq!(opened_readers() - before, 3);
2502    }
2503
2504    /// An idle connection must not outlive the store it was opened on, and
2505    /// must not hide a migration another build made while it sat idle.
2506    #[test]
2507    fn idle_readers_follow_store_replacement_and_schema_changes() {
2508        let directory = tempfile::tempdir().unwrap();
2509        let path = directory.path().join("mj.sqlite3");
2510        drop(open_writer(&path).unwrap());
2511        add_workspace(&path, "old store");
2512        drop(open_reader_strict(&path).unwrap());
2513
2514        for suffix in ["", "-wal", "-shm"] {
2515            let file = PathBuf::from(format!("{}{suffix}", path.display()));
2516            if file.exists() {
2517                fs::remove_file(file).unwrap();
2518            }
2519        }
2520        forget_verified_schema(&path);
2521        drop(open_writer(&path).unwrap());
2522        add_workspace(&path, "new store");
2523        let reader = open_reader_strict(&path).unwrap();
2524        assert!(has_workspace(&reader, "new store"));
2525        assert!(!has_workspace(&reader, "old store"));
2526        drop(reader);
2527
2528        drop(open_reader_strict(&path).unwrap());
2529        stamp_schema_version(&path, SCHEMA_VERSION + 1);
2530        let error = open_reader_strict(&path).unwrap_err();
2531        assert!(
2532            error
2533                .chain()
2534                .any(|cause| cause.downcast_ref::<StoreSchemaMismatch>().is_some()),
2535            "a reused connection checks compatibility like a new one: {error:#}"
2536        );
2537    }
2538}