1use super::*;
2use rusqlite::OpenFlags;
3
4const COMPATIBILITY_METADATA_VERSION: i64 = 30;
5
6pub(super) struct SchemaState {
7 pub(super) revision: i64,
8 minimum_compatible: Option<i64>,
9}
10
11impl SchemaState {
12 pub(super) fn ensure_supported(&self) -> Result<()> {
13 self.ensure_supported_by(SCHEMA_VERSION)
14 }
15
16 fn ensure_supported_by(&self, supported: i64) -> Result<()> {
17 let reason = if self.revision < supported {
18 StoreSchemaMismatchReason::NeedsMigration
19 } else if let Some(minimum_compatible) = self.minimum_compatible {
20 if minimum_compatible <= supported {
21 return Ok(());
22 }
23 StoreSchemaMismatchReason::Incompatible { minimum_compatible }
24 } else {
25 StoreSchemaMismatchReason::InvalidCompatibilityMetadata
26 };
27 Err(StoreSchemaMismatch {
28 found: self.revision,
29 supported,
30 reason,
31 }
32 .into())
33 }
34}
35
36pub(super) fn read_schema_state(connection: &Connection) -> Result<SchemaState> {
39 let snapshot = connection
40 .unchecked_transaction()
41 .context("start database compatibility snapshot")?;
42 let revision: i64 = snapshot
43 .query_row("PRAGMA user_version", [], |row| row.get(0))
44 .context("read database migration revision")?;
45 let minimum_compatible = if revision >= COMPATIBILITY_METADATA_VERSION {
46 let invalid = || StoreSchemaMismatch {
47 found: revision,
48 supported: SCHEMA_VERSION,
49 reason: StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
50 };
51 let (count, singleton, floor, recorded): (i64, Option<i64>, Option<i64>, Option<i64>) =
52 snapshot
53 .query_row(
54 "SELECT count(*), min(singleton), min(minimum_compatible_version),
55 (SELECT max(version) FROM schema_migrations)
56 FROM schema_compatibility",
57 [],
58 |row| Ok((row.get(0)?, row.get(1)?, row.get(2)?, row.get(3)?)),
59 )
60 .map_err(|error| {
61 let structural = match &error {
65 rusqlite::Error::SqliteFailure(code, _) => {
66 code.code == rusqlite::ErrorCode::Unknown
67 }
68 _ => true,
69 };
70 let error = anyhow::Error::new(error);
71 if structural {
72 error.context(invalid())
73 } else {
74 error.context("read database compatibility metadata")
75 }
76 })?;
77 if count != 1
78 || singleton != Some(1)
79 || recorded != Some(revision)
80 || !floor
81 .is_some_and(|floor| (COMPATIBILITY_METADATA_VERSION..=revision).contains(&floor))
82 {
83 return Err(invalid().into());
84 }
85 floor
86 } else {
87 None
88 };
89 snapshot
90 .commit()
91 .context("finish database compatibility snapshot")?;
92 Ok(SchemaState {
93 revision,
94 minimum_compatible,
95 })
96}
97
98pub fn database_path() -> PathBuf {
99 data_dir().join("mj.sqlite3")
100}
101
102pub fn check_read_compatibility() -> Result<()> {
106 open_reader_strict(&database_path()).map(drop)
107}
108
109pub(super) fn open_writer(path: &Path) -> Result<Connection> {
115 let mut connection = open_writable(path)?;
116 connection.set_transaction_behavior(rusqlite::TransactionBehavior::Immediate);
117 Ok(connection)
118}
119
120fn open_writable(path: &Path) -> Result<Connection> {
121 if let Some(store) = path.parent() {
124 mj_core::config::ensure_may_control_store(store, "open this database for writing")?;
125 }
126 if let Some(parent) = path.parent() {
127 fs::create_dir_all(parent)
128 .with_context(|| format!("create Mjolnir data directory {}", parent.display()))?;
129 }
130 let connection = Connection::open(path)
131 .with_context(|| format!("open Mjolnir database {}", path.display()))?;
132 connection.busy_timeout(Duration::from_secs(5))?;
133 connection.execute_batch(
134 "PRAGMA foreign_keys = ON;
135 PRAGMA journal_mode = WAL;
136 PRAGMA synchronous = FULL;",
137 )?;
138 verify_schema_once(path, &connection)?;
139 committed::observe_connection(&connection, path)?;
140 Ok(connection)
141}
142
143pub(super) fn open(path: &Path) -> Result<Connection> {
144 open_writer(path)
145}
146
147#[cfg(not(test))]
151pub(super) fn open_reader(path: &Path) -> Result<Connection> {
152 open_reader_strict(path)
153}
154
155#[cfg(test)]
156pub(super) fn open_reader(path: &Path) -> Result<Connection> {
157 open_writable(path)
163}
164
165#[cfg_attr(test, allow(dead_code))]
166fn open_reader_strict(path: &Path) -> Result<Connection> {
167 let connection = Connection::open_with_flags(
168 path,
169 OpenFlags::SQLITE_OPEN_READ_ONLY | OpenFlags::SQLITE_OPEN_NO_MUTEX,
170 )
171 .with_context(|| format!("open Mjolnir database read-only {}", path.display()))?;
172 connection.busy_timeout(Duration::from_secs(5))?;
173 connection.execute_batch(
174 "PRAGMA foreign_keys = ON;
175 PRAGMA query_only = ON;",
176 )?;
177 read_schema_state(&connection)?.ensure_supported()?;
178 Ok(connection)
179}
180
181fn verified_schemas() -> &'static Mutex<HashSet<PathBuf>> {
185 static VERIFIED: OnceLock<Mutex<HashSet<PathBuf>>> = OnceLock::new();
186 VERIFIED.get_or_init(|| Mutex::new(HashSet::new()))
187}
188
189fn schema_cache_key(path: &Path) -> PathBuf {
192 let Some(parent) = path
193 .parent()
194 .filter(|parent| !parent.as_os_str().is_empty())
195 else {
196 return path.to_owned();
197 };
198 match (fs::canonicalize(parent), path.file_name()) {
199 (Ok(canonical), Some(name)) => canonical.join(name),
200 _ => path.to_owned(),
201 }
202}
203
204fn verify_schema_once(path: &Path, connection: &Connection) -> Result<()> {
209 let key = schema_cache_key(path);
210 let mut verified = verified_schemas()
211 .lock()
212 .unwrap_or_else(PoisonError::into_inner);
213 let state = read_schema_state(connection)?;
214 if state.revision > SCHEMA_VERSION
215 || (state.revision == SCHEMA_VERSION && verified.contains(&key))
216 {
217 return state.ensure_supported();
219 }
220 migrate_schema(connection)?;
223 read_schema_state(connection)?.ensure_supported()?;
224 verified.insert(key);
225 Ok(())
226}
227
228#[cfg(test)]
232pub(super) fn forget_verified_schema(path: &Path) {
233 verified_schemas()
234 .lock()
235 .unwrap_or_else(PoisonError::into_inner)
236 .remove(&schema_cache_key(path));
237}
238
239const BASELINE_SCHEMA_VERSION: i64 = 33;
242
243const BASELINE_MINIMUM_COMPATIBLE_VERSION: i64 = 32;
246
247const ACCOUNTING_MIGRATION_SQL: &str = "
250 ALTER TABLE session_turn_usage RENAME TO old_session_turn_usage;
251 CREATE TABLE session_turn_usage (
252 session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
253 command_id TEXT NOT NULL,
254 completed_ordinal INTEGER NOT NULL,
255 turn_start_position INTEGER,
256 body TEXT NOT NULL,
257 PRIMARY KEY(session_id, command_id)
258 );
259 INSERT INTO session_turn_usage SELECT * FROM old_session_turn_usage;
260 DROP TABLE old_session_turn_usage;
261 CREATE INDEX session_turn_usage_order ON session_turn_usage(session_id, completed_ordinal);
262 ALTER TABLE session_provider_cost RENAME TO old_session_provider_cost;
263 CREATE TABLE session_provider_cost (
264 session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
265 body TEXT NOT NULL
266 );
267 INSERT INTO session_provider_cost SELECT * FROM old_session_provider_cost;
268 DROP TABLE old_session_provider_cost;
269 CREATE TABLE subagent_accounting (
270 child_session_id TEXT PRIMARY KEY REFERENCES session_contexts(session_id),
271 parent_session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
272 task_name TEXT NOT NULL,
273 CHECK(child_session_id <> parent_session_id)
274 ) STRICT;
275 CREATE INDEX subagent_accounting_parent ON subagent_accounting(parent_session_id);
276 INSERT INTO subagent_accounting SELECT child_session_id, parent_session_id,
277 json_extract(record_json, '$.task_name') FROM subagent_sessions;
278 CREATE TABLE session_turn_selections (
279 session_id TEXT NOT NULL REFERENCES session_contexts(session_id),
280 command_id TEXT NOT NULL,
281 model TEXT,
282 effort TEXT,
283 PRIMARY KEY(session_id, command_id)
284 ) STRICT;
285";
286
287fn migrate_schema(connection: &Connection) -> Result<()> {
288 let state = read_schema_state(connection)?;
289 let version = state.revision;
290 if version > SCHEMA_VERSION {
291 return state.ensure_supported();
292 }
293 if version == 0 {
294 create_baseline_schema(connection)?;
295 } else if version < BASELINE_SCHEMA_VERSION {
296 super::legacy_schema::migrate_to_baseline(connection)
297 .context("upgrade historical database schema")?;
298 }
299 if version < 34 {
307 connection.execute_batch(
308 "BEGIN IMMEDIATE;
309 CREATE TABLE IF NOT EXISTS session_mount_access (
310 session_id TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
311 source BLOB NOT NULL,
312 destination BLOB NOT NULL,
313 access TEXT NOT NULL CHECK(access IN ('rw')),
314 PRIMARY KEY(session_id, destination)
315 ) STRICT;
316 INSERT INTO schema_migrations(version, applied_at)
317 VALUES (34, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
318 PRAGMA user_version = 34;
319 COMMIT;",
320 )?;
321 }
322 if version < 35 {
329 connection.execute_batch(
330 "BEGIN IMMEDIATE;
331 ALTER TABLE sessions ADD COLUMN container_workspace TEXT;
332 INSERT INTO schema_migrations(version, applied_at)
333 VALUES (35, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
334 PRAGMA user_version = 35;
335 COMMIT;",
336 )?;
337 }
338 if version < 36 {
344 connection.execute_batch(
345 "BEGIN IMMEDIATE;
346 ALTER TABLE sessions ADD COLUMN build_cache_json TEXT;
347 INSERT INTO schema_migrations(version, applied_at)
348 VALUES (36, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
349 PRAGMA user_version = 36;
350 COMMIT;",
351 )?;
352 }
353 if version < 37 {
359 connection.execute_batch(
360 "BEGIN IMMEDIATE;
361 ALTER TABLE session_targets ADD COLUMN borrowed_from TEXT;
362 INSERT INTO schema_migrations(version, applied_at)
363 VALUES (37, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
364 PRAGMA user_version = 37;
365 COMMIT;",
366 )?;
367 }
368 if version < 38 {
375 connection.execute_batch(
376 "BEGIN IMMEDIATE;
377 CREATE TABLE IF NOT EXISTS workspace_layouts (
378 workspace_id TEXT PRIMARY KEY REFERENCES workspaces(workspace_id) ON DELETE CASCADE,
379 layout TEXT NOT NULL
380 ) STRICT;
381 INSERT INTO schema_migrations(version, applied_at)
382 VALUES (38, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
383 PRAGMA user_version = 38;
384 COMMIT;",
385 )?;
386 }
387 if version < 39 {
391 connection.execute_batch(
392 "BEGIN IMMEDIATE;
393 UPDATE schema_compatibility SET minimum_compatible_version = 39 WHERE singleton = 1;
394 INSERT INTO schema_migrations(version, applied_at)
395 VALUES (39, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
396 PRAGMA user_version = 39;
397 COMMIT;",
398 )?;
399 }
400 if version < 40 {
403 connection.execute_batch(
404 "BEGIN IMMEDIATE;
405 CREATE TABLE native_agents (
406 owner TEXT NOT NULL REFERENCES sessions(session_id) ON DELETE CASCADE,
407 child TEXT NOT NULL,
408 staging INTEGER NOT NULL CHECK(staging IN (0,1)),
409 body TEXT NOT NULL CHECK(json_valid(body)),
410 PRIMARY KEY(owner, child, staging)
411 ) STRICT;
412 CREATE TABLE native_agent_transcript (
413 owner TEXT NOT NULL,
414 child TEXT NOT NULL,
415 staging INTEGER NOT NULL,
416 stable_id TEXT NOT NULL,
417 position INTEGER NOT NULL,
418 body TEXT NOT NULL CHECK(json_valid(body)),
419 PRIMARY KEY(owner, child, staging, stable_id),
420 FOREIGN KEY(owner, child, staging) REFERENCES native_agents(owner, child, staging)
421 ON DELETE CASCADE ON UPDATE CASCADE
422 ) STRICT;
423 CREATE INDEX native_agent_transcript_position ON native_agent_transcript(owner, child, staging, position);
424 CREATE TABLE native_agent_replay (
425 owner TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE
426 ) STRICT;
427 UPDATE schema_compatibility SET minimum_compatible_version = 40 WHERE singleton = 1;
428 INSERT INTO schema_migrations(version, applied_at)
429 VALUES (40, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
430 PRAGMA user_version = 40;
431 COMMIT;",
432 )?;
433 }
434 if version < 41 {
437 connection.execute_batch(
438 "BEGIN IMMEDIATE;
439 UPDATE schema_compatibility SET minimum_compatible_version = 41 WHERE singleton = 1;
440 INSERT INTO schema_migrations(version, applied_at)
441 VALUES (41, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
442 PRAGMA user_version = 41;
443 COMMIT;",
444 )?;
445 }
446
447 if version < 42 {
449 connection.execute_batch(
450 "BEGIN IMMEDIATE;
451 UPDATE schema_compatibility SET minimum_compatible_version = 42 WHERE singleton = 1;
452 INSERT INTO schema_migrations(version, applied_at)
453 VALUES (42, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
454 PRAGMA user_version = 42;
455 COMMIT;",
456 )?;
457 }
458
459 if version < 43 {
462 connection.execute_batch("BEGIN IMMEDIATE;
463 UPDATE schema_compatibility SET minimum_compatible_version = 43 WHERE singleton = 1;
464 INSERT INTO schema_migrations(version, applied_at) VALUES (43, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
465 PRAGMA user_version = 43;
466 COMMIT;")?;
467 }
468
469 if version < 44 {
472 connection.execute_batch("BEGIN IMMEDIATE;
473 CREATE TABLE quota_reset_cache (identity TEXT PRIMARY KEY, body TEXT NOT NULL);
474 UPDATE schema_compatibility SET minimum_compatible_version = 44 WHERE singleton = 1;
475 INSERT INTO schema_migrations(version, applied_at) VALUES (44, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
476 PRAGMA user_version = 44;
477 COMMIT;")?;
478 }
479
480 if version < 45 {
487 let add_column =
491 match super::legacy_schema::table_has_column(connection, "sessions", "launch_base")? {
492 true => "",
493 false => "ALTER TABLE sessions ADD COLUMN launch_base TEXT;",
494 };
495 connection.execute_batch(&format!(
496 "BEGIN IMMEDIATE;
497 {add_column}
498 INSERT INTO schema_migrations(version, applied_at)
499 VALUES (45, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
500 PRAGMA user_version = 45;
501 COMMIT;"
502 ))?;
503 }
504
505 if version < 46 {
509 let add_column = if super::legacy_schema::table_has_column(
510 connection,
511 "sessions",
512 "target_runtime_json",
513 )? {
514 ""
515 } else {
516 "ALTER TABLE sessions ADD COLUMN target_runtime_json TEXT;"
517 };
518 connection.execute_batch(&format!(
519 "BEGIN IMMEDIATE;
520 {add_column}
521 UPDATE schema_compatibility SET minimum_compatible_version = 46 WHERE singleton = 1;
522 INSERT INTO schema_migrations(version, applied_at)
523 VALUES (46, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
524 PRAGMA user_version = 46;
525 COMMIT;"
526 ))?;
527 }
528
529 if version < 47 {
533 let add_branch =
534 if super::legacy_schema::table_has_column(connection, "sessions", "launch_branch")? {
535 ""
536 } else {
537 "ALTER TABLE sessions ADD COLUMN launch_branch TEXT;"
538 };
539 let add_publication = if super::legacy_schema::table_has_column(
540 connection,
541 "sessions",
542 "publication_json",
543 )? {
544 ""
545 } else {
546 "ALTER TABLE sessions ADD COLUMN publication_json TEXT;"
547 };
548 connection.execute_batch(&format!(
549 "BEGIN IMMEDIATE;
550 {add_branch}
551 {add_publication}
552 UPDATE schema_compatibility SET minimum_compatible_version = 47 WHERE singleton = 1;
553 INSERT INTO schema_migrations(version, applied_at)
554 VALUES (47, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
555 PRAGMA user_version = 47;
556 COMMIT;",
557 ))?;
558 }
559
560 if version < 48 {
564 connection.execute_batch(
565 "BEGIN IMMEDIATE;
566 UPDATE schema_compatibility SET minimum_compatible_version = 48 WHERE singleton = 1;
567 INSERT INTO schema_migrations(version, applied_at)
568 VALUES (48, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
569 PRAGMA user_version = 48;
570 COMMIT;",
571 )?;
572 }
573
574 if version < 49 {
579 connection.execute_batch(
580 "BEGIN IMMEDIATE;
581 CREATE TABLE IF NOT EXISTS subagent_handbacks (
582 child_session_id TEXT PRIMARY KEY,
583 handback_command_id TEXT,
584 handback_message TEXT,
585 handback_recorded_at_ms INTEGER,
586 reminder_command_id TEXT,
587 reminder_for_command_id TEXT,
588 reminder_sent_at_ms INTEGER,
589 reminder_failed_for_command_id TEXT
590 );
591 UPDATE schema_compatibility SET minimum_compatible_version = 49 WHERE singleton = 1;
592 INSERT INTO schema_migrations(version, applied_at)
593 VALUES (49, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
594 PRAGMA user_version = 49;
595 COMMIT;",
596 )?;
597 }
598
599 if version < 50 {
603 let add_column = if super::legacy_schema::table_has_column(
604 connection,
605 "subagent_handbacks",
606 "awaited_ordinal",
607 )? {
608 ""
609 } else {
610 "ALTER TABLE subagent_handbacks ADD COLUMN awaited_ordinal INTEGER;"
611 };
612 connection.execute_batch(&format!(
613 "BEGIN IMMEDIATE;
614 {add_column}
615 INSERT INTO schema_migrations(version, applied_at)
616 VALUES (50, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
617 PRAGMA user_version = 50;
618 COMMIT;"
619 ))?;
620 }
621
622 if version < 51 {
626 let add_column = if super::legacy_schema::table_has_column(
627 connection,
628 "subagent_handbacks",
629 "report_dir",
630 )? {
631 ""
632 } else {
633 "ALTER TABLE subagent_handbacks ADD COLUMN report_dir TEXT;"
634 };
635 connection.execute_batch(&format!(
636 "BEGIN IMMEDIATE;
637 {add_column}
638 INSERT INTO schema_migrations(version, applied_at)
639 VALUES (51, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
640 PRAGMA user_version = 51;
641 COMMIT;"
642 ))?;
643 }
644
645 if version < 52 {
652 connection.execute_batch(
653 "BEGIN IMMEDIATE;
654 CREATE TABLE IF NOT EXISTS stopped_subagents (
655 parent_session_id TEXT NOT NULL
656 REFERENCES sessions(session_id) ON DELETE CASCADE,
657 child_session_id TEXT NOT NULL,
658 record_json TEXT NOT NULL CHECK(json_valid(record_json)),
659 PRIMARY KEY(parent_session_id, child_session_id)
660 ) STRICT;
661 INSERT INTO schema_migrations(version, applied_at)
662 VALUES (52, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
663 PRAGMA user_version = 52;
664 COMMIT;",
665 )?;
666 }
667
668 if version < 53 {
673 migrate_parked_session_state(connection)?;
674 }
675
676 if version < 54 {
679 let add_column =
680 if super::legacy_schema::table_has_column(connection, "sessions", "checkout_json")? {
681 ""
682 } else {
683 "ALTER TABLE sessions ADD COLUMN checkout_json TEXT;"
684 };
685 connection.execute_batch(&format!(
686 "BEGIN IMMEDIATE;
687 {add_column}
688 UPDATE schema_compatibility SET minimum_compatible_version = 54 WHERE singleton = 1;
689 INSERT INTO schema_migrations(version, applied_at)
690 VALUES (54, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
691 PRAGMA user_version = 54;
692 COMMIT;"
693 ))?;
694 }
695
696 if version < 55 {
699 let add_column = if super::legacy_schema::table_has_column(
700 connection,
701 "sessions",
702 "expected_runtime_identity",
703 )? {
704 ""
705 } else {
706 "ALTER TABLE sessions ADD COLUMN expected_runtime_identity TEXT;"
707 };
708 connection.execute_batch(&format!(
709 "BEGIN IMMEDIATE;
710 {add_column}
711 UPDATE schema_compatibility SET minimum_compatible_version = 55 WHERE singleton = 1;
712 INSERT INTO schema_migrations(version, applied_at)
713 VALUES (55, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
714 PRAGMA user_version = 55;
715 COMMIT;"
716 ))?;
717 }
718
719 if version < 56 {
721 let add_column = if super::legacy_schema::table_has_column(
722 connection,
723 "sessions",
724 "subagents",
725 )? {
726 ""
727 } else {
728 "ALTER TABLE sessions ADD COLUMN subagents TEXT CHECK(subagents IS NULL OR json_valid(subagents));"
729 };
730 connection.execute_batch(&format!(
731 "BEGIN IMMEDIATE;
732 {add_column}
733 UPDATE sessions SET subagents = CASE WHEN mjolnir_subagents = 1
734 THEN '{{\"mode\":\"all_models\"}}' ELSE '{{\"mode\":\"native\"}}' END WHERE subagents IS NULL AND mjolnir_subagents IS NOT NULL;
735 CREATE TABLE IF NOT EXISTS subagent_preference (singleton INTEGER PRIMARY KEY CHECK(singleton = 1), policy TEXT NOT NULL CHECK(json_valid(policy)));
736 UPDATE schema_compatibility SET minimum_compatible_version = 56 WHERE singleton = 1;
737 INSERT INTO schema_migrations(version, applied_at) VALUES (56, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
738 PRAGMA user_version = 56;
739 COMMIT;"
740 ))?;
741 }
742
743 if version < 57 {
746 let tx = connection.unchecked_transaction()?;
747 tx.execute_batch("DROP TRIGGER IF EXISTS api_session_error_updated;
748 DROP TRIGGER IF EXISTS api_session_error_inserted;
749 CREATE TABLE IF NOT EXISTS checkpoint_operations (
750 session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
751 command_id TEXT NOT NULL UNIQUE,
752 related_command_ids TEXT NOT NULL DEFAULT '[]' CHECK(json_valid(related_command_ids))
753 ) STRICT;")?;
754 super::events::migrate_event_outcomes(&tx)?;
755 tx.execute_batch("UPDATE schema_compatibility SET minimum_compatible_version = 57 WHERE singleton = 1;
756 INSERT INTO schema_migrations(version, applied_at) VALUES (57, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
757 PRAGMA user_version = 57;")?;
758 tx.commit()?;
759 }
760
761 if version < 58 {
764 connection.execute_batch("BEGIN IMMEDIATE;
765 UPDATE schema_compatibility SET minimum_compatible_version = 58 WHERE singleton = 1;
766 INSERT INTO schema_migrations(version, applied_at) VALUES (58, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
767 PRAGMA user_version = 58;
768 COMMIT;")?;
769 }
770
771 if version < 59 {
774 connection.execute_batch("BEGIN IMMEDIATE;
775 CREATE TABLE IF NOT EXISTS startup_steps (
776 sequence INTEGER PRIMARY KEY AUTOINCREMENT,
777 session_id TEXT NOT NULL,
778 group_id TEXT,
779 command_id TEXT NOT NULL UNIQUE,
780 step_json TEXT NOT NULL,
781 phase TEXT NOT NULL DEFAULT 'pending'
782 CHECK (phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting', 'done', 'failed', 'dismissed')),
783 error TEXT,
784 accepted_ordinal INTEGER
785 ) STRICT;
786 CREATE INDEX IF NOT EXISTS startup_steps_session_sequence
787 ON startup_steps(session_id, sequence);
788 CREATE INDEX IF NOT EXISTS startup_steps_group_sequence
789 ON startup_steps(group_id, sequence) WHERE group_id IS NOT NULL;
790 CREATE INDEX IF NOT EXISTS startup_steps_pending
791 ON startup_steps(session_id, sequence)
792 WHERE phase IN ('pending', 'delivering', 'accepted', 'cancelling', 'rejecting');
793 UPDATE schema_compatibility SET minimum_compatible_version = 59 WHERE singleton = 1;
794 INSERT INTO schema_migrations(version, applied_at) VALUES (59, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
795 PRAGMA user_version = 59;
796 COMMIT;")?;
797 }
798 if version < 60 {
801 connection.execute_batch("BEGIN IMMEDIATE;
802 CREATE TABLE IF NOT EXISTS delegation_effects (
803 parent_session_id TEXT NOT NULL,
804 request_id TEXT NOT NULL,
805 phase TEXT NOT NULL,
806 prepared_json TEXT,
807 result_json TEXT,
808 receipt_pending INTEGER NOT NULL DEFAULT 0,
809 PRIMARY KEY(parent_session_id, request_id)
810 ) STRICT;
811 CREATE INDEX IF NOT EXISTS delegation_effects_receipt_cleanup
812 ON delegation_effects(parent_session_id, request_id) WHERE receipt_pending = 1;
813 UPDATE schema_compatibility SET minimum_compatible_version = 60 WHERE singleton = 1;
814 INSERT INTO schema_migrations(version, applied_at) VALUES (60, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
815 PRAGMA user_version = 60;
816 COMMIT;")?;
817 }
818 if version < 61 {
821 let add_column = if super::legacy_schema::table_has_column(
822 connection,
823 "turn_review_state",
824 "orchestration",
825 )? {
826 ""
827 } else {
828 "ALTER TABLE turn_review_state ADD COLUMN orchestration TEXT;"
829 };
830 connection.execute_batch(&format!("BEGIN IMMEDIATE;
831 {add_column}
832 UPDATE schema_compatibility SET minimum_compatible_version = 61 WHERE singleton = 1;
833 INSERT INTO schema_migrations(version, applied_at) VALUES (61, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
834 PRAGMA user_version = 61;
835 COMMIT;"))?;
836 }
837 if version < 62 {
840 connection.execute_batch("BEGIN IMMEDIATE;
841 CREATE TABLE IF NOT EXISTS worker_restart_intents (
842 session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
843 operation_id TEXT NOT NULL,
844 target_json TEXT NOT NULL,
845 desired_build TEXT NOT NULL,
846 phase TEXT NOT NULL CHECK(phase IN ('prepared', 'swapping', 'awaiting_readiness'))
847 ) STRICT;
848 UPDATE schema_compatibility SET minimum_compatible_version = 62 WHERE singleton = 1;
849 INSERT INTO schema_migrations(version, applied_at) VALUES (62, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
850 PRAGMA user_version = 62;
851 COMMIT;")?;
852 }
853
854 if version < 63 {
857 connection.execute_batch("BEGIN IMMEDIATE;
858 CREATE TABLE IF NOT EXISTS session_incarnations (
859 session_id TEXT PRIMARY KEY REFERENCES sessions(session_id) ON DELETE CASCADE,
860 identity TEXT NOT NULL
861 ) STRICT;
862 INSERT OR IGNORE INTO session_incarnations(session_id, identity)
863 SELECT session_id, lower(hex(randomblob(16))) FROM sessions;
864 CREATE TRIGGER IF NOT EXISTS session_incarnation_insert
865 AFTER INSERT ON sessions BEGIN
866 INSERT INTO session_incarnations(session_id, identity)
867 VALUES(NEW.session_id, lower(hex(randomblob(16))));
868 END;
869 CREATE TRIGGER IF NOT EXISTS session_incarnation_resume
870 AFTER UPDATE OF state ON sessions
871 WHEN (NEW.state = 'provisioning' AND OLD.state <> 'provisioning')
872 OR (NEW.state = 'running' AND OLD.state IN
873 ('stopped', 'parked', 'error', 'lost', 'destroyed-with-data-loss'))
874 BEGIN
875 UPDATE session_incarnations SET identity = lower(hex(randomblob(16)))
876 WHERE session_id = NEW.session_id;
877 END;
878 UPDATE schema_compatibility SET minimum_compatible_version = 63 WHERE singleton = 1;
879 INSERT INTO schema_migrations(version, applied_at) VALUES (63, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
880 PRAGMA user_version = 63;
881 COMMIT;")?;
882 }
883
884 if version < 64 {
887 connection.execute_batch("BEGIN IMMEDIATE;
888 CREATE TABLE IF NOT EXISTS retained_move_sources (
889 operation_id TEXT PRIMARY KEY,
890 session_id TEXT NOT NULL,
891 source_json TEXT NOT NULL,
892 exclusions_json TEXT NOT NULL,
893 created_at TEXT NOT NULL
894 ) STRICT;
895 UPDATE schema_compatibility SET minimum_compatible_version = 64 WHERE singleton = 1;
896 INSERT INTO schema_migrations(version, applied_at) VALUES (64, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
897 PRAGMA user_version = 64;
898 COMMIT;")?;
899 }
900
901 if version < 65 {
907 let drop_column = if super::legacy_schema::table_has_column(
908 connection,
909 "turn_review_state",
910 "orchestration",
911 )? {
912 "ALTER TABLE turn_review_state DROP COLUMN orchestration;"
913 } else {
914 ""
915 };
916 connection.execute_batch(&format!("BEGIN IMMEDIATE;
917 DROP TRIGGER IF EXISTS session_incarnation_insert;
918 DROP TRIGGER IF EXISTS session_incarnation_resume;
919 DROP TABLE IF EXISTS session_incarnations;
920 {drop_column}
921 UPDATE schema_compatibility SET minimum_compatible_version = 65 WHERE singleton = 1;
922 INSERT INTO schema_migrations(version, applied_at) VALUES (65, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
923 PRAGMA user_version = 65;
924 COMMIT;"))?;
925 }
926
927 if version < 66 {
930 let drop_column = if super::legacy_schema::table_has_column(
931 connection,
932 "sessions",
933 "expected_runtime_identity",
934 )? {
935 "ALTER TABLE sessions DROP COLUMN expected_runtime_identity;"
936 } else {
937 ""
938 };
939 connection.execute_batch(&format!("BEGIN IMMEDIATE;
940 {drop_column}
941 UPDATE schema_compatibility SET minimum_compatible_version = 66 WHERE singleton = 1;
942 INSERT INTO schema_migrations(version, applied_at) VALUES (66, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
943 PRAGMA user_version = 66;
944 COMMIT;"))?;
945 }
946
947 if version < 67 {
950 connection.execute_batch(&format!("BEGIN IMMEDIATE;
951 {ACCOUNTING_MIGRATION_SQL}
952 UPDATE schema_compatibility SET minimum_compatible_version = 67 WHERE singleton = 1;
953 INSERT INTO schema_migrations(version, applied_at) VALUES (67, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'));
954 PRAGMA user_version = 67;
955 COMMIT;"))?;
956 }
957
958 if version < 68 {
961 migrate_startup_cleanup_state(connection)?;
962 }
963
964 if version < 69 {
967 let has_accounting = connection
968 .prepare(
969 "SELECT 1 FROM sqlite_schema WHERE type='table' AND name='subagent_accounting'",
970 )?
971 .exists([])?;
972 let accounting = if has_accounting {
973 ""
974 } else {
975 ACCOUNTING_MIGRATION_SQL
976 };
977 let add_snapshot = if super::legacy_schema::table_has_column(
978 connection,
979 "sessions",
980 "project_json",
981 )? {
982 ""
983 } else {
984 "ALTER TABLE sessions ADD COLUMN project_json TEXT CHECK(project_json IS NULL OR json_valid(project_json));"
985 };
986 connection.execute_batch(&format!("BEGIN IMMEDIATE;
987 {accounting}
988 {add_snapshot}
989 CREATE TABLE IF NOT EXISTS project_catalog (
990 bundle_id TEXT PRIMARY KEY,
991 project_key TEXT NOT NULL UNIQUE,
992 snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
993 hidden INTEGER NOT NULL DEFAULT 0 CHECK(hidden IN (0,1))
994 ) STRICT;
995 CREATE TABLE IF NOT EXISTS project_aliases (
996 bundle_id TEXT PRIMARY KEY,
997 canonical_id TEXT NOT NULL REFERENCES project_catalog(bundle_id),
998 snapshot_json TEXT NOT NULL CHECK(json_valid(snapshot_json)),
999 config_pending INTEGER NOT NULL DEFAULT 0 CHECK(config_pending IN (0,1))
1000 ) STRICT;
1001 CREATE TABLE IF NOT EXISTS project_session_aliases (
1002 session_id TEXT NOT NULL REFERENCES session_contexts(session_id) ON DELETE CASCADE,
1003 bundle_id TEXT NOT NULL, PRIMARY KEY(session_id,bundle_id)
1004 ) STRICT;
1005 CREATE TABLE IF NOT EXISTS project_locations (
1006 host TEXT NOT NULL,
1007 directory BLOB NOT NULL,
1008 checkout_root BLOB NOT NULL,
1009 repository_root BLOB NOT NULL,
1010 identity_json TEXT NOT NULL CHECK(json_valid(identity_json)),
1011 seen_at TEXT NOT NULL,
1012 PRIMARY KEY(host, directory)
1013 ) STRICT;
1014 CREATE TABLE IF NOT EXISTS project_seed_homes (
1015 harness TEXT NOT NULL,
1016 home BLOB NOT NULL,
1017 PRIMARY KEY(harness, home)
1018 ) STRICT;
1019 CREATE TABLE IF NOT EXISTS project_seed_failures (
1020 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)),
1021 PRIMARY KEY(harness,home,directory)
1022 ) STRICT;
1023 CREATE TABLE IF NOT EXISTS project_discovery_changes (
1024 sequence INTEGER PRIMARY KEY AUTOINCREMENT,
1025 session_id TEXT NOT NULL,
1026 directory BLOB,
1027 managed_worktree TEXT,
1028 target_template_id TEXT NOT NULL
1029 ) STRICT;
1030 CREATE TABLE IF NOT EXISTS project_discovery_progress (
1031 singleton INTEGER PRIMARY KEY CHECK(singleton=1),
1032 sequence INTEGER NOT NULL DEFAULT 0
1033 ) STRICT;
1034 CREATE TABLE IF NOT EXISTS project_discovery_failures (
1035 sequence INTEGER PRIMARY KEY REFERENCES project_discovery_changes(sequence) ON DELETE CASCADE,
1036 error TEXT NOT NULL
1037 ) STRICT;
1038 INSERT OR IGNORE INTO project_discovery_progress(singleton) VALUES(1);
1039 INSERT INTO project_discovery_changes(session_id, directory, managed_worktree, target_template_id)
1040 SELECT session_id, project_directory, managed_worktree, target_template_id FROM sessions
1041 WHERE project_directory IS NOT NULL AND NOT EXISTS (
1042 SELECT 1 FROM project_discovery_changes d WHERE d.session_id=sessions.session_id
1043 );
1044 CREATE TRIGGER IF NOT EXISTS project_discovery_insert AFTER INSERT ON sessions
1045 WHEN NEW.project_directory IS NOT NULL BEGIN
1046 INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
1047 VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
1048 END;
1049 CREATE TRIGGER IF NOT EXISTS project_discovery_update AFTER UPDATE OF project_directory,managed_worktree,target_template_id ON sessions
1050 WHEN NEW.project_directory IS NOT NULL AND
1051 (NEW.project_directory IS NOT OLD.project_directory
1052 OR NEW.managed_worktree IS NOT OLD.managed_worktree
1053 OR NEW.target_template_id IS NOT OLD.target_template_id) BEGIN
1054 INSERT INTO project_discovery_changes(session_id,directory,managed_worktree,target_template_id)
1055 VALUES(NEW.session_id,NEW.project_directory,NEW.managed_worktree,NEW.target_template_id);
1056 END;
1057 UPDATE schema_compatibility SET minimum_compatible_version=69 WHERE singleton=1;
1058 INSERT INTO schema_migrations(version,applied_at) VALUES(69,strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1059 PRAGMA user_version=69;
1060 COMMIT;"))?;
1061 }
1062
1063 if version < 70 {
1066 let transaction = connection.unchecked_transaction()?;
1067 let titles: Vec<(String, String)> = transaction
1068 .prepare("SELECT session_id, acp_session_title FROM sessions WHERE length(acp_session_title) > ?1")?
1069 .query_map([mj_core::state::MAX_SESSION_TITLE_CHARS], |row| {
1070 Ok((row.get(0)?, row.get(1)?))
1071 })?
1072 .collect::<rusqlite::Result<_>>()?;
1073 for (session_id, title) in titles {
1074 transaction.execute(
1075 "UPDATE sessions SET acp_session_title = ?2 WHERE session_id = ?1",
1076 params![session_id, mj_core::state::normalize_session_title(&title)],
1077 )?;
1078 }
1079 transaction.execute_batch(
1080 "INSERT INTO schema_migrations(version, applied_at)
1081 VALUES (70, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1082 PRAGMA user_version = 70;",
1083 )?;
1084 transaction.commit()?;
1085 }
1086
1087 let recorded: Option<i64> =
1088 connection.query_row("SELECT max(version) FROM schema_migrations", [], |row| {
1089 row.get(0)
1090 })?;
1091 if recorded == Some(SCHEMA_VERSION) && version < SCHEMA_VERSION {
1092 tracing::info!(
1093 from_revision = version,
1094 to_revision = SCHEMA_VERSION,
1095 "database migrations applied"
1096 );
1097 }
1098 if recorded != Some(SCHEMA_VERSION) {
1099 bail!(
1100 "Mjolnir database migration ledger {:?} does not match schema {}",
1101 recorded,
1102 SCHEMA_VERSION
1103 );
1104 }
1105 Ok(())
1106}
1107
1108fn migrate_parked_session_state(connection: &Connection) -> Result<()> {
1115 const BEFORE: &str = "'stopped','lost',";
1116 const AFTER: &str = "'stopped','parked','lost',";
1117 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1118 let migration = (|| -> Result<()> {
1119 let transaction = connection.unchecked_transaction()?;
1120 let sql: String = transaction.query_row(
1121 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1122 [],
1123 |row| row.get(0),
1124 )?;
1125 let (_, definition) = sql
1126 .split_once('(')
1127 .context("missing sessions table definition")?;
1128 if !definition.contains(AFTER) {
1129 ensure!(
1130 definition.matches(BEFORE).count() == 1,
1131 "unexpected sessions state constraint"
1132 );
1133 let definition = definition.replace(BEFORE, AFTER);
1134 let objects: Vec<String> = transaction
1135 .prepare(
1136 "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1137 AND type IN ('index','trigger') AND sql IS NOT NULL",
1138 )?
1139 .query_map([], |row| row.get(0))?
1140 .collect::<rusqlite::Result<_>>()?;
1141 transaction.execute_batch(&format!(
1142 "CREATE TABLE sessions_parked_v53 ({definition};
1143 INSERT INTO sessions_parked_v53 SELECT * FROM sessions;
1144 DROP TABLE sessions;
1145 ALTER TABLE sessions_parked_v53 RENAME TO sessions;"
1146 ))?;
1147 for object in objects {
1148 transaction.execute_batch(&object)?;
1149 }
1150 ensure!(
1151 !transaction
1152 .prepare("PRAGMA foreign_key_check")?
1153 .exists([])?,
1154 "foreign key violation in the parked-state migration"
1155 );
1156 }
1157 transaction.execute_batch(
1158 "UPDATE schema_compatibility SET minimum_compatible_version = 53
1159 WHERE singleton = 1;
1160 INSERT INTO schema_migrations(version, applied_at)
1161 VALUES (53, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1162 PRAGMA user_version = 53;",
1163 )?;
1164 transaction.commit()?;
1165 Ok(())
1166 })();
1167 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1168 migration.context("migrate the sessions state constraint for parked sub-agents")?;
1169 restored.context("restore foreign key enforcement after the parked-state migration")?;
1170 Ok(())
1171}
1172
1173fn migrate_startup_cleanup_state(connection: &Connection) -> Result<()> {
1175 const BEFORE: &str = "'stopped','parked','lost',";
1176 const AFTER: &str = "'stopped','parked','startup-cleanup','lost',";
1177 connection.execute_batch("PRAGMA foreign_keys = OFF;")?;
1178 let migration = (|| -> Result<()> {
1179 let transaction = connection.unchecked_transaction()?;
1180 let sql: String = transaction.query_row(
1181 "SELECT sql FROM sqlite_schema WHERE type='table' AND name='sessions'",
1182 [],
1183 |row| row.get(0),
1184 )?;
1185 let (_, definition) = sql
1186 .split_once('(')
1187 .context("missing sessions table definition")?;
1188 if !definition.contains(AFTER) {
1189 ensure!(
1190 definition.matches(BEFORE).count() == 1,
1191 "unexpected sessions state constraint"
1192 );
1193 let definition = definition.replace(BEFORE, AFTER);
1194 let objects: Vec<String> = transaction
1195 .prepare(
1196 "SELECT sql FROM sqlite_schema WHERE tbl_name='sessions'
1197 AND type IN ('index','trigger') AND sql IS NOT NULL",
1198 )?
1199 .query_map([], |row| row.get(0))?
1200 .collect::<rusqlite::Result<_>>()?;
1201 transaction.execute_batch(&format!(
1202 "CREATE TABLE sessions_startup_cleanup_v68 ({definition};
1203 INSERT INTO sessions_startup_cleanup_v68 SELECT * FROM sessions;
1204 DROP TABLE sessions;
1205 ALTER TABLE sessions_startup_cleanup_v68 RENAME TO sessions;"
1206 ))?;
1207 for object in objects {
1208 transaction.execute_batch(&object)?;
1209 }
1210 ensure!(
1211 !transaction
1212 .prepare("PRAGMA foreign_key_check")?
1213 .exists([])?,
1214 "foreign key violation in the startup-cleanup migration"
1215 );
1216 }
1217 transaction.execute_batch(
1218 "UPDATE schema_compatibility SET minimum_compatible_version = 68
1219 WHERE singleton = 1;
1220 INSERT INTO schema_migrations(version, applied_at)
1221 VALUES (68, strftime('%Y-%m-%dT%H:%M:%fZ','now'));
1222 PRAGMA user_version = 68;",
1223 )?;
1224 transaction.commit()?;
1225 Ok(())
1226 })();
1227 let restored = connection.execute_batch("PRAGMA foreign_keys = ON;");
1228 migration.context("migrate the sessions state constraint for failed startup cleanup")?;
1229 restored.context("restore foreign key enforcement after the startup-cleanup migration")?;
1230 Ok(())
1231}
1232
1233fn create_baseline_schema(connection: &Connection) -> Result<()> {
1237 connection.execute_batch("BEGIN IMMEDIATE;")?;
1238 let created = (|| -> Result<()> {
1239 let version: i64 = connection.query_row("PRAGMA user_version", [], |row| row.get(0))?;
1240 if version != 0 {
1241 return Ok(());
1242 }
1243 connection.execute_batch(include_str!("baseline.sql"))?;
1244 connection.execute(
1245 "INSERT INTO schema_compatibility(singleton, minimum_compatible_version) VALUES (1, ?1)",
1246 [BASELINE_MINIMUM_COMPATIBLE_VERSION],
1247 )?;
1248 connection.execute(
1249 "INSERT INTO schema_migrations(version, applied_at)
1250 VALUES (?1, strftime('%Y-%m-%dT%H:%M:%fZ', 'now'))",
1251 [BASELINE_SCHEMA_VERSION],
1252 )?;
1253 connection.pragma_update(None, "user_version", BASELINE_SCHEMA_VERSION)?;
1254 Ok(())
1255 })();
1256 match created {
1257 Ok(()) => connection
1258 .execute_batch("COMMIT;")
1259 .context("commit baseline database schema"),
1260 Err(error) => {
1261 if let Err(rollback) = connection.execute_batch("ROLLBACK;") {
1262 tracing::warn!(%rollback, "could not roll back a failed baseline schema");
1263 }
1264 Err(error.context("create baseline database schema"))
1265 }
1266 }
1267}
1268
1269#[cfg(test)]
1270pub(super) fn advance_test_schema(path: &Path, revision: i64, minimum_compatible: i64) {
1271 let connection = Connection::open(path).unwrap();
1272 let transaction = connection.unchecked_transaction().unwrap();
1273 transaction
1274 .execute(
1275 "UPDATE schema_compatibility SET minimum_compatible_version = ?1",
1276 [minimum_compatible],
1277 )
1278 .unwrap();
1279 transaction
1280 .execute(
1281 "INSERT INTO schema_migrations(version, applied_at) VALUES (?1, 'test')",
1282 [revision],
1283 )
1284 .unwrap();
1285 transaction
1286 .pragma_update(None, "user_version", revision)
1287 .unwrap();
1288 transaction.commit().unwrap();
1289 forget_verified_schema(path);
1290}
1291
1292#[cfg(test)]
1293mod reader_tests {
1294 use super::*;
1295
1296 fn assert_divergent_history_upgrades(revision: i64, project_history: bool, interrupt: bool) {
1297 let directory = tempfile::tempdir().unwrap();
1298 let path = directory.path().join("divergent-history.sqlite");
1299 let connection = Connection::open(&path).unwrap();
1300 create_baseline_schema(&connection).unwrap();
1301 connection.execute_batch(&format!(
1302 "INSERT INTO session_contexts(session_id,bundle_id,created_at)
1303 VALUES ('kept','project','now');
1304 INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at,project_directory)
1305 VALUES ('kept','Keep my work','codex','codex','local','error','now',X'2F7265706F');
1306 INSERT INTO materialized_sessions(session_id) VALUES ('kept');
1307 INSERT INTO session_turn_usage VALUES ('kept','turn',1,1,'{{\"tokens\":42}}');
1308 INSERT INTO session_provider_cost VALUES ('kept','{{\"amount\":1}}');
1309 CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1310 WHEN NEW.version > {revision}
1311 BEGIN SELECT RAISE(ABORT,'fixture migration boundary'); END;"
1312 )).unwrap();
1313 assert!(migrate_schema(&connection).is_err());
1314 if !connection.is_autocommit() {
1315 connection.execute_batch("ROLLBACK").unwrap();
1316 }
1317 connection
1318 .execute_batch("DROP TRIGGER stop_at_revision")
1319 .unwrap();
1320 assert_eq!(read_schema_state(&connection).unwrap().revision, revision);
1321 if project_history {
1322 connection
1323 .execute_batch(include_str!("project_catalog_v67.sql"))
1324 .unwrap();
1325 connection
1326 .execute_batch(
1327 "INSERT INTO project_catalog VALUES ('project','key','{}',0);
1328 INSERT INTO project_aliases VALUES ('alias','project','{}',1);
1329 INSERT INTO project_session_aliases VALUES ('kept','alias');
1330 UPDATE sessions SET project_json='{\"kept\":true}';",
1331 )
1332 .unwrap();
1333 }
1334 if interrupt {
1335 connection
1336 .execute_batch(
1337 "CREATE TRIGGER interrupt_reconciliation BEFORE INSERT ON schema_migrations
1338 WHEN NEW.version=69 BEGIN SELECT RAISE(ABORT,'interrupted reconciliation'); END;",
1339 )
1340 .unwrap();
1341 assert!(migrate_schema(&connection).is_err());
1342 drop(connection);
1343 let connection = Connection::open(&path).unwrap();
1344 assert_eq!(read_schema_state(&connection).unwrap().revision, 68);
1345 connection
1346 .execute_batch("DROP TRIGGER interrupt_reconciliation")
1347 .unwrap();
1348 } else {
1349 drop(connection);
1350 }
1351
1352 let writer = open_writer(&path).unwrap();
1353 let state = read_schema_state(&writer).unwrap();
1354 assert_eq!(state.revision, SCHEMA_VERSION);
1355 assert!(state.ensure_supported_by(68).is_err());
1356 assert!(state.ensure_supported_by(67).is_err());
1357 if project_history {
1358 let snapshot: String = writer
1359 .query_row(
1360 "SELECT project_json FROM sessions WHERE session_id='kept'",
1361 [],
1362 |row| row.get(0),
1363 )
1364 .unwrap();
1365 assert_eq!(snapshot, r#"{"kept":true}"#);
1366 let alias: (String, i64) = writer.query_row(
1367 "SELECT canonical_id,config_pending FROM project_aliases WHERE bundle_id='alias'", [],
1368 |row| Ok((row.get(0)?, row.get(1)?))
1369 ).unwrap();
1370 assert_eq!(alias, ("project".to_owned(), 1));
1371 }
1372 writer.execute_batch(
1374 "UPDATE sessions SET state='startup-cleanup',project_directory=X'2F6E6577' WHERE session_id='kept';
1375 INSERT INTO session_turn_selections VALUES ('kept','turn','model','high');"
1376 ).unwrap();
1377 let changed: Vec<u8> = writer
1378 .query_row(
1379 "SELECT directory FROM project_discovery_changes ORDER BY sequence DESC LIMIT 1",
1380 [],
1381 |row| row.get(0),
1382 )
1383 .unwrap();
1384 assert_eq!(changed, b"/new");
1385 writer
1386 .execute("DELETE FROM sessions WHERE session_id='kept'", [])
1387 .unwrap();
1388 let usage: String = writer
1389 .query_row("SELECT body FROM session_turn_usage", [], |row| row.get(0))
1390 .unwrap();
1391 assert_eq!(usage, r#"{"tokens":42}"#);
1392 let cost: String = writer
1393 .query_row("SELECT body FROM session_provider_cost", [], |row| {
1394 row.get(0)
1395 })
1396 .unwrap();
1397 assert_eq!(cost, r#"{"amount":1}"#);
1398 assert!(
1399 !writer
1400 .prepare("PRAGMA foreign_key_check")
1401 .unwrap()
1402 .exists([])
1403 .unwrap()
1404 );
1405 drop(writer);
1406 forget_verified_schema(&path);
1407 assert_eq!(
1408 read_schema_state(&open_writer(&path).unwrap())
1409 .unwrap()
1410 .revision,
1411 SCHEMA_VERSION
1412 );
1413 }
1414
1415 #[test]
1416 fn divergent_accounting_revision_67_preserves_usage_and_adds_projects() {
1417 assert_divergent_history_upgrades(67, false, false);
1418 }
1419
1420 #[test]
1421 fn divergent_cleanup_revision_68_preserves_usage_and_adds_projects() {
1422 assert_divergent_history_upgrades(68, false, false);
1423 }
1424
1425 #[test]
1426 fn divergent_project_revision_67_preserves_aliases_snapshots_and_usage() {
1427 assert_divergent_history_upgrades(66, true, false);
1428 }
1429
1430 #[test]
1431 fn divergent_project_reconciliation_resumes_after_interruption() {
1432 assert_divergent_history_upgrades(66, true, true);
1433 }
1434
1435 #[test]
1436 fn move_ownership_upgrade_retains_sources_and_refuses_previous_daemons() {
1437 let directory = tempfile::tempdir().unwrap();
1438 let path = directory.path().join("move.sqlite");
1439 let connection = open_writer(&path).unwrap();
1440 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();
1441 drop(connection);
1442 forget_verified_schema(&path);
1443 let upgraded = open_writer(&path).unwrap();
1444 let schema = read_schema_state(&upgraded).unwrap();
1445 assert!(schema.ensure_supported_by(63).is_err());
1446 upgraded.execute("INSERT INTO retained_move_sources VALUES ('move-one','session-one','{}','[]','now')", []).unwrap();
1447 drop(upgraded);
1448 let reopened = open_writer(&path).unwrap();
1449 let count: i64 = reopened
1450 .query_row("SELECT count(*) FROM retained_move_sources", [], |row| {
1451 row.get(0)
1452 })
1453 .unwrap();
1454 assert_eq!(count, 1);
1455 }
1456
1457 #[test]
1458 fn title_migration_caps_old_titles_preserves_short_titles_and_is_compatible() {
1459 let directory = tempfile::tempdir().unwrap();
1460 let path = directory.path().join("title-migration.sqlite3");
1461 let titles = [
1462 ("long", Some("word ".repeat(20_000))),
1463 ("unicode", Some("界".repeat(257))),
1464 ("exact", Some("界".repeat(256))),
1465 ("short", Some(" Keep\nthis title ".into())),
1466 ("unset", None),
1467 ];
1468 for (id, _) in &titles {
1469 let mut session = super::super::tests::session(id, "project");
1470 session.state = SessionState::Stopped;
1471 save_session_to(&path, &session).unwrap();
1472 }
1473 stamp_schema_version(&path, 69);
1474 let old = Connection::open(&path).unwrap();
1475 for (id, title) in &titles {
1476 old.execute(
1477 "UPDATE sessions SET acp_session_title=?2 WHERE session_id=?1",
1478 params![id, title],
1479 )
1480 .unwrap();
1481 }
1482 drop(old);
1483
1484 let upgraded = open_writer(&path).unwrap();
1485 let state = read_schema_state(&upgraded).unwrap();
1486 assert_eq!(state.revision, 70);
1487 assert_eq!(state.minimum_compatible, Some(69));
1488 state.ensure_supported_by(69).unwrap();
1489 for (id, original) in &titles {
1490 let stored: Option<String> = upgraded
1491 .query_row(
1492 "SELECT acp_session_title FROM sessions WHERE session_id=?1",
1493 [id],
1494 |row| row.get(0),
1495 )
1496 .unwrap();
1497 let expected = match *id {
1498 "long" => Some(format!("{}word…", "word ".repeat(50))),
1499 "unicode" => Some(format!("{}…", "界".repeat(255))),
1500 _ => original.clone(),
1501 };
1502 assert_eq!(stored, expected, "session {id}");
1503 assert!(stored.is_none_or(|title| title.chars().count() <= 256));
1504 }
1505 drop(upgraded);
1506 forget_verified_schema(&path);
1507 drop(open_writer(&path).unwrap());
1508 }
1509
1510 #[test]
1511 fn interrupted_title_migration_rolls_back_titles_and_revision() {
1512 let directory = tempfile::tempdir().unwrap();
1513 let path = directory.path().join("title-migration-interrupted.sqlite3");
1514 save_session_to(&path, &super::super::tests::session("old", "project")).unwrap();
1515 stamp_schema_version(&path, 69);
1516 let original = "word ".repeat(20_000);
1517 let old = Connection::open(&path).unwrap();
1518 old.execute("UPDATE sessions SET acp_session_title=?1", [&original])
1519 .unwrap();
1520 old.execute_batch(
1521 "CREATE TRIGGER stop_title_migration BEFORE INSERT ON schema_migrations
1522 WHEN NEW.version=70 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;",
1523 )
1524 .unwrap();
1525 assert!(migrate_schema(&old).is_err());
1526 assert_eq!(read_schema_state(&old).unwrap().revision, 69);
1527 let stored: String = old
1528 .query_row("SELECT acp_session_title FROM sessions", [], |row| {
1529 row.get(0)
1530 })
1531 .unwrap();
1532 assert_eq!(stored, original);
1533 old.execute_batch("DROP TRIGGER stop_title_migration")
1534 .unwrap();
1535 drop(old);
1536 drop(open_writer(&path).unwrap());
1537 assert!(
1538 load_state_from(&path).unwrap().sessions["old"]
1539 .acp_session_title
1540 .as_ref()
1541 .unwrap()
1542 .chars()
1543 .count()
1544 <= 256
1545 );
1546 }
1547
1548 #[test]
1549 fn recent_revisions_upgrade_directly_and_preserve_user_data() {
1550 const TEST_REVISION_FLOOR: i64 = 47;
1553 for revision in TEST_REVISION_FLOOR..SCHEMA_VERSION {
1554 let directory = tempfile::tempdir().unwrap();
1555 let path = directory.path().join("mj.sqlite3");
1556 let connection = Connection::open(&path).unwrap();
1557 connection
1558 .execute_batch(include_str!("legacy_v1.sql"))
1559 .unwrap();
1560 connection.execute_batch(
1561 "INSERT INTO session_contexts VALUES ('old-session', 'project', '2026-01-01T00:00:00Z');
1562 INSERT INTO sessions(session_id, title, harness_kind, last_profile,
1563 target_template_id, state, updated_at, native_session_id)
1564 VALUES ('old-session', 'Keep my work', 'codex', 'codex', 'local', 'error',
1565 '2026-01-01T00:00:00Z', 'native-original');
1566 INSERT INTO prompt_history(session_id, event_sequence, submitted_at, text)
1567 VALUES ('old-session', 1, '2026-01-01T00:00:00Z', 'Keep my prompt');"
1568 ).unwrap();
1569 connection
1573 .execute_batch(&format!(
1574 "CREATE TRIGGER stop_at_revision BEFORE INSERT ON schema_migrations
1575 WHEN NEW.version > {revision}
1576 BEGIN SELECT RAISE(ABORT, 'fixture migration boundary'); END;"
1577 ))
1578 .unwrap();
1579 super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
1580 assert!(migrate_schema(&connection).is_err());
1581 drop(connection);
1582 let connection = Connection::open(&path).unwrap();
1583 let found: i64 = connection
1584 .query_row("PRAGMA user_version", [], |row| row.get(0))
1585 .unwrap();
1586 assert_eq!(found, revision);
1587 connection
1588 .execute_batch("DROP TRIGGER stop_at_revision")
1589 .unwrap();
1590 drop(connection);
1591
1592 let writer =
1593 open_writer(&path).unwrap_or_else(|error| panic!("revision {revision}: {error:#}"));
1594 assert_eq!(read_schema_state(&writer).unwrap().revision, SCHEMA_VERSION);
1595 assert_eq!(
1596 writer
1597 .query_row("PRAGMA integrity_check", [], |row| row.get::<_, String>(0))
1598 .unwrap(),
1599 "ok"
1600 );
1601 assert!(
1602 !writer
1603 .prepare("PRAGMA foreign_key_check")
1604 .unwrap()
1605 .exists([])
1606 .unwrap()
1607 );
1608 drop(writer);
1609 let reader = open_reader_strict(&path).unwrap();
1610 let prompt: String = reader
1611 .query_row("SELECT text FROM prompt_history", [], |row| row.get(0))
1612 .unwrap();
1613 assert_eq!(prompt, "Keep my prompt");
1614 let state = load_state_from(&path).unwrap();
1615 assert_eq!(state.sessions["old-session"].title, "Keep my work");
1616 assert_eq!(
1617 state.sessions["old-session"].native_session_id.as_deref(),
1618 Some("native-original")
1619 );
1620 drop(reader);
1621 forget_verified_schema(&path);
1623 drop(open_writer(&path).unwrap());
1624 }
1625 }
1626
1627 #[test]
1628 fn accounting_migration_preserves_usage_and_changes_its_deletion_owner() {
1629 let dir = tempfile::tempdir().unwrap();
1630 let path = dir.path().join("migration.sqlite");
1631 let connection = Connection::open(&path).unwrap();
1632 connection
1633 .execute_batch(include_str!("legacy_v1.sql"))
1634 .unwrap();
1635 connection.execute_batch("INSERT INTO session_contexts VALUES ('old-session','project','2026-01-01T00:00:00Z');
1636 INSERT INTO sessions(session_id,title,harness_kind,last_profile,target_template_id,state,updated_at)
1637 VALUES ('old-session','Retain usage','codex','codex','local','error','2026-01-01T00:00:00Z');
1638 CREATE TRIGGER stop_before_accounting BEFORE INSERT ON schema_migrations WHEN NEW.version=67
1639 BEGIN SELECT RAISE(ABORT,'fixture boundary'); END;").unwrap();
1640 super::super::legacy_schema::migrate_to_baseline(&connection).unwrap();
1641 assert!(migrate_schema(&connection).is_err());
1642 connection
1643 .execute_batch(
1644 "ROLLBACK; DROP TRIGGER stop_before_accounting;
1645 INSERT INTO session_turn_usage VALUES ('old-session','turn',7,1,'{}');
1646 INSERT INTO session_provider_cost VALUES ('old-session','{\"amount\":1}');",
1647 )
1648 .unwrap();
1649 drop(connection);
1650 let writer = open_writer(&path).unwrap();
1651 writer
1652 .execute("DELETE FROM sessions WHERE session_id='old-session'", [])
1653 .unwrap();
1654 for table in ["session_turn_usage", "session_provider_cost"] {
1655 assert_eq!(
1656 writer
1657 .query_row(&format!("SELECT COUNT(*) FROM {table}"), [], |row| row
1658 .get::<_, u64>(0))
1659 .unwrap(),
1660 1
1661 );
1662 }
1663 assert_eq!(
1664 writer
1665 .query_row("SELECT body FROM session_turn_usage", [], |row| row
1666 .get::<_, String>(0))
1667 .unwrap(),
1668 "{}"
1669 );
1670 assert!(
1671 !writer
1672 .prepare("PRAGMA foreign_key_check")
1673 .unwrap()
1674 .exists([])
1675 .unwrap()
1676 );
1677 assert_eq!(
1678 read_schema_state(&writer).unwrap().minimum_compatible,
1679 Some(MINIMUM_COMPATIBLE_VERSION)
1680 );
1681 }
1682
1683 const MINIMUM_COMPATIBLE_VERSION: i64 = 69;
1686
1687 fn stamp_schema_version(path: &Path, version: i64) {
1690 if version > SCHEMA_VERSION {
1691 advance_test_schema(path, version, version);
1692 return;
1693 }
1694 let connection = Connection::open(path).unwrap();
1695 if version < 67 {
1696 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();
1697 }
1698 connection
1699 .execute_batch(&format!("PRAGMA user_version = {version};"))
1700 .unwrap();
1701 connection
1702 .execute(
1703 "DELETE FROM schema_migrations WHERE version > ?1",
1704 [version],
1705 )
1706 .unwrap();
1707 if version == 30 {
1708 connection
1709 .execute(
1710 "UPDATE schema_compatibility SET minimum_compatible_version = 30 WHERE singleton = 1",
1711 [],
1712 )
1713 .unwrap();
1714 }
1715 drop(connection);
1716 forget_verified_schema(path);
1717 }
1718
1719 #[test]
1720 fn subagent_policy_migration_preserves_legacy_choices_and_refuses_old_writers() {
1721 use mj_core::subagent::SubagentPolicy;
1722 let directory = tempfile::tempdir().unwrap();
1723 let path = directory.path().join("subagent-policy.sqlite3");
1724 for id in ["all", "native", "unset"] {
1725 save_session_to(&path, &super::super::tests::session(id, "project")).unwrap();
1726 }
1727 let connection = Connection::open(&path).unwrap();
1728 connection
1729 .execute_batch(
1730 "UPDATE sessions SET mjolnir_subagents = 1 WHERE session_id = 'all';
1731 UPDATE sessions SET mjolnir_subagents = 0 WHERE session_id = 'native';
1732 ALTER TABLE sessions DROP COLUMN subagents;
1733 DROP TABLE subagent_preference;
1734 DELETE FROM schema_migrations WHERE version >= 56;
1735 UPDATE schema_compatibility SET minimum_compatible_version = 55;
1736 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;",
1737 )
1738 .unwrap();
1739 drop(connection);
1740 forget_verified_schema(&path);
1741 let upgraded = open_writer(&path).unwrap();
1742 assert!(
1743 read_schema_state(&upgraded)
1744 .unwrap()
1745 .ensure_supported_by(55)
1746 .is_err()
1747 );
1748 drop(upgraded);
1749 let state = load_state_from(&path).unwrap();
1750 assert_eq!(
1751 state.sessions["all"].subagents,
1752 Some(SubagentPolicy::AllModels)
1753 );
1754 assert_eq!(
1755 state.sessions["native"].subagents,
1756 Some(SubagentPolicy::Native)
1757 );
1758 assert_eq!(state.sessions["unset"].subagents, None);
1759 assert_eq!(state.last_subagent_policy, SubagentPolicy::Native);
1760 }
1761
1762 #[test]
1763 fn steering_migration_refuses_builds_that_cannot_read_returned_steers() {
1764 let directory = tempfile::tempdir().unwrap();
1765 let path = directory.path().join("steering-migration.sqlite3");
1766 let record = super::super::tests::session("steered-session", "project");
1767 save_session_to(&path, &record).unwrap();
1768 let connection = Connection::open(&path).unwrap();
1769 connection
1770 .execute_batch(
1771 "DELETE FROM schema_migrations WHERE version >= 48;
1772 UPDATE schema_compatibility SET minimum_compatible_version = 47;
1773 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;",
1774 )
1775 .unwrap();
1776 drop(connection);
1777 forget_verified_schema(&path);
1778 let upgraded = open_writer(&path).unwrap();
1779 let schema = read_schema_state(&upgraded).unwrap();
1780 assert_eq!(schema.revision, SCHEMA_VERSION);
1781 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1782 let error = schema.ensure_supported_by(47).unwrap_err();
1783 assert!(matches!(
1784 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
1785 StoreSchemaMismatchReason::Incompatible {
1786 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
1787 }
1788 ));
1789 drop(upgraded);
1790 assert_eq!(load_state_from(&path).unwrap().sessions[&record.id], record);
1791 }
1792
1793 #[test]
1794 fn durable_target_migration_preserves_sessions_and_refuses_previous_builds() {
1795 let directory = tempfile::tempdir().unwrap();
1796 let path = directory.path().join("target-migration.sqlite3");
1797 let record = super::super::tests::session("preserved-session", "project");
1798 save_session_to(&path, &record).unwrap();
1799 let connection = Connection::open(&path).unwrap();
1800 connection
1801 .execute_batch(
1802 "ALTER TABLE sessions DROP COLUMN target_runtime_json;
1803 DELETE FROM schema_migrations WHERE version >= 46;
1804 UPDATE schema_compatibility SET minimum_compatible_version = 44;
1805 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;",
1806 )
1807 .unwrap();
1808 forget_verified_schema(&path);
1809 let upgraded = open_writer(&path).unwrap();
1810 let schema = read_schema_state(&upgraded).unwrap();
1811 assert_eq!(schema.revision, SCHEMA_VERSION);
1812 assert_eq!(schema.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1813 let error = schema.ensure_supported_by(45).unwrap_err();
1814 assert!(matches!(
1815 error.downcast_ref::<StoreSchemaMismatch>().unwrap().reason,
1816 StoreSchemaMismatchReason::Incompatible {
1817 minimum_compatible: MINIMUM_COMPATIBLE_VERSION
1818 }
1819 ));
1820 drop(upgraded);
1821 let restored = load_state_from(&path).unwrap();
1822 assert_eq!(restored.sessions[&record.id], record);
1823 }
1824
1825 #[test]
1826 fn native_agents_and_unstructured_input_raise_the_store_compatibility_floor() {
1827 let directory = tempfile::tempdir().unwrap();
1828 let path = directory.path().join("mj.sqlite3");
1829 let connection = open_writer(&path).unwrap();
1830 connection
1831 .execute_batch(
1832 "BEGIN IMMEDIATE;
1833 DROP TABLE quota_reset_cache;
1834 DROP TABLE native_agent_transcript;
1835 DROP TABLE native_agents;
1836 DROP TABLE native_agent_replay;
1837 DELETE FROM schema_migrations WHERE version >= 39;
1838 UPDATE schema_compatibility SET minimum_compatible_version = 32;
1839 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;
1840 COMMIT;",
1841 )
1842 .unwrap();
1843 migrate_schema(&connection).unwrap();
1844 let state = read_schema_state(&connection).unwrap();
1845 assert_eq!(state.revision, SCHEMA_VERSION);
1846 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1847 let event = ApiEventData::InputRequired {
1848 request: None,
1849 turn_id: Some(1),
1850 };
1851 #[derive(serde::Deserialize)]
1852 struct LegacyInputEvent {
1853 #[serde(rename = "request")]
1854 _request: mj_core::elicitation::ElicitationRequest,
1855 }
1856 let encoded = serde_json::to_value(&event).unwrap();
1857 assert!(serde_json::from_value::<LegacyInputEvent>(encoded["data"].clone()).is_err());
1858 }
1859
1860 #[test]
1861 fn older_readers_and_reopened_writers_preserve_a_compatible_future_schema() {
1862 let directory = tempfile::tempdir().unwrap();
1863 let path = directory.path().join("mj.sqlite3");
1864 let connection = open_writer(&path).unwrap();
1865 connection
1866 .execute_batch(
1867 "CREATE TABLE future_feature(value TEXT NOT NULL);
1868 INSERT INTO future_feature VALUES ('preserve me');",
1869 )
1870 .unwrap();
1871 drop(connection);
1872 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1873
1874 let reader = open_reader_strict(&path).unwrap();
1875 assert_eq!(
1876 reader
1877 .query_row("SELECT value FROM future_feature", [], |row| row
1878 .get::<_, String>(0))
1879 .unwrap(),
1880 "preserve me"
1881 );
1882 assert!(reader.execute("DELETE FROM future_feature", []).is_err());
1883 drop(reader);
1884
1885 let raw = Connection::open(&path).unwrap();
1888 raw.execute_batch("DROP TRIGGER session_contexts_workspace_update;")
1889 .unwrap();
1890 drop(raw);
1891 let writer = open_writer(&path).unwrap();
1892 assert!(!writer.query_row("SELECT EXISTS(SELECT 1 FROM sqlite_schema WHERE name = 'session_contexts_workspace_update')", [], |row| row.get::<_, bool>(0)).unwrap());
1893 assert_eq!(
1894 writer
1895 .query_row("SELECT value FROM future_feature", [], |row| row
1896 .get::<_, String>(0))
1897 .unwrap(),
1898 "preserve me"
1899 );
1900 let state = read_schema_state(&writer).unwrap();
1901 assert_eq!(state.revision, SCHEMA_VERSION + 1);
1902 assert_eq!(state.minimum_compatible, Some(SCHEMA_VERSION));
1903 }
1904
1905 #[test]
1906 fn invalid_compatibility_metadata_refuses_readers_and_writers() {
1907 for alteration in [
1908 "DROP TABLE schema_compatibility",
1909 "DELETE FROM schema_compatibility",
1910 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET minimum_compatible_version = 0",
1911 "UPDATE schema_compatibility SET minimum_compatible_version = 99999",
1912 "PRAGMA ignore_check_constraints = ON; UPDATE schema_compatibility SET singleton = 2",
1913 "PRAGMA ignore_check_constraints = ON; INSERT INTO schema_compatibility VALUES (2, 30)",
1914 "DROP TABLE schema_compatibility; CREATE TABLE schema_compatibility(singleton, minimum_compatible_version); INSERT INTO schema_compatibility VALUES (1, 'invalid')",
1915 "DELETE FROM schema_migrations WHERE version = (SELECT max(version) FROM schema_migrations)",
1916 ] {
1917 for future in [false, true] {
1918 let directory = tempfile::tempdir().unwrap();
1919 let path = directory.path().join("mj.sqlite3");
1920 drop(open_writer(&path).unwrap());
1921 if future {
1922 advance_test_schema(&path, SCHEMA_VERSION + 1, SCHEMA_VERSION);
1923 }
1924 let raw = Connection::open(&path).unwrap();
1925 raw.execute_batch(alteration).unwrap();
1926 let before: i64 = raw
1927 .query_row("PRAGMA schema_version", [], |row| row.get(0))
1928 .unwrap();
1929 for error in [
1931 open_reader_strict(&path).unwrap_err(),
1932 open_writer(&path).unwrap_err(),
1933 ] {
1934 let mismatch = error.downcast_ref::<StoreSchemaMismatch>().unwrap();
1935 assert_eq!(
1936 mismatch.reason,
1937 StoreSchemaMismatchReason::InvalidCompatibilityMetadata,
1938 "{alteration}"
1939 );
1940 }
1941 forget_verified_schema(&path);
1942 assert!(open_writer(&path).is_err(), "{alteration}");
1943 let after: i64 = raw
1944 .query_row("PRAGMA schema_version", [], |row| row.get(0))
1945 .unwrap();
1946 assert_eq!(
1947 before, after,
1948 "a rejected open repaired schema: {alteration}"
1949 );
1950 }
1951 }
1952 }
1953
1954 #[test]
1955 fn a_failed_baseline_leaves_an_empty_store_that_a_retry_creates() {
1956 let directory = tempfile::tempdir().unwrap();
1957 let path = directory.path().join("mj.sqlite3");
1958 let connection = Connection::open(&path).unwrap();
1959 connection
1961 .execute_batch("CREATE TABLE workspaces(conflict TEXT)")
1962 .unwrap();
1963
1964 let error = migrate_schema(&connection).unwrap_err();
1965
1966 assert!(format!("{error:#}").contains("create baseline database schema"));
1967 assert!(
1968 connection.is_autocommit(),
1969 "the failed baseline left a transaction open"
1970 );
1971 assert_eq!(read_schema_state(&connection).unwrap().revision, 0);
1972 let tables: i64 = connection
1973 .query_row(
1974 "SELECT count(*) FROM sqlite_schema WHERE type = 'table'",
1975 [],
1976 |row| row.get(0),
1977 )
1978 .unwrap();
1979 assert_eq!(tables, 1, "only the conflicting table remains");
1980
1981 connection.execute_batch("DROP TABLE workspaces").unwrap();
1982 drop(connection);
1983 let writer = open_writer(&path).unwrap();
1984 let state = read_schema_state(&writer).unwrap();
1985 assert_eq!(state.revision, SCHEMA_VERSION);
1986 assert_eq!(state.minimum_compatible, Some(MINIMUM_COMPATIBLE_VERSION));
1987 }
1988
1989 #[test]
1993 fn strict_reader_reports_a_newer_store_without_blaming_the_daemon() {
1994 let directory = tempfile::tempdir().unwrap();
1995 let path = directory.path().join("mj.sqlite3");
1996 drop(open_writer(&path).unwrap());
1997 stamp_schema_version(&path, SCHEMA_VERSION + 1);
1998
1999 let error = open_reader_strict(&path).unwrap_err();
2000
2001 let mismatch = error
2002 .chain()
2003 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2004 .expect("the reader reports the mismatch as a typed cause");
2005 assert_eq!(mismatch.found, SCHEMA_VERSION + 1);
2006 assert_eq!(mismatch.supported, SCHEMA_VERSION);
2007 let message = mismatch.to_string();
2008 assert!(message.contains("upgrade Mjolnir"), "got {message}");
2009 assert!(
2010 !message.contains("start the Mjolnir daemon"),
2011 "got {message}"
2012 );
2013 }
2014
2015 #[test]
2018 fn strict_reader_keeps_the_migrate_advice_when_the_store_is_behind() {
2019 let directory = tempfile::tempdir().unwrap();
2020 let path = directory.path().join("mj.sqlite3");
2021 drop(open_writer(&path).unwrap());
2022 let raw = Connection::open(&path).unwrap();
2023 raw.execute_batch(&format!(
2024 "UPDATE schema_compatibility SET minimum_compatible_version = {0};
2025 DELETE FROM schema_migrations WHERE version > {0};
2026 INSERT OR IGNORE INTO schema_migrations(version, applied_at) VALUES ({0}, 'test');
2027 PRAGMA user_version = {0};",
2028 SCHEMA_VERSION - 1
2029 ))
2030 .unwrap();
2031 drop(raw);
2032
2033 let error = open_reader_strict(&path).unwrap_err();
2034
2035 let mismatch = error
2036 .chain()
2037 .find_map(|cause| cause.downcast_ref::<StoreSchemaMismatch>())
2038 .expect("the reader reports the mismatch as a typed cause");
2039 assert_eq!(
2040 mismatch.to_string(),
2041 format!(
2042 "Mjolnir database schema {} is not the supported schema {SCHEMA_VERSION}; \
2043 start the Mjolnir daemon to migrate it",
2044 SCHEMA_VERSION - 1
2045 )
2046 );
2047 }
2048
2049 #[test]
2050 fn strict_reader_rejects_mutation() {
2051 let directory = tempfile::tempdir().unwrap();
2052 let path = directory.path().join("mj.sqlite3");
2053 drop(open_writer(&path).unwrap());
2054
2055 let reader = open_reader_strict(&path).unwrap();
2056 let error = reader
2057 .execute("CREATE TABLE forbidden(value TEXT)", [])
2058 .unwrap_err();
2059 assert!(
2060 matches!(
2061 error.sqlite_error_code(),
2062 Some(rusqlite::ErrorCode::ReadOnly)
2063 ),
2064 "unexpected mutation error: {error}"
2065 );
2066 }
2067}