Skip to main content

mj_controller/database/
schema.rs

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