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