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
127const 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
181pub 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
199pub 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
208pub 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 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}