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