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