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