Skip to main content

aft/db/
bash_tasks.rs

1use std::io::ErrorKind;
2use std::path::{Path, PathBuf};
3use std::time::Duration;
4
5use rusqlite::types::Value;
6use rusqlite::{params, params_from_iter, Connection, OptionalExtension, Row};
7
8use crate::bash_background::persistence::{
9    resolve_task_layout, session_tasks_dir, uninitialized_layout_is_recent,
10};
11
12pub const TERMINAL_ROW_RETENTION_AGE_MS: i64 = 30 * 24 * 60 * 60 * 1000;
13const MAX_TERMINAL_PRUNE_ROWS: usize = 500;
14const LAYOUT_CREATION_GRACE: Duration = Duration::from_secs(5 * 60);
15
16const TERMINAL_PRUNE_PREDICATE: &str = "
17    status IN ('completed', 'failed', 'killed', 'timed_out')
18    AND completed_at IS NOT NULL
19    AND completed_at < ?1
20    AND completion_delivered = 1
21    AND NOT EXISTS (
22        SELECT 1 FROM bash_pattern_watches AS watch
23        WHERE watch.harness = bash_tasks.harness
24          AND watch.session_id = bash_tasks.session_id
25          AND watch.task_id = bash_tasks.task_id
26    )";
27
28#[derive(Debug, Clone, Copy, PartialEq, Eq)]
29pub struct TerminalRowsPrune {
30    pub removed: usize,
31    /// Rows left in the bounded `limit + 1` probe. Reaching the probe limit
32    /// means additional SQL candidates may remain beyond this count.
33    pub remaining_candidates: usize,
34}
35
36#[derive(Debug)]
37struct TerminalPruneCandidate {
38    harness: String,
39    session_id: String,
40    task_id: String,
41    stdout_path: Option<String>,
42    stderr_path: Option<String>,
43}
44
45#[derive(Debug)]
46pub(crate) struct TerminalPrunePlan {
47    storage_root: Option<PathBuf>,
48    candidates: Vec<TerminalPruneCandidate>,
49    probed_candidates: usize,
50    cutoff: i64,
51}
52
53#[derive(Debug)]
54pub(crate) struct PreparedTerminalPrune {
55    identities: Vec<TerminalPruneCandidate>,
56    probed_candidates: usize,
57    cutoff: i64,
58}
59
60#[derive(Debug, Clone)]
61pub struct BashTaskRow {
62    pub harness: String,
63    pub session_id: String,
64    pub task_id: String,
65    pub project_key: String,
66    pub command: String,
67    pub cwd: String,
68    pub status: String,
69    pub exit_code: Option<i32>,
70    pub pid: Option<i64>,
71    pub pgid: Option<i64>,
72    pub started_at: i64,
73    pub completed_at: Option<i64>,
74    pub stdout_path: Option<String>,
75    pub stderr_path: Option<String>,
76    pub compressed: bool,
77    pub timeout_ms: Option<i64>,
78    pub completion_delivered: bool,
79    pub output_bytes: Option<i64>,
80    pub metadata: String,
81}
82
83pub fn upsert_bash_task(conn: &Connection, row: &BashTaskRow) -> rusqlite::Result<()> {
84    conn.execute(
85        "INSERT INTO bash_tasks (
86            harness, session_id, task_id, project_key, command, cwd, status,
87            exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
88            compressed, timeout_ms, completion_delivered, output_bytes, metadata
89         ) VALUES (
90            ?1, ?2, ?3, ?4, ?5, ?6, ?7,
91            ?8, ?9, ?10, ?11, ?12, ?13, ?14,
92            ?15, ?16, ?17, ?18, ?19
93         )
94         ON CONFLICT(harness, session_id, task_id) DO UPDATE SET
95            project_key = excluded.project_key,
96            command = excluded.command,
97            cwd = excluded.cwd,
98            status = excluded.status,
99            exit_code = excluded.exit_code,
100            pid = excluded.pid,
101            pgid = excluded.pgid,
102            started_at = excluded.started_at,
103            completed_at = excluded.completed_at,
104            stdout_path = excluded.stdout_path,
105            stderr_path = excluded.stderr_path,
106            compressed = excluded.compressed,
107            timeout_ms = excluded.timeout_ms,
108            completion_delivered = excluded.completion_delivered,
109            output_bytes = excluded.output_bytes,
110            metadata = excluded.metadata",
111        params![
112            row.harness,
113            row.session_id,
114            row.task_id,
115            row.project_key,
116            row.command,
117            row.cwd,
118            row.status,
119            row.exit_code,
120            row.pid,
121            row.pgid,
122            row.started_at,
123            row.completed_at,
124            row.stdout_path,
125            row.stderr_path,
126            row.compressed,
127            row.timeout_ms,
128            row.completion_delivered,
129            row.output_bytes,
130            row.metadata,
131        ],
132    )?;
133    Ok(())
134}
135
136pub fn delete_delivered_terminal_bash_task(
137    conn: &Connection,
138    harness: &str,
139    session_id: &str,
140    task_id: &str,
141    reason: &str,
142) -> rusqlite::Result<usize> {
143    let deleted = conn.execute(
144        "DELETE FROM bash_tasks
145         WHERE harness = ?1 AND session_id = ?2 AND task_id = ?3
146           AND completion_delivered = 1
147           AND status IN ('completed', 'failed', 'killed', 'timed_out', 'fate_unknown')",
148        params![harness, session_id, task_id],
149    )?;
150    // A row can produce this warning only once: retries affect zero rows after
151    // the first successful DELETE, preventing a cleanup loop from flooding logs.
152    if deleted > 0 {
153        crate::slog_warn!("bash task row deleted: task_id={task_id} reason={reason}");
154    }
155    Ok(deleted)
156}
157
158pub fn delete_bash_task(
159    conn: &Connection,
160    harness: &str,
161    session_id: &str,
162    task_id: &str,
163) -> rusqlite::Result<usize> {
164    conn.execute(
165        "DELETE FROM bash_tasks
166         WHERE harness = ?1 AND session_id = ?2 AND task_id = ?3",
167        params![harness, session_id, task_id],
168    )
169}
170
171/// Remove acknowledged terminal rows only after the retention age has passed
172/// and the task layout is absent. The filesystem check uses the same lookup and
173/// initialization grace as persisted-task garbage collection, so a task being
174/// created concurrently is not mistaken for a missing task.
175pub fn prune_terminal_rows(
176    conn: &Connection,
177    now_ms: i64,
178    limit: usize,
179) -> rusqlite::Result<TerminalRowsPrune> {
180    prune_terminal_rows_guarded(conn, now_ms, limit, |_| false)
181}
182
183pub(crate) fn prune_terminal_rows_guarded(
184    conn: &Connection,
185    now_ms: i64,
186    limit: usize,
187    is_registered_in_process: impl Fn(&str) -> bool,
188) -> rusqlite::Result<TerminalRowsPrune> {
189    let plan = select_terminal_prune_candidates(conn, now_ms, limit)?;
190    let prepared = prepare_terminal_prune(plan, is_registered_in_process);
191    delete_prepared_terminal_rows(conn, prepared)
192}
193
194pub(crate) fn select_terminal_prune_candidates(
195    conn: &Connection,
196    now_ms: i64,
197    limit: usize,
198) -> rusqlite::Result<TerminalPrunePlan> {
199    let cutoff = now_ms.saturating_sub(TERMINAL_ROW_RETENTION_AGE_MS);
200    let bounded_limit = limit.min(MAX_TERMINAL_PRUNE_ROWS);
201    let probe_limit = bounded_limit.saturating_add(1);
202    let mut candidates = conn
203        .prepare(&format!(
204            "SELECT harness, session_id, task_id, stdout_path, stderr_path
205             FROM bash_tasks
206             WHERE {TERMINAL_PRUNE_PREDICATE}
207             LIMIT ?2"
208        ))?
209        .query_map(
210            params![cutoff, i64::try_from(probe_limit).unwrap_or(501)],
211            |row| {
212                Ok(TerminalPruneCandidate {
213                    harness: row.get(0)?,
214                    session_id: row.get(1)?,
215                    task_id: row.get(2)?,
216                    stdout_path: row.get(3)?,
217                    stderr_path: row.get(4)?,
218                })
219            },
220        )?
221        .collect::<rusqlite::Result<Vec<_>>>()?;
222    let probed_candidates = candidates.len();
223    candidates.truncate(bounded_limit);
224    let storage_root = conn
225        .path()
226        .and_then(|path| Path::new(path).parent())
227        .filter(|path| !path.as_os_str().is_empty())
228        .map(Path::to_path_buf);
229    Ok(TerminalPrunePlan {
230        storage_root,
231        candidates,
232        probed_candidates,
233        cutoff,
234    })
235}
236
237pub(crate) fn prepare_terminal_prune(
238    plan: TerminalPrunePlan,
239    is_registered_in_process: impl Fn(&str) -> bool,
240) -> PreparedTerminalPrune {
241    prepare_terminal_prune_observed(plan, is_registered_in_process, || {})
242}
243
244pub(crate) fn prepare_terminal_prune_observed(
245    plan: TerminalPrunePlan,
246    is_registered_in_process: impl Fn(&str) -> bool,
247    observe_stat_phase: impl FnOnce(),
248) -> PreparedTerminalPrune {
249    observe_stat_phase();
250    let identities = plan
251        .candidates
252        .into_iter()
253        .filter(|candidate| {
254            !is_registered_in_process(&candidate.task_id)
255                && task_layout_is_gone(plan.storage_root.as_deref(), candidate)
256        })
257        .collect();
258    PreparedTerminalPrune {
259        identities,
260        probed_candidates: plan.probed_candidates,
261        cutoff: plan.cutoff,
262    }
263}
264
265pub(crate) fn cap_prepared_terminal_rows(prepared: &mut PreparedTerminalPrune, limit: usize) {
266    prepared.identities.truncate(limit);
267}
268
269pub(crate) fn delete_prepared_terminal_rows(
270    conn: &Connection,
271    prepared: PreparedTerminalPrune,
272) -> rusqlite::Result<TerminalRowsPrune> {
273    let removed = if prepared.identities.is_empty() {
274        0
275    } else {
276        let mut values = Vec::with_capacity(1 + prepared.identities.len() * 3);
277        values.push(Value::Integer(prepared.cutoff));
278        let identities = prepared
279            .identities
280            .into_iter()
281            .enumerate()
282            .map(|(index, candidate)| {
283                let parameter = 2 + index * 3;
284                values.extend([
285                    Value::Text(candidate.harness),
286                    Value::Text(candidate.session_id),
287                    Value::Text(candidate.task_id),
288                ]);
289                format!("(?{parameter}, ?{}, ?{})", parameter + 1, parameter + 2)
290            })
291            .collect::<Vec<_>>()
292            .join(", ");
293        conn.execute(
294            &format!(
295                "DELETE FROM bash_tasks
296                 WHERE {TERMINAL_PRUNE_PREDICATE}
297                   AND (harness, session_id, task_id) IN (VALUES {identities})"
298            ),
299            params_from_iter(values),
300        )?
301    };
302
303    Ok(TerminalRowsPrune {
304        removed,
305        remaining_candidates: prepared.probed_candidates.saturating_sub(removed),
306    })
307}
308
309fn task_layout_is_gone(storage_root: Option<&Path>, candidate: &TerminalPruneCandidate) -> bool {
310    let session_dir = candidate_session_dir(candidate).or_else(|| {
311        storage_root.map(|storage_root| session_tasks_dir(storage_root, &candidate.session_id))
312    });
313    let Some(session_dir) = session_dir else {
314        return false;
315    };
316    let task_id = &candidate.task_id;
317    let directory_layout = session_dir.join(task_id);
318    let flat_layout = session_dir.join(format!("{task_id}.json"));
319    match resolve_task_layout(&session_dir, task_id) {
320        Ok(_) => false,
321        Err(error) if error.kind() == ErrorKind::NotFound => {
322            if !directory_layout.exists() && !flat_layout.exists() {
323                return true;
324            }
325            !uninitialized_layout_is_recent(&session_dir, task_id, LAYOUT_CREATION_GRACE)
326                .unwrap_or(true)
327        }
328        // Releases before the 64-bit random suffix used eight hexadecimal
329        // digits. Current layout validation intentionally rejects those IDs,
330        // so retain any surviving legacy artifact and prune only full absence.
331        Err(error) if error.kind() == ErrorKind::InvalidInput && is_legacy_task_id(task_id) => {
332            !directory_layout.exists()
333                && !flat_layout.exists()
334                && !candidate
335                    .stdout_path
336                    .iter()
337                    .chain(&candidate.stderr_path)
338                    .any(|path| Path::new(path).exists())
339        }
340        Err(_) => false,
341    }
342}
343
344fn candidate_session_dir(candidate: &TerminalPruneCandidate) -> Option<PathBuf> {
345    candidate
346        .stdout_path
347        .iter()
348        .chain(&candidate.stderr_path)
349        .find_map(|path| {
350            let path = Path::new(path);
351            let parent = path.parent()?;
352            let file_name = path.file_name()?.to_str()?;
353            if file_name.starts_with(&format!("{}.", candidate.task_id)) {
354                return Some(parent.to_path_buf());
355            }
356            let task_dir = parent.parent()?;
357            (task_dir.file_name()?.to_str()? == candidate.task_id)
358                .then(|| task_dir.parent().map(Path::to_path_buf))
359                .flatten()
360        })
361}
362
363fn is_legacy_task_id(task_id: &str) -> bool {
364    task_id.strip_prefix("bash-").is_some_and(|suffix| {
365        suffix.len() == 8 && suffix.bytes().all(|byte| byte.is_ascii_hexdigit())
366    })
367}
368
369pub fn get_bash_task(
370    conn: &Connection,
371    harness: &str,
372    session_id: &str,
373    task_id: &str,
374) -> rusqlite::Result<Option<BashTaskRow>> {
375    conn.query_row(
376        "SELECT harness, session_id, task_id, project_key, command, cwd, status,
377                exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
378                compressed, timeout_ms, completion_delivered, output_bytes, metadata
379         FROM bash_tasks
380         WHERE harness = ?1 AND session_id = ?2 AND task_id = ?3",
381        params![harness, session_id, task_id],
382        map_bash_task_row,
383    )
384    .optional()
385}
386
387const SESSION_TASKS_SQL: &str =
388    "SELECT harness, session_id, task_id, project_key, command, cwd, status,
389                exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
390                compressed, timeout_ms, completion_delivered, output_bytes, metadata
391         FROM bash_tasks
392         WHERE harness = ?1 AND session_id = ?2";
393
394pub fn list_bash_tasks_for_session(
395    conn: &Connection,
396    harness: &str,
397    session_id: &str,
398) -> rusqlite::Result<Vec<BashTaskRow>> {
399    let mut stmt = conn.prepare(SESSION_TASKS_SQL)?;
400    let mut rows = stmt
401        .query_map(params![harness, session_id], map_bash_task_row)?
402        .collect::<rusqlite::Result<Vec<_>>>()?;
403    // Task rows carry large command/metadata payloads. Sorting them in SQLite
404    // spills entire rows to a temporary file for long-lived sessions. Keep the
405    // legacy integer/BINARY order, but sort the already-required result vector.
406    rows.sort_by(|a, b| (a.started_at, &a.task_id).cmp(&(b.started_at, &b.task_id)));
407    Ok(rows)
408}
409
410pub fn list_bash_tasks_by_id(
411    conn: &Connection,
412    harness: &str,
413    task_id: &str,
414) -> rusqlite::Result<Vec<BashTaskRow>> {
415    let mut stmt = conn.prepare(
416        "SELECT harness, session_id, task_id, project_key, command, cwd, status,
417                exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
418                compressed, timeout_ms, completion_delivered, output_bytes, metadata
419         FROM bash_tasks
420         WHERE harness = ?1 AND task_id = ?2
421         ORDER BY started_at DESC",
422    )?;
423    let rows = stmt
424        .query_map(params![harness, task_id], map_bash_task_row)?
425        .collect();
426    rows
427}
428
429pub fn list_replayable_bash_tasks_for_project(
430    conn: &Connection,
431    harness: &str,
432    project_key: &str,
433) -> rusqlite::Result<Vec<BashTaskRow>> {
434    let mut stmt = conn.prepare(
435        "SELECT harness, session_id, task_id, project_key, command, cwd, status,
436                exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
437                compressed, timeout_ms, completion_delivered, output_bytes, metadata
438         FROM bash_tasks
439         WHERE harness = ?1 AND project_key = ?2
440           AND (status NOT IN ('completed', 'failed', 'killed', 'timed_out', 'fate_unknown')
441                OR completion_delivered = 0)
442         ORDER BY started_at ASC, task_id ASC",
443    )?;
444    let rows = stmt
445        .query_map(params![harness, project_key], map_bash_task_row)?
446        .collect();
447    rows
448}
449
450pub fn find_bash_task_for_project(
451    conn: &Connection,
452    harness: &str,
453    project_key: &str,
454    task_id: &str,
455) -> rusqlite::Result<Option<BashTaskRow>> {
456    conn.query_row(
457        "SELECT harness, session_id, task_id, project_key, command, cwd, status,
458                exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
459                compressed, timeout_ms, completion_delivered, output_bytes, metadata
460         FROM bash_tasks
461         WHERE harness = ?1 AND project_key = ?2 AND task_id = ?3
462         ORDER BY started_at DESC
463         LIMIT 1",
464        params![harness, project_key, task_id],
465        map_bash_task_row,
466    )
467    .optional()
468}
469
470fn map_bash_task_row(row: &Row<'_>) -> rusqlite::Result<BashTaskRow> {
471    Ok(BashTaskRow {
472        harness: row.get(0)?,
473        session_id: row.get(1)?,
474        task_id: row.get(2)?,
475        project_key: row.get(3)?,
476        command: row.get(4)?,
477        cwd: row.get(5)?,
478        status: row.get(6)?,
479        exit_code: row.get(7)?,
480        pid: row.get(8)?,
481        pgid: row.get(9)?,
482        started_at: row.get(10)?,
483        completed_at: row.get(11)?,
484        stdout_path: row.get(12)?,
485        stderr_path: row.get(13)?,
486        compressed: row.get::<_, i64>(14)? != 0,
487        timeout_ms: row.get(15)?,
488        completion_delivered: row.get::<_, i64>(16)? != 0,
489        output_bytes: row.get(17)?,
490        metadata: row.get::<_, Option<String>>(18)?.unwrap_or_default(),
491    })
492}
493
494#[cfg(test)]
495mod tests {
496    use super::*;
497
498    #[test]
499    fn session_history_preserves_sqlite_order_without_a_temp_sort() {
500        let temp = tempfile::tempdir().unwrap();
501        let conn = crate::db::open(&temp.path().join("aft.db")).unwrap();
502        for (task, started, status) in [
503            ("é", 4, "running"),
504            ("a", 4, "completed"),
505            ("z", -1, "failed"),
506            ("A", 4, "failed"),
507            ("aa", 4, "running"),
508            ("first", i64::MIN, "completed"),
509        ] {
510            conn.execute("INSERT INTO bash_tasks
511                (harness, session_id, task_id, project_key, command, cwd, status, started_at, metadata)
512                VALUES ('opencode', 'session', ?1, 'project', ?2, '.', ?3, ?4, ?5)",
513                params![task, "command".repeat(8192), status, started, format!("metadata-{task}")]).unwrap();
514        }
515        let legacy = conn
516            .prepare(
517                "SELECT harness, session_id, task_id, project_key, command, cwd, status,
518            exit_code, pid, pgid, started_at, completed_at, stdout_path, stderr_path,
519            compressed, timeout_ms, completion_delivered, output_bytes, metadata
520            FROM bash_tasks WHERE harness = ?1 AND session_id = ?2
521            ORDER BY started_at ASC, task_id ASC",
522            )
523            .unwrap()
524            .query_map(params!["opencode", "session"], map_bash_task_row)
525            .unwrap()
526            .collect::<rusqlite::Result<Vec<_>>>()
527            .unwrap();
528        let actual = list_bash_tasks_for_session(&conn, "opencode", "session").unwrap();
529        assert_eq!(format!("{actual:?}"), format!("{legacy:?}"));
530        assert_eq!(
531            actual
532                .iter()
533                .map(|r| r.task_id.as_str())
534                .collect::<Vec<_>>(),
535            ["first", "z", "A", "a", "aa", "é"]
536        );
537        assert!(list_bash_tasks_for_session(&conn, "other", "session")
538            .unwrap()
539            .is_empty());
540        let plan = conn
541            .prepare(&format!("EXPLAIN QUERY PLAN {SESSION_TASKS_SQL}"))
542            .unwrap()
543            .query_map(params!["opencode", "session"], |r| r.get::<_, String>(3))
544            .unwrap()
545            .collect::<rusqlite::Result<Vec<_>>>()
546            .unwrap()
547            .join("\n");
548        assert!(
549            !plan.contains("TEMP B-TREE"),
550            "session history must not spill task rows: {plan}"
551        );
552    }
553}