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 if version < 74 {
1285 migrate_opencode_harness_kind(connection)?;
1286 }
1287
1288 if version < 75 {
1293 connection.execute_batch(
1294 "BEGIN IMMEDIATE;
1295 CREATE TABLE IF NOT EXISTS github_watch_cursors (
1296 owner TEXT NOT NULL,
1297 repo TEXT NOT NULL,
1298 items_watermark_at TEXT,
1299 items_watermark_id INTEGER,
1300 items_etag TEXT,
1301 comments_cursor TEXT,
1302 PRIMARY KEY(owner, repo)
1303 ) STRICT;
1304 CREATE TABLE IF NOT EXISTS github_watch_classifications (
1305 owner TEXT NOT NULL,
1306 repo TEXT NOT NULL,
1307 number INTEGER NOT NULL,
1308 session_id TEXT NOT NULL,
1309 kind TEXT NOT NULL CHECK(kind IN ('issue', 'pull_request')),
1310 title TEXT NOT NULL,
1311 url TEXT NOT NULL,
1312 created_at TEXT NOT NULL,
1313 interested INTEGER NOT NULL CHECK(interested IN (0, 1)),
1314 created INTEGER NOT NULL CHECK(created IN (0, 1)),
1315 PRIMARY KEY(owner, repo, number, session_id)
1316 ) STRICT;
1317 CREATE TABLE IF NOT EXISTS github_watch_items (
1318 owner TEXT NOT NULL,
1319 repo TEXT NOT NULL,
1320 number INTEGER NOT NULL,
1321 creator_session_id TEXT NOT NULL,
1322 kind TEXT NOT NULL CHECK(kind IN ('issue', 'pull_request')),
1323 title TEXT NOT NULL,
1324 url TEXT NOT NULL,
1325 created_at TEXT NOT NULL,
1326 pull_request_etag TEXT,
1327 pull_request_state TEXT CHECK(pull_request_state IS NULL OR pull_request_state IN ('open', 'closed')),
1328 PRIMARY KEY(owner, repo, number)
1329 ) STRICT;
1330 CREATE TABLE IF NOT EXISTS mailbox_outbox (
1331 event_key TEXT PRIMARY KEY,
1332 target_session_id TEXT NOT NULL,
1333 event_json TEXT NOT NULL,
1334 wake INTEGER NOT NULL CHECK(wake IN (0, 1)),
1335 unpark INTEGER NOT NULL DEFAULT 0 CHECK(unpark IN (0, 1)),
1336 created_at TEXT NOT NULL,
1337 accepted_at TEXT,
1338 accepted_command_id TEXT,
1339 accepted_ordinal INTEGER
1340 ) STRICT;
1341 CREATE INDEX IF NOT EXISTS mailbox_outbox_pending
1342 ON mailbox_outbox(accepted_at, created_at);
1343 INSERT INTO schema_migrations(version, applied_at)
1344 VALUES (75, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
1345 PRAGMA user_version = 75;
1346 COMMIT;",
1347 )?;
1348 }
1349
1350 if version < 76 {
1355 let add_column =
1356 if super::legacy_schema::table_has_column(connection, "mailbox_outbox", "failure")? {
1357 ""
1358 } else {
1359 "ALTER TABLE mailbox_outbox ADD COLUMN failure TEXT;"
1360 };
1361 connection.execute_batch(&format!(
1362 "BEGIN IMMEDIATE;
1363 {add_column}
1364 INSERT INTO schema_migrations(version, applied_at)
1365 VALUES (76, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
1366 PRAGMA user_version = 76;
1367 COMMIT;"
1368 ))?;
1369 }
1370
1371 if version < 77 {
1377 connection.execute_batch(
1378 "BEGIN IMMEDIATE;
1379 CREATE TABLE prompt_history_nullable_v77 (
1380 history_id INTEGER PRIMARY KEY,
1381 session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
1382 event_ordinal INTEGER CHECK(event_ordinal IS NULL OR event_ordinal >= 0),
1383 submitted_at TEXT NOT NULL,
1384 text TEXT NOT NULL CHECK(length(trim(text)) > 0),
1385 UNIQUE(session_id, event_ordinal)
1386 ) STRICT;
1387 INSERT INTO prompt_history_nullable_v77
1388 SELECT history_id, session_id, event_ordinal, submitted_at, text
1389 FROM prompt_history;
1390 DROP TABLE prompt_history;
1391 ALTER TABLE prompt_history_nullable_v77 RENAME TO prompt_history;
1392 CREATE INDEX prompt_history_session_recent
1393 ON prompt_history(session_id, history_id DESC);
1394 CREATE INDEX prompt_history_recent
1395 ON prompt_history(history_id DESC);
1396 INSERT INTO schema_migrations(version, applied_at)
1397 VALUES (77, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
1398 PRAGMA user_version = 77;
1399 COMMIT;",
1400 )?;
1401 }
1402
1403 let recorded: Option<i64> =
1404 connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
1405 row.get(0)
1406 })?;
1407 if recorded == Some(SCHEMA_VERSION) && version < SCHEMA_VERSION {
1408 tracing::info!(
1409 from_revision = version,
1410 to_revision = SCHEMA_VERSION,
1411 "database migrations applied"
1412 );
1413 }
1414 if recorded != Some(SCHEMA_VERSION) {
1415 bail!(
1416 "Mjolnir database migration ledger {:?} does not match schema {}",
1417 recorded,
1418 SCHEMA_VERSION
1419 );
1420 }
1421 Ok(())
1422}
1423
1424fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
1431 const BEFORE: &str = "'stopped','lost',";
1432 const AFTER: &str = "'stopped','parked','lost',";
1433 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1434 let migration = (|| -> Result<()> {
1435 let transaction = connection.unchecked_transaction()?;
1436 let sql: String = transaction.query_row(
1437 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1438 [],
1439 |row| row.get(0),
1440 )?;
1441 let (_, definition) = sql
1442 .split_once('(')
1443 .context("missing sessions table definition")?;
1444 if !definition.contains(AFTER) {
1445 ensure!(
1446 definition.matches(BEFORE).count() == 1,
1447 "unexpected sessions state constraint"
1448 );
1449 let definition = definition.replace(BEFORE, AFTER);
1450 let objects: Vec<String> = transaction
1451 .prepare(
1452 "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1453 AND type IN ('index','trigger') AND sql IS NOT NULL",
1454 )?
1455 .query_map([], |row| row.get(0))?
1456 .collect::<rusqlite::Result<_>>()?;
1457 transaction.execute_batch(&format!(
1458 "CREATE TABLE sessions_parked_v53 ({definition};
1459 INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
1460 DROP TABLE sessions;
1461 ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
1462 ))?;
1463 for object in objects {
1464 transaction.execute_batch(&object)?;
1465 }
1466 ensure!(
1467 !transaction
1468 .prepare("PRAGMA foreign_key_check")?
1469 .exists([])?,
1470 "foreign key violation in the parked-state migration"
1471 );
1472 }
1473 transaction.execute_batch(
1474 "UPDATE schema_compatibility SET minimum_compatible_version = 53
1475 WHERE singleton = 1;
1476 INSERT INTO schema_migrations(version, applied_at)
1477 VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1478 PRAGMA user_version = 53;",
1479 )?;
1480 transaction.commit()?;
1481 Ok(())
1482 })();
1483 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1484 migration.context("migrate the sessions state constraint for parked sub-agents")?;
1485 restored.context("restore foreign key enforcement after the parked-state migration")?;
1486 Ok(())
1487}
1488
1489fn migrate_opencode_harness_kind(connection: &Connection) -> Result<()> {
1498 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1499 let migration = (|| -> Result<()> {
1500 let transaction = connection.unchecked_transaction()?;
1501 for table in ["sessions", "hidden_native_sessions"] {
1502 let sql: String = transaction.query_row(
1503 "SELECT sql FROM sqlite_schema WHERE type='table' AND name=?1",
1504 [table],
1505 |row| row.get(0),
1506 )?;
1507 let (_, definition) = sql
1508 .split_once('(')
1509 .context("missing harness table definition")?;
1510 if definition.contains("'opencode')") {
1511 continue;
1512 }
1513 ensure!(
1514 definition.contains("'zcode')"),
1515 "unexpected {table} harness constraint"
1516 );
1517 let definition = definition.replace("'zcode')", "'zcode','opencode')");
1518 let objects: Vec<String> = transaction
1519 .prepare(
1520 "SELECT sql FROM sqlite_schema WHERE tbl_name=?1
1521 AND type IN ('index','trigger') AND sql IS NOT NULL",
1522 )?
1523 .query_map([table], |row| row.get(0))?
1524 .collect::<rusqlite::Result<_>>()?;
1525 transaction.execute_batch(&format!(
1526 "CREATE TABLE {table}_opencode_v74 ({definition};
1527 INSERT INTO {table}_opencode_v74 SELECT * FROM {table};
1528 DROP TABLE {table};
1529 ALTER TABLE {table}_opencode_v74 RENAME TO {table};"
1530 ))?;
1531 for object in objects {
1532 transaction.execute_batch(&object)?;
1533 }
1534 }
1535 ensure!(
1536 !transaction
1537 .prepare("PRAGMA foreign_key_check")?
1538 .exists([])?,
1539 "foreign key violation in OpenCode migration"
1540 );
1541 transaction.execute_batch(
1542 "UPDATE schema_compatibility SET minimum_compatible_version = 74
1543 WHERE singleton = 1;
1544 INSERT INTO schema_migrations(version, applied_at)
1545 VALUES (74, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1546 PRAGMA user_version = 74;",
1547 )?;
1548 transaction.commit()?;
1549 Ok(())
1550 })();
1551 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1552 migration.context("migrate OpenCode harness constraints")?;
1553 restored.context("restore foreign key enforcement after OpenCode migration")?;
1554 Ok(())
1555}
1556
1557fn migrate_startup_cleanup_state(connection: &Connection) -> Result<()> {
1559 const BEFORE: &str = "'stopped','parked','lost',";
1560 const AFTER: &str = "'stopped','parked','startup-cleanup','lost',";
1561 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1562 let migration = (|| -> Result<()> {
1563 let transaction = connection.unchecked_transaction()?;
1564 let sql: String = transaction.query_row(
1565 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1566 [],
1567 |row| row.get(0),
1568 )?;
1569 let (_, definition) = sql
1570 .split_once('(')
1571 .context("missing sessions table definition")?;
1572 if !definition.contains(AFTER) {
1573 ensure!(
1574 definition.matches(BEFORE).count() == 1,
1575 "unexpected sessions state constraint"
1576 );
1577 let definition = definition.replace(BEFORE, AFTER);
1578 let objects: Vec<String> = transaction
1579 .prepare(
1580 "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1581 AND type IN ('index','trigger') AND sql IS NOT NULL",
1582 )?
1583 .query_map([], |row| row.get(0))?
1584 .collect::<rusqlite::Result<_>>()?;
1585 transaction.execute_batch(&format!(
1586 "CREATE TABLE sessions_startup_cleanup_v68 ({definition};
1587 INSERT INTO sessions_startup_cleanup_v68 SELECT * FROM sessions;
1588 DROP TABLE sessions;
1589 ALTER TABLE sessions_startup_cleanup_v68 RENAME TO sessions;"
1590 ))?;
1591 for object in objects {
1592 transaction.execute_batch(&object)?;
1593 }
1594 ensure!(
1595 !transaction
1596 .prepare("PRAGMA foreign_key_check")?
1597 .exists([])?,
1598 "foreign key violation in the startup-cleanup migration"
1599 );
1600 }
1601 transaction.execute_batch(
1602 "UPDATE schema_compatibility SET minimum_compatible_version = 68
1603 WHERE singleton = 1;
1604 INSERT INTO schema_migrations(version, applied_at)
1605 VALUES (68, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1606 PRAGMA user_version = 68;",
1607 )?;
1608 transaction.commit()?;
1609 Ok(())
1610 })();
1611 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1612 migration.context("migrate the sessions state constraint for failed startup cleanup")?;
1613 restored.context("restore foreign key enforcement after the startup-cleanup migration")?;
1614 Ok(())
1615}
1616
1617fn create_baseline_schema(connection: &Connection) -> Result<()> {
1621 connection.execute_batch("BEGIN IMMEDIATE;")?;
1622 let created = (|| -> Result<()> {
1623 let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
1624 if version != 0 {
1625 return Ok(());
1626 }
1627 connection.execute_batch(include_str!("baseline.sql"))?;
1628 connection.execute(
1629 "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
1630 [BASELINE_MINIMUM_COMPATIBLE_VERSION],
1631 )?;
1632 connection.execute(
1633 "INSERT INTO schema_migrations(version, applied_at)
1634 VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
1635 [BASELINE_SCHEMA_VERSION],
1636 )?;
1637 connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
1638 Ok(())
1639 })();
1640 match created {
1641 Ok(()) => connection
1642 .execute_batch("COMMIT;")
1643 .context("commit baseline database schema"),
1644 Err(error) => {
1645 if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
1646 tracing::warn!(%rollback, "could not roll back a failed baseline schema");
1647 }
1648 Err(error.context("create baseline database schema"))
1649 }
1650 }
1651}
1652
1653#[cfg(test)]
1654pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
1655 let connection = Connection::open(path).unwrap();
1656 let transaction = connection.unchecked_transaction().unwrap();
1657 transaction
1658 .execute(
1659 "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
1660 [minimum_compatible],
1661 )
1662 .unwrap();
1663 transaction
1664 .execute(
1665 "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
1666 [revision],
1667 )
1668 .unwrap();
1669 transaction
1670 .pragma_update(None, "user_version", revision)
1671 .unwrap();
1672 transaction.commit().unwrap();
1673 forget_verified_schema(path);
1674}
1675
1676#[cfg(test)]
1677mod reader_tests {
1678 use super::*;
1679
1680 #[test]
1681 fn restart_intent_schema_matches_for_fresh_and_upgraded_stores() {
1682 let directory = tempfile::tempdir().unwrap();
1683 let fresh = Connection::open(directory.path().join("fresh.sqlite3")).unwrap();
1684 migrate_schema(&fresh).unwrap();
1685
1686 let upgraded = Connection::open(directory.path().join("upgraded.sqlite3")).unwrap();
1687 create_baseline_schema(&upgraded).unwrap();
1688 upgraded
1689 .execute_batch(
1690 "CREATE TRIGGER stop_before_restart_migration BEFORE INSERT ON schema_migrations
1691 WHEN NEW.version > 72
1692 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;",
1693 )
1694 .unwrap();
1695 assert!(migrate_schema(&upgraded).is_err());
1696 if !upgraded.is_autocommit() {
1697 upgraded.execute_batch("ROLLBACK").unwrap();
1698 }
1699 upgraded
1700 .execute_batch("DROP TRIGGER stop_before_restart_migration")
1701 .unwrap();
1702 assert_eq!(read_schema_state(&upgraded).unwrap().revision, 72);
1703 assert!(migrate_schema(&upgraded).is_ok());
1704
1705 let table_sql = |connection: &Connection| {
1706 connection
1707 .query_row(
1708 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='session_restart_intents'",
1709 [],
1710 |row| row.get::<_, String>(0),
1711 )
1712 .unwrap()
1713 };
1714 assert_eq!(table_sql(&fresh), table_sql(&upgraded));
1715 assert_eq!(read_schema_state(&fresh).unwrap().revision, SCHEMA_VERSION);
1716 assert_eq!(
1717 read_schema_state(&upgraded).unwrap().revision,
1718 SCHEMA_VERSION
1719 );
1720 }
1721
1722 fn assert_divergent_history_upgrades(revision: i64, project_history: bool, interrupt: bool) {
1723 let directory = tempfile::tempdir().unwrap();
1724 let path = directory.path().join("divergent-history.sqlite");
1725 let connection = Connection::open(&path).unwrap();
1726 create_baseline_schema(&connection).unwrap();
1727 connection.execute_batch(&format!(
1728 "INSERT INTO session_contexts(session_id,bundle_id,created_at)
1729 VALUES ('kept','project','now');
1730 INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at,project_directory)
1731 VALUES ('kept','Keep my work','codex','codex','local','error','now',X'2F7265706F');
1732 INSERT INTO materialized_sessions(session_id) VALUES ('kept');
1733 INSERT INTO session_turn_usage VALUES ('kept','turn',1,1,'{{\"tokens\":42}}');
1734 INSERT INTO session_provider_cost VALUES ('kept','{{\"amount\":1}}');
1735 CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1736 WHEN NEW.version > {revision}
1737 BEGIN SELECT RAISE(ABORT,'fixture migration boundary'); END;"
1738 )).unwrap();
1739 assert!(migrate_schema(&connection).is_err());
1740 if !connection.is_autocommit() {
1741 connection.execute_batch("ROLLBACK").unwrap();
1742 }
1743 connection
1744 .execute_batch("DROP TRIGGER stop_at_revision")
1745 .unwrap();
1746 assert_eq!(read_schema_state(&connection).unwrap().revision, revision);
1747 if project_history {
1748 connection
1749 .execute_batch(include_str!("project_catalog_v67.sql"))
1750 .unwrap();
1751 connection
1752 .execute_batch(
1753 "INSERT INTO project_catalog VALUES ('project','key','{}',0);
1754 INSERT INTO project_aliases VALUES ('alias','project','{}',1);
1755 INSERT INTO project_session_aliases VALUES ('kept','alias');
1756 UPDATE sessions SET project_json='{\"kept\":true}';",
1757 )
1758 .unwrap();
1759 }
1760 if interrupt {
1761 connection
1762 .execute_batch(
1763 "CREATE TRIGGER interrupt_reconciliation BEFORE INSERT ON schema_migrations
1764 WHEN NEW.version=69 BEGIN SELECT RAISE(ABORT,'interrupted reconciliation'); END;",
1765 )
1766 .unwrap();
1767 assert!(migrate_schema(&connection).is_err());
1768 drop(connection);
1769 let connection = Connection::open(&path).unwrap();
1770 assert_eq!(read_schema_state(&connection).unwrap().revision, 68);
1771 connection
1772 .execute_batch("DROP TRIGGER interrupt_reconciliation")
1773 .unwrap();
1774 } else {
1775 drop(connection);
1776 }
1777
1778 let writer = open_writer(&path).unwrap();
1779 let state = read_schema_state(&writer).unwrap();
1780 assert_eq!(state.revision, SCHEMA_VERSION);
1781 assert!(state.ensure_supported_by(68).is_err());
1782 assert!(state.ensure_supported_by(67).is_err());
1783 if project_history {
1784 let snapshot: String = writer
1785 .query_row(
1786 "SELECT project_json FROM sessions WHERE session_id='kept'",
1787 [],
1788 |row| row.get(0),
1789 )
1790 .unwrap();
1791 assert_eq!(snapshot, r#"{"kept":true}"#);
1792 let alias: (String, i64) = writer.query_row(
1793 "SELECT canonical_id,config_pending FROM project_aliases WHERE bundle_id='alias'", [],
1794 |row| Ok((row.get(0)?, row.get(1)?))
1795 ).unwrap();
1796 assert_eq!(alias, ("project".to_owned(), 1));
1797 }
1798 writer.execute_batch(
1800 "UPDATE sessions SET state='startup-cleanup',project_directory=X'2F6E6577' WHERE session_id='kept';
1801 INSERT INTO session_turn_selections VALUES ('kept','turn','model','high');"
1802 ).unwrap();
1803 let changed: Vec<u8> = writer
1804 .query_row(
1805 "SELECT directory FROM project_discovery_changes ORDER BY sequence DESC LIMIT 1",
1806 [],
1807 |row| row.get(0),
1808 )
1809 .unwrap();
1810 assert_eq!(changed, b"/new");
1811 writer
1812 .execute("DELETE FROM sessions WHERE session_id='kept'", [])
1813 .unwrap();
1814 let usage: String = writer
1815 .query_row("SELECT body FROM session_turn_usage", [], |row| row.get(0))
1816 .unwrap();
1817 assert_eq!(usage, r#"{"tokens":42}"#);
1818 let cost: String = writer
1819 .query_row("SELECT body FROM session_provider_cost", [], |row| {
1820 row.get(0)
1821 })
1822 .unwrap();
1823 assert_eq!(cost, r#"{"amount":1}"#);
1824 assert!(
1825 !writer
1826 .prepare("PRAGMA foreign_key_check")
1827 .unwrap()
1828 .exists([])
1829 .unwrap()
1830 );
1831 drop(writer);
1832 forget_verified_schema(&path);
1833 assert_eq!(
1834 read_schema_state(&open_writer(&path).unwrap())
1835 .unwrap()
1836 .revision,
1837 SCHEMA_VERSION
1838 );
1839 }
1840
1841 #[test]
1842 fn divergent_accounting_revision_67_preserves_usage_and_adds_projects() {
1843 assert_divergent_history_upgrades(67, false, false);
1844 }
1845
1846 #[test]
1847 fn divergent_cleanup_revision_68_preserves_usage_and_adds_projects() {
1848 assert_divergent_history_upgrades(68, false, false);
1849 }
1850
1851 #[test]
1852 fn divergent_project_revision_67_preserves_aliases_snapshots_and_usage() {
1853 assert_divergent_history_upgrades(66, true, false);
1854 }
1855
1856 #[test]
1857 fn divergent_project_reconciliation_resumes_after_interruption() {
1858 assert_divergent_history_upgrades(66, true, true);
1859 }
1860
1861 #[test]
1862 fn move_ownership_upgrade_retains_sources_and_refuses_previous_daemons() {
1863 let directory = tempfile::tempdir().unwrap();
1864 let path = directory.path().join("move.sqlite");
1865 let connection = open_writer(&path).unwrap();
1866 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; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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();
1867 drop(connection);
1868 forget_verified_schema(&path);
1869 let upgraded = open_writer(&path).unwrap();
1870 let schema = read_schema_state(&upgraded).unwrap();
1871 assert!(schema.ensure_supported_by(63).is_err());
1872 upgraded.execute("INSERT INTO retained_move_sources VALUES ('move-one','session-one','{}','[]','now')", []).unwrap();
1873 drop(upgraded);
1874 let reopened = open_writer(&path).unwrap();
1875 let count: i64 = reopened
1876 .query_row("SELECT count(*) FROM retained_move_sources", [], |row| {
1877 row.get(0)
1878 })
1879 .unwrap();
1880 assert_eq!(count, 1);
1881 }
1882
1883 #[test]
1884 fn ec2_move_ownership_upgrade_is_atomic_and_preserves_existing_moves() {
1885 let directory = tempfile::tempdir().unwrap();
1886 let path = directory.path().join("ec2-move-migration.sqlite3");
1887 save_session_to(&path, &super::super::tests::session("source", "project")).unwrap();
1888 stamp_schema_version(&path, 70);
1889 let old = Connection::open(&path).unwrap();
1890 old.execute_batch(
1891 "UPDATE schema_compatibility SET minimum_compatible_version=69;
1892 INSERT INTO session_moves(session_id,operation_id,operation_json)
1893 VALUES('source','existing-move','{\"operation_id\":\"existing-move\"}');
1894 CREATE TRIGGER stop_ec2_move_migration BEFORE INSERT ON schema_migrations
1895 WHEN NEW.version=71 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1896 )
1897 .unwrap();
1898 assert!(migrate_schema(&old).is_err());
1899 let state = read_schema_state(&old).unwrap();
1900 assert_eq!(state.revision, 70);
1901 assert_eq!(state.minimum_compatible, Some(69));
1902 old.execute_batch("DROP TRIGGER stop_ec2_move_migration")
1903 .unwrap();
1904 drop(old);
1905 let upgraded = open_writer(&path).unwrap();
1906 let state = read_schema_state(&upgraded).unwrap();
1907 assert_eq!(state.revision, SCHEMA_VERSION);
1908 assert_eq!(state.minimum_compatible, Some(74));
1909 assert!(state.ensure_supported_by(70).is_err());
1910 let preserved: String = upgraded
1911 .query_row(
1912 "SELECT operation_json FROM session_moves WHERE session_id='source'",
1913 [],
1914 |row| row.get(0),
1915 )
1916 .unwrap();
1917 assert_eq!(preserved, r#"{"operation_id":"existing-move"}"#);
1918 assert_eq!(
1919 upgraded
1920 .query_row(
1921 "SELECT count(*) FROM sessions WHERE session_id='source'",
1922 [],
1923 |row| row.get::<_, i64>(0)
1924 )
1925 .unwrap(),
1926 1
1927 );
1928 }
1929
1930 #[test]
1931 fn title_migration_caps_old_titles_preserves_short_titles_and_is_compatible() {
1932 let directory = tempfile::tempdir().unwrap();
1933 let path = directory.path().join("title-migration.sqlite3");
1934 let titles = [
1935 ("long", Some("word ".repeat(20_000))),
1936 ("unicode", Some("界".repeat(257))),
1937 ("exact", Some("界".repeat(256))),
1938 ("short", Some(" Keep\nthis title ".into())),
1939 ("unset", None),
1940 ];
1941 for (id, _) in &titles {
1942 let mut session = super::super::tests::session(id, "project");
1943 session.state = SessionState::Stopped;
1944 save_session_to(&path, &session).unwrap();
1945 }
1946 stamp_schema_version(&path, 69);
1947 let old = Connection::open(&path).unwrap();
1948 old.execute_batch(
1949 "UPDATE schema_compatibility SET minimum_compatible_version=69;
1950 CREATE TRIGGER stop_after_title_migration BEFORE INSERT ON schema_migrations
1951 WHEN NEW.version=71 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1952 )
1953 .unwrap();
1954 for (id, title) in &titles {
1955 old.execute(
1956 "UPDATE sessions SET acp_session_title=?2 WHERE session_id=?1",
1957 params![id, title],
1958 )
1959 .unwrap();
1960 }
1961 assert!(migrate_schema(&old).is_err());
1962 let upgraded = old;
1963 let state = read_schema_state(&upgraded).unwrap();
1964 assert_eq!(state.revision, 70);
1965 assert_eq!(state.minimum_compatible, Some(69));
1966 state.ensure_supported_by(69).unwrap();
1967 for (id, original) in &titles {
1968 let stored: Option<String> = upgraded
1969 .query_row(
1970 "SELECT acp_session_title FROM sessions WHERE session_id=?1",
1971 [id],
1972 |row| row.get(0),
1973 )
1974 .unwrap();
1975 let expected = match *id {
1976 "long" => Some(format!("{}word…", "word ".repeat(50))),
1977 "unicode" => Some(format!("{}…", "界".repeat(255))),
1978 _ => original.clone(),
1979 };
1980 assert_eq!(stored, expected, "session {id}");
1981 assert!(stored.is_none_or(|title| title.chars().count() <= 256));
1982 }
1983 upgraded
1984 .execute_batch("DROP TRIGGER stop_after_title_migration")
1985 .unwrap();
1986 drop(upgraded);
1987 forget_verified_schema(&path);
1988 drop(open_writer(&path).unwrap());
1989 }
1990
1991 #[test]
1992 fn interrupted_title_migration_rolls_back_titles_and_revision() {
1993 let directory = tempfile::tempdir().unwrap();
1994 let path = directory.path().join("title-migration-interrupted.sqlite3");
1995 save_session_to(&path, &super::super::tests::session("old", "project")).unwrap();
1996 stamp_schema_version(&path, 69);
1997 let original = "word ".repeat(20_000);
1998 let old = Connection::open(&path).unwrap();
1999 old.execute("UPDATE sessions SET acp_session_title=?1", [&original])
2000 .unwrap();
2001 old.execute_batch(
2002 "CREATE TRIGGER stop_title_migration BEFORE INSERT ON schema_migrations
2003 WHEN NEW.version=70 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
2004 )
2005 .unwrap();
2006 assert!(migrate_schema(&old).is_err());
2007 assert_eq!(read_schema_state(&old).unwrap().revision, 69);
2008 let stored: String = old
2009 .query_row("SELECT acp_session_title FROM sessions", [], |row| {
2010 row.get(0)
2011 })
2012 .unwrap();
2013 assert_eq!(stored, original);
2014 old.execute_batch("DROP TRIGGER stop_title_migration")
2015 .unwrap();
2016 drop(old);
2017 drop(open_writer(&path).unwrap());
2018 assert!(
2019 load_state_from(&path).unwrap().sessions["old"]
2020 .acp_session_title
2021 .as_ref()
2022 .unwrap()
2023 .chars()
2024 .count()
2025 <= 256
2026 );
2027 }
2028
2029 #[test]
2030 fn recent_revisions_upgrade_directly_and_preserve_user_data() {
2031 const TEST_REVISION_FLOOR: i64 = 47;
2034 for revision in TEST_REVISION_FLOOR..SCHEMA_VERSION {
2035 let directory = tempfile::tempdir().unwrap();
2036 let path = directory.path().join("mj.sqlite3");
2037 let connection = Connection::open(&path).unwrap();
2038 connection
2039 .execute_batch(include_str!("legacy_v1.sql"))
2040 .unwrap();
2041 connection.execute_batch(
2042 "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
2043 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
2044 target_template_id, state, updated_at, native_session_id)
2045 VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
2046 '2026-01-01T00:00:00Z', 'native-original');
2047 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
2048 VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
2049 ).unwrap();
2050 connection
2054 .execute_batch(&format!(
2055 "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
2056 WHEN NEW.version > {revision}
2057 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
2058 ))
2059 .unwrap();
2060 super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
2061 assert!(migrate_schema(&connection).is_err());
2062 drop(connection);
2063 let connection = Connection::open(&path).unwrap();
2064 let found: i64 = connection
2065 .query_row("PRAGMA user_version", [], |row| row.get(0))
2066 .unwrap();
2067 assert_eq!(found, revision);
2068 connection
2069 .execute_batch("DROP TRIGGER stop_at_revision")
2070 .unwrap();
2071 drop(connection);
2072
2073 let writer =
2074 open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
2075 assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
2076 assert_eq!(
2077 writer
2078 .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
2079 .unwrap(),
2080 "ok"
2081 );
2082 assert!(
2083 !writer
2084 .prepare("PRAGMA foreign_key_check")
2085 .unwrap()
2086 .exists([])
2087 .unwrap()
2088 );
2089 drop(writer);
2090 let reader = open_reader_strict(&path).unwrap();
2091 let prompt: String = reader
2092 .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
2093 .unwrap();
2094 assert_eq!(prompt, "Keep my prompt");
2095 let state = load_state_from(&path).unwrap();
2096 assert_eq!(state.sessions["old-session"].title, "Keep my work");
2097 assert_eq!(
2098 state.sessions["old-session"].native_session_id.as_deref(),
2099 Some("native-original")
2100 );
2101 drop(reader);
2102 forget_verified_schema(&path);
2104 drop(open_writer(&path).unwrap());
2105 }
2106 }
2107
2108 #[test]
2109 fn revision_77_makes_prompt_ordinals_nullable_and_preserves_history() {
2110 let directory = tempfile::tempdir().unwrap();
2111 let path = directory.path().join("mj.sqlite3");
2112 let connection = Connection::open(&path).unwrap();
2113 create_baseline_schema(&connection).unwrap();
2114 migrate_schema(&connection).unwrap();
2115 connection
2116 .execute(
2117 "INSERT INTO session_contexts(session_id, bundle_id, created_at)
2118 VALUES ('old-session', 'project', '2026-01-01T00:00:00Z')",
2119 [],
2120 )
2121 .unwrap();
2122 connection
2123 .execute(
2124 "INSERT INTO prompt_history(session_id, event_ordinal, submitted_at, text)
2125 VALUES ('old-session', 11, '2026-01-01T00:00:00Z', 'kept prompt')",
2126 [],
2127 )
2128 .unwrap();
2129 connection
2132 .execute_batch(
2133 "CREATE TABLE prompt_history_v76 (
2134 history_id INTEGER PRIMARY KEY,
2135 session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
2136 event_ordinal INTEGER NOT NULL CHECK(event_ordinal >= 0),
2137 submitted_at TEXT NOT NULL,
2138 text TEXT NOT NULL CHECK(length(trim(text)) > 0),
2139 UNIQUE(session_id, event_ordinal)
2140 ) STRICT;
2141 INSERT INTO prompt_history_v76
2142 SELECT history_id, session_id, event_ordinal, submitted_at, text
2143 FROM prompt_history;
2144 DROP TABLE prompt_history;
2145 ALTER TABLE prompt_history_v76 RENAME TO prompt_history;
2146 CREATE INDEX prompt_history_session_recent
2147 ON prompt_history(session_id, history_id DESC);
2148 CREATE INDEX prompt_history_recent ON prompt_history(history_id DESC);
2149 DELETE FROM schema_migrations WHERE version = 77;
2150 PRAGMA user_version = 76;",
2151 )
2152 .unwrap();
2153 drop(connection);
2154 forget_verified_schema(&path);
2155
2156 let migrated = open_writer(&path).unwrap();
2157 assert_eq!(read_schema_state(&migrated).unwrap().revision, 77);
2158 let preserved: (i64, i64, String) = migrated
2159 .query_row(
2160 "SELECT history_id, event_ordinal, text FROM prompt_history",
2161 [],
2162 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?)),
2163 )
2164 .unwrap();
2165 assert_eq!(preserved, (1, 11, "kept prompt".into()));
2166 let ordinal_not_null: i64 = migrated
2167 .query_row(
2168 "SELECT \"notnull\" FROM pragma_table_info('prompt_history')
2169 WHERE name = 'event_ordinal'",
2170 [],
2171 |row| row.get(0),
2172 )
2173 .unwrap();
2174 assert_eq!(ordinal_not_null, 0);
2175 let indexes: Vec<String> = migrated
2176 .prepare("PRAGMA index_list('prompt_history')")
2177 .unwrap()
2178 .query_map([], |row| row.get(1))
2179 .unwrap()
2180 .collect::<rusqlite::Result<_>>()
2181 .unwrap();
2182 assert!(indexes.contains(&"prompt_history_session_recent".into()));
2183 assert!(indexes.contains(&"prompt_history_recent".into()));
2184 }
2185
2186 #[test]
2187 fn accounting_migration_preserves_usage_and_changes_its_deletion_owner() {
2188 let dir = tempfile::tempdir().unwrap();
2189 let path = dir.path().join("migration.sqlite");
2190 let connection = Connection::open(&path).unwrap();
2191 connection
2192 .execute_batch(include_str!("legacy_v1.sql"))
2193 .unwrap();
2194 connection.execute_batch("INSERT INTO session_contexts VALUES ('old-session','project','2026-01-01T00:00:00Z');
2195 INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at)
2196 VALUES ('old-session','Retain usage','codex','codex','local','error','2026-01-01T00:00:00Z');
2197 CREATE TRIGGER stop_before_accounting BEFORE INSERT ON schema_migrations WHEN NEW.version=67
2198 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;").unwrap();
2199 super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
2200 assert!(migrate_schema(&connection).is_err());
2201 connection
2202 .execute_batch(
2203 "ROLLBACK; DROP TRIGGER stop_before_accounting;
2204 INSERT INTO session_turn_usage VALUES ('old-session','turn',7,1,'{}');
2205 INSERT INTO session_provider_cost VALUES ('old-session','{\"amount\":1}');",
2206 )
2207 .unwrap();
2208 drop(connection);
2209 let writer = open_writer(&path).unwrap();
2210 writer
2211 .execute("DELETE FROM sessions WHERE session_id='old-session'", [])
2212 .unwrap();
2213 for table in ["session_turn_usage", "session_provider_cost"] {
2214 assert_eq!(
2215 writer
2216 .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| row
2217 .get::<_, u64>(0))
2218 .unwrap(),
2219 1
2220 );
2221 }
2222 assert_eq!(
2223 writer
2224 .query_row("SELECT body FROM session_turn_usage", [], |row| row
2225 .get::<_, String>(0))
2226 .unwrap(),
2227 "{}"
2228 );
2229 assert!(
2230 !writer
2231 .prepare("PRAGMA foreign_key_check")
2232 .unwrap()
2233 .exists([])
2234 .unwrap()
2235 );
2236 assert_eq!(
2237 read_schema_state(&writer).unwrap().minimum_compatible,
2238 Some(MINIMUM_COMPATIBLE_VERSION)
2239 );
2240 }
2241
2242 const MINIMUM_COMPATIBLE_VERSION: i64 = 74;
2246
2247 fn stamp_schema_version(path: &Path, version: i64) {
2250 if version > SCHEMA_VERSION {
2251 advance_test_schema(path, version, version);
2252 return;
2253 }
2254 let connection = Connection::open(path).unwrap();
2255 if version < 67 {
2256 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();
2257 }
2258 if version < 76
2260 && legacy_schema::table_has_column(&connection, "mailbox_outbox", "failure").unwrap()
2261 {
2262 connection
2263 .execute_batch("ALTER TABLE mailbox_outbox DROP COLUMN failure;")
2264 .unwrap();
2265 }
2266 connection
2267 .execute_batch(&format!("PRAGMA user_version = {version};"))
2268 .unwrap();
2269 connection
2270 .execute(
2271 "DELETE FROM schema_migrations WHERE version > ?1",
2272 [version],
2273 )
2274 .unwrap();
2275 if version == 30 {
2276 connection
2277 .execute(
2278 "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
2279 [],
2280 )
2281 .unwrap();
2282 }
2283 if matches!(version, 69 | 70) {
2284 connection.execute(
2285 "UPDATE schema_compatibility SET minimum_compatible_version=69 WHERE singleton=1", [],
2286 ).unwrap();
2287 }
2288 drop(connection);
2289 forget_verified_schema(path);
2290 }
2291
2292 #[test]
2293 fn subagent_policy_migration_preserves_legacy_choices_and_refuses_old_writers() {
2294 use mj_core::subagent::SubagentPolicy;
2295 let directory = tempfile::tempdir().unwrap();
2296 let path = directory.path().join("subagent-policy.sqlite3");
2297 for id in ["all", "native", "unset"] {
2298 save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
2299 }
2300 let connection = Connection::open(&path).unwrap();
2301 connection
2302 .execute_batch(
2303 "UPDATE sessions SET mjolnir_subagents = 1 WHERE session_id = 'all';
2304 UPDATE sessions SET mjolnir_subagents = 0 WHERE session_id = 'native';
2305 ALTER TABLE sessions DROP COLUMN subagents;
2306 DROP TABLE subagent_preference;
2307 DELETE FROM schema_migrations WHERE version >= 56;
2308 UPDATE schema_compatibility SET minimum_compatible_version = 55;
2309 DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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;",
2310 )
2311 .unwrap();
2312 drop(connection);
2313 forget_verified_schema(&path);
2314 let upgraded = open_writer(&path).unwrap();
2315 assert!(
2316 read_schema_state(&upgraded)
2317 .unwrap()
2318 .ensure_supported_by(55)
2319 .is_err()
2320 );
2321 drop(upgraded);
2322 let state = load_state_from(&path).unwrap();
2323 assert_eq!(
2324 state.sessions["all"].subagents,
2325 Some(SubagentPolicy::AllModels)
2326 );
2327 assert_eq!(
2328 state.sessions["native"].subagents,
2329 Some(SubagentPolicy::Native)
2330 );
2331 assert_eq!(state.sessions["unset"].subagents, None);
2332 assert_eq!(state.last_subagent_policy, SubagentPolicy::Native);
2333 }
2334
2335 #[test]
2336 fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
2337 let directory = tempfile::tempdir().unwrap();
2338 let path = directory.path().join("steering-migration.sqlite3");
2339 let record = super::super::tests::session("steered-session", "project");
2340 save_session_to(&path, &record).unwrap();
2341 let connection = Connection::open(&path).unwrap();
2342 connection
2343 .execute_batch(
2344 "DELETE FROM schema_migrations WHERE version >= 48;
2345 UPDATE schema_compatibility SET minimum_compatible_version = 47;
2346 DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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;",
2347 )
2348 .unwrap();
2349 drop(connection);
2350 forget_verified_schema(&path);
2351 let upgraded = open_writer(&path).unwrap();
2352 let schema = read_schema_state(&upgraded).unwrap();
2353 assert_eq!(schema.revision, SCHEMA_VERSION);
2354 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2355 let error = schema.ensure_supported_by(47).unwrap_err();
2356 assert!(matches!(
2357 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
2358 StoreSchemaMismatchReason::Incompatible {
2359 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
2360 }
2361 ));
2362 drop(upgraded);
2363 assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
2364 }
2365
2366 #[test]
2367 fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
2368 let directory = tempfile::tempdir().unwrap();
2369 let path = directory.path().join("target-migration.sqlite3");
2370 let record = super::super::tests::session("preserved-session", "project");
2371 save_session_to(&path, &record).unwrap();
2372 let connection = Connection::open(&path).unwrap();
2373 connection
2374 .execute_batch(
2375 "ALTER TABLE sessions DROP COLUMN target_runtime_json;
2376 DELETE FROM schema_migrations WHERE version >= 46;
2377 UPDATE schema_compatibility SET minimum_compatible_version = 44;
2378 DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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;",
2379 )
2380 .unwrap();
2381 forget_verified_schema(&path);
2382 let upgraded = open_writer(&path).unwrap();
2383 let schema = read_schema_state(&upgraded).unwrap();
2384 assert_eq!(schema.revision, SCHEMA_VERSION);
2385 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2386 let error = schema.ensure_supported_by(45).unwrap_err();
2387 assert!(matches!(
2388 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
2389 StoreSchemaMismatchReason::Incompatible {
2390 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
2391 }
2392 ));
2393 drop(upgraded);
2394 let restored = load_state_from(&path).unwrap();
2395 assert_eq!(restored.sessions[&record.id], record);
2396 }
2397
2398 #[test]
2399 fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
2400 let directory = tempfile::tempdir().unwrap();
2401 let path = directory.path().join("mj.sqlite3");
2402 let connection = open_writer(&path).unwrap();
2403 connection
2404 .execute_batch(
2405 "BEGIN IMMEDIATE;
2406 DROP TABLE quota_reset_cache;
2407 DROP TABLE native_agent_transcript;
2408 DROP TABLE native_agents;
2409 DROP TABLE native_agent_replay;
2410 DELETE FROM schema_migrations WHERE version >= 39;
2411 UPDATE schema_compatibility SET minimum_compatible_version = 32;
2412 DROP TABLE IF EXISTS subagent_accounting; DROP TABLE IF EXISTS session_turn_selections; ALTER TABLE mailbox_outbox DROP COLUMN failure; 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;
2413 COMMIT;",
2414 )
2415 .unwrap();
2416 migrate_schema(&connection).unwrap();
2417 let state = read_schema_state(&connection).unwrap();
2418 assert_eq!(state.revision, SCHEMA_VERSION);
2419 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2420 let event = ApiEventData::InputRequired {
2421 request: None,
2422 turn_id: Some(1),
2423 };
2424 #[derive(serde::Deserialize)]
2425 struct LegacyInputEvent {
2426 #[serde(rename = "request")]
2427 _request: mj_core::elicitation::ElicitationRequest,
2428 }
2429 let encoded = serde_json::to_value(&event).unwrap();
2430 assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
2431 }
2432
2433 #[test]
2434 fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
2435 let directory = tempfile::tempdir().unwrap();
2436 let path = directory.path().join("mj.sqlite3");
2437 let connection = open_writer(&path).unwrap();
2438 connection
2439 .execute_batch(
2440 "CREATE TABLE future_feature(value TEXT NOT NULL);
2441 INSERT INTO future_feature VALUES ('preserve me');",
2442 )
2443 .unwrap();
2444 drop(connection);
2445 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
2446
2447 let reader = open_reader_strict(&path).unwrap();
2448 assert_eq!(
2449 reader
2450 .query_row("SELECT value FROM future_feature", [], |row| row
2451 .get::<_, String>(0))
2452 .unwrap(),
2453 "preserve me"
2454 );
2455 assert!(reader.execute("DELETE FROM future_feature", []).is_err());
2456 drop(reader);
2457
2458 let raw = Connection::open(&path).unwrap();
2461 raw.execute_batch("DROP TRIGGER session_contexts_workspace_update;")
2462 .unwrap();
2463 drop(raw);
2464 let writer = open_writer(&path).unwrap();
2465 assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'session_contexts_workspace_update')", [], |row| row.get::<_, bool>(0)).unwrap());
2466 assert_eq!(
2467 writer
2468 .query_row("SELECT value FROM future_feature", [], |row| row
2469 .get::<_, String>(0))
2470 .unwrap(),
2471 "preserve me"
2472 );
2473 let state = read_schema_state(&writer).unwrap();
2474 assert_eq!(state.revision, SCHEMA_VERSION + 1);
2475 assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
2476 }
2477
2478 #[test]
2479 fn invalid_compatibility_metadata_refuses_readers_and_writers() {
2480 for alteration in [
2481 "DROP TABLE schema_compatibility",
2482 "DELETE FROM schema_compatibility",
2483 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
2484 "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
2485 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
2486 "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
2487 "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
2488 "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
2489 ] {
2490 for future in [false, true] {
2491 let directory = tempfile::tempdir().unwrap();
2492 let path = directory.path().join("mj.sqlite3");
2493 drop(open_writer(&path).unwrap());
2494 if future {
2495 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
2496 }
2497 let raw = Connection::open(&path).unwrap();
2498 raw.execute_batch(alteration).unwrap();
2499 let before: i64 = raw
2500 .query_row("PRAGMA schema_version", [], |row| row.get(0))
2501 .unwrap();
2502 for error in [
2504 open_reader_strict(&path).unwrap_err(),
2505 open_writer(&path).unwrap_err(),
2506 ] {
2507 let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
2508 assert_eq!(
2509 mismatch.reason,
2510 StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
2511 "{alteration}"
2512 );
2513 }
2514 forget_verified_schema(&path);
2515 assert!(open_writer(&path).is_err(), "{alteration}");
2516 let after: i64 = raw
2517 .query_row("PRAGMA schema_version", [], |row| row.get(0))
2518 .unwrap();
2519 assert_eq!(
2520 before, after,
2521 "a rejected open repaired schema: {alteration}"
2522 );
2523 }
2524 }
2525 }
2526
2527 #[test]
2528 fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
2529 let directory = tempfile::tempdir().unwrap();
2530 let path = directory.path().join("mj.sqlite3");
2531 let connection = Connection::open(&path).unwrap();
2532 connection
2534 .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
2535 .unwrap();
2536
2537 let error = migrate_schema(&connection).unwrap_err();
2538
2539 assert!(format!("{error:#}").contains("create baseline database schema"));
2540 assert!(
2541 connection.is_autocommit(),
2542 "the failed baseline left a transaction open"
2543 );
2544 assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
2545 let tables: i64 = connection
2546 .query_row(
2547 "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
2548 [],
2549 |row| row.get(0),
2550 )
2551 .unwrap();
2552 assert_eq!(tables, 1, "only the conflicting table remains");
2553
2554 connection.execute_batch("DROP TABLE workspaces").unwrap();
2555 drop(connection);
2556 let writer = open_writer(&path).unwrap();
2557 let state = read_schema_state(&writer).unwrap();
2558 assert_eq!(state.revision, SCHEMA_VERSION);
2559 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
2560 }
2561
2562 #[test]
2566 fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
2567 let directory = tempfile::tempdir().unwrap();
2568 let path = directory.path().join("mj.sqlite3");
2569 drop(open_writer(&path).unwrap());
2570 stamp_schema_version(&path, SCHEMA_VERSION + 1);
2571
2572 let error = open_reader_strict(&path).unwrap_err();
2573
2574 let mismatch = error
2575 .chain()
2576 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2577 .expect("the reader reports the mismatch as a typed cause");
2578 assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
2579 assert_eq!(mismatch.supported, SCHEMA_VERSION);
2580 let message = mismatch.to_string();
2581 assert!(message.contains("upgrade Mjolnir"), "got {message}");
2582 assert!(
2583 !message.contains("start the Mjolnir daemon"),
2584 "got {message}"
2585 );
2586 }
2587
2588 #[test]
2591 fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
2592 let directory = tempfile::tempdir().unwrap();
2593 let path = directory.path().join("mj.sqlite3");
2594 drop(open_writer(&path).unwrap());
2595 let raw = Connection::open(&path).unwrap();
2596 raw.execute_batch(&format!(
2597 "UPDATE schema_compatibility SET minimum_compatible_version = {0};
2598 DELETE FROM schema_migrations WHERE version > {0};
2599 INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
2600 PRAGMA user_version = {0};",
2601 SCHEMA_VERSION - 1
2602 ))
2603 .unwrap();
2604 drop(raw);
2605
2606 let error = open_reader_strict(&path).unwrap_err();
2607
2608 let mismatch = error
2609 .chain()
2610 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2611 .expect("the reader reports the mismatch as a typed cause");
2612 assert_eq!(
2613 mismatch.to_string(),
2614 format!(
2615 "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
2616 start the Mjolnir daemon to migrate it",
2617 SCHEMA_VERSION - 1
2618 )
2619 );
2620 }
2621
2622 #[test]
2623 fn strict_reader_rejects_mutation() {
2624 let directory = tempfile::tempdir().unwrap();
2625 let path = directory.path().join("mj.sqlite3");
2626 drop(open_writer(&path).unwrap());
2627
2628 let reader = open_reader_strict(&path).unwrap();
2629 let error = reader
2630 .execute("CREATE TABLE forbidden(value TEXT)", [])
2631 .unwrap_err();
2632 assert!(
2633 matches!(
2634 error.sqlite_error_code(),
2635 Some(rusqlite::ErrorCode::ReadOnly)
2636 ),
2637 "unexpected mutation error: {error}"
2638 );
2639 }
2640
2641 fn opened_readers() -> usize {
2642 OPENED_READERS.with(std::cell::Cell::get)
2643 }
2644
2645 fn has_workspace(reader: &Connection, name: &str) -> bool {
2646 reader
2647 .query_row(
2648 "SELECT EXISTS(SELECT 1 FROM workspaces WHERE name = ?1)",
2649 [name],
2650 |row| row.get(0),
2651 )
2652 .unwrap()
2653 }
2654
2655 fn add_workspace(path: &Path, name: &str) {
2656 open_writer(path)
2657 .unwrap()
2658 .execute(
2659 "INSERT INTO workspaces(workspace_id, name, name_key, created_at, last_opened_at)
2660 VALUES (?1, ?1, ?1, 'now', 'now')",
2661 [name],
2662 )
2663 .unwrap();
2664 }
2665
2666 #[test]
2669 fn strict_readers_are_reused_and_see_later_commits() {
2670 let directory = tempfile::tempdir().unwrap();
2671 let path = directory.path().join("mj.sqlite3");
2672 drop(open_writer(&path).unwrap());
2673 let before = opened_readers();
2674 assert!(!has_workspace(&open_reader_strict(&path).unwrap(), "later"));
2675 add_workspace(&path, "later");
2676 for _ in 0..20 {
2677 assert!(has_workspace(&open_reader_strict(&path).unwrap(), "later"));
2678 }
2679 assert_eq!(
2680 opened_readers() - before,
2681 1,
2682 "one connection served every read"
2683 );
2684
2685 let first = open_reader_strict(&path).unwrap();
2687 let second = open_reader_strict(&path).unwrap();
2688 drop((first, second));
2689 assert_eq!(opened_readers() - before, 2);
2690 drop(open_reader_strict(&path).unwrap());
2691 drop(open_reader_strict(&path).unwrap());
2692 assert_eq!(opened_readers() - before, 2);
2693
2694 let reader = open_reader_strict(&path).unwrap();
2696 reader.execute_batch("BEGIN").unwrap();
2697 drop(reader);
2698 let (first, second) = (
2699 open_reader_strict(&path).unwrap(),
2700 open_reader_strict(&path).unwrap(),
2701 );
2702 assert!(first.is_autocommit() && second.is_autocommit());
2703 assert_eq!(opened_readers() - before, 3);
2704 }
2705
2706 #[test]
2709 fn idle_readers_follow_store_replacement_and_schema_changes() {
2710 let directory = tempfile::tempdir().unwrap();
2711 let path = directory.path().join("mj.sqlite3");
2712 drop(open_writer(&path).unwrap());
2713 add_workspace(&path, "old store");
2714 drop(open_reader_strict(&path).unwrap());
2715
2716 for suffix in ["", "-wal", "-shm"] {
2717 let file = PathBuf::from(format!("{}{suffix}", path.display()));
2718 if file.exists() {
2719 fs::remove_file(file).unwrap();
2720 }
2721 }
2722 forget_verified_schema(&path);
2723 drop(open_writer(&path).unwrap());
2724 add_workspace(&path, "new store");
2725 let reader = open_reader_strict(&path).unwrap();
2726 assert!(has_workspace(&reader, "new store"));
2727 assert!(!has_workspace(&reader, "old store"));
2728 drop(reader);
2729
2730 drop(open_reader_strict(&path).unwrap());
2731 stamp_schema_version(&path, SCHEMA_VERSION + 1);
2732 let error = open_reader_strict(&path).unwrap_err();
2733 assert!(
2734 error
2735 .chain()
2736 .any(|cause| cause.downcast_ref::<StoreSchemaMismatch>().is_some()),
2737 "a reused connection checks compatibility like a new one: {error:#}"
2738 );
2739 }
2740}