1use rusqlite::{Connection, OpenFlags, TransactionBehavior};
2
3pub mod lifecycle;
4pub use lifecycle::{
5 connection_snapshot, SqliteConnectionSnapshot, SqliteStore, SqliteStoreCount, TrackedConnection,
6};
7use std::fmt;
8use std::fs;
9use std::path::Path;
10use std::time::Duration;
11
12pub mod backups;
13pub mod bash_tasks;
14pub mod bash_watches;
15pub mod compression_events;
16pub mod github_read_cache;
17pub mod removal;
18pub mod standing_roots;
19pub mod state;
20
21pub const CURRENT_SCHEMA_VERSION: u32 = 10;
22
23const MIGRATION_V10: &str = r#"
24CREATE TABLE IF NOT EXISTS compression_event_rollups (
25 harness TEXT NOT NULL,
26 project_key TEXT NOT NULL,
27 session_is_null INTEGER NOT NULL,
28 session_id TEXT NOT NULL,
29 events INTEGER NOT NULL,
30 original_tokens INTEGER NOT NULL,
31 compressed_tokens INTEGER NOT NULL,
32 PRIMARY KEY (harness, project_key, session_is_null, session_id)
33);
34CREATE INDEX IF NOT EXISTS idx_compression_created ON compression_events(created_at, id);
35CREATE TABLE IF NOT EXISTS compression_retention_cursor (
36 singleton INTEGER PRIMARY KEY CHECK (singleton = 1),
37 created_at INTEGER NOT NULL,
38 event_id INTEGER NOT NULL
39);
40"#;
41
42const MIGRATION_V1: &str = r#"
43CREATE TABLE IF NOT EXISTS schema_version (
44 version INTEGER NOT NULL PRIMARY KEY
45);
46
47CREATE TABLE IF NOT EXISTS bash_tasks (
48 harness TEXT NOT NULL,
49 session_id TEXT NOT NULL,
50 task_id TEXT NOT NULL,
51 project_key TEXT NOT NULL,
52 command TEXT NOT NULL,
53 cwd TEXT NOT NULL,
54 status TEXT NOT NULL,
55 exit_code INTEGER,
56 pid INTEGER,
57 pgid INTEGER,
58 started_at INTEGER NOT NULL,
59 completed_at INTEGER,
60 stdout_path TEXT,
61 stderr_path TEXT,
62 compressed INTEGER NOT NULL DEFAULT 1,
63 timeout_ms INTEGER,
64 completion_delivered INTEGER NOT NULL DEFAULT 0,
65 output_bytes INTEGER,
66 metadata TEXT,
67 PRIMARY KEY (harness, session_id, task_id)
68);
69CREATE INDEX IF NOT EXISTS idx_bash_tasks_project_key ON bash_tasks(project_key);
70CREATE INDEX IF NOT EXISTS idx_bash_tasks_status ON bash_tasks(status);
71CREATE INDEX IF NOT EXISTS idx_bash_tasks_session_status ON bash_tasks(harness, session_id, status);
72
73CREATE TABLE IF NOT EXISTS compression_events (
74 id INTEGER PRIMARY KEY AUTOINCREMENT,
75 harness TEXT NOT NULL,
76 session_id TEXT,
77 project_key TEXT NOT NULL,
78 tool TEXT NOT NULL,
79 task_id TEXT,
80 command TEXT,
81 compressor TEXT NOT NULL,
82 original_bytes INTEGER NOT NULL,
83 compressed_bytes INTEGER NOT NULL,
84 original_tokens INTEGER NOT NULL,
85 compressed_tokens INTEGER NOT NULL,
86 created_at INTEGER NOT NULL
87);
88CREATE INDEX IF NOT EXISTS idx_compression_session ON compression_events(harness, session_id);
89CREATE INDEX IF NOT EXISTS idx_compression_session_created ON compression_events(harness, session_id, created_at);
90CREATE INDEX IF NOT EXISTS idx_compression_project_key ON compression_events(project_key);
91
92CREATE TABLE IF NOT EXISTS backups (
93 id INTEGER PRIMARY KEY AUTOINCREMENT,
94 backup_id TEXT,
95 harness TEXT NOT NULL,
96 session_id TEXT NOT NULL,
97 project_key TEXT NOT NULL,
98 op_id TEXT,
99 order_blob BLOB NOT NULL,
100 file_path TEXT NOT NULL,
101 path_hash TEXT NOT NULL,
102 backup_path TEXT,
103 kind TEXT NOT NULL,
104 description TEXT,
105 created_at INTEGER NOT NULL,
106 is_tombstone INTEGER NOT NULL DEFAULT 0
107);
108CREATE INDEX IF NOT EXISTS idx_backups_session_path ON backups(harness, session_id, path_hash);
109CREATE INDEX IF NOT EXISTS idx_backups_session_op ON backups(harness, session_id, op_id) WHERE op_id IS NOT NULL;
110CREATE INDEX IF NOT EXISTS idx_backups_session_order ON backups(harness, session_id, order_blob DESC);
111CREATE INDEX IF NOT EXISTS idx_backups_session_path_order ON backups(harness, session_id, path_hash, order_blob DESC);
112
113CREATE TABLE IF NOT EXISTS harness_state (
114 harness TEXT NOT NULL,
115 key TEXT NOT NULL,
116 value TEXT NOT NULL,
117 updated_at INTEGER NOT NULL,
118 PRIMARY KEY (harness, key)
119);
120
121CREATE TABLE IF NOT EXISTS host_state (
122 key TEXT NOT NULL PRIMARY KEY,
123 value TEXT NOT NULL,
124 updated_at INTEGER NOT NULL
125);
126"#;
127
128const MIGRATION_V2: &str = r#"
129DELETE FROM compression_events
130WHERE id NOT IN (
131 SELECT MIN(id)
132 FROM compression_events
133 GROUP BY
134 harness,
135 COALESCE(session_id, char(0)),
136 project_key,
137 tool,
138 COALESCE(task_id, char(0))
139);
140
141CREATE UNIQUE INDEX IF NOT EXISTS idx_compression_event_identity
142ON compression_events (
143 harness,
144 COALESCE(session_id, char(0)),
145 project_key,
146 tool,
147 COALESCE(task_id, char(0))
148);
149"#;
150
151const MIGRATION_V3: &str = r#"
152CREATE INDEX IF NOT EXISTS idx_bash_tasks_project_lookup
153ON bash_tasks (harness, project_key, task_id, started_at DESC);
154"#;
155
156const MIGRATION_V4: &str = r#"
159ALTER TABLE backups ADD COLUMN restore_meta TEXT;
160"#;
161
162const MIGRATION_V5: &str = r#"
165CREATE TABLE IF NOT EXISTS bash_pattern_watches (
166 harness TEXT NOT NULL,
167 session_id TEXT NOT NULL,
168 task_id TEXT NOT NULL,
169 watch_id TEXT NOT NULL,
170 pattern_kind TEXT NOT NULL,
171 pattern TEXT NOT NULL,
172 once INTEGER NOT NULL DEFAULT 1,
173 created_at INTEGER NOT NULL,
174 stdout_offset INTEGER NOT NULL DEFAULT 0,
175 stderr_offset INTEGER NOT NULL DEFAULT 0,
176 pty_offset INTEGER NOT NULL DEFAULT 0,
177 scanning INTEGER NOT NULL DEFAULT 1,
178 pending_match INTEGER NOT NULL DEFAULT 0,
179 match_text TEXT,
180 match_offset INTEGER,
181 match_context TEXT,
182 PRIMARY KEY (harness, session_id, task_id, watch_id)
183);
184CREATE INDEX IF NOT EXISTS idx_bash_pattern_watches_session
185 ON bash_pattern_watches (harness, session_id);
186CREATE INDEX IF NOT EXISTS idx_bash_pattern_watches_task
187 ON bash_pattern_watches (harness, session_id, task_id);
188"#;
189
190const MIGRATION_V6: &str = r#"
193CREATE INDEX IF NOT EXISTS idx_bash_tasks_started_activity
194 ON bash_tasks (started_at, project_key, harness, session_id);
195CREATE INDEX IF NOT EXISTS idx_bash_tasks_non_terminal_pid
196 ON bash_tasks (pid)
197 WHERE status NOT IN ('completed', 'failed', 'killed', 'timed_out');
198CREATE INDEX IF NOT EXISTS idx_backups_created_activity
199 ON backups (created_at, project_key, harness, session_id);
200"#;
201
202const MIGRATION_V7: &str = r#"
205CREATE TABLE IF NOT EXISTS standing_roots (
206 literal_path TEXT NOT NULL PRIMARY KEY,
207 resolved_target TEXT NOT NULL,
208 resolved_git_toplevel TEXT,
209 scoped_relative_path TEXT
210);
211
212CREATE TABLE IF NOT EXISTS standing_root_freshness (
213 literal_path TEXT NOT NULL,
214 index_kind TEXT NOT NULL CHECK (index_kind IN ('search', 'semantic', 'callgraph')),
215 needs_strict_verify INTEGER NOT NULL CHECK (needs_strict_verify IN (0, 1)),
216 strict_verified_at INTEGER,
217 PRIMARY KEY (literal_path, index_kind),
218 FOREIGN KEY (literal_path) REFERENCES standing_roots(literal_path) ON DELETE CASCADE
219);
220CREATE INDEX IF NOT EXISTS idx_standing_root_freshness_needs_verify
221 ON standing_root_freshness (needs_strict_verify, literal_path);
222"#;
223
224const MIGRATION_V8: &str = r#"
227DROP INDEX IF EXISTS idx_bash_tasks_non_terminal_pid;
228CREATE INDEX idx_bash_tasks_non_terminal_pid
229 ON bash_tasks (pid)
230 WHERE status NOT IN ('completed', 'failed', 'killed', 'timed_out', 'fate_unknown');
231"#;
232
233const MIGRATION_V9: &str = r#"
238DROP INDEX IF EXISTS idx_bash_pattern_watches_session;
239DROP INDEX IF EXISTS idx_bash_pattern_watches_task;
240ALTER TABLE bash_pattern_watches RENAME TO bash_pattern_watches_without_task_fk;
241CREATE TABLE bash_pattern_watches (
242 harness TEXT NOT NULL,
243 session_id TEXT NOT NULL,
244 task_id TEXT NOT NULL,
245 watch_id TEXT NOT NULL,
246 pattern_kind TEXT NOT NULL,
247 pattern TEXT NOT NULL,
248 once INTEGER NOT NULL DEFAULT 1,
249 created_at INTEGER NOT NULL,
250 stdout_offset INTEGER NOT NULL DEFAULT 0,
251 stderr_offset INTEGER NOT NULL DEFAULT 0,
252 pty_offset INTEGER NOT NULL DEFAULT 0,
253 scanning INTEGER NOT NULL DEFAULT 1,
254 pending_match INTEGER NOT NULL DEFAULT 0,
255 match_text TEXT,
256 match_offset INTEGER,
257 match_context TEXT,
258 PRIMARY KEY (harness, session_id, task_id, watch_id),
259 FOREIGN KEY (harness, session_id, task_id)
260 REFERENCES bash_tasks (harness, session_id, task_id) ON DELETE CASCADE
261);
262INSERT INTO bash_pattern_watches (
263 harness, session_id, task_id, watch_id, pattern_kind, pattern, once,
264 created_at, stdout_offset, stderr_offset, pty_offset, scanning,
265 pending_match, match_text, match_offset, match_context
266)
267SELECT
268 watch.harness, watch.session_id, watch.task_id, watch.watch_id,
269 watch.pattern_kind, watch.pattern, watch.once, watch.created_at,
270 watch.stdout_offset, watch.stderr_offset, watch.pty_offset, watch.scanning,
271 watch.pending_match, watch.match_text, watch.match_offset, watch.match_context
272FROM bash_pattern_watches_without_task_fk AS watch
273WHERE EXISTS (
274 SELECT 1
275 FROM bash_tasks AS task
276 WHERE task.harness = watch.harness
277 AND task.session_id = watch.session_id
278 AND task.task_id = watch.task_id
279);
280DROP TABLE bash_pattern_watches_without_task_fk;
281CREATE INDEX idx_bash_pattern_watches_session
282 ON bash_pattern_watches (harness, session_id);
283CREATE INDEX idx_bash_pattern_watches_task
284 ON bash_pattern_watches (harness, session_id, task_id);
285"#;
286
287#[derive(Debug)]
288pub enum OpenError {
289 Io(std::io::Error),
290 Sqlite(rusqlite::Error),
291 DowngradeRefused {
292 db_version: u32,
293 supported: u32,
294 },
295 MigrationFailed {
296 from: u32,
297 to: u32,
298 error: rusqlite::Error,
299 },
300}
301
302impl fmt::Display for OpenError {
303 fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
304 match self {
305 OpenError::Io(error) => write!(f, "database I/O error: {error}"),
306 OpenError::Sqlite(error) => write!(f, "sqlite error: {error}"),
307 OpenError::DowngradeRefused {
308 db_version,
309 supported,
310 } => write!(
311 f,
312 "database schema version {db_version} is newer than supported version {supported}"
313 ),
314 OpenError::MigrationFailed { from, to, error } => {
315 write!(f, "database migration {from}->{to} failed: {error}")
316 }
317 }
318 }
319}
320
321impl std::error::Error for OpenError {}
322
323impl From<std::io::Error> for OpenError {
324 fn from(error: std::io::Error) -> Self {
325 OpenError::Io(error)
326 }
327}
328
329impl From<rusqlite::Error> for OpenError {
330 fn from(error: rusqlite::Error) -> Self {
331 OpenError::Sqlite(error)
332 }
333}
334
335pub fn open(path: &Path) -> Result<TrackedConnection, OpenError> {
341 if let Some(parent) = path.parent() {
342 if !parent.as_os_str().is_empty() {
343 fs::create_dir_all(parent)?;
344 }
345 }
346
347 let mut conn = TrackedConnection::open(path, SqliteStore::AftDb)?;
348 apply_pragmas(&conn)?;
349 run_migrations(&mut conn)?;
350 Ok(conn)
351}
352
353pub fn open_readonly(path: &Path) -> Result<TrackedConnection, OpenError> {
358 let conn = TrackedConnection::open_path_with_flags(
359 path,
360 OpenFlags::SQLITE_OPEN_READ_ONLY,
361 SqliteStore::AftDb,
362 )?;
363 conn.busy_timeout(Duration::from_secs(5))?;
364 Ok(conn)
365}
366
367pub fn apply_pragmas(conn: &Connection) -> Result<(), rusqlite::Error> {
369 conn.pragma_update(None, "foreign_keys", "ON")?;
370 conn.pragma_update(None, "busy_timeout", 5000)?;
373 conn.pragma_update(None, "journal_mode", "WAL")?;
374 conn.pragma_update(None, "synchronous", "NORMAL")?;
375 Ok(())
376}
377
378pub fn run_migrations(conn: &mut Connection) -> Result<u32, OpenError> {
383 conn.execute_batch(
384 "CREATE TABLE IF NOT EXISTS schema_version (version INTEGER NOT NULL PRIMARY KEY);",
385 )?;
386
387 let db_version = current_schema_version(conn)?;
388 if db_version == CURRENT_SCHEMA_VERSION {
389 return Ok(db_version);
390 }
391 if db_version > CURRENT_SCHEMA_VERSION {
392 return Err(OpenError::DowngradeRefused {
393 db_version,
394 supported: CURRENT_SCHEMA_VERSION,
395 });
396 }
397
398 for version in (db_version + 1)..=CURRENT_SCHEMA_VERSION {
402 apply_migration(conn, version)?;
403 }
404
405 Ok(CURRENT_SCHEMA_VERSION)
406}
407
408fn current_schema_version(conn: &Connection) -> Result<u32, rusqlite::Error> {
409 conn.query_row(
410 "SELECT COALESCE(MAX(version), 0) FROM schema_version",
411 [],
412 |row| row.get::<_, u32>(0),
413 )
414}
415
416fn apply_migration(conn: &mut Connection, version: u32) -> Result<(), OpenError> {
417 let planned_from = version - 1;
418 let tx = conn
419 .transaction_with_behavior(TransactionBehavior::Immediate)
420 .map_err(|error| OpenError::MigrationFailed {
421 from: planned_from,
422 to: version,
423 error,
424 })?;
425 let db_version = current_schema_version(&tx).map_err(|error| OpenError::MigrationFailed {
426 from: planned_from,
427 to: version,
428 error,
429 })?;
430 if db_version > CURRENT_SCHEMA_VERSION {
431 return Err(OpenError::DowngradeRefused {
432 db_version,
433 supported: CURRENT_SCHEMA_VERSION,
434 });
435 }
436 if db_version >= version {
437 return tx.commit().map_err(|error| OpenError::MigrationFailed {
438 from: planned_from,
439 to: version,
440 error,
441 });
442 }
443
444 let from = db_version;
445 let already_applied =
446 migration_already_applied(&tx, version).map_err(|error| OpenError::MigrationFailed {
447 from,
448 to: version,
449 error,
450 })?;
451 let result = if already_applied {
452 Ok(())
453 } else {
454 apply_migration_statements(&tx, version)
455 }
456 .and_then(|()| {
457 tx.execute("DELETE FROM schema_version", [])?;
458 tx.execute(
459 "INSERT OR REPLACE INTO schema_version (version) VALUES (?1)",
460 [version],
461 )?;
462 tx.commit()
463 });
464
465 result.map_err(|error| OpenError::MigrationFailed {
466 from,
467 to: version,
468 error,
469 })?;
470 if already_applied && from == 9 && version == 10 {
471 log::warn!("aft.db: schema 9 with v10 objects present; recorded 10 (issue #312 recovery)");
472 }
473 Ok(())
474}
475
476fn migration_already_applied(conn: &Connection, version: u32) -> rusqlite::Result<bool> {
477 match version {
478 10 => conn
479 .query_row(
480 "SELECT COUNT(*) FROM sqlite_master
481 WHERE (type = 'table' AND name IN (
482 'compression_event_rollups',
483 'compression_retention_cursor'
484 )) OR (type = 'index' AND name = 'idx_compression_created')",
485 [],
486 |row| row.get::<_, u32>(0),
487 )
488 .map(|object_count| object_count == 3),
489 _ => Ok(false),
490 }
491}
492
493fn apply_migration_statements(conn: &Connection, version: u32) -> rusqlite::Result<()> {
494 match version {
495 1 => conn.execute_batch(MIGRATION_V1),
496 2 => conn.execute_batch(MIGRATION_V2),
497 3 => conn.execute_batch(MIGRATION_V3),
498 4 => apply_migration_v4(conn),
499 5 => conn.execute_batch(MIGRATION_V5),
500 6 => conn.execute_batch(MIGRATION_V6),
501 7 => conn.execute_batch(MIGRATION_V7),
502 8 => conn.execute_batch(MIGRATION_V8),
503 9 => conn.execute_batch(MIGRATION_V9),
504 10 => conn.execute_batch(MIGRATION_V10),
505 _ => Ok(()),
506 }
507}
508
509fn apply_migration_v4(conn: &Connection) -> rusqlite::Result<()> {
510 let mut stmt = conn.prepare("PRAGMA table_info(backups)")?;
511 let columns = stmt
512 .query_map([], |row| row.get::<_, String>(1))?
513 .collect::<rusqlite::Result<Vec<_>>>()?;
514 drop(stmt);
515
516 if !columns.iter().any(|column| column == "restore_meta") {
517 conn.execute_batch(MIGRATION_V4)?;
518 }
519 Ok(())
520}
521
522#[cfg(test)]
523mod tests {
524 use super::*;
525 use rusqlite::params;
526 use tempfile::tempdir;
527
528 const EXPECTED_TABLES: &[&str] = &[
529 "schema_version",
530 "bash_tasks",
531 "bash_pattern_watches",
532 "compression_events",
533 "backups",
534 "harness_state",
535 "host_state",
536 "standing_roots",
537 "standing_root_freshness",
538 ];
539
540 const EXPECTED_INDEXES: &[&str] = &[
541 "idx_bash_tasks_project_key",
542 "idx_bash_tasks_status",
543 "idx_bash_tasks_session_status",
544 "idx_bash_tasks_project_lookup",
545 "idx_bash_pattern_watches_session",
546 "idx_bash_pattern_watches_task",
547 "idx_bash_tasks_started_activity",
548 "idx_bash_tasks_non_terminal_pid",
549 "idx_backups_created_activity",
550 "idx_compression_session",
551 "idx_compression_session_created",
552 "idx_compression_project_key",
553 "idx_compression_event_identity",
554 "idx_backups_session_path",
555 "idx_backups_session_op",
556 "idx_backups_session_order",
557 "idx_backups_session_path_order",
558 "idx_standing_root_freshness_needs_verify",
559 ];
560
561 #[test]
562 fn open_fresh_db_creates_all_tables() {
563 let dir = tempdir().unwrap();
564 let conn = open(&dir.path().join("aft.db")).unwrap();
565
566 let tables = sqlite_names(&conn, "table");
567 for table in EXPECTED_TABLES {
568 assert!(tables.contains(&table.to_string()), "missing table {table}");
569 }
570 }
571
572 #[test]
573 fn open_fresh_db_creates_all_indexes() {
574 let dir = tempdir().unwrap();
575 let conn = open(&dir.path().join("aft.db")).unwrap();
576
577 let indexes = sqlite_names(&conn, "index");
578 for index in EXPECTED_INDEXES {
579 assert!(
580 indexes.contains(&index.to_string()),
581 "missing index {index}"
582 );
583 }
584 }
585
586 #[test]
587 fn open_existing_db_is_idempotent() {
588 let dir = tempdir().unwrap();
589 let path = dir.path().join("aft.db");
590
591 let conn = open(&path).unwrap();
592 let first_version = schema_version(&conn);
593 drop(conn);
594
595 let conn = open(&path).unwrap();
596 assert_eq!(schema_version(&conn), first_version);
597 }
598
599 #[test]
600 fn pragmas_applied_correctly() {
601 let dir = tempdir().unwrap();
602 let conn = open(&dir.path().join("aft.db")).unwrap();
603
604 let foreign_keys: i64 = conn
605 .query_row("PRAGMA foreign_keys", [], |row| row.get(0))
606 .unwrap();
607 let journal_mode: String = conn
608 .query_row("PRAGMA journal_mode", [], |row| row.get(0))
609 .unwrap();
610 let busy_timeout: i64 = conn
611 .query_row("PRAGMA busy_timeout", [], |row| row.get(0))
612 .unwrap();
613 let synchronous: i64 = conn
614 .query_row("PRAGMA synchronous", [], |row| row.get(0))
615 .unwrap();
616
617 assert_eq!(foreign_keys, 1);
618 assert_eq!(journal_mode, "wal");
619 assert_eq!(busy_timeout, 5000);
620 assert_eq!(synchronous, 1);
621 }
622
623 #[test]
624 fn downgrade_refused() {
625 let dir = tempdir().unwrap();
626 let path = dir.path().join("aft.db");
627 let conn = open(&path).unwrap();
628 conn.execute("INSERT OR REPLACE INTO schema_version VALUES (999)", [])
629 .unwrap();
630 drop(conn);
631
632 match open(&path).unwrap_err() {
633 OpenError::DowngradeRefused {
634 db_version,
635 supported,
636 } => {
637 assert_eq!(db_version, 999);
638 assert_eq!(supported, CURRENT_SCHEMA_VERSION);
639 }
640 error => panic!("expected downgrade refusal, got {error:?}"),
641 }
642 }
643
644 #[test]
645 fn migration_runner_advances_version() {
646 let dir = tempdir().unwrap();
647 let conn = open(&dir.path().join("aft.db")).unwrap();
648
649 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
650 }
651
652 #[test]
653 fn migration_v10_probe_recognizes_the_reported_version_nine_wedge() {
654 let mut conn = Connection::open_in_memory().unwrap();
655 run_migrations(&mut conn).unwrap();
656 conn.execute_batch(MIGRATION_V9).unwrap();
657 conn.execute_batch(
658 "DELETE FROM schema_version;
659 INSERT INTO schema_version (version) VALUES (9);",
660 )
661 .unwrap();
662
663 assert_eq!(schema_version(&conn), 9);
664 assert!(migration_already_applied(&conn, 10).unwrap());
665 }
666
667 #[test]
668 fn every_migration_is_safe_to_apply_twice() {
669 let conn = Connection::open_in_memory().unwrap();
670
671 for version in 1..=CURRENT_SCHEMA_VERSION {
672 apply_migration_statements(&conn, version).unwrap_or_else(|error| {
673 panic!("migration V{version} failed on its first application: {error}")
674 });
675 apply_migration_statements(&conn, version).unwrap_or_else(|error| {
676 panic!("migration V{version} failed when applied a second time: {error}")
677 });
678 }
679 }
680
681 #[test]
682 fn migration_v2_deduplicates_compression_events_and_adds_unique_index() {
683 let dir = tempdir().unwrap();
684 let path = dir.path().join("aft.db");
685
686 let conn = Connection::open(&path).unwrap();
687 conn.execute_batch(MIGRATION_V1).unwrap();
688 conn.execute("DELETE FROM schema_version", []).unwrap();
689 conn.execute("INSERT INTO schema_version (version) VALUES (1)", [])
690 .unwrap();
691 insert_compression_event(
692 &conn,
693 1,
694 "opencode",
695 Some("session-1"),
696 "project-key",
697 "bash",
698 Some("task-1"),
699 )
700 .unwrap();
701 insert_compression_event(
702 &conn,
703 2,
704 "opencode",
705 Some("session-1"),
706 "project-key",
707 "bash",
708 Some("task-1"),
709 )
710 .unwrap();
711 insert_compression_event(&conn, 3, "opencode", None, "project-key", "bash", None).unwrap();
712 drop(conn);
713
714 let conn = open(&path).unwrap();
715
716 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
717 let ids = compression_event_ids(&conn);
718 assert_eq!(ids, vec![1, 3]);
719 let indexes = sqlite_names(&conn, "index");
720 assert!(
721 indexes.contains(&"idx_compression_event_identity".to_string()),
722 "missing v2 unique compression event identity index"
723 );
724 assert_unique_constraint(insert_compression_event(
725 &conn,
726 4,
727 "opencode",
728 Some("session-1"),
729 "project-key",
730 "bash",
731 Some("task-1"),
732 ));
733 }
734
735 #[test]
736 fn migration_v3_upgrades_existing_v2_database() {
737 let dir = tempdir().unwrap();
738 let path = dir.path().join("aft.db");
739 let conn = Connection::open(&path).unwrap();
740 conn.execute_batch(MIGRATION_V1).unwrap();
741 conn.execute_batch(MIGRATION_V2).unwrap();
742 conn.execute("DELETE FROM schema_version", []).unwrap();
743 conn.execute("INSERT INTO schema_version (version) VALUES (2)", [])
744 .unwrap();
745 drop(conn);
746
747 let conn = open(&path).unwrap();
748
749 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
752 assert!(sqlite_names(&conn, "index").contains(&"idx_bash_tasks_project_lookup".to_string()));
753 }
754
755 #[test]
756 fn bash_task_project_lookup_uses_composite_filter_and_order_index() {
757 let dir = tempdir().unwrap();
758 let conn = open(&dir.path().join("aft.db")).unwrap();
759 let mut statement = conn
760 .prepare(
761 "EXPLAIN QUERY PLAN
762 SELECT harness, session_id, task_id, project_key, command, cwd, status,
763 exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
764 compressed, timeout_ms, completion_delivered, output_bytes, metadata
765 FROM bash_tasks
766 WHERE harness = ?1 AND project_key = ?2 AND task_id = ?3
767 ORDER BY started_at DESC
768 LIMIT 1",
769 )
770 .unwrap();
771 let plan = statement
772 .query_map(params!["opencode", "project-key", "bash-task"], |row| {
773 row.get::<_, String>(3)
774 })
775 .unwrap()
776 .collect::<Result<Vec<_>, _>>()
777 .unwrap();
778
779 assert!(
780 plan.iter()
781 .any(|detail| detail.contains("idx_bash_tasks_project_lookup")),
782 "lookup plan did not use the composite index: {plan:?}"
783 );
784 assert!(
785 plan.iter()
786 .all(|detail| !detail.contains("USE TEMP B-TREE FOR ORDER BY")),
787 "lookup plan still sorts into a temporary B-tree: {plan:?}"
788 );
789 }
790
791 #[test]
792 fn migration_v4_adds_restore_metadata_to_v2_and_v3_databases() {
793 for initial_version in [2, 3] {
794 let dir = tempdir().unwrap();
795 let path = dir.path().join(format!("aft-v{initial_version}.db"));
796 let conn = Connection::open(&path).unwrap();
797 conn.execute_batch(MIGRATION_V1).unwrap();
798 conn.execute_batch(MIGRATION_V2).unwrap();
799 conn.execute("DELETE FROM schema_version", []).unwrap();
800 conn.execute(
801 "INSERT INTO schema_version (version) VALUES (?1)",
802 [initial_version],
803 )
804 .unwrap();
805 insert_backup(&conn, "legacy", &order_blob(1)).unwrap();
806 drop(conn);
807
808 let conn = open(&path).unwrap();
809
810 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
811 assert!(table_columns(&conn, "backups").contains(&"restore_meta".to_string()));
812 let restore_meta: Option<String> = conn
813 .query_row(
814 "SELECT restore_meta FROM backups WHERE backup_id = 'legacy'",
815 [],
816 |row| row.get(0),
817 )
818 .unwrap();
819 assert_eq!(restore_meta, None, "legacy rows stay nullable");
820 }
821 }
822
823 #[test]
824 fn migration_v4_is_idempotent_when_column_already_exists() {
825 let dir = tempdir().unwrap();
826 let path = dir.path().join("aft.db");
827 let conn = Connection::open(&path).unwrap();
828 conn.execute_batch(MIGRATION_V1).unwrap();
829 conn.execute_batch(MIGRATION_V2).unwrap();
830 conn.execute_batch(MIGRATION_V4).unwrap();
831 conn.execute("DELETE FROM schema_version", []).unwrap();
832 conn.execute("INSERT INTO schema_version (version) VALUES (3)", [])
833 .unwrap();
834 drop(conn);
835
836 let conn = open(&path).unwrap();
837
838 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
839 assert_eq!(
840 table_columns(&conn, "backups")
841 .iter()
842 .filter(|column| column.as_str() == "restore_meta")
843 .count(),
844 1
845 );
846 }
847
848 #[test]
849 fn migration_v5_adds_bash_pattern_watches_table() {
850 let dir = tempdir().unwrap();
851 let path = dir.path().join("aft.db");
852 let conn = Connection::open(&path).unwrap();
853 conn.execute_batch(MIGRATION_V1).unwrap();
854 conn.execute_batch(MIGRATION_V2).unwrap();
855 conn.execute_batch(MIGRATION_V3).unwrap();
856 conn.execute_batch(MIGRATION_V4).unwrap();
857 conn.execute("DELETE FROM schema_version", []).unwrap();
858 conn.execute("INSERT INTO schema_version (version) VALUES (4)", [])
859 .unwrap();
860 drop(conn);
861
862 let conn = open(&path).unwrap();
863
864 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
865 assert!(sqlite_names(&conn, "table").contains(&"bash_pattern_watches".to_string()));
866 assert!(sqlite_names(&conn, "index").contains(&"idx_bash_pattern_watches_task".to_string()));
867 }
868
869 #[test]
870 fn migration_v9_adds_task_cascade_and_drops_existing_orphan_watches() {
871 let dir = tempdir().unwrap();
872 let path = dir.path().join("aft.db");
873 let conn = Connection::open(&path).unwrap();
874 for migration in [
875 MIGRATION_V1,
876 MIGRATION_V2,
877 MIGRATION_V3,
878 MIGRATION_V4,
879 MIGRATION_V5,
880 MIGRATION_V6,
881 MIGRATION_V7,
882 MIGRATION_V8,
883 ] {
884 conn.execute_batch(migration).unwrap();
885 }
886 insert_bash_task(&conn, "pi", "session", "bash-attached").unwrap();
887 insert_bash_pattern_watch(&conn, "pi", "session", "bash-attached", "watch-attached")
888 .unwrap();
889 insert_bash_pattern_watch(&conn, "pi", "session", "bash-orphan", "watch-orphan").unwrap();
890 conn.execute("DELETE FROM schema_version", []).unwrap();
891 conn.execute("INSERT INTO schema_version (version) VALUES (8)", [])
892 .unwrap();
893 drop(conn);
894
895 let conn = open(&path).unwrap();
896 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
897 let attached: i64 = conn
898 .query_row(
899 "SELECT COUNT(*) FROM bash_pattern_watches WHERE task_id = 'bash-attached'",
900 [],
901 |row| row.get(0),
902 )
903 .unwrap();
904 let orphaned: i64 = conn
905 .query_row(
906 "SELECT COUNT(*) FROM bash_pattern_watches WHERE task_id = 'bash-orphan'",
907 [],
908 |row| row.get(0),
909 )
910 .unwrap();
911 assert_eq!(attached, 1, "migration lost a watch with a live task row");
912 assert_eq!(
913 orphaned, 0,
914 "migration retained a pre-existing orphan watch"
915 );
916
917 conn.execute(
918 "DELETE FROM bash_tasks WHERE harness = 'pi' AND session_id = 'session' AND task_id = 'bash-attached'",
919 [],
920 )
921 .unwrap();
922 let watches: i64 = conn
923 .query_row("SELECT COUNT(*) FROM bash_pattern_watches", [], |row| {
924 row.get(0)
925 })
926 .unwrap();
927 assert_eq!(watches, 0, "task deletion did not cascade to its watch");
928 }
929
930 #[test]
931 fn migration_v7_adds_machine_scoped_standing_root_tables() {
932 let dir = tempdir().unwrap();
933 let path = dir.path().join("aft.db");
934 let conn = Connection::open(&path).unwrap();
935 conn.execute_batch(MIGRATION_V1).unwrap();
936 conn.execute_batch(MIGRATION_V2).unwrap();
937 conn.execute_batch(MIGRATION_V3).unwrap();
938 conn.execute_batch(MIGRATION_V4).unwrap();
939 conn.execute_batch(MIGRATION_V5).unwrap();
940 conn.execute_batch(MIGRATION_V6).unwrap();
941 conn.execute("DELETE FROM schema_version", []).unwrap();
942 conn.execute("INSERT INTO schema_version (version) VALUES (6)", [])
943 .unwrap();
944 drop(conn);
945
946 let conn = open(&path).unwrap();
947 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
948 assert!(sqlite_names(&conn, "table").contains(&"standing_roots".to_string()));
949 assert!(sqlite_names(&conn, "table").contains(&"standing_root_freshness".to_string()));
950 assert!(sqlite_names(&conn, "index")
951 .contains(&"idx_standing_root_freshness_needs_verify".to_string()));
952 }
953
954 #[test]
955 fn migration_v6_adds_removal_health_indexes() {
956 let dir = tempdir().unwrap();
957 let path = dir.path().join("aft.db");
958 let conn = Connection::open(&path).unwrap();
959 conn.execute_batch(MIGRATION_V1).unwrap();
960 conn.execute_batch(MIGRATION_V2).unwrap();
961 conn.execute_batch(MIGRATION_V3).unwrap();
962 conn.execute_batch(MIGRATION_V4).unwrap();
963 conn.execute_batch(MIGRATION_V5).unwrap();
964 conn.execute("DELETE FROM schema_version", []).unwrap();
965 conn.execute("INSERT INTO schema_version (version) VALUES (5)", [])
966 .unwrap();
967 drop(conn);
968
969 let conn = open(&path).unwrap();
970
971 let indexes = sqlite_names(&conn, "index");
972 for index in [
973 "idx_bash_tasks_started_activity",
974 "idx_bash_tasks_non_terminal_pid",
975 "idx_backups_created_activity",
976 ] {
977 assert!(
978 indexes.contains(&index.to_string()),
979 "missing v6 index {index}"
980 );
981 }
982 }
983
984 #[test]
985 fn open_readonly_does_not_create_a_missing_database() {
986 let dir = tempdir().unwrap();
987 let path = dir.path().join("missing-aft.db");
988
989 assert!(open_readonly(&path).is_err());
990 assert!(!path.exists());
991 }
992
993 #[test]
994 fn migration_runner_no_op_when_current() {
995 let dir = tempdir().unwrap();
996 let path = dir.path().join("aft.db");
997
998 let conn = open(&path).unwrap();
999 assert_eq!(schema_version_row_count(&conn), 1);
1000 drop(conn);
1001
1002 let conn = open(&path).unwrap();
1003 assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
1004 assert_eq!(schema_version_row_count(&conn), 1);
1005 }
1006
1007 #[test]
1008 fn harness_state_compound_pk_works() {
1009 let dir = tempdir().unwrap();
1010 let conn = open(&dir.path().join("aft.db")).unwrap();
1011
1012 conn.execute(
1013 "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
1014 params!["opencode", "warned_tools", "{}", 1_i64],
1015 )
1016 .unwrap();
1017 let duplicate = conn.execute(
1018 "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
1019 params!["opencode", "warned_tools", "{}", 2_i64],
1020 );
1021 assert_unique_constraint(duplicate);
1022
1023 conn.execute(
1024 "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
1025 params!["pi", "warned_tools", "{}", 3_i64],
1026 )
1027 .unwrap();
1028 }
1029
1030 #[test]
1031 fn host_state_simple_pk_works() {
1032 let dir = tempdir().unwrap();
1033 let conn = open(&dir.path().join("aft.db")).unwrap();
1034
1035 conn.execute(
1036 "INSERT INTO host_state (key, value, updated_at) VALUES (?1, ?2, ?3)",
1037 params!["trusted_filter_projects", "[]", 1_i64],
1038 )
1039 .unwrap();
1040 let duplicate = conn.execute(
1041 "INSERT INTO host_state (key, value, updated_at) VALUES (?1, ?2, ?3)",
1042 params!["trusted_filter_projects", "[]", 2_i64],
1043 );
1044 assert_unique_constraint(duplicate);
1045 }
1046
1047 #[test]
1048 fn bash_tasks_compound_pk_works() {
1049 let dir = tempdir().unwrap();
1050 let conn = open(&dir.path().join("aft.db")).unwrap();
1051
1052 insert_bash_task(&conn, "opencode", "session-1", "bash-12345678").unwrap();
1053 let duplicate = insert_bash_task(&conn, "opencode", "session-1", "bash-12345678");
1054 assert_unique_constraint(duplicate);
1055
1056 insert_bash_task(&conn, "pi", "session-1", "bash-12345678").unwrap();
1057 }
1058
1059 #[test]
1060 fn backups_order_blob_sort() {
1061 let dir = tempdir().unwrap();
1062 let conn = open(&dir.path().join("aft.db")).unwrap();
1063
1064 let one = order_blob(1);
1065 let two = order_blob(2);
1066 let max = [0xFF; 16];
1067
1068 insert_backup(&conn, "one", &one).unwrap();
1069 insert_backup(&conn, "two", &two).unwrap();
1070 insert_backup(&conn, "max", &max).unwrap();
1071
1072 assert_eq!(backup_ids_ordered(&conn, "ASC"), vec!["one", "two", "max"]);
1073 assert_eq!(backup_ids_ordered(&conn, "DESC"), vec!["max", "two", "one"]);
1074 }
1075
1076 fn sqlite_names(conn: &Connection, kind: &str) -> Vec<String> {
1077 let sql = match kind {
1078 "table" => "SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name",
1079 "index" => "SELECT name FROM sqlite_master WHERE type='index' AND name NOT LIKE 'sqlite_%' ORDER BY name",
1080 _ => panic!("unsupported sqlite_master kind: {kind}"),
1081 };
1082 let mut stmt = conn.prepare(sql).unwrap();
1083 stmt.query_map([], |row| row.get::<_, String>(0))
1084 .unwrap()
1085 .collect::<Result<Vec<_>, _>>()
1086 .unwrap()
1087 }
1088
1089 fn table_columns(conn: &Connection, table: &str) -> Vec<String> {
1090 let mut stmt = conn
1091 .prepare(&format!("PRAGMA table_info({table})"))
1092 .unwrap();
1093 stmt.query_map([], |row| row.get::<_, String>(1))
1094 .unwrap()
1095 .collect::<Result<Vec<_>, _>>()
1096 .unwrap()
1097 }
1098
1099 fn schema_version(conn: &Connection) -> u32 {
1100 conn.query_row("SELECT version FROM schema_version", [], |row| row.get(0))
1101 .unwrap()
1102 }
1103
1104 fn schema_version_row_count(conn: &Connection) -> i64 {
1105 conn.query_row("SELECT COUNT(*) FROM schema_version", [], |row| row.get(0))
1106 .unwrap()
1107 }
1108
1109 fn assert_unique_constraint(result: rusqlite::Result<usize>) {
1110 let error = result.expect_err("expected a unique constraint violation");
1111 assert!(
1112 error.to_string().contains("UNIQUE constraint failed"),
1113 "expected UNIQUE constraint failure, got {error}"
1114 );
1115 }
1116
1117 fn insert_bash_task(
1118 conn: &Connection,
1119 harness: &str,
1120 session_id: &str,
1121 task_id: &str,
1122 ) -> rusqlite::Result<usize> {
1123 conn.execute(
1124 "INSERT INTO bash_tasks (
1125 harness, session_id, task_id, project_key, command, cwd, status, started_at
1126 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
1127 params![
1128 harness,
1129 session_id,
1130 task_id,
1131 "project-key",
1132 "echo ok",
1133 "/tmp",
1134 "running",
1135 1_i64
1136 ],
1137 )
1138 }
1139
1140 fn insert_bash_pattern_watch(
1141 conn: &Connection,
1142 harness: &str,
1143 session_id: &str,
1144 task_id: &str,
1145 watch_id: &str,
1146 ) -> rusqlite::Result<usize> {
1147 conn.execute(
1148 "INSERT INTO bash_pattern_watches (
1149 harness, session_id, task_id, watch_id, pattern_kind, pattern, created_at
1150 ) VALUES (?1, ?2, ?3, ?4, 'substring', 'needle', 1)",
1151 params![harness, session_id, task_id, watch_id],
1152 )
1153 }
1154
1155 fn insert_compression_event(
1156 conn: &Connection,
1157 id: i64,
1158 harness: &str,
1159 session_id: Option<&str>,
1160 project_key: &str,
1161 tool: &str,
1162 task_id: Option<&str>,
1163 ) -> rusqlite::Result<usize> {
1164 conn.execute(
1165 "INSERT INTO compression_events (
1166 id, harness, session_id, project_key, tool, task_id, command, compressor,
1167 original_bytes, compressed_bytes, original_tokens, compressed_tokens, created_at
1168 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
1169 params![
1170 id,
1171 harness,
1172 session_id,
1173 project_key,
1174 tool,
1175 task_id,
1176 "echo ok",
1177 "test-compressor",
1178 100_i64,
1179 50_i64,
1180 20_i64,
1181 10_i64,
1182 id
1183 ],
1184 )
1185 }
1186
1187 fn compression_event_ids(conn: &Connection) -> Vec<i64> {
1188 let mut stmt = conn
1189 .prepare("SELECT id FROM compression_events ORDER BY id")
1190 .unwrap();
1191 stmt.query_map([], |row| row.get::<_, i64>(0))
1192 .unwrap()
1193 .collect::<Result<Vec<_>, _>>()
1194 .unwrap()
1195 }
1196
1197 fn insert_backup(
1198 conn: &Connection,
1199 backup_id: &str,
1200 order_blob: &[u8],
1201 ) -> rusqlite::Result<usize> {
1202 conn.execute(
1203 "INSERT INTO backups (
1204 backup_id, harness, session_id, project_key, order_blob, file_path,
1205 path_hash, kind, created_at
1206 ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1207 params![
1208 backup_id,
1209 "opencode",
1210 "session-1",
1211 "project-key",
1212 order_blob,
1213 "/tmp/file.txt",
1214 "path-hash",
1215 "content",
1216 1_i64
1217 ],
1218 )
1219 }
1220
1221 fn order_blob(value: u128) -> [u8; 16] {
1222 value.to_be_bytes()
1223 }
1224
1225 fn backup_ids_ordered(conn: &Connection, direction: &str) -> Vec<String> {
1226 let sql = match direction {
1227 "ASC" => "SELECT backup_id FROM backups ORDER BY order_blob ASC",
1228 "DESC" => "SELECT backup_id FROM backups ORDER BY order_blob DESC",
1229 _ => panic!("unsupported order direction: {direction}"),
1230 };
1231 let mut stmt = conn.prepare(sql).unwrap();
1232 stmt.query_map([], |row| row.get::<_, String>(0))
1233 .unwrap()
1234 .collect::<Result<Vec<_>, _>>()
1235 .unwrap()
1236 }
1237}