Skip to main content

aft/db/
mod.rs

1use rusqlite::{Connection, TransactionBehavior};
2use std::fmt;
3use std::fs;
4use std::path::Path;
5
6pub mod backups;
7pub mod bash_tasks;
8pub mod bash_watches;
9pub mod compression_events;
10pub mod state;
11
12pub const CURRENT_SCHEMA_VERSION: u32 = 5;
13
14const MIGRATION_V1: &str = r#"
15CREATE TABLE IF NOT EXISTS schema_version (
16  version INTEGER NOT NULL PRIMARY KEY
17);
18
19CREATE TABLE IF NOT EXISTS bash_tasks (
20  harness      TEXT NOT NULL,
21  session_id   TEXT NOT NULL,
22  task_id      TEXT NOT NULL,
23  project_key  TEXT NOT NULL,
24  command      TEXT NOT NULL,
25  cwd          TEXT NOT NULL,
26  status       TEXT NOT NULL,
27  exit_code    INTEGER,
28  pid          INTEGER,
29  pgid         INTEGER,
30  started_at   INTEGER NOT NULL,
31  completed_at INTEGER,
32  stdout_path  TEXT,
33  stderr_path  TEXT,
34  compressed   INTEGER NOT NULL DEFAULT 1,
35  timeout_ms   INTEGER,
36  completion_delivered INTEGER NOT NULL DEFAULT 0,
37  output_bytes INTEGER,
38  metadata     TEXT,
39  PRIMARY KEY (harness, session_id, task_id)
40);
41CREATE INDEX IF NOT EXISTS idx_bash_tasks_project_key ON bash_tasks(project_key);
42CREATE INDEX IF NOT EXISTS idx_bash_tasks_status      ON bash_tasks(status);
43CREATE INDEX IF NOT EXISTS idx_bash_tasks_session_status ON bash_tasks(harness, session_id, status);
44
45CREATE TABLE IF NOT EXISTS compression_events (
46  id                INTEGER PRIMARY KEY AUTOINCREMENT,
47  harness           TEXT NOT NULL,
48  session_id        TEXT,
49  project_key       TEXT NOT NULL,
50  tool              TEXT NOT NULL,
51  task_id           TEXT,
52  command           TEXT,
53  compressor        TEXT NOT NULL,
54  original_bytes    INTEGER NOT NULL,
55  compressed_bytes  INTEGER NOT NULL,
56  original_tokens   INTEGER NOT NULL,
57  compressed_tokens INTEGER NOT NULL,
58  created_at        INTEGER NOT NULL
59);
60CREATE INDEX IF NOT EXISTS idx_compression_session         ON compression_events(harness, session_id);
61CREATE INDEX IF NOT EXISTS idx_compression_session_created ON compression_events(harness, session_id, created_at);
62CREATE INDEX IF NOT EXISTS idx_compression_project_key     ON compression_events(project_key);
63
64CREATE TABLE IF NOT EXISTS backups (
65  id            INTEGER PRIMARY KEY AUTOINCREMENT,
66  backup_id     TEXT,
67  harness       TEXT NOT NULL,
68  session_id    TEXT NOT NULL,
69  project_key   TEXT NOT NULL,
70  op_id         TEXT,
71  order_blob    BLOB NOT NULL,
72  file_path     TEXT NOT NULL,
73  path_hash     TEXT NOT NULL,
74  backup_path   TEXT,
75  kind          TEXT NOT NULL,
76  description   TEXT,
77  created_at    INTEGER NOT NULL,
78  is_tombstone  INTEGER NOT NULL DEFAULT 0
79);
80CREATE INDEX IF NOT EXISTS idx_backups_session_path  ON backups(harness, session_id, path_hash);
81CREATE INDEX IF NOT EXISTS idx_backups_session_op    ON backups(harness, session_id, op_id) WHERE op_id IS NOT NULL;
82CREATE INDEX IF NOT EXISTS idx_backups_session_order ON backups(harness, session_id, order_blob DESC);
83CREATE INDEX IF NOT EXISTS idx_backups_session_path_order ON backups(harness, session_id, path_hash, order_blob DESC);
84
85CREATE TABLE IF NOT EXISTS harness_state (
86  harness    TEXT NOT NULL,
87  key        TEXT NOT NULL,
88  value      TEXT NOT NULL,
89  updated_at INTEGER NOT NULL,
90  PRIMARY KEY (harness, key)
91);
92
93CREATE TABLE IF NOT EXISTS host_state (
94  key        TEXT NOT NULL PRIMARY KEY,
95  value      TEXT NOT NULL,
96  updated_at INTEGER NOT NULL
97);
98"#;
99
100const MIGRATION_V2: &str = r#"
101DELETE FROM compression_events
102WHERE id NOT IN (
103  SELECT MIN(id)
104  FROM compression_events
105  GROUP BY
106    harness,
107    COALESCE(session_id, char(0)),
108    project_key,
109    tool,
110    COALESCE(task_id, char(0))
111);
112
113CREATE UNIQUE INDEX IF NOT EXISTS idx_compression_event_identity
114ON compression_events (
115  harness,
116  COALESCE(session_id, char(0)),
117  project_key,
118  tool,
119  COALESCE(task_id, char(0))
120);
121"#;
122
123const MIGRATION_V3: &str = r#"
124CREATE INDEX IF NOT EXISTS idx_bash_tasks_project_lookup
125ON bash_tasks (harness, project_key, task_id, started_at DESC);
126"#;
127
128// V4 adds the restore_meta column to backups (Unix mode / created_dirs /
129// link_target for DB-fallback restores when the meta.json sidecar is gone).
130const MIGRATION_V4: &str = r#"
131ALTER TABLE backups ADD COLUMN restore_meta TEXT;
132"#;
133
134// V5 persists async bash_notify / bash_watch pattern registrations so a
135// bridge/daemon restart can re-arm watches and deliver gap matches.
136const MIGRATION_V5: &str = r#"
137CREATE TABLE IF NOT EXISTS bash_pattern_watches (
138  harness        TEXT NOT NULL,
139  session_id     TEXT NOT NULL,
140  task_id        TEXT NOT NULL,
141  watch_id       TEXT NOT NULL,
142  pattern_kind   TEXT NOT NULL,
143  pattern        TEXT NOT NULL,
144  once           INTEGER NOT NULL DEFAULT 1,
145  created_at     INTEGER NOT NULL,
146  stdout_offset  INTEGER NOT NULL DEFAULT 0,
147  stderr_offset  INTEGER NOT NULL DEFAULT 0,
148  pty_offset     INTEGER NOT NULL DEFAULT 0,
149  scanning       INTEGER NOT NULL DEFAULT 1,
150  pending_match  INTEGER NOT NULL DEFAULT 0,
151  match_text     TEXT,
152  match_offset   INTEGER,
153  match_context  TEXT,
154  PRIMARY KEY (harness, session_id, task_id, watch_id)
155);
156CREATE INDEX IF NOT EXISTS idx_bash_pattern_watches_session
157  ON bash_pattern_watches (harness, session_id);
158CREATE INDEX IF NOT EXISTS idx_bash_pattern_watches_task
159  ON bash_pattern_watches (harness, session_id, task_id);
160"#;
161
162#[derive(Debug)]
163pub enum OpenError {
164    Io(std::io::Error),
165    Sqlite(rusqlite::Error),
166    DowngradeRefused {
167        db_version: u32,
168        supported: u32,
169    },
170    MigrationFailed {
171        from: u32,
172        to: u32,
173        error: rusqlite::Error,
174    },
175}
176
177impl fmt::Display for OpenError {
178    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
179        match self {
180            OpenError::Io(error) => write!(f, "database I/O error: {error}"),
181            OpenError::Sqlite(error) => write!(f, "sqlite error: {error}"),
182            OpenError::DowngradeRefused {
183                db_version,
184                supported,
185            } => write!(
186                f,
187                "database schema version {db_version} is newer than supported version {supported}"
188            ),
189            OpenError::MigrationFailed { from, to, error } => {
190                write!(f, "database migration {from}->{to} failed: {error}")
191            }
192        }
193    }
194}
195
196impl std::error::Error for OpenError {}
197
198impl From<std::io::Error> for OpenError {
199    fn from(error: std::io::Error) -> Self {
200        OpenError::Io(error)
201    }
202}
203
204impl From<rusqlite::Error> for OpenError {
205    fn from(error: rusqlite::Error) -> Self {
206        OpenError::Sqlite(error)
207    }
208}
209
210/// Open or create the AFT SQLite database at the given path.
211///
212/// Applies per-connection PRAGMAs, runs schema migrations from the DB's
213/// current schema version up to [`CURRENT_SCHEMA_VERSION`], and returns the
214/// configured connection.
215pub fn open(path: &Path) -> Result<Connection, OpenError> {
216    if let Some(parent) = path.parent() {
217        if !parent.as_os_str().is_empty() {
218            fs::create_dir_all(parent)?;
219        }
220    }
221
222    let mut conn = Connection::open(path)?;
223    apply_pragmas(&conn)?;
224    run_migrations(&mut conn)?;
225    Ok(conn)
226}
227
228/// Apply the per-connection PRAGMAs required for every AFT SQLite connection.
229pub fn apply_pragmas(conn: &Connection) -> Result<(), rusqlite::Error> {
230    conn.pragma_update(None, "foreign_keys", "ON")?;
231    conn.pragma_update(None, "journal_mode", "WAL")?;
232    conn.pragma_update(None, "busy_timeout", 5000)?;
233    conn.pragma_update(None, "synchronous", "NORMAL")?;
234    Ok(())
235}
236
237/// Run forward-only migrations up to [`CURRENT_SCHEMA_VERSION`].
238///
239/// Returns the post-migration schema version. Refuses to open databases created
240/// by newer AFT versions.
241pub fn run_migrations(conn: &mut Connection) -> Result<u32, OpenError> {
242    conn.execute_batch(
243        "CREATE TABLE IF NOT EXISTS schema_version (version INTEGER NOT NULL PRIMARY KEY);",
244    )?;
245
246    let db_version = current_schema_version(conn)?;
247    if db_version > CURRENT_SCHEMA_VERSION {
248        return Err(OpenError::DowngradeRefused {
249            db_version,
250            supported: CURRENT_SCHEMA_VERSION,
251        });
252    }
253
254    for version in (db_version + 1)..=CURRENT_SCHEMA_VERSION {
255        apply_migration(conn, version)?;
256    }
257
258    Ok(current_schema_version(conn)?)
259}
260
261fn current_schema_version(conn: &Connection) -> Result<u32, rusqlite::Error> {
262    conn.query_row(
263        "SELECT COALESCE(MAX(version), 0) FROM schema_version",
264        [],
265        |row| row.get::<_, u32>(0),
266    )
267}
268
269fn apply_migration(conn: &mut Connection, version: u32) -> Result<(), OpenError> {
270    let from = version - 1;
271    let tx = conn
272        .transaction_with_behavior(TransactionBehavior::Immediate)
273        .map_err(|error| OpenError::MigrationFailed {
274            from,
275            to: version,
276            error,
277        })?;
278
279    let result = match version {
280        1 => tx.execute_batch(MIGRATION_V1),
281        2 => tx.execute_batch(MIGRATION_V2),
282        3 => tx.execute_batch(MIGRATION_V3),
283        4 => apply_migration_v4(&tx),
284        5 => tx.execute_batch(MIGRATION_V5),
285        _ => Ok(()),
286    }
287    .and_then(|()| {
288        tx.execute("DELETE FROM schema_version", [])?;
289        tx.execute(
290            "INSERT OR REPLACE INTO schema_version (version) VALUES (?1)",
291            [version],
292        )?;
293        tx.commit()
294    });
295
296    result.map_err(|error| OpenError::MigrationFailed {
297        from,
298        to: version,
299        error,
300    })
301}
302
303fn apply_migration_v4(conn: &Connection) -> rusqlite::Result<()> {
304    let mut stmt = conn.prepare("PRAGMA table_info(backups)")?;
305    let columns = stmt
306        .query_map([], |row| row.get::<_, String>(1))?
307        .collect::<rusqlite::Result<Vec<_>>>()?;
308    drop(stmt);
309
310    if !columns.iter().any(|column| column == "restore_meta") {
311        conn.execute_batch(MIGRATION_V4)?;
312    }
313    Ok(())
314}
315
316#[cfg(test)]
317mod tests {
318    use super::*;
319    use rusqlite::params;
320    use tempfile::tempdir;
321
322    const EXPECTED_TABLES: &[&str] = &[
323        "schema_version",
324        "bash_tasks",
325        "bash_pattern_watches",
326        "compression_events",
327        "backups",
328        "harness_state",
329        "host_state",
330    ];
331
332    const EXPECTED_INDEXES: &[&str] = &[
333        "idx_bash_tasks_project_key",
334        "idx_bash_tasks_status",
335        "idx_bash_tasks_session_status",
336        "idx_bash_tasks_project_lookup",
337        "idx_bash_pattern_watches_session",
338        "idx_bash_pattern_watches_task",
339        "idx_compression_session",
340        "idx_compression_session_created",
341        "idx_compression_project_key",
342        "idx_compression_event_identity",
343        "idx_backups_session_path",
344        "idx_backups_session_op",
345        "idx_backups_session_order",
346        "idx_backups_session_path_order",
347    ];
348
349    #[test]
350    fn open_fresh_db_creates_all_tables() {
351        let dir = tempdir().unwrap();
352        let conn = open(&dir.path().join("aft.db")).unwrap();
353
354        let tables = sqlite_names(&conn, "table");
355        for table in EXPECTED_TABLES {
356            assert!(tables.contains(&table.to_string()), "missing table {table}");
357        }
358    }
359
360    #[test]
361    fn open_fresh_db_creates_all_indexes() {
362        let dir = tempdir().unwrap();
363        let conn = open(&dir.path().join("aft.db")).unwrap();
364
365        let indexes = sqlite_names(&conn, "index");
366        for index in EXPECTED_INDEXES {
367            assert!(
368                indexes.contains(&index.to_string()),
369                "missing index {index}"
370            );
371        }
372    }
373
374    #[test]
375    fn open_existing_db_is_idempotent() {
376        let dir = tempdir().unwrap();
377        let path = dir.path().join("aft.db");
378
379        let conn = open(&path).unwrap();
380        let first_version = schema_version(&conn);
381        drop(conn);
382
383        let conn = open(&path).unwrap();
384        assert_eq!(schema_version(&conn), first_version);
385    }
386
387    #[test]
388    fn pragmas_applied_correctly() {
389        let dir = tempdir().unwrap();
390        let conn = open(&dir.path().join("aft.db")).unwrap();
391
392        let foreign_keys: i64 = conn
393            .query_row("PRAGMA foreign_keys", [], |row| row.get(0))
394            .unwrap();
395        let journal_mode: String = conn
396            .query_row("PRAGMA journal_mode", [], |row| row.get(0))
397            .unwrap();
398        let busy_timeout: i64 = conn
399            .query_row("PRAGMA busy_timeout", [], |row| row.get(0))
400            .unwrap();
401        let synchronous: i64 = conn
402            .query_row("PRAGMA synchronous", [], |row| row.get(0))
403            .unwrap();
404
405        assert_eq!(foreign_keys, 1);
406        assert_eq!(journal_mode, "wal");
407        assert_eq!(busy_timeout, 5000);
408        assert_eq!(synchronous, 1);
409    }
410
411    #[test]
412    fn downgrade_refused() {
413        let dir = tempdir().unwrap();
414        let path = dir.path().join("aft.db");
415        let conn = open(&path).unwrap();
416        conn.execute("INSERT OR REPLACE INTO schema_version VALUES (999)", [])
417            .unwrap();
418        drop(conn);
419
420        match open(&path).unwrap_err() {
421            OpenError::DowngradeRefused {
422                db_version,
423                supported,
424            } => {
425                assert_eq!(db_version, 999);
426                assert_eq!(supported, CURRENT_SCHEMA_VERSION);
427            }
428            error => panic!("expected downgrade refusal, got {error:?}"),
429        }
430    }
431
432    #[test]
433    fn migration_runner_advances_version() {
434        let dir = tempdir().unwrap();
435        let conn = open(&dir.path().join("aft.db")).unwrap();
436
437        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
438    }
439
440    #[test]
441    fn migration_v2_deduplicates_compression_events_and_adds_unique_index() {
442        let dir = tempdir().unwrap();
443        let path = dir.path().join("aft.db");
444
445        let conn = Connection::open(&path).unwrap();
446        conn.execute_batch(MIGRATION_V1).unwrap();
447        conn.execute("DELETE FROM schema_version", []).unwrap();
448        conn.execute("INSERT INTO schema_version (version) VALUES (1)", [])
449            .unwrap();
450        insert_compression_event(
451            &conn,
452            1,
453            "opencode",
454            Some("session-1"),
455            "project-key",
456            "bash",
457            Some("task-1"),
458        )
459        .unwrap();
460        insert_compression_event(
461            &conn,
462            2,
463            "opencode",
464            Some("session-1"),
465            "project-key",
466            "bash",
467            Some("task-1"),
468        )
469        .unwrap();
470        insert_compression_event(&conn, 3, "opencode", None, "project-key", "bash", None).unwrap();
471        drop(conn);
472
473        let conn = open(&path).unwrap();
474
475        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
476        let ids = compression_event_ids(&conn);
477        assert_eq!(ids, vec![1, 3]);
478        let indexes = sqlite_names(&conn, "index");
479        assert!(
480            indexes.contains(&"idx_compression_event_identity".to_string()),
481            "missing v2 unique compression event identity index"
482        );
483        assert_unique_constraint(insert_compression_event(
484            &conn,
485            4,
486            "opencode",
487            Some("session-1"),
488            "project-key",
489            "bash",
490            Some("task-1"),
491        ));
492    }
493
494    #[test]
495    fn migration_v3_upgrades_existing_v2_database() {
496        let dir = tempdir().unwrap();
497        let path = dir.path().join("aft.db");
498        let conn = Connection::open(&path).unwrap();
499        conn.execute_batch(MIGRATION_V1).unwrap();
500        conn.execute_batch(MIGRATION_V2).unwrap();
501        conn.execute("DELETE FROM schema_version", []).unwrap();
502        conn.execute("INSERT INTO schema_version (version) VALUES (2)", [])
503            .unwrap();
504        drop(conn);
505
506        let conn = open(&path).unwrap();
507
508        // A v2 database migrates all the way to the current version; V3 creates
509        // the bash-task lookup index on the way.
510        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
511        assert!(sqlite_names(&conn, "index").contains(&"idx_bash_tasks_project_lookup".to_string()));
512    }
513
514    #[test]
515    fn bash_task_project_lookup_uses_composite_filter_and_order_index() {
516        let dir = tempdir().unwrap();
517        let conn = open(&dir.path().join("aft.db")).unwrap();
518        let mut statement = conn
519            .prepare(
520                "EXPLAIN QUERY PLAN
521                 SELECT harness, session_id, task_id, project_key, command, cwd, status,
522                        exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
523                        compressed, timeout_ms, completion_delivered, output_bytes, metadata
524                 FROM bash_tasks
525                 WHERE harness = ?1 AND project_key = ?2 AND task_id = ?3
526                 ORDER BY started_at DESC
527                 LIMIT 1",
528            )
529            .unwrap();
530        let plan = statement
531            .query_map(params!["opencode", "project-key", "bash-task"], |row| {
532                row.get::<_, String>(3)
533            })
534            .unwrap()
535            .collect::<Result<Vec<_>, _>>()
536            .unwrap();
537
538        assert!(
539            plan.iter()
540                .any(|detail| detail.contains("idx_bash_tasks_project_lookup")),
541            "lookup plan did not use the composite index: {plan:?}"
542        );
543        assert!(
544            plan.iter()
545                .all(|detail| !detail.contains("USE TEMP B-TREE FOR ORDER BY")),
546            "lookup plan still sorts into a temporary B-tree: {plan:?}"
547        );
548    }
549
550    #[test]
551    fn migration_v4_adds_restore_metadata_to_v2_and_v3_databases() {
552        for initial_version in [2, 3] {
553            let dir = tempdir().unwrap();
554            let path = dir.path().join(format!("aft-v{initial_version}.db"));
555            let conn = Connection::open(&path).unwrap();
556            conn.execute_batch(MIGRATION_V1).unwrap();
557            conn.execute_batch(MIGRATION_V2).unwrap();
558            conn.execute("DELETE FROM schema_version", []).unwrap();
559            conn.execute(
560                "INSERT INTO schema_version (version) VALUES (?1)",
561                [initial_version],
562            )
563            .unwrap();
564            insert_backup(&conn, "legacy", &order_blob(1)).unwrap();
565            drop(conn);
566
567            let conn = open(&path).unwrap();
568
569            assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
570            assert!(table_columns(&conn, "backups").contains(&"restore_meta".to_string()));
571            let restore_meta: Option<String> = conn
572                .query_row(
573                    "SELECT restore_meta FROM backups WHERE backup_id = 'legacy'",
574                    [],
575                    |row| row.get(0),
576                )
577                .unwrap();
578            assert_eq!(restore_meta, None, "legacy rows stay nullable");
579        }
580    }
581
582    #[test]
583    fn migration_v4_is_idempotent_when_column_already_exists() {
584        let dir = tempdir().unwrap();
585        let path = dir.path().join("aft.db");
586        let conn = Connection::open(&path).unwrap();
587        conn.execute_batch(MIGRATION_V1).unwrap();
588        conn.execute_batch(MIGRATION_V2).unwrap();
589        conn.execute_batch(MIGRATION_V4).unwrap();
590        conn.execute("DELETE FROM schema_version", []).unwrap();
591        conn.execute("INSERT INTO schema_version (version) VALUES (3)", [])
592            .unwrap();
593        drop(conn);
594
595        let conn = open(&path).unwrap();
596
597        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
598        assert_eq!(
599            table_columns(&conn, "backups")
600                .iter()
601                .filter(|column| column.as_str() == "restore_meta")
602                .count(),
603            1
604        );
605    }
606
607    #[test]
608    fn migration_v5_adds_bash_pattern_watches_table() {
609        let dir = tempdir().unwrap();
610        let path = dir.path().join("aft.db");
611        let conn = Connection::open(&path).unwrap();
612        conn.execute_batch(MIGRATION_V1).unwrap();
613        conn.execute_batch(MIGRATION_V2).unwrap();
614        conn.execute_batch(MIGRATION_V3).unwrap();
615        conn.execute_batch(MIGRATION_V4).unwrap();
616        conn.execute("DELETE FROM schema_version", []).unwrap();
617        conn.execute("INSERT INTO schema_version (version) VALUES (4)", [])
618            .unwrap();
619        drop(conn);
620
621        let conn = open(&path).unwrap();
622
623        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
624        assert!(sqlite_names(&conn, "table").contains(&"bash_pattern_watches".to_string()));
625        assert!(sqlite_names(&conn, "index").contains(&"idx_bash_pattern_watches_task".to_string()));
626    }
627
628    #[test]
629    fn migration_runner_no_op_when_current() {
630        let dir = tempdir().unwrap();
631        let path = dir.path().join("aft.db");
632
633        let conn = open(&path).unwrap();
634        assert_eq!(schema_version_row_count(&conn), 1);
635        drop(conn);
636
637        let conn = open(&path).unwrap();
638        assert_eq!(schema_version(&conn), CURRENT_SCHEMA_VERSION);
639        assert_eq!(schema_version_row_count(&conn), 1);
640    }
641
642    #[test]
643    fn harness_state_compound_pk_works() {
644        let dir = tempdir().unwrap();
645        let conn = open(&dir.path().join("aft.db")).unwrap();
646
647        conn.execute(
648            "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
649            params!["opencode", "warned_tools", "{}", 1_i64],
650        )
651        .unwrap();
652        let duplicate = conn.execute(
653            "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
654            params!["opencode", "warned_tools", "{}", 2_i64],
655        );
656        assert_unique_constraint(duplicate);
657
658        conn.execute(
659            "INSERT INTO harness_state (harness, key, value, updated_at) VALUES (?1, ?2, ?3, ?4)",
660            params!["pi", "warned_tools", "{}", 3_i64],
661        )
662        .unwrap();
663    }
664
665    #[test]
666    fn host_state_simple_pk_works() {
667        let dir = tempdir().unwrap();
668        let conn = open(&dir.path().join("aft.db")).unwrap();
669
670        conn.execute(
671            "INSERT INTO host_state (key, value, updated_at) VALUES (?1, ?2, ?3)",
672            params!["trusted_filter_projects", "[]", 1_i64],
673        )
674        .unwrap();
675        let duplicate = conn.execute(
676            "INSERT INTO host_state (key, value, updated_at) VALUES (?1, ?2, ?3)",
677            params!["trusted_filter_projects", "[]", 2_i64],
678        );
679        assert_unique_constraint(duplicate);
680    }
681
682    #[test]
683    fn bash_tasks_compound_pk_works() {
684        let dir = tempdir().unwrap();
685        let conn = open(&dir.path().join("aft.db")).unwrap();
686
687        insert_bash_task(&conn, "opencode", "session-1", "bash-12345678").unwrap();
688        let duplicate = insert_bash_task(&conn, "opencode", "session-1", "bash-12345678");
689        assert_unique_constraint(duplicate);
690
691        insert_bash_task(&conn, "pi", "session-1", "bash-12345678").unwrap();
692    }
693
694    #[test]
695    fn backups_order_blob_sort() {
696        let dir = tempdir().unwrap();
697        let conn = open(&dir.path().join("aft.db")).unwrap();
698
699        let one = order_blob(1);
700        let two = order_blob(2);
701        let max = [0xFF; 16];
702
703        insert_backup(&conn, "one", &one).unwrap();
704        insert_backup(&conn, "two", &two).unwrap();
705        insert_backup(&conn, "max", &max).unwrap();
706
707        assert_eq!(backup_ids_ordered(&conn, "ASC"), vec!["one", "two", "max"]);
708        assert_eq!(backup_ids_ordered(&conn, "DESC"), vec!["max", "two", "one"]);
709    }
710
711    fn sqlite_names(conn: &Connection, kind: &str) -> Vec<String> {
712        let sql = match kind {
713            "table" => "SELECT name FROM sqlite_master WHERE type='table' AND name NOT LIKE 'sqlite_%' ORDER BY name",
714            "index" => "SELECT name FROM sqlite_master WHERE type='index' AND name NOT LIKE 'sqlite_%' ORDER BY name",
715            _ => panic!("unsupported sqlite_master kind: {kind}"),
716        };
717        let mut stmt = conn.prepare(sql).unwrap();
718        stmt.query_map([], |row| row.get::<_, String>(0))
719            .unwrap()
720            .collect::<Result<Vec<_>, _>>()
721            .unwrap()
722    }
723
724    fn table_columns(conn: &Connection, table: &str) -> Vec<String> {
725        let mut stmt = conn
726            .prepare(&format!("PRAGMA table_info({table})"))
727            .unwrap();
728        stmt.query_map([], |row| row.get::<_, String>(1))
729            .unwrap()
730            .collect::<Result<Vec<_>, _>>()
731            .unwrap()
732    }
733
734    fn schema_version(conn: &Connection) -> u32 {
735        conn.query_row("SELECT version FROM schema_version", [], |row| row.get(0))
736            .unwrap()
737    }
738
739    fn schema_version_row_count(conn: &Connection) -> i64 {
740        conn.query_row("SELECT COUNT(*) FROM schema_version", [], |row| row.get(0))
741            .unwrap()
742    }
743
744    fn assert_unique_constraint(result: rusqlite::Result<usize>) {
745        let error = result.expect_err("expected a unique constraint violation");
746        assert!(
747            error.to_string().contains("UNIQUE constraint failed"),
748            "expected UNIQUE constraint failure, got {error}"
749        );
750    }
751
752    fn insert_bash_task(
753        conn: &Connection,
754        harness: &str,
755        session_id: &str,
756        task_id: &str,
757    ) -> rusqlite::Result<usize> {
758        conn.execute(
759            "INSERT INTO bash_tasks (
760                harness, session_id, task_id, project_key, command, cwd, status, started_at
761             ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8)",
762            params![
763                harness,
764                session_id,
765                task_id,
766                "project-key",
767                "echo ok",
768                "/tmp",
769                "running",
770                1_i64
771            ],
772        )
773    }
774
775    fn insert_compression_event(
776        conn: &Connection,
777        id: i64,
778        harness: &str,
779        session_id: Option<&str>,
780        project_key: &str,
781        tool: &str,
782        task_id: Option<&str>,
783    ) -> rusqlite::Result<usize> {
784        conn.execute(
785            "INSERT INTO compression_events (
786                id, harness, session_id, project_key, tool, task_id, command, compressor,
787                original_bytes, compressed_bytes, original_tokens, compressed_tokens, created_at
788             ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10, ?11, ?12, ?13)",
789            params![
790                id,
791                harness,
792                session_id,
793                project_key,
794                tool,
795                task_id,
796                "echo ok",
797                "test-compressor",
798                100_i64,
799                50_i64,
800                20_i64,
801                10_i64,
802                id
803            ],
804        )
805    }
806
807    fn compression_event_ids(conn: &Connection) -> Vec<i64> {
808        let mut stmt = conn
809            .prepare("SELECT id FROM compression_events ORDER BY id")
810            .unwrap();
811        stmt.query_map([], |row| row.get::<_, i64>(0))
812            .unwrap()
813            .collect::<Result<Vec<_>, _>>()
814            .unwrap()
815    }
816
817    fn insert_backup(
818        conn: &Connection,
819        backup_id: &str,
820        order_blob: &[u8],
821    ) -> rusqlite::Result<usize> {
822        conn.execute(
823            "INSERT INTO backups (
824                backup_id, harness, session_id, project_key, order_blob, file_path,
825                path_hash, kind, created_at
826             ) VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
827            params![
828                backup_id,
829                "opencode",
830                "session-1",
831                "project-key",
832                order_blob,
833                "/tmp/file.txt",
834                "path-hash",
835                "content",
836                1_i64
837            ],
838        )
839    }
840
841    fn order_blob(value: u128) -> [u8; 16] {
842        value.to_be_bytes()
843    }
844
845    fn backup_ids_ordered(conn: &Connection, direction: &str) -> Vec<String> {
846        let sql = match direction {
847            "ASC" => "SELECT backup_id FROM backups ORDER BY order_blob ASC",
848            "DESC" => "SELECT backup_id FROM backups ORDER BY order_blob DESC",
849            _ => panic!("unsupported order direction: {direction}"),
850        };
851        let mut stmt = conn.prepare(sql).unwrap();
852        stmt.query_map([], |row| row.get::<_, String>(0))
853            .unwrap()
854            .collect::<Result<Vec<_>, _>>()
855            .unwrap()
856    }
857}