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