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 self.ensure_supported_by(SCHEMA_VERSION)
14 }
15
16 fn ensure_supported_by(&self, supported: i64) -> Result<()> {
17 let reason = if self.revision < supported {
18 StoreSchemaMismatchReason::NeedsMigration
19 } else if let Some(minimum_compatible) = self.minimum_compatible {
20 if minimum_compatible <= supported {
21 return Ok(());
22 }
23 StoreSchemaMismatchReason::Incompatible { minimum_compatible }
24 } else {
25 StoreSchemaMismatchReason::InvalidCompatibilityMetadata
26 };
27 Err(StoreSchemaMismatch {
28 found: self.revision,
29 supported,
30 reason,
31 }
32 .into())
33 }
34}
35
36pub(super) fn read_schema_state(connection: &Connection) -> Result<SchemaState> {
39 let snapshot = connection
40 .unchecked_transaction()
41 .context("start database compatibility snapshot")?;
42 let revision: i64 = snapshot
43 .query_row("PRAGMA user_version", [], |row| row.get(0))
44 .context("read database migration revision")?;
45 let minimum_compatible = if revision >= COMPATIBILITY_METADATA_VERSION {
46 let invalid = || StoreSchemaMismatch {
47 found: revision,
48 supported: SCHEMA_VERSION,
49 reason: StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
50 };
51 let (count, singleton, floor, recorded): (i64, Option<i64>, Option<i64>, Option<i64>) =
52 snapshot
53 .query_row(
54 "SELECT count(*), min(singleton), min(minimum_compatible_version),
55 (SELECT max(version) FROM schema_migrations)
56 FROM schema_compatibility",
57 [],
58 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
59 )
60 .map_err(|error| {
61 let structural = match &error {
65 rusqlite::Error::SqliteFailure(code, _) => {
66 code.code == rusqlite::ErrorCode::Unknown
67 }
68 _ => true,
69 };
70 let error = anyhow::Error::new(error);
71 if structural {
72 error.context(invalid())
73 } else {
74 error.context("read database compatibility metadata")
75 }
76 })?;
77 if count != 1
78 || singleton != Some(1)
79 || recorded != Some(revision)
80 || !floor
81 .is_some_and(|floor| (COMPATIBILITY_METADATA_VERSION..=revision).contains(&floor))
82 {
83 return Err(invalid().into());
84 }
85 floor
86 } else {
87 None
88 };
89 snapshot
90 .commit()
91 .context("finish database compatibility snapshot")?;
92 Ok(SchemaState {
93 revision,
94 minimum_compatible,
95 })
96}
97
98pub fn database_path() -> PathBuf {
99 data_dir().join("mj.sqlite3")
100}
101
102pub fn check_read_compatibility() -> Result<()> {
106 open_reader_strict(&database_path()).map(drop)
107}
108
109pub(super) fn open_writer(path: &Path) -> Result<Connection> {
115 let mut connection = open_writable(path)?;
116 connection.set_transaction_behavior(rusqlite::TransactionBehavior::Immediate);
117 Ok(connection)
118}
119
120fn open_writable(path: &Path) -> Result<Connection> {
121 if let Some(parent) = path.parent() {
122 fs::create_dir_all(parent)
123 .with_context(|| format!("create Mjolnir data directory {}", parent.display()))?;
124 }
125 let connection = Connection::open(path)
126 .with_context(|| format!("open Mjolnir database {}", path.display()))?;
127 connection.busy_timeout(Duration::from_secs(5))?;
128 connection.execute_batch(
129 "PRAGMA foreign_keys = ON;
130 PRAGMA journal_mode = WAL;
131 PRAGMA synchronous = FULL;",
132 )?;
133 verify_schema_once(path, &connection)?;
134 Ok(connection)
135}
136
137pub(super) fn open(path: &Path) -> Result<Connection> {
138 open_writer(path)
139}
140
141#[cfg(not(test))]
145pub(super) fn open_reader(path: &Path) -> Result<Connection> {
146 open_reader_strict(path)
147}
148
149#[cfg(test)]
150pub(super) fn open_reader(path: &Path) -> Result<Connection> {
151 open_writable(path)
157}
158
159#[cfg_attr(test, allow(dead_code))]
160fn open_reader_strict(path: &Path) -> Result<Connection> {
161 let connection = Connection::open_with_flags(
162 path,
163 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
164 )
165 .with_context(|| format!("open Mjolnir database read-only {}", path.display()))?;
166 connection.busy_timeout(Duration::from_secs(5))?;
167 connection.execute_batch(
168 "PRAGMA foreign_keys = ON;
169 PRAGMA query_only = ON;",
170 )?;
171 read_schema_state(&connection)?.ensure_supported()?;
172 Ok(connection)
173}
174
175fn verified_schemas() -> &'static Mutex<HashSet<PathBuf>> {
179 static VERIFIED: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
180 VERIFIED.get_or_init(|| Mutex::new(HashSet::new()))
181}
182
183fn schema_cache_key(path: &Path) -> PathBuf {
186 let Some(parent) = path
187 .parent()
188 .filter(|parent| !parent.as_os_str().is_empty())
189 else {
190 return path.to_owned();
191 };
192 match (fs::canonicalize(parent), path.file_name()) {
193 (Ok(canonical), Some(name)) => canonical.join(name),
194 _ => path.to_owned(),
195 }
196}
197
198fn verify_schema_once(path: &Path, connection: &Connection) -> Result<()> {
203 let key = schema_cache_key(path);
204 let mut verified = verified_schemas()
205 .lock()
206 .unwrap_or_else(PoisonError::into_inner);
207 let state = read_schema_state(connection)?;
208 if state.revision > SCHEMA_VERSION
209 || (state.revision == SCHEMA_VERSION && verified.contains(&key))
210 {
211 return state.ensure_supported();
213 }
214 migrate_schema(connection)?;
217 read_schema_state(connection)?.ensure_supported()?;
218 verified.insert(key);
219 Ok(())
220}
221
222#[cfg(test)]
226pub(super) fn forget_verified_schema(path: &Path) {
227 verified_schemas()
228 .lock()
229 .unwrap_or_else(PoisonError::into_inner)
230 .remove(&schema_cache_key(path));
231}
232
233const BASELINE_SCHEMA_VERSION: i64 = 33;
236
237const BASELINE_MINIMUM_COMPATIBLE_VERSION: i64 = 32;
240
241fn migrate_schema(connection: &Connection) -> Result<()> {
242 let state = read_schema_state(connection)?;
243 let version = state.revision;
244 if version > SCHEMA_VERSION {
245 return state.ensure_supported();
246 }
247 if version == 0 {
248 create_baseline_schema(connection)?;
249 } else if version < BASELINE_SCHEMA_VERSION {
250 super::legacy_schema::migrate_to_baseline(connection)
251 .context("upgrade historical database schema")?;
252 }
253 if version < 34 {
261 connection.execute_batch(
262 "BEGIN IMMEDIATE;
263 CREATE TABLE IF NOT EXISTS session_mount_access (
264 session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
265 source BLOB NOT NULL,
266 destination BLOB NOT NULL,
267 access TEXT NOT NULL CHECK(access IN ('rw')),
268 PRIMARY KEY(session_id, destination)
269 ) STRICT;
270 INSERT INTO schema_migrations(version, applied_at)
271 VALUES (34, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
272 PRAGMA user_version = 34;
273 COMMIT;",
274 )?;
275 }
276 if version < 35 {
283 connection.execute_batch(
284 "BEGIN IMMEDIATE;
285 ALTER TABLE sessions ADD COLUMN container_workspace TEXT;
286 INSERT INTO schema_migrations(version, applied_at)
287 VALUES (35, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
288 PRAGMA user_version = 35;
289 COMMIT;",
290 )?;
291 }
292 if version < 36 {
298 connection.execute_batch(
299 "BEGIN IMMEDIATE;
300 ALTER TABLE sessions ADD COLUMN build_cache_json TEXT;
301 INSERT INTO schema_migrations(version, applied_at)
302 VALUES (36, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
303 PRAGMA user_version = 36;
304 COMMIT;",
305 )?;
306 }
307 if version < 37 {
313 connection.execute_batch(
314 "BEGIN IMMEDIATE;
315 ALTER TABLE session_targets ADD COLUMN borrowed_from TEXT;
316 INSERT INTO schema_migrations(version, applied_at)
317 VALUES (37, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
318 PRAGMA user_version = 37;
319 COMMIT;",
320 )?;
321 }
322 if version < 38 {
329 connection.execute_batch(
330 "BEGIN IMMEDIATE;
331 CREATE TABLE IF NOT EXISTS workspace_layouts (
332 workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
333 layout TEXT NOT NULL
334 ) STRICT;
335 INSERT INTO schema_migrations(version, applied_at)
336 VALUES (38, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
337 PRAGMA user_version = 38;
338 COMMIT;",
339 )?;
340 }
341 if version < 39 {
345 connection.execute_batch(
346 "BEGIN IMMEDIATE;
347 UPDATE schema_compatibility SET minimum_compatible_version = 39 WHERE singleton = 1;
348 INSERT INTO schema_migrations(version, applied_at)
349 VALUES (39, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
350 PRAGMA user_version = 39;
351 COMMIT;",
352 )?;
353 }
354 if version < 40 {
357 connection.execute_batch(
358 "BEGIN IMMEDIATE;
359 CREATE TABLE native_agents (
360 owner TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
361 child TEXT NOT NULL,
362 staging INTEGER NOT NULL CHECK(staging IN (0,1)),
363 body TEXT NOT NULL CHECK(json_valid(body)),
364 PRIMARY KEY(owner, child, staging)
365 ) STRICT;
366 CREATE TABLE native_agent_transcript (
367 owner TEXT NOT NULL,
368 child TEXT NOT NULL,
369 staging INTEGER NOT NULL,
370 stable_id TEXT NOT NULL,
371 position INTEGER NOT NULL,
372 body TEXT NOT NULL CHECK(json_valid(body)),
373 PRIMARY KEY(owner, child, staging, stable_id),
374 FOREIGN KEY(owner, child, staging) REFERENCES native_agents(owner, child, staging)
375 ON DELETE CASCADE ON UPDATE CASCADE
376 ) STRICT;
377 CREATE INDEX native_agent_transcript_position ON native_agent_transcript(owner, child, staging, position);
378 CREATE TABLE native_agent_replay (
379 owner TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE
380 ) STRICT;
381 UPDATE schema_compatibility SET minimum_compatible_version = 40 WHERE singleton = 1;
382 INSERT INTO schema_migrations(version, applied_at)
383 VALUES (40, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
384 PRAGMA user_version = 40;
385 COMMIT;",
386 )?;
387 }
388 if version < 41 {
391 connection.execute_batch(
392 "BEGIN IMMEDIATE;
393 UPDATE schema_compatibility SET minimum_compatible_version = 41 WHERE singleton = 1;
394 INSERT INTO schema_migrations(version, applied_at)
395 VALUES (41, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
396 PRAGMA user_version = 41;
397 COMMIT;",
398 )?;
399 }
400
401 if version < 42 {
403 connection.execute_batch(
404 "BEGIN IMMEDIATE;
405 UPDATE schema_compatibility SET minimum_compatible_version = 42 WHERE singleton = 1;
406 INSERT INTO schema_migrations(version, applied_at)
407 VALUES (42, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
408 PRAGMA user_version = 42;
409 COMMIT;",
410 )?;
411 }
412
413 if version < 43 {
416 connection.execute_batch("BEGIN IMMEDIATE;
417 UPDATE schema_compatibility SET minimum_compatible_version = 43 WHERE singleton = 1;
418 INSERT INTO schema_migrations(version, applied_at) VALUES (43, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
419 PRAGMA user_version = 43;
420 COMMIT;")?;
421 }
422
423 if version < 44 {
426 connection.execute_batch("BEGIN IMMEDIATE;
427 CREATE TABLE quota_reset_cache (identity TEXT PRIMARY KEY, body TEXT NOT NULL);
428 UPDATE schema_compatibility SET minimum_compatible_version = 44 WHERE singleton = 1;
429 INSERT INTO schema_migrations(version, applied_at) VALUES (44, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
430 PRAGMA user_version = 44;
431 COMMIT;")?;
432 }
433
434 if version < 45 {
441 let add_column =
445 match super::legacy_schema::table_has_column(connection, "sessions", "launch_base")? {
446 true => "",
447 false => "ALTER TABLE sessions ADD COLUMN launch_base TEXT;",
448 };
449 connection.execute_batch(&format!(
450 "BEGIN IMMEDIATE;
451 {add_column}
452 INSERT INTO schema_migrations(version, applied_at)
453 VALUES (45, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
454 PRAGMA user_version = 45;
455 COMMIT;"
456 ))?;
457 }
458
459 if version < 46 {
463 let add_column = if super::legacy_schema::table_has_column(
464 connection,
465 "sessions",
466 "target_runtime_json",
467 )? {
468 ""
469 } else {
470 "ALTER TABLE sessions ADD COLUMN target_runtime_json TEXT;"
471 };
472 connection.execute_batch(&format!(
473 "BEGIN IMMEDIATE;
474 {add_column}
475 UPDATE schema_compatibility SET minimum_compatible_version = 46 WHERE singleton = 1;
476 INSERT INTO schema_migrations(version, applied_at)
477 VALUES (46, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
478 PRAGMA user_version = 46;
479 COMMIT;"
480 ))?;
481 }
482
483 if version < 47 {
487 let add_branch =
488 if super::legacy_schema::table_has_column(connection, "sessions", "launch_branch")? {
489 ""
490 } else {
491 "ALTER TABLE sessions ADD COLUMN launch_branch TEXT;"
492 };
493 let add_publication = if super::legacy_schema::table_has_column(
494 connection,
495 "sessions",
496 "publication_json",
497 )? {
498 ""
499 } else {
500 "ALTER TABLE sessions ADD COLUMN publication_json TEXT;"
501 };
502 connection.execute_batch(&format!(
503 "BEGIN IMMEDIATE;
504 {add_branch}
505 {add_publication}
506 UPDATE schema_compatibility SET minimum_compatible_version = 47 WHERE singleton = 1;
507 INSERT INTO schema_migrations(version, applied_at)
508 VALUES (47, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
509 PRAGMA user_version = 47;
510 COMMIT;",
511 ))?;
512 }
513
514 if version < 48 {
518 connection.execute_batch(
519 "BEGIN IMMEDIATE;
520 UPDATE schema_compatibility SET minimum_compatible_version = 48 WHERE singleton = 1;
521 INSERT INTO schema_migrations(version, applied_at)
522 VALUES (48, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
523 PRAGMA user_version = 48;
524 COMMIT;",
525 )?;
526 }
527
528 if version < 49 {
533 connection.execute_batch(
534 "BEGIN IMMEDIATE;
535 CREATE TABLE IF NOT EXISTS subagent_handbacks (
536 child_session_id TEXT PRIMARY KEY,
537 handback_command_id TEXT,
538 handback_message TEXT,
539 handback_recorded_at_ms INTEGER,
540 reminder_command_id TEXT,
541 reminder_for_command_id TEXT,
542 reminder_sent_at_ms INTEGER,
543 reminder_failed_for_command_id TEXT
544 );
545 UPDATE schema_compatibility SET minimum_compatible_version = 49 WHERE singleton = 1;
546 INSERT INTO schema_migrations(version, applied_at)
547 VALUES (49, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
548 PRAGMA user_version = 49;
549 COMMIT;",
550 )?;
551 }
552
553 if version < 50 {
557 let add_column = if super::legacy_schema::table_has_column(
558 connection,
559 "subagent_handbacks",
560 "awaited_ordinal",
561 )? {
562 ""
563 } else {
564 "ALTER TABLE subagent_handbacks ADD COLUMN awaited_ordinal INTEGER;"
565 };
566 connection.execute_batch(&format!(
567 "BEGIN IMMEDIATE;
568 {add_column}
569 INSERT INTO schema_migrations(version, applied_at)
570 VALUES (50, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
571 PRAGMA user_version = 50;
572 COMMIT;"
573 ))?;
574 }
575
576 if version < 51 {
580 let add_column = if super::legacy_schema::table_has_column(
581 connection,
582 "subagent_handbacks",
583 "report_dir",
584 )? {
585 ""
586 } else {
587 "ALTER TABLE subagent_handbacks ADD COLUMN report_dir TEXT;"
588 };
589 connection.execute_batch(&format!(
590 "BEGIN IMMEDIATE;
591 {add_column}
592 INSERT INTO schema_migrations(version, applied_at)
593 VALUES (51, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
594 PRAGMA user_version = 51;
595 COMMIT;"
596 ))?;
597 }
598
599 if version < 52 {
606 connection.execute_batch(
607 "BEGIN IMMEDIATE;
608 CREATE TABLE IF NOT EXISTS stopped_subagents (
609 parent_session_id TEXT NOT NULL
610 REFERENCES sessions(session_id) ON DELETE CASCADE,
611 child_session_id TEXT NOT NULL,
612 record_json TEXT NOT NULL CHECK(json_valid(record_json)),
613 PRIMARY KEY(parent_session_id, child_session_id)
614 ) STRICT;
615 INSERT INTO schema_migrations(version, applied_at)
616 VALUES (52, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
617 PRAGMA user_version = 52;
618 COMMIT;",
619 )?;
620 }
621
622 if version < 53 {
627 migrate_parked_session_state(connection)?;
628 }
629
630 if version < 54 {
633 let add_column =
634 if super::legacy_schema::table_has_column(connection, "sessions", "checkout_json")? {
635 ""
636 } else {
637 "ALTER TABLE sessions ADD COLUMN checkout_json TEXT;"
638 };
639 connection.execute_batch(&format!(
640 "BEGIN IMMEDIATE;
641 {add_column}
642 UPDATE schema_compatibility SET minimum_compatible_version = 54 WHERE singleton = 1;
643 INSERT INTO schema_migrations(version, applied_at)
644 VALUES (54, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
645 PRAGMA user_version = 54;
646 COMMIT;"
647 ))?;
648 }
649
650 if version < 55 {
653 let add_column = if super::legacy_schema::table_has_column(
654 connection,
655 "sessions",
656 "expected_runtime_identity",
657 )? {
658 ""
659 } else {
660 "ALTER TABLE sessions ADD COLUMN expected_runtime_identity TEXT;"
661 };
662 connection.execute_batch(&format!(
663 "BEGIN IMMEDIATE;
664 {add_column}
665 UPDATE schema_compatibility SET minimum_compatible_version = 55 WHERE singleton = 1;
666 INSERT INTO schema_migrations(version, applied_at)
667 VALUES (55, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
668 PRAGMA user_version = 55;
669 COMMIT;"
670 ))?;
671 }
672
673 let recorded: Option<i64> =
674 connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
675 row.get(0)
676 })?;
677 if recorded != Some(SCHEMA_VERSION) {
678 bail!(
679 "Mjolnir database migration ledger {:?} does not match schema {}",
680 recorded,
681 SCHEMA_VERSION
682 );
683 }
684 Ok(())
685}
686
687fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
694 const BEFORE: &str = "'stopped','lost',";
695 const AFTER: &str = "'stopped','parked','lost',";
696 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
697 let migration = (|| -> Result<()> {
698 let transaction = connection.unchecked_transaction()?;
699 let sql: String = transaction.query_row(
700 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
701 [],
702 |row| row.get(0),
703 )?;
704 let (_, definition) = sql
705 .split_once('(')
706 .context("missing sessions table definition")?;
707 if !definition.contains(AFTER) {
708 ensure!(
709 definition.matches(BEFORE).count() == 1,
710 "unexpected sessions state constraint"
711 );
712 let definition = definition.replace(BEFORE, AFTER);
713 let objects: Vec<String> = transaction
714 .prepare(
715 "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
716 AND type IN ('index','trigger') AND sql IS NOT NULL",
717 )?
718 .query_map([], |row| row.get(0))?
719 .collect::<rusqlite::Result<_>>()?;
720 transaction.execute_batch(&format!(
721 "CREATE TABLE sessions_parked_v53 ({definition};
722 INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
723 DROP TABLE sessions;
724 ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
725 ))?;
726 for object in objects {
727 transaction.execute_batch(&object)?;
728 }
729 ensure!(
730 !transaction
731 .prepare("PRAGMA foreign_key_check")?
732 .exists([])?,
733 "foreign key violation in the parked-state migration"
734 );
735 }
736 transaction.execute_batch(
737 "UPDATE schema_compatibility SET minimum_compatible_version = 53
738 WHERE singleton = 1;
739 INSERT INTO schema_migrations(version, applied_at)
740 VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
741 PRAGMA user_version = 53;",
742 )?;
743 transaction.commit()?;
744 Ok(())
745 })();
746 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
747 migration.context("migrate the sessions state constraint for parked sub-agents")?;
748 restored.context("restore foreign key enforcement after the parked-state migration")?;
749 Ok(())
750}
751
752fn create_baseline_schema(connection: &Connection) -> Result<()> {
756 connection.execute_batch("BEGIN IMMEDIATE;")?;
757 let created = (|| -> Result<()> {
758 let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
759 if version != 0 {
760 return Ok(());
761 }
762 connection.execute_batch(include_str!("baseline.sql"))?;
763 connection.execute(
764 "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
765 [BASELINE_MINIMUM_COMPATIBLE_VERSION],
766 )?;
767 connection.execute(
768 "INSERT INTO schema_migrations(version, applied_at)
769 VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
770 [BASELINE_SCHEMA_VERSION],
771 )?;
772 connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
773 Ok(())
774 })();
775 match created {
776 Ok(()) => connection
777 .execute_batch("COMMIT;")
778 .context("commit baseline database schema"),
779 Err(error) => {
780 if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
781 tracing::warn!(%rollback, "could not roll back a failed baseline schema");
782 }
783 Err(error.context("create baseline database schema"))
784 }
785 }
786}
787
788#[cfg(test)]
789pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
790 let connection = Connection::open(path).unwrap();
791 let transaction = connection.unchecked_transaction().unwrap();
792 transaction
793 .execute(
794 "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
795 [minimum_compatible],
796 )
797 .unwrap();
798 transaction
799 .execute(
800 "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
801 [revision],
802 )
803 .unwrap();
804 transaction
805 .pragma_update(None, "user_version", revision)
806 .unwrap();
807 transaction.commit().unwrap();
808 forget_verified_schema(path);
809}
810
811#[cfg(test)]
812mod reader_tests {
813 use super::*;
814
815 #[test]
816 fn every_historical_revision_upgrades_directly_and_preserves_user_data() {
817 for revision in 1..SCHEMA_VERSION {
818 let directory = tempfile::tempdir().unwrap();
819 let path = directory.path().join("mj.sqlite3");
820 let connection = Connection::open(&path).unwrap();
821 connection
822 .execute_batch(include_str!("legacy_v1.sql"))
823 .unwrap();
824 connection.execute_batch(
825 "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
826 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
827 target_template_id, state, updated_at, native_session_id)
828 VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
829 '2026-01-01T00:00:00Z', 'native-original');
830 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
831 VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
832 ).unwrap();
833 connection
837 .execute_batch(&format!(
838 "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
839 WHEN NEW.version > {revision}
840 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
841 ))
842 .unwrap();
843 if revision < BASELINE_SCHEMA_VERSION {
844 assert!(super::super::legacy_schema::migrate_to_baseline(&connection).is_err());
845 } else {
846 super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
847 assert!(migrate_schema(&connection).is_err());
848 }
849 drop(connection);
850 let connection = Connection::open(&path).unwrap();
851 let found: i64 = connection
852 .query_row("PRAGMA user_version", [], |row| row.get(0))
853 .unwrap();
854 assert_eq!(found, revision);
855 connection
856 .execute_batch("DROP TRIGGER stop_at_revision")
857 .unwrap();
858 drop(connection);
859
860 let writer =
861 open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
862 assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
863 assert_eq!(
864 writer
865 .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
866 .unwrap(),
867 "ok"
868 );
869 assert!(
870 !writer
871 .prepare("PRAGMA foreign_key_check")
872 .unwrap()
873 .exists([])
874 .unwrap()
875 );
876 drop(writer);
877 let reader = open_reader_strict(&path).unwrap();
878 let prompt: String = reader
879 .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
880 .unwrap();
881 assert_eq!(prompt, "Keep my prompt");
882 let state = load_state_from(&path).unwrap();
883 assert_eq!(state.sessions["old-session"].title, "Keep my work");
884 assert_eq!(
885 state.sessions["old-session"].native_session_id.as_deref(),
886 Some("native-original")
887 );
888 drop(reader);
889 forget_verified_schema(&path);
891 drop(open_writer(&path).unwrap());
892 }
893 }
894
895 const MINIMUM_COMPATIBLE_VERSION: i64 = 55;
898
899 fn stamp_schema_version(path: &Path, version: i64) {
902 if version > SCHEMA_VERSION {
903 advance_test_schema(path, version, version);
904 return;
905 }
906 let connection = Connection::open(path).unwrap();
907 connection
908 .execute_batch(&format!("PRAGMA user_version = {version};"))
909 .unwrap();
910 connection
911 .execute(
912 "DELETE FROM schema_migrations WHERE version > ?1",
913 [version],
914 )
915 .unwrap();
916 if version == 30 {
917 connection
918 .execute(
919 "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
920 [],
921 )
922 .unwrap();
923 }
924 drop(connection);
925 forget_verified_schema(path);
926 }
927
928 #[test]
929 fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
930 let directory = tempfile::tempdir().unwrap();
931 let path = directory.path().join("steering-migration.sqlite3");
932 let record = super::super::tests::session("steered-session", "project");
933 save_session_to(&path, &record).unwrap();
934 let connection = Connection::open(&path).unwrap();
935 connection
936 .execute_batch(
937 "DELETE FROM schema_migrations WHERE version >= 48;
938 UPDATE schema_compatibility SET minimum_compatible_version = 47;
939 PRAGMA user_version = 47;",
940 )
941 .unwrap();
942 drop(connection);
943 forget_verified_schema(&path);
944 let upgraded = open_writer(&path).unwrap();
945 let schema = read_schema_state(&upgraded).unwrap();
946 assert_eq!(schema.revision, SCHEMA_VERSION);
947 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
948 let error = schema.ensure_supported_by(47).unwrap_err();
949 assert!(matches!(
950 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
951 StoreSchemaMismatchReason::Incompatible {
952 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
953 }
954 ));
955 drop(upgraded);
956 assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
957 }
958
959 #[test]
960 fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
961 let directory = tempfile::tempdir().unwrap();
962 let path = directory.path().join("target-migration.sqlite3");
963 let record = super::super::tests::session("preserved-session", "project");
964 save_session_to(&path, &record).unwrap();
965 let connection = Connection::open(&path).unwrap();
966 connection
967 .execute_batch(
968 "ALTER TABLE sessions DROP COLUMN target_runtime_json;
969 DELETE FROM schema_migrations WHERE version >= 46;
970 UPDATE schema_compatibility SET minimum_compatible_version = 44;
971 PRAGMA user_version = 45;",
972 )
973 .unwrap();
974 forget_verified_schema(&path);
975 let upgraded = open_writer(&path).unwrap();
976 let schema = read_schema_state(&upgraded).unwrap();
977 assert_eq!(schema.revision, SCHEMA_VERSION);
978 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
979 let error = schema.ensure_supported_by(45).unwrap_err();
980 assert!(matches!(
981 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
982 StoreSchemaMismatchReason::Incompatible {
983 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
984 }
985 ));
986 drop(upgraded);
987 let restored = load_state_from(&path).unwrap();
988 assert_eq!(restored.sessions[&record.id], record);
989 }
990
991 #[test]
992 fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
993 let directory = tempfile::tempdir().unwrap();
994 let path = directory.path().join("mj.sqlite3");
995 let connection = open_writer(&path).unwrap();
996 connection
997 .execute_batch(
998 "BEGIN IMMEDIATE;
999 DROP TABLE quota_reset_cache;
1000 DROP TABLE native_agent_transcript;
1001 DROP TABLE native_agents;
1002 DROP TABLE native_agent_replay;
1003 DELETE FROM schema_migrations WHERE version >= 39;
1004 UPDATE schema_compatibility SET minimum_compatible_version = 32;
1005 PRAGMA user_version = 38;
1006 COMMIT;",
1007 )
1008 .unwrap();
1009 migrate_schema(&connection).unwrap();
1010 let state = read_schema_state(&connection).unwrap();
1011 assert_eq!(state.revision, SCHEMA_VERSION);
1012 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1013 let event = ApiEventData::InputRequired {
1014 request: None,
1015 turn_id: Some(1),
1016 };
1017 #[derive(serde::Deserialize)]
1018 struct LegacyInputEvent {
1019 #[serde(rename = "request")]
1020 _request: mj_core::elicitation::ElicitationRequest,
1021 }
1022 let encoded = serde_json::to_value(&event).unwrap();
1023 assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
1024 }
1025
1026 #[test]
1027 fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
1028 let directory = tempfile::tempdir().unwrap();
1029 let path = directory.path().join("mj.sqlite3");
1030 let connection = open_writer(&path).unwrap();
1031 connection
1032 .execute_batch(
1033 "CREATE TABLE future_feature(value TEXT NOT NULL);
1034 INSERT INTO future_feature VALUES ('preserve me');",
1035 )
1036 .unwrap();
1037 drop(connection);
1038 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1039
1040 let reader = open_reader_strict(&path).unwrap();
1041 assert_eq!(
1042 reader
1043 .query_row("SELECT value FROM future_feature", [], |row| row
1044 .get::<_, String>(0))
1045 .unwrap(),
1046 "preserve me"
1047 );
1048 assert!(reader.execute("DELETE FROM future_feature", []).is_err());
1049 drop(reader);
1050
1051 let raw = Connection::open(&path).unwrap();
1054 raw.execute_batch("DROP TRIGGER api_session_error_updated;")
1055 .unwrap();
1056 drop(raw);
1057 let writer = open_writer(&path).unwrap();
1058 assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'api_session_error_updated')", [], |row| row.get::<_, bool>(0)).unwrap());
1059 assert_eq!(
1060 writer
1061 .query_row("SELECT value FROM future_feature", [], |row| row
1062 .get::<_, String>(0))
1063 .unwrap(),
1064 "preserve me"
1065 );
1066 let state = read_schema_state(&writer).unwrap();
1067 assert_eq!(state.revision, SCHEMA_VERSION + 1);
1068 assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
1069 }
1070
1071 #[test]
1072 fn invalid_compatibility_metadata_refuses_readers_and_writers() {
1073 for alteration in [
1074 "DROP TABLE schema_compatibility",
1075 "DELETE FROM schema_compatibility",
1076 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
1077 "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
1078 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
1079 "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
1080 "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
1081 "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
1082 ] {
1083 for future in [false, true] {
1084 let directory = tempfile::tempdir().unwrap();
1085 let path = directory.path().join("mj.sqlite3");
1086 drop(open_writer(&path).unwrap());
1087 if future {
1088 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1089 }
1090 let raw = Connection::open(&path).unwrap();
1091 raw.execute_batch(alteration).unwrap();
1092 let before: i64 = raw
1093 .query_row("PRAGMA schema_version", [], |row| row.get(0))
1094 .unwrap();
1095 for error in [
1097 open_reader_strict(&path).unwrap_err(),
1098 open_writer(&path).unwrap_err(),
1099 ] {
1100 let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
1101 assert_eq!(
1102 mismatch.reason,
1103 StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
1104 "{alteration}"
1105 );
1106 }
1107 forget_verified_schema(&path);
1108 assert!(open_writer(&path).is_err(), "{alteration}");
1109 let after: i64 = raw
1110 .query_row("PRAGMA schema_version", [], |row| row.get(0))
1111 .unwrap();
1112 assert_eq!(
1113 before, after,
1114 "a rejected open repaired schema: {alteration}"
1115 );
1116 }
1117 }
1118 }
1119
1120 #[test]
1121 fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
1122 let directory = tempfile::tempdir().unwrap();
1123 let path = directory.path().join("mj.sqlite3");
1124 let connection = Connection::open(&path).unwrap();
1125 connection
1127 .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
1128 .unwrap();
1129
1130 let error = migrate_schema(&connection).unwrap_err();
1131
1132 assert!(format!("{error:#}").contains("create baseline database schema"));
1133 assert!(
1134 connection.is_autocommit(),
1135 "the failed baseline left a transaction open"
1136 );
1137 assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
1138 let tables: i64 = connection
1139 .query_row(
1140 "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
1141 [],
1142 |row| row.get(0),
1143 )
1144 .unwrap();
1145 assert_eq!(tables, 1, "only the conflicting table remains");
1146
1147 connection.execute_batch("DROP TABLE workspaces").unwrap();
1148 drop(connection);
1149 let writer = open_writer(&path).unwrap();
1150 let state = read_schema_state(&writer).unwrap();
1151 assert_eq!(state.revision, SCHEMA_VERSION);
1152 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1153 }
1154
1155 #[test]
1159 fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
1160 let directory = tempfile::tempdir().unwrap();
1161 let path = directory.path().join("mj.sqlite3");
1162 drop(open_writer(&path).unwrap());
1163 stamp_schema_version(&path, SCHEMA_VERSION + 1);
1164
1165 let error = open_reader_strict(&path).unwrap_err();
1166
1167 let mismatch = error
1168 .chain()
1169 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
1170 .expect("the reader reports the mismatch as a typed cause");
1171 assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
1172 assert_eq!(mismatch.supported, SCHEMA_VERSION);
1173 let message = mismatch.to_string();
1174 assert!(message.contains("upgrade Mjolnir"), "got {message}");
1175 assert!(
1176 !message.contains("start the Mjolnir daemon"),
1177 "got {message}"
1178 );
1179 }
1180
1181 #[test]
1184 fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
1185 let directory = tempfile::tempdir().unwrap();
1186 let path = directory.path().join("mj.sqlite3");
1187 drop(open_writer(&path).unwrap());
1188 let raw = Connection::open(&path).unwrap();
1189 raw.execute_batch(&format!(
1190 "UPDATE schema_compatibility SET minimum_compatible_version = {0};
1191 DELETE FROM schema_migrations WHERE version > {0};
1192 INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
1193 PRAGMA user_version = {0};",
1194 SCHEMA_VERSION - 1
1195 ))
1196 .unwrap();
1197 drop(raw);
1198
1199 let error = open_reader_strict(&path).unwrap_err();
1200
1201 let mismatch = error
1202 .chain()
1203 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
1204 .expect("the reader reports the mismatch as a typed cause");
1205 assert_eq!(
1206 mismatch.to_string(),
1207 format!(
1208 "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
1209 start the Mjolnir daemon to migrate it",
1210 SCHEMA_VERSION - 1
1211 )
1212 );
1213 }
1214
1215 #[test]
1216 fn strict_reader_rejects_mutation() {
1217 let directory = tempfile::tempdir().unwrap();
1218 let path = directory.path().join("mj.sqlite3");
1219 drop(open_writer(&path).unwrap());
1220
1221 let reader = open_reader_strict(&path).unwrap();
1222 let error = reader
1223 .execute("CREATE TABLE forbidden(value TEXT)", [])
1224 .unwrap_err();
1225 assert!(
1226 matches!(
1227 error.sqlite_error_code(),
1228 Some(rusqlite::ErrorCode::ReadOnly)
1229 ),
1230 "unexpected mutation error: {error}"
1231 );
1232 }
1233}