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