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