Skip to main content

mj_controller/database/
schema.rs

1use super::*;
2use rusqlite::OpenFlags;
3
4const COMPATIBILITY_METADATA_VERSION: i64 = 30;
5
6pub(super) struct SchemaState {
7    pub(super) revision: i64,
8    minimum_compatible: Option<i64>,
9}
10
11impl SchemaState {
12    pub(super) fn ensure_supported(&self) -> Result<()> {
13        let reason = if self.revision < SCHEMA_VERSION {
14            StoreSchemaMismatchReason::NeedsMigration
15        } else if let Some(minimum_compatible) = self.minimum_compatible {
16            if minimum_compatible <= SCHEMA_VERSION {
17                return Ok(());
18            }
19            StoreSchemaMismatchReason::Incompatible { minimum_compatible }
20        } else {
21            StoreSchemaMismatchReason::InvalidCompatibilityMetadata
22        };
23        Err(StoreSchemaMismatch {
24            found: self.revision,
25            supported: SCHEMA_VERSION,
26            reason,
27        }
28        .into())
29    }
30}
31
32/// The revision, ledger, and compatibility floor must describe one snapshot.
33/// A missing floor is only legitimate before compatibility was introduced.
34pub(super) fn read_schema_state(connection: &Connection) -> Result<SchemaState> {
35    let snapshot = connection
36        .unchecked_transaction()
37        .context("start database compatibility snapshot")?;
38    let revision: i64 = snapshot
39        .query_row("PRAGMA user_version", [], |row| row.get(0))
40        .context("read database migration revision")?;
41    let minimum_compatible = if revision >= COMPATIBILITY_METADATA_VERSION {
42        let invalid = || StoreSchemaMismatch {
43            found: revision,
44            supported: SCHEMA_VERSION,
45            reason: StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
46        };
47        let (count, singleton, floor, recorded): (i64, Option<i64>, Option<i64>, Option<i64>) =
48            snapshot
49                .query_row(
50                    "SELECT count(*), min(singleton), min(minimum_compatible_version),
51                    (SELECT max(version) FROM schema_migrations)
52             FROM schema_compatibility",
53                    [],
54                    |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
55                )
56                .map_err(|error| {
57                    // Missing tables/columns and invalid field types are
58                    // structural. Busy, I/O, and interruption errors are not
59                    // evidence of an incompatible migration.
60                    let structural = match &error {
61                        rusqlite::Error::SqliteFailure(code, _) => {
62                            code.code == rusqlite::ErrorCode::Unknown
63                        }
64                        _ => true,
65                    };
66                    let error = anyhow::Error::new(error);
67                    if structural {
68                        error.context(invalid())
69                    } else {
70                        error.context("read database compatibility metadata")
71                    }
72                })?;
73        if count != 1
74            || singleton != Some(1)
75            || recorded != Some(revision)
76            || !floor
77                .is_some_and(|floor| (COMPATIBILITY_METADATA_VERSION..=revision).contains(&floor))
78        {
79            return Err(invalid().into());
80        }
81        floor
82    } else {
83        None
84    };
85    snapshot
86        .commit()
87        .context("finish database compatibility snapshot")?;
88    Ok(SchemaState {
89        revision,
90        minimum_compatible,
91    })
92}
93
94pub fn database_path() -> PathBuf {
95    data_dir().join("mj.sqlite3")
96}
97
98/// Verify that this client can read the daemon's store without creating or
99/// migrating it. Startup must pass this gate before handing out a connection,
100/// even when an older daemon happens to speak the same wire protocol.
101pub fn check_read_compatibility() -> Result<()> {
102    open_reader_strict(&database_path()).map(drop)
103}
104
105/// A writer-capable connection whose transactions take the WAL write lock at
106/// `BEGIN`, where the busy handler applies. A DEFERRED transaction that has
107/// already read cannot wait: SQLite only calls the busy handler when the
108/// connection holds no transaction, so the upgrade to a write returns
109/// `SQLITE_BUSY` at once (issue 1117).
110pub(super) fn open_writer(path: &Path) -> Result<Connection> {
111    let mut connection = open_writable(path)?;
112    connection.set_transaction_behavior(rusqlite::TransactionBehavior::Immediate);
113    Ok(connection)
114}
115
116fn open_writable(path: &Path) -> Result<Connection> {
117    if let Some(parent) = path.parent() {
118        fs::create_dir_all(parent)
119            .with_context(|| format!("create Mjolnir data directory {}", parent.display()))?;
120    }
121    let connection = Connection::open(path)
122        .with_context(|| format!("open Mjolnir database {}", path.display()))?;
123    connection.busy_timeout(Duration::from_secs(5))?;
124    connection.execute_batch(
125        "PRAGMA foreign_keys = ON;
126         PRAGMA journal_mode = WAL;
127         PRAGMA synchronous = FULL;",
128    )?;
129    verify_schema_once(path, &connection)?;
130    Ok(connection)
131}
132
133pub(super) fn open(path: &Path) -> Result<Connection> {
134    open_writer(path)
135}
136
137/// Open an existing database without permitting schema or data mutation.
138/// Client processes use this path so an accidental write fails locally
139/// instead of competing with the daemon's writer.
140#[cfg(not(test))]
141pub(super) fn open_reader(path: &Path) -> Result<Connection> {
142    open_reader_strict(path)
143}
144
145#[cfg(test)]
146pub(super) fn open_reader(path: &Path) -> Result<Connection> {
147    // Path-taking database helpers are migration fixtures in unit tests: they
148    // intentionally open old or not-yet-created schemas. Production query
149    // entry points compile against the strict reader above. The connection is
150    // writable but keeps SQLite's DEFERRED default, so a fixture read does not
151    // take the write lock.
152    open_writable(path)
153}
154
155#[cfg_attr(test, allow(dead_code))]
156fn open_reader_strict(path: &Path) -> Result<Connection> {
157    let connection = Connection::open_with_flags(
158        path,
159        OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
160    )
161    .with_context(|| format!("open Mjolnir database read-only {}", path.display()))?;
162    connection.busy_timeout(Duration::from_secs(5))?;
163    connection.execute_batch(
164        "PRAGMA foreign_keys = ON;
165         PRAGMA query_only = ON;",
166    )?;
167    read_schema_state(&connection)?.ensure_supported()?;
168    Ok(connection)
169}
170
171/// Databases this process has already migrated. A controller owns its store
172/// exclusively (`ControllerStoreGuard`), so a schema verified once stays
173/// verified and later connections skip the migration probes entirely.
174fn verified_schemas() -> &'static Mutex<HashSet<PathBuf>> {
175    static VERIFIED: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
176    VERIFIED.get_or_init(|| Mutex::new(HashSet::new()))
177}
178
179/// Stable cache identity for a database. The file itself may not exist yet, so
180/// the canonicalized parent directory carries the identity.
181fn schema_cache_key(path: &Path) -> PathBuf {
182    let Some(parent) = path
183        .parent()
184        .filter(|parent| !parent.as_os_str().is_empty())
185    else {
186        return path.to_owned();
187    };
188    match (fs::canonicalize(parent), path.file_name()) {
189        (Ok(canonical), Some(name)) => canonical.join(name),
190        _ => path.to_owned(),
191    }
192}
193
194/// Run the migration ladder the first time this process opens a database.
195/// Later opens confirm compatibility without repeating schema repairs. A
196/// database behind this build is migrated again, so a recreated file under a
197/// reused path still converges. Compatible future stores are never repaired.
198fn verify_schema_once(path: &Path, connection: &Connection) -> Result<()> {
199    let key = schema_cache_key(path);
200    let mut verified = verified_schemas()
201        .lock()
202        .unwrap_or_else(PoisonError::into_inner);
203    let state = read_schema_state(connection)?;
204    if state.revision > SCHEMA_VERSION
205        || (state.revision == SCHEMA_VERSION && verified.contains(&key))
206    {
207        // An older build must never run its repairs against a newer schema.
208        return state.ensure_supported();
209    }
210    // Holding the lock across the ladder keeps two first opens of the same
211    // database from running the additive migration steps against each other.
212    migrate_schema(connection)?;
213    read_schema_state(connection)?.ensure_supported()?;
214    verified.insert(key);
215    Ok(())
216}
217
218/// Forget that this process verified a database's schema. Only tests need it:
219/// they simulate a store written by an older build by editing the schema of a
220/// database this process has already opened, which no controller can do.
221#[cfg(test)]
222pub(super) fn forget_verified_schema(path: &Path) {
223    verified_schemas()
224        .lock()
225        .unwrap_or_else(PoisonError::into_inner)
226        .remove(&schema_cache_key(path));
227}
228
229/// New stores start at the revision Mjolnir 2.7.2 shipped. Existing stores
230/// retain the complete migration path, even when users skip many releases.
231const BASELINE_SCHEMA_VERSION: i64 = 33;
232
233/// The compatibility floor a baseline store records. Migration 32 (ZCode) was
234/// the last breaking change before the baseline.
235const BASELINE_MINIMUM_COMPATIBLE_VERSION: i64 = 32;
236
237fn migrate_schema(connection: &Connection) -> Result<()> {
238    let state = read_schema_state(connection)?;
239    let version = state.revision;
240    if version > SCHEMA_VERSION {
241        return state.ensure_supported();
242    }
243    if version == 0 {
244        create_baseline_schema(connection)?;
245    } else if version < BASELINE_SCHEMA_VERSION {
246        super::legacy_schema::migrate_to_baseline(connection)
247            .context("upgrade historical database schema")?;
248    }
249    // Compatible: adds one table. Older readers ignore it and treat read-write
250    // mounts as copy-on-write, a behaviour difference rather than lost data.
251    // Older writers rewrite `session_mounts` but never touch this table, so its
252    // rows survive their updates; a row only applies while a mount with the same
253    // source and destination is still not read-only, so an older build that
254    // makes the mount read-only or removes it keeps that choice. The
255    // compatibility floor stays where it is.
256    if version < 34 {
257        connection.execute_batch(
258            "BEGIN IMMEDIATE;
259             CREATE TABLE IF NOT EXISTS session_mount_access (
260                 session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
261                 source BLOB NOT NULL,
262                 destination BLOB NOT NULL,
263                 access TEXT NOT NULL CHECK(access IN ('rw')),
264                 PRIMARY KEY(session_id, destination)
265             ) STRICT;
266             INSERT INTO schema_migrations(version, applied_at)
267                 VALUES (34, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
268             PRAGMA user_version = 34;
269             COMMIT;",
270        )?;
271    }
272    // Compatible: adds one nullable column. Older readers ignore it, and the
273    // older writer's session upsert lists columns explicitly, so it preserves
274    // the value. An older executable launching such a session uses the shared
275    // `/workspace` instead of the recorded per-session path, which is a
276    // behaviour difference, not data loss. The compatibility floor stays where
277    // it is.
278    if version < 35 {
279        connection.execute_batch(
280            "BEGIN IMMEDIATE;
281             ALTER TABLE sessions ADD COLUMN container_workspace TEXT;
282             INSERT INTO schema_migrations(version, applied_at)
283                 VALUES (35, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
284             PRAGMA user_version = 35;
285             COMMIT;",
286        )?;
287    }
288    // Compatible: adds one nullable column. Older readers ignore it, and the
289    // older writer's session upsert lists columns explicitly, so it preserves
290    // the value. An older executable launching such a session runs it without
291    // the mbx build cache, which is a behaviour difference, not data loss. The
292    // compatibility floor stays where it is.
293    if version < 36 {
294        connection.execute_batch(
295            "BEGIN IMMEDIATE;
296             ALTER TABLE sessions ADD COLUMN build_cache_json TEXT;
297             INSERT INTO schema_migrations(version, applied_at)
298                 VALUES (36, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
299             PRAGMA user_version = 36;
300             COMMIT;",
301        )?;
302    }
303    // Compatible: adds one nullable column to `session_targets`. Only a
304    // container sub-agent child row ever carries a value, and older builds
305    // could never start such a child, so an older update that rewrites the row
306    // without the column loses nothing usable. Older readers ignore it. The
307    // compatibility floor stays where it is.
308    if version < 37 {
309        connection.execute_batch(
310            "BEGIN IMMEDIATE;
311             ALTER TABLE session_targets ADD COLUMN borrowed_from TEXT;
312             INSERT INTO schema_migrations(version, applied_at)
313                 VALUES (37, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
314             PRAGMA user_version = 37;
315             COMMIT;",
316        )?;
317    }
318    // Compatible: adds one table holding the dashboard's conversation pane
319    // arrangement per workspace. Older readers never select from it and older
320    // writers never touch it, so their updates leave its rows intact; a
321    // workspace deleted by an older build still removes them through the
322    // foreign key. Losing the table only means the conversation area opens as
323    // a single pane. The compatibility floor stays where it is.
324    if version < 38 {
325        connection.execute_batch(
326            "BEGIN IMMEDIATE;
327             CREATE TABLE IF NOT EXISTS workspace_layouts (
328                 workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
329                 layout TEXT NOT NULL
330             ) STRICT;
331             INSERT INTO schema_migrations(version, applied_at)
332                 VALUES (38, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
333             PRAGMA user_version = 38;
334             COMMIT;",
335        )?;
336    }
337    // Breaking: persisted input_required API events can now omit the structured
338    // request. Older readers require it and fail to deserialize the event log;
339    // the shared read/write compatibility floor must advance with the revision.
340    if version < 39 {
341        connection.execute_batch(
342            "BEGIN IMMEDIATE;
343             UPDATE schema_compatibility SET minimum_compatible_version = 39 WHERE singleton = 1;
344             INSERT INTO schema_migrations(version, applied_at)
345                 VALUES (39, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
346             PRAGMA user_version = 39;
347             COMMIT;",
348        )?;
349    }
350    // Breaking: native child projections and relay observations must be preserved
351    // by every reader/writer; older builds cannot interpret their lifecycle.
352    if version < 40 {
353        connection.execute_batch(
354            "BEGIN IMMEDIATE;
355             CREATE TABLE native_agents (
356                 owner TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
357                 child TEXT NOT NULL,
358                 staging INTEGER NOT NULL CHECK(staging IN (0,1)),
359                 body TEXT NOT NULL CHECK(json_valid(body)),
360                 PRIMARY KEY(owner, child, staging)
361             ) STRICT;
362             CREATE TABLE native_agent_transcript (
363                 owner TEXT NOT NULL,
364                 child TEXT NOT NULL,
365                 staging INTEGER NOT NULL,
366                 stable_id TEXT NOT NULL,
367                 position INTEGER NOT NULL,
368                 body TEXT NOT NULL CHECK(json_valid(body)),
369                 PRIMARY KEY(owner, child, staging, stable_id),
370                 FOREIGN KEY(owner, child, staging) REFERENCES native_agents(owner, child, staging)
371                     ON DELETE CASCADE ON UPDATE CASCADE
372             ) STRICT;
373             CREATE INDEX native_agent_transcript_position ON native_agent_transcript(owner, child, staging, position);
374             CREATE TABLE native_agent_replay (
375                 owner TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE
376             ) STRICT;
377             UPDATE schema_compatibility SET minimum_compatible_version = 40 WHERE singleton = 1;
378             INSERT INTO schema_migrations(version, applied_at)
379                 VALUES (40, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
380             PRAGMA user_version = 40;
381             COMMIT;",
382        )?;
383    }
384    // Breaking: clear-context relay commands/outcomes are persisted in event
385    // JSON. Older readers cannot decode them or honor the context boundary.
386    if version < 41 {
387        connection.execute_batch(
388            "BEGIN IMMEDIATE;
389             UPDATE schema_compatibility SET minimum_compatible_version = 41 WHERE singleton = 1;
390             INSERT INTO schema_migrations(version, applied_at)
391                 VALUES (41, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
392             PRAGMA user_version = 41;
393             COMMIT;",
394        )?;
395    }
396
397    // Breaking: older layout writers discard Browse identity and pin badges.
398    if version < 42 {
399        connection.execute_batch(
400            "BEGIN IMMEDIATE;
401             UPDATE schema_compatibility SET minimum_compatible_version = 42 WHERE singleton = 1;
402             INSERT INTO schema_migrations(version, applied_at)
403                 VALUES (42, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
404             PRAGMA user_version = 42;
405             COMMIT;",
406        )?;
407    }
408
409    // Breaking: durable steering commands/observations and native availability
410    // cannot be interpreted or preserved by older readers and writers.
411    if version < 43 {
412        connection.execute_batch("BEGIN IMMEDIATE;
413            UPDATE schema_compatibility SET minimum_compatible_version = 43 WHERE singleton = 1;
414            INSERT INTO schema_migrations(version, applied_at) VALUES (43, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
415            PRAGMA user_version = 43;
416            COMMIT;")?;
417    }
418
419    // Breaking: new quota recovery commands in stored relay JSON cannot be
420    // read or preserved by older binaries, even though the cache is additive.
421    if version < 44 {
422        connection.execute_batch("BEGIN IMMEDIATE;
423            CREATE TABLE quota_reset_cache (identity TEXT PRIMARY KEY, body TEXT NOT NULL);
424            UPDATE schema_compatibility SET minimum_compatible_version = 44 WHERE singleton = 1;
425            INSERT INTO schema_migrations(version, applied_at) VALUES (44, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
426            PRAGMA user_version = 44;
427            COMMIT;")?;
428    }
429
430    // Compatible: adds one nullable column. Older readers ignore it, and the
431    // older writer's session upsert lists columns explicitly, so it preserves
432    // the value. An older executable relaunching such a session starts it at
433    // HEAD or the remote default branch instead of the recorded revision,
434    // which is a behaviour difference, not data loss. The compatibility floor
435    // stays where it is.
436    if version < 45 {
437        // The column is added only when it is absent, the way migration 34
438        // creates its table only when absent: a store rolled back to an older
439        // revision still carries the column, and a second ALTER would refuse.
440        let add_column =
441            match super::legacy_schema::table_has_column(connection, "sessions", "launch_base")? {
442                true => "",
443                false => "ALTER TABLE sessions ADD COLUMN launch_base TEXT;",
444            };
445        connection.execute_batch(&format!(
446            "BEGIN IMMEDIATE;
447             {add_column}
448             INSERT INTO schema_migrations(version, applied_at)
449                 VALUES (45, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
450             PRAGMA user_version = 45;
451             COMMIT;"
452        ))?;
453    }
454
455    let recorded: Option<i64> =
456        connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
457            row.get(0)
458        })?;
459    if recorded != Some(SCHEMA_VERSION) {
460        bail!(
461            "Mjolnir database migration ledger {:?} does not match schema {}",
462            recorded,
463            SCHEMA_VERSION
464        );
465    }
466    Ok(())
467}
468
469/// Create an empty store at the baseline revision in one immediate transaction.
470/// The revision is read again under the write lock, so a second process that
471/// raced to create the same store finds it already created.
472fn create_baseline_schema(connection: &Connection) -> Result<()> {
473    connection.execute_batch("BEGIN IMMEDIATE;")?;
474    let created = (|| -> Result<()> {
475        let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
476        if version != 0 {
477            return Ok(());
478        }
479        connection.execute_batch(include_str!("baseline.sql"))?;
480        connection.execute(
481            "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
482            [BASELINE_MINIMUM_COMPATIBLE_VERSION],
483        )?;
484        connection.execute(
485            "INSERT INTO schema_migrations(version, applied_at)
486             VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
487            [BASELINE_SCHEMA_VERSION],
488        )?;
489        connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
490        Ok(())
491    })();
492    match created {
493        Ok(()) => connection
494            .execute_batch("COMMIT;")
495            .context("commit baseline database schema"),
496        Err(error) => {
497            if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
498                tracing::warn!(%rollback, "could not roll back a failed baseline schema");
499            }
500            Err(error.context("create baseline database schema"))
501        }
502    }
503}
504
505#[cfg(test)]
506pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
507    let connection = Connection::open(path).unwrap();
508    let transaction = connection.unchecked_transaction().unwrap();
509    transaction
510        .execute(
511            "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
512            [minimum_compatible],
513        )
514        .unwrap();
515    transaction
516        .execute(
517            "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
518            [revision],
519        )
520        .unwrap();
521    transaction
522        .pragma_update(None, "user_version", revision)
523        .unwrap();
524    transaction.commit().unwrap();
525    forget_verified_schema(path);
526}
527
528#[cfg(test)]
529mod reader_tests {
530    use super::*;
531
532    #[test]
533    fn every_historical_revision_upgrades_directly_and_preserves_user_data() {
534        for revision in 1..SCHEMA_VERSION {
535            let directory = tempfile::tempdir().unwrap();
536            let path = directory.path().join("mj.sqlite3");
537            let connection = Connection::open(&path).unwrap();
538            connection
539                .execute_batch(include_str!("legacy_v1.sql"))
540                .unwrap();
541            connection.execute_batch(
542                "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
543                 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
544                     target_template_id, state, updated_at, native_session_id)
545                     VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
546                         '2026-01-01T00:00:00Z', 'native-original');
547                 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
548                     VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
549            ).unwrap();
550            // Interrupt at each historical transaction boundary. Closing the
551            // connection rolls back the interrupted step, exactly as a killed
552            // updater would. Reopening must resume without manual cleanup.
553            connection
554                .execute_batch(&format!(
555                    "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
556                 WHEN NEW.version > {revision}
557                 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
558                ))
559                .unwrap();
560            if revision < BASELINE_SCHEMA_VERSION {
561                assert!(super::super::legacy_schema::migrate_to_baseline(&connection).is_err());
562            } else {
563                super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
564                assert!(migrate_schema(&connection).is_err());
565            }
566            drop(connection);
567            let connection = Connection::open(&path).unwrap();
568            let found: i64 = connection
569                .query_row("PRAGMA user_version", [], |row| row.get(0))
570                .unwrap();
571            assert_eq!(found, revision);
572            connection
573                .execute_batch("DROP TRIGGER stop_at_revision")
574                .unwrap();
575            drop(connection);
576
577            let writer =
578                open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
579            assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
580            assert_eq!(
581                writer
582                    .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
583                    .unwrap(),
584                "ok"
585            );
586            assert!(
587                !writer
588                    .prepare("PRAGMA foreign_key_check")
589                    .unwrap()
590                    .exists([])
591                    .unwrap()
592            );
593            drop(writer);
594            let reader = open_reader_strict(&path).unwrap();
595            let prompt: String = reader
596                .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
597                .unwrap();
598            assert_eq!(prompt, "Keep my prompt");
599            let state = load_state_from(&path).unwrap();
600            assert_eq!(state.sessions["old-session"].title, "Keep my work");
601            assert_eq!(
602                state.sessions["old-session"].native_session_id.as_deref(),
603                Some("native-original")
604            );
605            drop(reader);
606            // A fresh writer process must not apply any migration twice.
607            forget_verified_schema(&path);
608            drop(open_writer(&path).unwrap());
609        }
610    }
611
612    /// The oldest executable revision that can still read and write a store at
613    /// `SCHEMA_VERSION`. Migration 44 adds durable quota recovery commands.
614    const MINIMUM_COMPATIBLE_VERSION: i64 = 44;
615
616    /// Rewrites a store's recorded schema version the way another build's
617    /// migration ladder would, and forgets that this process verified it.
618    fn stamp_schema_version(path: &Path, version: i64) {
619        if version > SCHEMA_VERSION {
620            advance_test_schema(path, version, version);
621            return;
622        }
623        let connection = Connection::open(path).unwrap();
624        connection
625            .execute_batch(&format!("PRAGMA user_version = {version};"))
626            .unwrap();
627        connection
628            .execute(
629                "DELETE FROM schema_migrations WHERE version > ?1",
630                [version],
631            )
632            .unwrap();
633        if version == 30 {
634            connection
635                .execute(
636                    "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
637                    [],
638                )
639                .unwrap();
640        }
641        drop(connection);
642        forget_verified_schema(path);
643    }
644
645    #[test]
646    fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
647        let directory = tempfile::tempdir().unwrap();
648        let path = directory.path().join("mj.sqlite3");
649        let connection = open_writer(&path).unwrap();
650        connection
651            .execute_batch(
652                "BEGIN IMMEDIATE;
653             DROP TABLE quota_reset_cache;
654             DROP TABLE native_agent_transcript;
655             DROP TABLE native_agents;
656             DROP TABLE native_agent_replay;
657             DELETE FROM schema_migrations WHERE version >= 39;
658             UPDATE schema_compatibility SET minimum_compatible_version = 32;
659             PRAGMA user_version = 38;
660             COMMIT;",
661            )
662            .unwrap();
663        migrate_schema(&connection).unwrap();
664        let state = read_schema_state(&connection).unwrap();
665        assert_eq!(state.revision, SCHEMA_VERSION);
666        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
667        let event = ApiEventData::InputRequired {
668            request: None,
669            turn_id: Some(1),
670        };
671        #[derive(serde::Deserialize)]
672        struct LegacyInputEvent {
673            #[serde(rename = "request")]
674            _request: mj_core::elicitation::ElicitationRequest,
675        }
676        let encoded = serde_json::to_value(&event).unwrap();
677        assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
678    }
679
680    #[test]
681    fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
682        let directory = tempfile::tempdir().unwrap();
683        let path = directory.path().join("mj.sqlite3");
684        let connection = open_writer(&path).unwrap();
685        connection
686            .execute_batch(
687                "CREATE TABLE future_feature(value TEXT NOT NULL);
688                 INSERT INTO future_feature VALUES ('preserve me');",
689            )
690            .unwrap();
691        drop(connection);
692        advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
693
694        let reader = open_reader_strict(&path).unwrap();
695        assert_eq!(
696            reader
697                .query_row("SELECT value FROM future_feature", [], |row| row
698                    .get::<_, String>(0))
699                .unwrap(),
700            "preserve me"
701        );
702        assert!(reader.execute("DELETE FROM future_feature", []).is_err());
703        drop(reader);
704
705        // A repair would recreate this deliberately removed trigger. A future
706        // schema is authoritative even when it differs from our own repairs.
707        let raw = Connection::open(&path).unwrap();
708        raw.execute_batch("DROP TRIGGER api_session_error_updated;")
709            .unwrap();
710        drop(raw);
711        let writer = open_writer(&path).unwrap();
712        assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'api_session_error_updated')", [], |row| row.get::<_, bool>(0)).unwrap());
713        assert_eq!(
714            writer
715                .query_row("SELECT value FROM future_feature", [], |row| row
716                    .get::<_, String>(0))
717                .unwrap(),
718            "preserve me"
719        );
720        let state = read_schema_state(&writer).unwrap();
721        assert_eq!(state.revision, SCHEMA_VERSION + 1);
722        assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
723    }
724
725    #[test]
726    fn invalid_compatibility_metadata_refuses_readers_and_writers() {
727        for alteration in [
728            "DROP TABLE schema_compatibility",
729            "DELETE FROM schema_compatibility",
730            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
731            "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
732            "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
733            "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
734            "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
735            "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
736        ] {
737            for future in [false, true] {
738                let directory = tempfile::tempdir().unwrap();
739                let path = directory.path().join("mj.sqlite3");
740                drop(open_writer(&path).unwrap());
741                if future {
742                    advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
743                }
744                let raw = Connection::open(&path).unwrap();
745                raw.execute_batch(alteration).unwrap();
746                let before: i64 = raw
747                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
748                    .unwrap();
749                // Exercise the cached path as well as a fresh writer open.
750                for error in [
751                    open_reader_strict(&path).unwrap_err(),
752                    open_writer(&path).unwrap_err(),
753                ] {
754                    let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
755                    assert_eq!(
756                        mismatch.reason,
757                        StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
758                        "{alteration}"
759                    );
760                }
761                forget_verified_schema(&path);
762                assert!(open_writer(&path).is_err(), "{alteration}");
763                let after: i64 = raw
764                    .query_row("PRAGMA schema_version", [], |row| row.get(0))
765                    .unwrap();
766                assert_eq!(
767                    before, after,
768                    "a rejected open repaired schema: {alteration}"
769                );
770            }
771        }
772    }
773
774    #[test]
775    fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
776        let directory = tempfile::tempdir().unwrap();
777        let path = directory.path().join("mj.sqlite3");
778        let connection = Connection::open(&path).unwrap();
779        // A table the baseline also creates makes its batch fail part way.
780        connection
781            .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
782            .unwrap();
783
784        let error = migrate_schema(&connection).unwrap_err();
785
786        assert!(format!("{error:#}").contains("create baseline database schema"));
787        assert!(
788            connection.is_autocommit(),
789            "the failed baseline left a transaction open"
790        );
791        assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
792        let tables: i64 = connection
793            .query_row(
794                "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
795                [],
796                |row| row.get(0),
797            )
798            .unwrap();
799        assert_eq!(tables, 1, "only the conflicting table remains");
800
801        connection.execute_batch("DROP TABLE workspaces").unwrap();
802        drop(connection);
803        let writer = open_writer(&path).unwrap();
804        let state = read_schema_state(&writer).unwrap();
805        assert_eq!(state.revision, SCHEMA_VERSION);
806        assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
807    }
808
809    /// A store ahead of this build cannot be fixed by starting a daemon of
810    /// this build, so the reader must not say so. This is the message the
811    /// incident in #24 printed twice a second for an hour.
812    #[test]
813    fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
814        let directory = tempfile::tempdir().unwrap();
815        let path = directory.path().join("mj.sqlite3");
816        drop(open_writer(&path).unwrap());
817        stamp_schema_version(&path, SCHEMA_VERSION + 1);
818
819        let error = open_reader_strict(&path).unwrap_err();
820
821        let mismatch = error
822            .chain()
823            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
824            .expect("the reader reports the mismatch as a typed cause");
825        assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
826        assert_eq!(mismatch.supported, SCHEMA_VERSION);
827        let message = mismatch.to_string();
828        assert!(message.contains("upgrade Mjolnir"), "got {message}");
829        assert!(
830            !message.contains("start the Mjolnir daemon"),
831            "got {message}"
832        );
833    }
834
835    /// A store behind this build keeps the advice that works, verbatim, so
836    /// existing log greps and runbooks keep matching.
837    #[test]
838    fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
839        let directory = tempfile::tempdir().unwrap();
840        let path = directory.path().join("mj.sqlite3");
841        drop(open_writer(&path).unwrap());
842        let raw = Connection::open(&path).unwrap();
843        raw.execute_batch(&format!(
844            "UPDATE schema_compatibility SET minimum_compatible_version = {0};
845             DELETE FROM schema_migrations WHERE version > {0};
846             INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
847             PRAGMA user_version = {0};",
848            SCHEMA_VERSION - 1
849        ))
850        .unwrap();
851        drop(raw);
852
853        let error = open_reader_strict(&path).unwrap_err();
854
855        let mismatch = error
856            .chain()
857            .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
858            .expect("the reader reports the mismatch as a typed cause");
859        assert_eq!(
860            mismatch.to_string(),
861            format!(
862                "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
863                 start the Mjolnir daemon to migrate it",
864                SCHEMA_VERSION - 1
865            )
866        );
867    }
868
869    #[test]
870    fn strict_reader_rejects_mutation() {
871        let directory = tempfile::tempdir().unwrap();
872        let path = directory.path().join("mj.sqlite3");
873        drop(open_writer(&path).unwrap());
874
875        let reader = open_reader_strict(&path).unwrap();
876        let error = reader
877            .execute("CREATE TABLE forbidden(value TEXT)", [])
878            .unwrap_err();
879        assert!(
880            matches!(
881                error.sqlite_error_code(),
882                Some(rusqlite::ErrorCode::ReadOnly)
883            ),
884            "unexpected mutation error: {error}"
885        );
886    }
887}