Skip to main content

task_runs/
store.rs

1//! Per-camp Turso store for TaskRun bytes (Tier 1).
2//!
3//! One `TaskStore` per daemon instance; callers share it via `Arc<TaskStore>`.
4//! Backed by `turso` (in-process, async) per W195 §Engine. Concurrency model:
5//! each `TaskStore` method gets a fresh `turso::Connection` from the shared
6//! `Database` and drops it at the end of the call. A single `Connection`
7//! cannot be used concurrently in turso 0.6.x — the SDK's `ConcurrentGuard`
8//! returns `Misuse("concurrent use forbidden")` rather than serializing —
9//! so we don't share one across the reader thread, the tail-poll loop, and
10//! the GC sweep. `Database::connect()` is cheap (an `Arc` clone plus a
11//! per-connection state struct).
12//!
13//! Independent connections mean independent write-lock contenders, so every
14//! connection carries a `busy_timeout` and every write goes through
15//! [`TaskStore::exec_retry`], which retries the `Busy` class with backoff. A
16//! write that returns `Busy` is a *lost* write, not a slow one — R617-B9 was
17//! exactly that: the lifecycle's terminal `update_status` losing the race
18//! against chunk/event writers and leaving the run `Running` forever.
19//!
20//! Storage contract (W195 §3 / Shape 1): this store owns
21//! `.yah/db/task-runs.turso` under the camp daemon.
22
23use crate::types::{
24    BeholderStatus, ChunkRef, Diagnostic, Event, EventSource, Initiator, KeepRange, Level,
25    OutputChunk, RunStatus, SeqRange, Stream, TaskRunId, TaskRunMeta, Triage,
26};
27use serde_json::Value as JsonValue;
28use std::collections::HashMap;
29use std::path::Path;
30use std::str::FromStr;
31use thiserror::Error;
32use tokio::sync::Mutex;
33use turso::{params, params_from_iter, Builder, Connection, Database, Value};
34
35/// Total in-turso backoff before a statement gives up and reports `Busy`.
36const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
37/// First outer-retry delay for a write that still came back `Busy`.
38const BUSY_RETRY_BASE: std::time::Duration = std::time::Duration::from_millis(10);
39/// Once the doubling delay passes this, the write gives up (≈10s total).
40const BUSY_RETRY_MAX: std::time::Duration = std::time::Duration::from_secs(2);
41
42/// `true` when the error means "someone else holds the write lock" — the
43/// retryable class, as opposed to a schema/constraint/misuse failure.
44fn is_busy(e: &turso::Error) -> bool {
45    matches!(e, turso::Error::Busy(_) | turso::Error::BusySnapshot(_))
46}
47
48// ─── Error ────────────────────────────────────────────────────────────────────
49
50#[derive(Debug, Error)]
51pub enum StoreError {
52    #[error("turso: {0}")]
53    Sql(#[from] turso::Error),
54    #[error("json: {0}")]
55    Json(#[from] serde_json::Error),
56    #[error("run not found: {0}")]
57    NotFound(String),
58    #[error("invalid stream value: {0}")]
59    InvalidStream(String),
60    #[error("io: {0}")]
61    Io(#[from] std::io::Error),
62    #[error("invalid field path: {0}")]
63    InvalidFieldPath(String),
64}
65
66// ─── Schema ───────────────────────────────────────────────────────────────────
67
68const SCHEMA: &str = r#"
69CREATE TABLE IF NOT EXISTS runs (
70    id              TEXT PRIMARY KEY,
71    command         TEXT NOT NULL,
72    cwd             TEXT NOT NULL,
73    env_json        TEXT NOT NULL,
74    started_at      INTEGER NOT NULL,
75    ended_at        INTEGER,
76    exit_code       INTEGER,
77    signal          INTEGER,
78    status          TEXT NOT NULL,
79    status_detail   TEXT,
80    label           TEXT,
81    initiator       TEXT NOT NULL,
82    beholder_status TEXT,
83    archived_at     INTEGER,
84    pinned          INTEGER NOT NULL DEFAULT 0,
85    origin          TEXT,
86    host_pid        INTEGER
87);
88
89CREATE TABLE IF NOT EXISTS chunks (
90    run_id      TEXT NOT NULL,
91    seq         INTEGER NOT NULL,
92    offset_ms   INTEGER NOT NULL,
93    stream      TEXT NOT NULL,
94    bytes       BLOB NOT NULL,
95    PRIMARY KEY (run_id, seq)
96);
97CREATE INDEX IF NOT EXISTS chunks_by_offset ON chunks(run_id, offset_ms);
98
99CREATE TABLE IF NOT EXISTS events (
100    run_id      TEXT NOT NULL,
101    seq         INTEGER NOT NULL,
102    offset_ms   INTEGER NOT NULL,
103    level       TEXT NOT NULL,
104    target      TEXT NOT NULL,
105    msg         TEXT NOT NULL,
106    fields_json TEXT NOT NULL,
107    anchor_seq  INTEGER,
108    source_kind TEXT NOT NULL,
109    source_name TEXT NOT NULL,
110    PRIMARY KEY (run_id, seq)
111);
112CREATE INDEX IF NOT EXISTS events_by_target ON events(run_id, target, level);
113CREATE INDEX IF NOT EXISTS events_by_offset ON events(run_id, offset_ms);
114
115CREATE TABLE IF NOT EXISTS triages (
116    run_id         TEXT PRIMARY KEY,
117    synopsis       TEXT NOT NULL,
118    keep_json      TEXT NOT NULL,
119    primary_lo     INTEGER NOT NULL,
120    primary_hi     INTEGER NOT NULL,
121    model          TEXT NOT NULL,
122    prompt_version INTEGER NOT NULL,
123    cached_at      INTEGER NOT NULL,
124    partial        INTEGER NOT NULL DEFAULT 0
125);
126
127CREATE TABLE IF NOT EXISTS _event_field_indexes (
128    field_path  TEXT PRIMARY KEY,
129    index_name  TEXT NOT NULL,
130    created_at  INTEGER NOT NULL
131);
132"#;
133
134// ─── Filter types ─────────────────────────────────────────────────────────────
135
136#[derive(Debug, Default)]
137pub struct RunFilter {
138    pub since: Option<u64>,
139    pub label: Option<String>,
140    pub status: Option<String>,
141    pub limit: Option<usize>,
142    pub archived: Option<bool>,
143    /// Filter by provenance tag (e.g. `Some("terminal")` to list only
144    /// terminal-session runs). `None` returns runs of every origin.
145    pub origin: Option<String>,
146}
147
148#[derive(Debug, Default)]
149pub struct ChunkFilter {
150    pub stream: Option<Stream>,
151    pub seq_range: Option<(u32, u32)>,
152    pub offset_range: Option<(u32, u32)>,
153    pub limit: Option<usize>,
154}
155
156// ─── Durability ───────────────────────────────────────────────────────────────
157
158/// Put this connection's syncs where they stop corruption, and take them off
159/// the commit path where they only cost.
160///
161/// On macOS `fsync(2)` returns once the data reaches the drive — it does not
162/// wait for the drive to flush its own write cache, and `fcntl(F_FULLFSYNC)`
163/// is the documented barrier. `turso_core` implements both but defaults to
164/// plain `fsync` on Apple, so the window between "the checkpoint synced the
165/// database file" and "the drive actually wrote it" is real. A panic inside
166/// that window takes the pages the checkpoint just wrote while the WAL reset
167/// that followed survives, and the result is committed data that exists in
168/// neither file. That is what happened to `task-runs.turso` (R914, 118 pages)
169/// and `gnome_queue.turso` (R914-B4, 38 rows) two days apart on one machine.
170///
171/// `F_FULLFSYNC` is a device-wide cache flush — 3.86 ms against `fsync`'s
172/// 0.029 ms, measured on the machine that lost both stores — so it ships
173/// alongside `synchronous = NORMAL`, which stops turso syncing on every
174/// commit and leaves the barrier on the checkpoint and WAL-restart path that
175/// actually failed. turso's own comment says so at `storage/pager.rs:4312`:
176/// "NORMAL mode skips fsync on WAL commit (but still fsyncs on checkpoint and
177/// wal restart)". This store writes a row per PTY chunk, so the pairing is
178/// load-bearing. Measured over 200 single-row commits: turso's default
179/// (`FULL` + `fsync`) 0.090 ms/commit, `FULL` + `F_FULLFSYNC` 3.612, and the
180/// shipped `NORMAL` + `F_FULLFSYNC` 0.225. The barrier is not free here —
181/// 2.5x the default — but `fullfsync` on its own would have been 40x.
182///
183/// Full writeup, including the forensics, lives in the yah monorepo at
184/// `crates/yah/camp-service/src/store_durability.rs`. It is restated rather
185/// than shared because `oss/qed` is an independent workspace that must build
186/// standalone.
187async fn harden_sync(conn: &Connection) -> Result<(), StoreError> {
188    conn.execute("PRAGMA synchronous = NORMAL", ()).await?;
189    if cfg!(target_vendor = "apple") {
190        // `PRAGMA fullfsync` is Apple-gated inside `turso_core`; plain
191        // `fsync` is already a barrier elsewhere.
192        conn.execute("PRAGMA fullfsync = ON", ()).await?;
193    }
194    Ok(())
195}
196
197// ─── TaskStore ────────────────────────────────────────────────────────────────
198
199struct SeqCounters {
200    next_seq: HashMap<String, u32>,
201    next_event_seq: HashMap<String, u32>,
202}
203
204pub struct TaskStore {
205    db: Database,
206    seq: Mutex<SeqCounters>,
207}
208
209impl TaskStore {
210    /// Open (or create) the per-camp task-runs database at `path`.
211    pub async fn open(path: &Path) -> Result<Self, StoreError> {
212        if let Some(parent) = path.parent() {
213            std::fs::create_dir_all(parent)?;
214        }
215        let db = Builder::new_local(path.to_string_lossy().as_ref())
216            .build()
217            .await?;
218        let conn = db.connect()?;
219        harden_sync(&conn).await?;
220        conn.execute_batch(SCHEMA).await?;
221        /* `CREATE TABLE IF NOT EXISTS` won't add a column to a runs table that
222           predates `origin`, so add it idempotently for already-created DBs.
223           A duplicate-column error means an up-to-date schema — swallow it. */
224        let _ = conn.execute("ALTER TABLE runs ADD COLUMN origin TEXT", ()).await;
225        /* Same idempotent add for `host_pid` (R617-F6). A row that predates it
226           reads back `None`, which the stale-run policy treats as "owner
227           unknown" and therefore tombstones — i.e. an old DB keeps exactly the
228           old Lost-on-disappear behaviour. */
229        let _ = conn
230            .execute("ALTER TABLE runs ADD COLUMN host_pid INTEGER", ())
231            .await;
232        Ok(TaskStore {
233            db,
234            seq: Mutex::new(SeqCounters {
235                next_seq: HashMap::new(),
236                next_event_seq: HashMap::new(),
237            }),
238        })
239    }
240
241    /// Open an in-memory store (tests only).
242    #[cfg(test)]
243    pub async fn open_in_memory() -> Result<Self, StoreError> {
244        let db = Builder::new_local(":memory:").build().await?;
245        db.connect()?.execute_batch(SCHEMA).await?;
246        Ok(TaskStore {
247            db,
248            seq: Mutex::new(SeqCounters {
249                next_seq: HashMap::new(),
250                next_event_seq: HashMap::new(),
251            }),
252        })
253    }
254
255    /// Open a fresh connection to the underlying database. Each call returns
256    /// an independent `Connection` that the caller may use within one logical
257    /// operation and drop. Never share a `Connection` across awaiting tasks —
258    /// `turso` 0.6.x rejects concurrent use on the same handle.
259    ///
260    /// Every connection carries a busy timeout: a live run drives three
261    /// concurrent writers (PTY reader chunks, shim-FIFO events, lifecycle
262    /// status), and without it the loser of a write-lock race gets an
263    /// immediate `Busy` instead of waiting its turn.
264    /// Every connection is also hardened against the macOS `fsync` gap
265    /// (R914-B4). `PRAGMA fullfsync` lives on the *pager*, which turso builds
266    /// per `Connection`, so it cannot be set once at open: auto-checkpoint
267    /// fires on whichever connection happens to cross the WAL threshold, and
268    /// an un-hardened one there is exactly how this store lost 118 committed
269    /// pages on 2026-09-15. The full mechanism is documented in the yah
270    /// monorepo at `crates/yah/camp-service/src/store_durability.rs`; it is
271    /// restated rather than shared because `oss/qed` is an independent
272    /// workspace that must build standalone.
273    async fn conn(&self) -> Result<Connection, StoreError> {
274        let conn = self.db.connect()?;
275        let _ = conn.busy_timeout(BUSY_TIMEOUT);
276        harden_sync(&conn).await?;
277        Ok(conn)
278    }
279
280    /// Execute a write statement, retrying while turso reports the database
281    /// is locked.
282    ///
283    /// `busy_timeout` alone is not enough: turso caps its internal backoff at
284    /// the configured total and then hands `Busy` back to the caller. Losing a
285    /// write here is not a slow path but a *wrong* one — a dropped terminal
286    /// `update_status` leaves the run `Running` forever — so writes get an
287    /// outer retry with exponential backoff on top.
288    async fn exec_retry(
289        &self,
290        sql: &str,
291        params: impl turso::params::IntoParams,
292    ) -> Result<u64, StoreError> {
293        let params = params.into_params()?;
294        let mut delay = BUSY_RETRY_BASE;
295        let mut last: StoreError;
296        loop {
297            match self.conn().await?.execute(sql, params.clone()).await {
298                Ok(n) => return Ok(n),
299                Err(e) if is_busy(&e) => last = StoreError::Sql(e),
300                Err(e) => return Err(StoreError::Sql(e)),
301            }
302            if delay > BUSY_RETRY_MAX {
303                return Err(last);
304            }
305            tokio::time::sleep(delay).await;
306            delay *= 2;
307        }
308    }
309
310    pub async fn insert_run(&self, meta: &TaskRunMeta) -> Result<(), StoreError> {
311        let (status, ended_at, exit_code, signal, detail) = status_columns(&meta.status);
312        let initiator_json = serde_json::to_string(&meta.initiator)?;
313        let env_json = serde_json::to_string(&meta.env)?;
314        let beholder_json = meta
315            .beholder_status
316            .as_ref()
317            .map(|b| serde_json::to_string(b))
318            .transpose()?;
319        let cwd = meta.cwd.to_string_lossy().to_string();
320
321        self.exec_retry(
322                "INSERT INTO runs \
323                 (id, command, cwd, env_json, started_at, ended_at, exit_code, signal, \
324                  status, status_detail, label, initiator, beholder_status, pinned, origin, \
325                  host_pid) \
326                 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10,?11,?12,?13,?14,?15,?16)",
327                params![
328                    meta.id.to_string(),
329                    meta.command.clone(),
330                    cwd,
331                    env_json,
332                    meta.started_at as i64,
333                    ended_at,
334                    exit_code.map(|c| c as i64),
335                    signal.map(|s| s as i64),
336                    status.to_string(),
337                    detail,
338                    meta.label.clone(),
339                    initiator_json,
340                    beholder_json,
341                    meta.pinned as i64,
342                    meta.origin.clone(),
343                    meta.host_pid.map(|p| p as i64),
344                ],
345            )
346            .await?;
347
348        let mut seq = self.seq.lock().await;
349        seq.next_seq.entry(meta.id.to_string()).or_insert(0);
350        seq.next_event_seq.entry(meta.id.to_string()).or_insert(0);
351        Ok(())
352    }
353
354    pub async fn update_beholder_status(
355        &self,
356        id: &TaskRunId,
357        status: &BeholderStatus,
358    ) -> Result<(), StoreError> {
359        let json = serde_json::to_string(status)?;
360        self.exec_retry(
361            "UPDATE runs SET beholder_status = ?1 WHERE id = ?2",
362            params![json, id.to_string()],
363        )
364        .await?;
365        Ok(())
366    }
367
368    pub async fn update_status(
369        &self,
370        id: &TaskRunId,
371        status: &RunStatus,
372    ) -> Result<(), StoreError> {
373        let (status_str, ended_at, exit_code, signal, detail) = status_columns(status);
374        self.exec_retry(
375            "UPDATE runs SET status=?1, status_detail=?2, ended_at=?3, exit_code=?4, signal=?5 \
376             WHERE id=?6",
377            params![
378                status_str.to_string(),
379                detail,
380                ended_at,
381                exit_code.map(|c| c as i64),
382                signal.map(|s| s as i64),
383                id.to_string(),
384            ],
385        )
386        .await?;
387        Ok(())
388    }
389
390    /// Append an output chunk. Returns the assigned `seq` number.
391    pub async fn append_chunk(
392        &self,
393        run_id: &TaskRunId,
394        offset_ms: u32,
395        stream: Stream,
396        bytes: &[u8],
397    ) -> Result<u32, StoreError> {
398        let key = run_id.to_string();
399        let seq = {
400            let mut g = self.seq.lock().await;
401            if let Some(s) = g.next_seq.get_mut(&key) {
402                let v = *s;
403                *s += 1;
404                v
405            } else {
406                drop(g);
407                let max = self.max_seq("chunks", &key).await?;
408                let next = max.map(|m| m + 1).unwrap_or(0);
409                let mut g = self.seq.lock().await;
410                g.next_seq.insert(key.clone(), next + 1);
411                next
412            }
413        };
414
415        self.exec_retry(
416                "INSERT INTO chunks (run_id, seq, offset_ms, stream, bytes) VALUES (?1,?2,?3,?4,?5)",
417                params![
418                    key,
419                    seq as i64,
420                    offset_ms as i64,
421                    stream.as_str().to_string(),
422                    bytes.to_vec(),
423                ],
424            )
425            .await?;
426        Ok(seq)
427    }
428
429    async fn max_seq(&self, table: &str, run_id: &str) -> Result<Option<u32>, StoreError> {
430        let sql = format!("SELECT MAX(seq) FROM {table} WHERE run_id = ?1");
431        let mut rows = self.conn().await?.query(&sql, params![run_id.to_string()]).await?;
432        match rows.next().await? {
433            Some(row) => {
434                let v: Option<i64> = row.get(0)?;
435                Ok(v.map(|n| n as u32))
436            }
437            None => Ok(None),
438        }
439    }
440
441    pub async fn get_run(&self, id: &TaskRunId) -> Result<Option<TaskRunMeta>, StoreError> {
442        let mut rows = self
443            .conn().await?
444            .query(
445                "SELECT id, command, cwd, env_json, started_at, ended_at, exit_code, signal, \
446                        status, status_detail, label, initiator, beholder_status, pinned, origin, \
447                        host_pid \
448                 FROM runs WHERE id = ?1",
449                params![id.to_string()],
450            )
451            .await?;
452        match rows.next().await? {
453            Some(row) => Ok(Some(row_to_meta(&row)?)),
454            None => Ok(None),
455        }
456    }
457
458    pub async fn chunk_count(&self, run_id: &TaskRunId) -> Result<u32, StoreError> {
459        let mut rows = self
460            .conn().await?
461            .query(
462                "SELECT COUNT(*) FROM chunks WHERE run_id = ?1",
463                params![run_id.to_string()],
464            )
465            .await?;
466        let row = rows.next().await?.expect("COUNT(*) always returns one row");
467        let n: i64 = row.get(0)?;
468        Ok(n as u32)
469    }
470
471    pub async fn list_runs(&self, filter: &RunFilter) -> Result<Vec<TaskRunMeta>, StoreError> {
472        let limit = filter.limit.unwrap_or(usize::MAX) as i64;
473        let since = filter.since.map(|s| s as i64).unwrap_or(0);
474
475        let archived_clause = match filter.archived {
476            Some(true) => "AND archived_at IS NOT NULL",
477            _ => "AND archived_at IS NULL",
478        };
479
480        let mut where_extra = String::new();
481        let mut p: Vec<Value> = vec![Value::Integer(since), Value::Integer(limit)];
482        // ?1 = since, ?2 = limit, then extra bindings starting at ?3
483        let mut next_param = 3;
484        if let Some(ref l) = filter.label {
485            where_extra.push_str(&format!(" AND label = ?{next_param}"));
486            p.push(Value::Text(l.clone()));
487            next_param += 1;
488        }
489        if let Some(ref s) = filter.status {
490            where_extra.push_str(&format!(" AND status = ?{next_param}"));
491            p.push(Value::Text(s.clone()));
492            next_param += 1;
493        }
494        if let Some(ref o) = filter.origin {
495            where_extra.push_str(&format!(" AND origin = ?{next_param}"));
496            p.push(Value::Text(o.clone()));
497        }
498
499        let sql = format!(
500            "SELECT id, command, cwd, env_json, started_at, ended_at, exit_code, signal, \
501                    status, status_detail, label, initiator, beholder_status, pinned, origin, \
502                    host_pid \
503             FROM runs \
504             WHERE started_at >= ?1 {} {} \
505             ORDER BY started_at DESC \
506             LIMIT ?2",
507            archived_clause, where_extra
508        );
509
510        let mut rows = self.conn().await?.query(&sql, params_from_iter(p)).await?;
511        let mut out = Vec::new();
512        while let Some(row) = rows.next().await? {
513            out.push(row_to_meta(&row)?);
514        }
515        Ok(out)
516    }
517
518    pub async fn archive_run(&self, id: &TaskRunId) -> Result<(), StoreError> {
519        let now = unix_now() as i64;
520        let count = self
521            .conn().await?
522            .execute(
523                "UPDATE runs SET archived_at = ?1 WHERE id = ?2 AND archived_at IS NULL",
524                params![now, id.to_string()],
525            )
526            .await?;
527        if count == 0 {
528            return Err(StoreError::NotFound(id.to_string()));
529        }
530        Ok(())
531    }
532
533    pub async fn pin_run(&self, id: &TaskRunId, pinned: bool) -> Result<(), StoreError> {
534        let count = self
535            .conn().await?
536            .execute(
537                "UPDATE runs SET pinned = ?1 WHERE id = ?2",
538                params![pinned as i64, id.to_string()],
539            )
540            .await?;
541        if count == 0 {
542            return Err(StoreError::NotFound(id.to_string()));
543        }
544        Ok(())
545    }
546
547    pub async fn gc_sweep(&self, config: &GcConfig) -> Result<GcResult, StoreError> {
548        let now = unix_now();
549        let warm_cutoff = (now as i64).saturating_sub(config.warm_secs as i64);
550        let mut result = GcResult::default();
551
552        result.archived_runs_cleaned =
553            self.count_query("SELECT COUNT(*) FROM runs WHERE archived_at IS NOT NULL", vec![])
554                .await?;
555
556        result.warm_rolloff_runs = self
557            .count_query(
558                "SELECT COUNT(*) FROM runs \
559                 WHERE archived_at IS NULL AND pinned = 0 AND started_at < ?1",
560                vec![Value::Integer(warm_cutoff)],
561            )
562            .await?;
563
564        result.chunks_deleted += self
565            .conn().await?
566            .execute(
567                "DELETE FROM chunks \
568                 WHERE run_id IN (SELECT id FROM runs WHERE archived_at IS NOT NULL)",
569                (),
570            )
571            .await?;
572
573        result.events_deleted += self
574            .conn().await?
575            .execute(
576                "DELETE FROM events \
577                 WHERE run_id IN (SELECT id FROM runs WHERE archived_at IS NOT NULL)",
578                (),
579            )
580            .await?;
581
582        result.chunks_deleted += self
583            .conn().await?
584            .execute(
585                "DELETE FROM chunks WHERE run_id IN \
586                 (SELECT id FROM runs WHERE archived_at IS NULL AND pinned = 0 AND started_at < ?1)",
587                params![warm_cutoff],
588            )
589            .await?;
590
591        result.events_deleted += self
592            .conn().await?
593            .execute(
594                "DELETE FROM events WHERE run_id IN \
595                 (SELECT id FROM runs WHERE archived_at IS NULL AND pinned = 0 AND started_at < ?1)",
596                params![warm_cutoff],
597            )
598            .await?;
599
600        Ok(result)
601    }
602
603    async fn count_query(&self, sql: &str, params: Vec<Value>) -> Result<u64, StoreError> {
604        let mut rows = self.conn().await?.query(sql, params_from_iter(params)).await?;
605        let row = rows.next().await?.expect("COUNT(*) returns one row");
606        let n: i64 = row.get(0)?;
607        Ok(n as u64)
608    }
609
610    pub async fn get_chunks(
611        &self,
612        run_id: &TaskRunId,
613        filter: &ChunkFilter,
614    ) -> Result<Vec<OutputChunk>, StoreError> {
615        let key = run_id.to_string();
616
617        let mut conditions = vec!["run_id = ?1".to_string()];
618        let mut p: Vec<Value> = vec![Value::Text(key)];
619        if let Some(s) = filter.stream {
620            conditions.push("stream = ?2".to_string());
621            p.push(Value::Text(s.as_str().to_string()));
622        }
623        if let Some((lo, hi)) = filter.seq_range {
624            conditions.push(format!("seq >= {} AND seq <= {}", lo, hi));
625        }
626        if let Some((from, to)) = filter.offset_range {
627            conditions.push(format!("offset_ms >= {} AND offset_ms <= {}", from, to));
628        }
629        let where_clause = conditions.join(" AND ");
630        let limit_clause = filter
631            .limit
632            .map(|l| format!("LIMIT {l}"))
633            .unwrap_or_default();
634
635        let sql = format!(
636            "SELECT seq, offset_ms, stream, bytes FROM chunks WHERE {} ORDER BY seq {}",
637            where_clause, limit_clause
638        );
639
640        let mut rows = self.conn().await?.query(&sql, params_from_iter(p)).await?;
641        let mut out = Vec::new();
642        while let Some(row) = rows.next().await? {
643            let seq: i64 = row.get(0)?;
644            let offset_ms: i64 = row.get(1)?;
645            let stream_str: String = row.get(2)?;
646            let bytes: Vec<u8> = row.get(3)?;
647            let stream = stream_str
648                .parse::<Stream>()
649                .map_err(StoreError::InvalidStream)?;
650            out.push(OutputChunk {
651                run_id: run_id.clone(),
652                seq: seq as u32,
653                offset_ms: offset_ms as u32,
654                stream,
655                bytes,
656            });
657        }
658        Ok(out)
659    }
660
661    // ─── Events ───────────────────────────────────────────────────────────────
662
663    pub async fn append_event(
664        &self,
665        run_id: &TaskRunId,
666        offset_ms: u32,
667        level: Level,
668        target: &str,
669        msg: &str,
670        fields: &serde_json::Value,
671        anchor_seq: Option<u32>,
672        source: &EventSource,
673    ) -> Result<u32, StoreError> {
674        let key = run_id.to_string();
675        let seq = {
676            let mut g = self.seq.lock().await;
677            if let Some(s) = g.next_event_seq.get_mut(&key) {
678                let v = *s;
679                *s += 1;
680                v
681            } else {
682                drop(g);
683                let max = self.max_seq("events", &key).await?;
684                let next = max.map(|m| m + 1).unwrap_or(0);
685                let mut g = self.seq.lock().await;
686                g.next_event_seq.insert(key.clone(), next + 1);
687                next
688            }
689        };
690
691        let fields_json = serde_json::to_string(fields)?;
692
693        self.exec_retry(
694                "INSERT INTO events \
695                 (run_id, seq, offset_ms, level, target, msg, fields_json, anchor_seq, source_kind, source_name) \
696                 VALUES (?1,?2,?3,?4,?5,?6,?7,?8,?9,?10)",
697                params![
698                    key,
699                    seq as i64,
700                    offset_ms as i64,
701                    level.as_str().to_string(),
702                    target.to_string(),
703                    msg.to_string(),
704                    fields_json,
705                    anchor_seq.map(|s| s as i64),
706                    source.kind_str().to_string(),
707                    source.name_str().to_string(),
708                ],
709            )
710            .await?;
711        Ok(seq)
712    }
713
714    pub async fn query_events(
715        &self,
716        run_id: &TaskRunId,
717        filter: &EventFilter,
718    ) -> Result<Vec<Event>, StoreError> {
719        let run_id_str = run_id.to_string();
720        let mut conditions = vec!["run_id = ?1".to_string()];
721        let mut p: Vec<Value> = vec![Value::Text(run_id_str)];
722        let mut next_param = 2usize;
723
724        if let Some(ref target) = filter.target {
725            conditions.push(format!("target = ?{next_param}"));
726            p.push(Value::Text(target.clone()));
727            next_param += 1;
728        }
729
730        if let Some(min_level) = filter.min_level {
731            let levels: Vec<String> = ALL_LEVELS
732                .iter()
733                .copied()
734                .filter(|&(l, _)| l >= min_level)
735                .map(|(_, s)| s.to_string())
736                .collect();
737            let placeholders: String = (0..levels.len())
738                .map(|i| format!("?{}", next_param + i))
739                .collect::<Vec<_>>()
740                .join(", ");
741            conditions.push(format!("level IN ({placeholders})"));
742            for level_str in levels {
743                p.push(Value::Text(level_str));
744                next_param += 1;
745            }
746        }
747
748        if let Some(ref range) = filter.seq_range {
749            conditions.push(format!("seq >= ?{next_param}"));
750            p.push(Value::Integer(range.start as i64));
751            next_param += 1;
752            conditions.push(format!("seq < ?{next_param}"));
753            p.push(Value::Integer(range.end as i64));
754            next_param += 1;
755        }
756
757        if let Some((from_ms, to_ms)) = filter.offset_range {
758            conditions.push(format!("offset_ms >= ?{next_param}"));
759            p.push(Value::Integer(from_ms as i64));
760            next_param += 1;
761            conditions.push(format!("offset_ms <= ?{next_param}"));
762            p.push(Value::Integer(to_ms as i64));
763            next_param += 1;
764        }
765
766        if let Some(ref ff) = filter.field_filter {
767            validate_field_path(&ff.path)?;
768            let escaped = ff.path.replace('\'', "''");
769            conditions.push(format!(
770                "json_extract(fields_json, '{escaped}') = ?{next_param}"
771            ));
772            match &ff.value {
773                JsonValue::String(s) => p.push(Value::Text(s.clone())),
774                JsonValue::Number(n) => {
775                    if let Some(i) = n.as_i64() {
776                        p.push(Value::Integer(i));
777                    } else if let Some(f) = n.as_f64() {
778                        p.push(Value::Real(f));
779                    } else {
780                        p.push(Value::Text(n.to_string()));
781                    }
782                }
783                JsonValue::Bool(b) => p.push(Value::Integer(*b as i64)),
784                JsonValue::Null => p.push(Value::Null),
785                other => p.push(Value::Text(other.to_string())),
786            }
787            let _ = next_param;
788        }
789
790        let where_clause = conditions.join(" AND ");
791        let limit_clause = filter.limit.map(|n| format!("LIMIT {n}")).unwrap_or_default();
792        let sql = format!(
793            "SELECT run_id, seq, offset_ms, level, target, msg, fields_json, anchor_seq, \
794             source_kind, source_name \
795             FROM events WHERE {where_clause} ORDER BY seq {limit_clause}"
796        );
797
798        let mut rows = self.conn().await?.query(&sql, params_from_iter(p)).await?;
799        let mut events = Vec::new();
800        while let Some(row) = rows.next().await? {
801            events.push(row_to_event(&row)?);
802        }
803        Ok(events)
804    }
805
806    pub async fn event_count(
807        &self,
808        run_id: &TaskRunId,
809        min_level: Option<Level>,
810    ) -> Result<u64, StoreError> {
811        let run_id_str = run_id.to_string();
812        if let Some(min_level) = min_level {
813            let levels: Vec<String> = ALL_LEVELS
814                .iter()
815                .copied()
816                .filter(|&(l, _)| l >= min_level)
817                .map(|(_, s)| s.to_string())
818                .collect();
819            let placeholders: String = levels
820                .iter()
821                .enumerate()
822                .map(|(i, _)| format!("?{}", i + 2))
823                .collect::<Vec<_>>()
824                .join(", ");
825            let sql = format!(
826                "SELECT COUNT(*) FROM events WHERE run_id = ?1 AND level IN ({placeholders})"
827            );
828            let mut p: Vec<Value> = vec![Value::Text(run_id_str)];
829            for s in levels {
830                p.push(Value::Text(s));
831            }
832            self.count_query(&sql, p).await
833        } else {
834            self.count_query(
835                "SELECT COUNT(*) FROM events WHERE run_id = ?1",
836                vec![Value::Text(run_id_str)],
837            )
838            .await
839        }
840    }
841
842    pub async fn query_diagnostics(
843        &self,
844        run_id: &TaskRunId,
845    ) -> Result<Vec<Diagnostic>, StoreError> {
846        let events = self
847            .query_events(
848                run_id,
849                &EventFilter {
850                    min_level: Some(Level::Warn),
851                    ..EventFilter::default()
852                },
853            )
854            .await?;
855        Ok(events.into_iter().map(event_to_diagnostic).collect())
856    }
857
858    // ─── json_extract hooks ───────────────────────────────────────────────────
859
860    pub async fn ensure_field_index(&self, field_path: &str) -> Result<(), StoreError> {
861        validate_field_path(field_path)?;
862
863        let exists: bool = {
864            let mut rows = self
865                .conn().await?
866                .query(
867                    "SELECT 1 FROM _event_field_indexes WHERE field_path = ?1",
868                    params![field_path.to_string()],
869                )
870                .await?;
871            rows.next().await?.is_some()
872        };
873
874        if exists {
875            return Ok(());
876        }
877
878        let index_name = field_path_to_index_name(field_path);
879        let escaped = field_path.replace('\'', "''");
880        let sql = format!(
881            "CREATE INDEX IF NOT EXISTS {index_name} \
882             ON events(run_id, json_extract(fields_json, '{escaped}')) \
883             WHERE json_extract(fields_json, '{escaped}') IS NOT NULL"
884        );
885        self.conn().await?.execute_batch(&sql).await?;
886
887        let now = unix_now();
888        self.conn().await?
889            .execute(
890                "INSERT OR IGNORE INTO _event_field_indexes (field_path, index_name, created_at) \
891                 VALUES (?1, ?2, ?3)",
892                params![field_path.to_string(), index_name, now as i64],
893            )
894            .await?;
895
896        Ok(())
897    }
898
899    pub async fn timeline_ticks(
900        &self,
901        run_id: &TaskRunId,
902        since_offset_ms: Option<u32>,
903        limit: Option<usize>,
904    ) -> Result<Vec<TimelineTick>, StoreError> {
905        let since = since_offset_ms.unwrap_or(0);
906        let cap = limit.unwrap_or(usize::MAX);
907
908        let mut chunks = self
909            .get_chunks(
910                run_id,
911                &ChunkFilter {
912                    offset_range: Some((since, u32::MAX)),
913                    ..Default::default()
914                },
915            )
916            .await?;
917        let mut events = self
918            .query_events(
919                run_id,
920                &EventFilter {
921                    offset_range: Some((since, u32::MAX)),
922                    ..Default::default()
923                },
924            )
925            .await?;
926
927        chunks.sort_by_key(|c| c.offset_ms);
928        events.sort_by_key(|e| e.offset_ms);
929
930        let mut ticks: Vec<TimelineTick> =
931            Vec::with_capacity((chunks.len() + events.len()).min(cap));
932        let mut ci = 0usize;
933        let mut ei = 0usize;
934
935        while ticks.len() < cap && (ci < chunks.len() || ei < events.len()) {
936            match (chunks.get(ci), events.get(ei)) {
937                (Some(c), Some(e)) if c.offset_ms <= e.offset_ms => {
938                    ticks.push(TimelineTick::Chunk(chunks[ci].clone()));
939                    ci += 1;
940                    let _ = e;
941                }
942                (Some(_), Some(_)) => {
943                    ticks.push(TimelineTick::Event(events[ei].clone()));
944                    ei += 1;
945                }
946                (Some(_), None) => {
947                    ticks.push(TimelineTick::Chunk(chunks[ci].clone()));
948                    ci += 1;
949                }
950                (None, Some(_)) => {
951                    ticks.push(TimelineTick::Event(events[ei].clone()));
952                    ei += 1;
953                }
954                (None, None) => break,
955            }
956        }
957
958        Ok(ticks)
959    }
960
961    pub async fn aggregate_events(
962        &self,
963        filter: &AggregateFilter,
964    ) -> Result<Vec<AggregateBucket>, StoreError> {
965        let limit = filter.limit.unwrap_or(100) as i64;
966        let since = filter.since.map(|s| s as i64).unwrap_or(0);
967        let group_by = filter.group_by.unwrap_or(AggregateGroupBy::Target);
968
969        let key_expr = match group_by {
970            AggregateGroupBy::Target => "e.target".to_string(),
971            AggregateGroupBy::Level => "e.level".to_string(),
972            AggregateGroupBy::ErrorCode => {
973                "json_extract(e.fields_json, '$.error.code')".to_string()
974            }
975        };
976
977        let mut where_parts = vec!["r.started_at >= ?1".to_string()];
978        let mut p: Vec<Value> = vec![Value::Integer(since)];
979        let mut next_param = 2usize;
980
981        if let Some(ref label) = filter.label {
982            where_parts.push(format!("r.label = ?{next_param}"));
983            p.push(Value::Text(label.clone()));
984            next_param += 1;
985        }
986
987        if matches!(group_by, AggregateGroupBy::ErrorCode) {
988            where_parts.push("json_extract(e.fields_json, '$.error.code') IS NOT NULL".to_string());
989        }
990
991        if let Some(ref ff) = filter.field_filter {
992            validate_field_path(&ff.path)?;
993            let escaped = ff.path.replace('\'', "''");
994            where_parts.push(format!(
995                "json_extract(e.fields_json, '{escaped}') = ?{next_param}"
996            ));
997            match &ff.value {
998                JsonValue::String(s) => p.push(Value::Text(s.clone())),
999                JsonValue::Number(n) => {
1000                    if let Some(i) = n.as_i64() {
1001                        p.push(Value::Integer(i));
1002                    } else if let Some(f) = n.as_f64() {
1003                        p.push(Value::Real(f));
1004                    } else {
1005                        p.push(Value::Text(n.to_string()));
1006                    }
1007                }
1008                JsonValue::Bool(b) => p.push(Value::Integer(*b as i64)),
1009                JsonValue::Null => p.push(Value::Null),
1010                other => p.push(Value::Text(other.to_string())),
1011            }
1012            let _ = next_param;
1013        }
1014
1015        let where_clause = where_parts.join(" AND ");
1016        let limit_param = p.len() + 1;
1017        let sql = format!(
1018            "SELECT {key_expr} as key, COUNT(*) as cnt \
1019             FROM events e \
1020             JOIN runs r ON r.id = e.run_id \
1021             WHERE {where_clause} \
1022             GROUP BY {key_expr} \
1023             ORDER BY cnt DESC \
1024             LIMIT ?{limit_param}"
1025        );
1026        p.push(Value::Integer(limit));
1027
1028        let mut rows = self.conn().await?.query(&sql, params_from_iter(p)).await?;
1029        let mut buckets = Vec::new();
1030        while let Some(row) = rows.next().await? {
1031            let key: Option<String> = row.get(0)?;
1032            let count: i64 = row.get(1)?;
1033            buckets.push(AggregateBucket {
1034                key: key.unwrap_or_default(),
1035                count: count as u64,
1036            });
1037        }
1038        Ok(buckets)
1039    }
1040
1041    // ── Tier 1.75 — triage ───────────────────────────────────────────────────
1042
1043    pub async fn upsert_triage(&self, triage: &Triage) -> Result<(), StoreError> {
1044        let keep_json = serde_json::to_string(&triage.keep)?;
1045        self.conn().await?
1046            .execute(
1047                "INSERT OR REPLACE INTO triages \
1048                 (run_id, synopsis, keep_json, primary_lo, primary_hi, \
1049                  model, prompt_version, cached_at, partial) \
1050                 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9)",
1051                params![
1052                    triage.run_id.to_string(),
1053                    triage.synopsis.clone(),
1054                    keep_json,
1055                    triage.primary.lo as i64,
1056                    triage.primary.hi as i64,
1057                    triage.model.clone(),
1058                    triage.prompt_version as i64,
1059                    triage.cached_at as i64,
1060                    triage.partial as i64,
1061                ],
1062            )
1063            .await?;
1064        Ok(())
1065    }
1066
1067    pub async fn get_triage(&self, run_id: &TaskRunId) -> Result<Option<Triage>, StoreError> {
1068        let mut rows = self
1069            .conn().await?
1070            .query(
1071                "SELECT run_id, synopsis, keep_json, primary_lo, primary_hi, \
1072                 model, prompt_version, cached_at, partial \
1073                 FROM triages WHERE run_id = ?1",
1074                params![run_id.to_string()],
1075            )
1076            .await?;
1077        let Some(row) = rows.next().await? else {
1078            return Ok(None);
1079        };
1080        let id: String = row.get(0)?;
1081        let synopsis: String = row.get(1)?;
1082        let keep_json: String = row.get(2)?;
1083        let primary_lo: i64 = row.get(3)?;
1084        let primary_hi: i64 = row.get(4)?;
1085        let model: String = row.get(5)?;
1086        let prompt_version: i64 = row.get(6)?;
1087        let cached_at: i64 = row.get(7)?;
1088        let partial: i64 = row.get(8)?;
1089        let keep: Vec<KeepRange> = serde_json::from_str(&keep_json)?;
1090        Ok(Some(Triage {
1091            run_id: id.parse().map_err(|_| StoreError::NotFound(id.clone()))?,
1092            synopsis,
1093            keep,
1094            primary: SeqRange {
1095                lo: primary_lo as u32,
1096                hi: primary_hi as u32,
1097            },
1098            model,
1099            prompt_version: prompt_version as u32,
1100            cached_at: cached_at as u64,
1101            partial: partial != 0,
1102        }))
1103    }
1104
1105    pub async fn list_field_indexes(&self) -> Result<Vec<FieldIndexInfo>, StoreError> {
1106        let mut rows = self
1107            .conn().await?
1108            .query(
1109                "SELECT field_path, index_name, created_at \
1110                 FROM _event_field_indexes ORDER BY created_at",
1111                (),
1112            )
1113            .await?;
1114        let mut out = Vec::new();
1115        while let Some(row) = rows.next().await? {
1116            let field_path: String = row.get(0)?;
1117            let index_name: String = row.get(1)?;
1118            let created_at: i64 = row.get(2)?;
1119            out.push(FieldIndexInfo {
1120                field_path,
1121                index_name,
1122                created_at: created_at as u64,
1123            });
1124        }
1125        Ok(out)
1126    }
1127}
1128
1129// ─── Timeline + aggregate ─────────────────────────────────────────────────────
1130
1131pub enum TimelineTick {
1132    Chunk(OutputChunk),
1133    Event(Event),
1134}
1135
1136impl TimelineTick {
1137    pub fn offset_ms(&self) -> u32 {
1138        match self {
1139            TimelineTick::Chunk(c) => c.offset_ms,
1140            TimelineTick::Event(e) => e.offset_ms,
1141        }
1142    }
1143}
1144
1145#[derive(Debug, Clone, Copy, PartialEq, Eq)]
1146pub enum AggregateGroupBy {
1147    Target,
1148    Level,
1149    ErrorCode,
1150}
1151
1152#[derive(Debug, Default)]
1153pub struct AggregateFilter {
1154    pub label: Option<String>,
1155    pub since: Option<u64>,
1156    pub group_by: Option<AggregateGroupBy>,
1157    pub field_filter: Option<FieldFilter>,
1158    pub limit: Option<u32>,
1159}
1160
1161#[derive(Debug, Clone)]
1162pub struct AggregateBucket {
1163    pub key: String,
1164    pub count: u64,
1165}
1166
1167// ─── Event filter types ───────────────────────────────────────────────────────
1168
1169#[derive(Debug, Default)]
1170pub struct EventFilter {
1171    pub target: Option<String>,
1172    pub min_level: Option<Level>,
1173    pub seq_range: Option<std::ops::Range<u32>>,
1174    pub offset_range: Option<(u32, u32)>,
1175    pub field_filter: Option<FieldFilter>,
1176    pub limit: Option<u32>,
1177}
1178
1179#[derive(Debug)]
1180pub struct FieldFilter {
1181    pub path: String,
1182    pub value: serde_json::Value,
1183}
1184
1185#[derive(Debug, Clone)]
1186pub struct FieldIndexInfo {
1187    pub field_path: String,
1188    pub index_name: String,
1189    pub created_at: u64,
1190}
1191
1192// ─── GC types ─────────────────────────────────────────────────────────────────
1193
1194#[derive(Debug, Clone)]
1195pub struct GcConfig {
1196    pub warm_secs: u64,
1197}
1198
1199impl Default for GcConfig {
1200    fn default() -> Self {
1201        Self {
1202            warm_secs: 30 * 24 * 3600,
1203        }
1204    }
1205}
1206
1207#[derive(Debug, Default)]
1208pub struct GcResult {
1209    pub archived_runs_cleaned: u64,
1210    pub warm_rolloff_runs: u64,
1211    pub chunks_deleted: u64,
1212    pub events_deleted: u64,
1213}
1214
1215// ─── Helpers ─────────────────────────────────────────────────────────────────
1216
1217const ALL_LEVELS: &[(Level, &str)] = &[
1218    (Level::Trace, "trace"),
1219    (Level::Debug, "debug"),
1220    (Level::Info, "info"),
1221    (Level::Warn, "warn"),
1222    (Level::Error, "error"),
1223    (Level::Fatal, "fatal"),
1224];
1225
1226pub fn validate_field_path(path: &str) -> Result<(), StoreError> {
1227    if path.is_empty() || !path.starts_with('$') {
1228        return Err(StoreError::InvalidFieldPath(path.to_string()));
1229    }
1230    let valid = path[1..]
1231        .chars()
1232        .all(|c| matches!(c, 'a'..='z' | 'A'..='Z' | '0'..='9' | '.' | '_' | '[' | ']'));
1233    if !valid {
1234        return Err(StoreError::InvalidFieldPath(path.to_string()));
1235    }
1236    Ok(())
1237}
1238
1239fn field_path_to_index_name(path: &str) -> String {
1240    let ident: String = path
1241        .chars()
1242        .skip(1)
1243        .map(|c| if c.is_alphanumeric() || c == '_' { c } else { '_' })
1244        .collect::<String>()
1245        .trim_matches('_')
1246        .to_string();
1247    format!("events_field_{ident}")
1248}
1249
1250fn row_to_event(r: &turso::Row) -> Result<Event, StoreError> {
1251    let run_id_str: String = r.get(0)?;
1252    let run_id = run_id_str
1253        .parse::<TaskRunId>()
1254        .map_err(|_| StoreError::NotFound(run_id_str.clone()))?;
1255
1256    let level_str: String = r.get(3)?;
1257    let level = Level::from_str(&level_str).map_err(StoreError::InvalidFieldPath)?;
1258
1259    let fields_json: String = r.get(6)?;
1260    let fields: JsonValue = serde_json::from_str(&fields_json)?;
1261
1262    let anchor_seq: Option<i64> = r.get(7)?;
1263    let anchor = anchor_seq.map(|s| ChunkRef { seq: s as u32 });
1264
1265    let source_kind: String = r.get(8)?;
1266    let source_name: String = r.get(9)?;
1267    let source = match source_kind.as_str() {
1268        "beholder" => EventSource::Beholder {
1269            name: source_name,
1270            version: "unknown".to_string(),
1271        },
1272        "shim" => EventSource::Shim {
1273            lib: source_name,
1274            version: "unknown".to_string(),
1275        },
1276        _ => EventSource::Synth,
1277    };
1278
1279    let seq: i64 = r.get(1)?;
1280    let offset_ms: i64 = r.get(2)?;
1281    let target: String = r.get(4)?;
1282    let msg: String = r.get(5)?;
1283
1284    Ok(Event {
1285        run_id,
1286        seq: seq as u32,
1287        offset_ms: offset_ms as u32,
1288        level,
1289        target,
1290        msg,
1291        fields,
1292        anchor,
1293        source,
1294    })
1295}
1296
1297fn event_to_diagnostic(e: Event) -> Diagnostic {
1298    let file = e
1299        .fields
1300        .get("file")
1301        .and_then(|f| f.get("path"))
1302        .and_then(|v| v.as_str())
1303        .map(str::to_owned);
1304    let line = e
1305        .fields
1306        .get("file")
1307        .and_then(|f| f.get("line"))
1308        .and_then(|v| v.as_u64())
1309        .map(|n| n as u32);
1310    let col = e
1311        .fields
1312        .get("file")
1313        .and_then(|f| f.get("col"))
1314        .and_then(|v| v.as_u64())
1315        .map(|n| n as u32);
1316    let code = e
1317        .fields
1318        .get("error")
1319        .and_then(|f| f.get("code"))
1320        .and_then(|v| v.as_str())
1321        .map(str::to_owned);
1322
1323    Diagnostic {
1324        severity: e.level,
1325        file,
1326        line,
1327        col,
1328        code,
1329        message: e.msg,
1330        source: format!("{}/{}", e.source.kind_str(), e.source.name_str()),
1331        run_id: e.run_id,
1332        event_seq: e.seq,
1333    }
1334}
1335
1336fn unix_now() -> u64 {
1337    std::time::SystemTime::now()
1338        .duration_since(std::time::UNIX_EPOCH)
1339        .unwrap_or_default()
1340        .as_secs()
1341}
1342
1343fn status_columns(
1344    status: &RunStatus,
1345) -> (&'static str, Option<i64>, Option<i32>, Option<i32>, Option<String>) {
1346    match status {
1347        RunStatus::Pending => ("pending", None, None, None, None),
1348        RunStatus::Running => ("running", None, None, None, None),
1349        RunStatus::Done { exit_code, ended_at } => {
1350            ("done", Some(*ended_at as i64), Some(*exit_code), None, None)
1351        }
1352        RunStatus::Killed { signal, ended_at } => {
1353            ("killed", Some(*ended_at as i64), None, Some(*signal), None)
1354        }
1355        RunStatus::Lost { reason } => ("lost", None, None, None, Some(reason.clone())),
1356    }
1357}
1358
1359fn row_to_meta(row: &turso::Row) -> Result<TaskRunMeta, StoreError> {
1360    let id_str: String = row.get(0)?;
1361    let command: String = row.get(1)?;
1362    let cwd: String = row.get(2)?;
1363    let env_json: String = row.get(3)?;
1364    let started_at: i64 = row.get(4)?;
1365    let ended_at: Option<i64> = row.get(5)?;
1366    let exit_code: Option<i64> = row.get(6)?;
1367    let signal: Option<i64> = row.get(7)?;
1368    let status_str: String = row.get(8)?;
1369    let status_detail: Option<String> = row.get(9)?;
1370    let label: Option<String> = row.get(10)?;
1371    let initiator_json: String = row.get(11)?;
1372    let beholder_json: Option<String> = row.get(12)?;
1373    let pinned: i64 = row.get(13).unwrap_or(0);
1374    let origin: Option<String> = row.get(14).ok().flatten();
1375    // `ok().flatten()` (not `?`): a DB written before the column existed has
1376    // no index 15 at all, and that must read as "owner unknown", not an error.
1377    let host_pid: Option<i64> = row.get(15).ok().flatten();
1378
1379    let status = reconstruct_status(
1380        &status_str,
1381        exit_code.map(|c| c as i32),
1382        signal.map(|s| s as i32),
1383        ended_at,
1384        status_detail,
1385    );
1386    let initiator: Initiator = serde_json::from_str(&initiator_json)?;
1387    let env: Vec<(String, String)> = serde_json::from_str(&env_json)?;
1388    let beholder_status: Option<BeholderStatus> = beholder_json
1389        .as_deref()
1390        .map(serde_json::from_str)
1391        .transpose()?;
1392
1393    Ok(TaskRunMeta {
1394        id: id_str
1395            .parse()
1396            .map_err(|e: uuid::Error| StoreError::NotFound(e.to_string()))?,
1397        command,
1398        cwd: cwd.into(),
1399        env,
1400        started_at: started_at as u64,
1401        status,
1402        label,
1403        initiator,
1404        beholder_status,
1405        pinned: pinned != 0,
1406        origin,
1407        host_pid: host_pid.map(|p| p as u32),
1408    })
1409}
1410
1411fn reconstruct_status(
1412    s: &str,
1413    exit_code: Option<i32>,
1414    signal: Option<i32>,
1415    ended_at: Option<i64>,
1416    detail: Option<String>,
1417) -> RunStatus {
1418    let ended_at = ended_at.unwrap_or(0) as u64;
1419    match s {
1420        "pending" => RunStatus::Pending,
1421        "running" => RunStatus::Running,
1422        "done" => RunStatus::Done {
1423            exit_code: exit_code.unwrap_or(0),
1424            ended_at,
1425        },
1426        "killed" => RunStatus::Killed {
1427            signal: signal.unwrap_or(15),
1428            ended_at,
1429        },
1430        "lost" => RunStatus::Lost {
1431            reason: detail.unwrap_or_default(),
1432        },
1433        other => RunStatus::Lost {
1434            reason: format!("unknown status in db: {other}"),
1435        },
1436    }
1437}
1438
1439// ─── Tests ────────────────────────────────────────────────────────────────────
1440
1441#[cfg(test)]
1442mod tests {
1443    use super::*;
1444    use crate::types::{Initiator, RunStatus, Stream, TaskRunId, TaskRunMeta};
1445    use std::path::PathBuf;
1446
1447    fn make_meta(id: TaskRunId, cmd: &str) -> TaskRunMeta {
1448        TaskRunMeta {
1449            id,
1450            command: cmd.to_string(),
1451            cwd: PathBuf::from("/tmp"),
1452            env: vec![],
1453            started_at: 1_000_000,
1454            status: RunStatus::Running,
1455            label: None,
1456            initiator: Initiator::Human {
1457                camp: "local".to_string(),
1458            },
1459            beholder_status: None,
1460            pinned: false,
1461            origin: None,
1462            host_pid: None,
1463        }
1464    }
1465
1466    #[tokio::test]
1467    async fn open_creates_schema() {
1468        let dir = tempfile::tempdir().unwrap();
1469        let path = dir.path().join("task-runs.turso");
1470        let store = TaskStore::open(&path).await.unwrap();
1471        assert!(path.exists());
1472        let n = store
1473            .count_query("SELECT COUNT(*) FROM runs", vec![])
1474            .await
1475            .unwrap();
1476        assert_eq!(n, 0);
1477    }
1478
1479    /// R617-B9: a live run has three concurrent writers (PTY chunks, shim
1480    /// events, lifecycle status). Before `exec_retry`, the terminal
1481    /// `update_status` lost the write-lock race and came back
1482    /// `Busy("database is locked")` — the caller's `let _ =` swallowed it and
1483    /// the run stayed `Running` forever. Writes must survive contention.
1484    #[tokio::test(flavor = "multi_thread", worker_threads = 4)]
1485    async fn concurrent_writers_do_not_lose_the_terminal_status() {
1486        let dir = tempfile::tempdir().unwrap();
1487        let store =
1488            std::sync::Arc::new(TaskStore::open(&dir.path().join("tr.turso")).await.unwrap());
1489        let id = TaskRunId::new();
1490        store.insert_run(&make_meta(id.clone(), "contended")).await.unwrap();
1491
1492        // Two background writers hammering the same DB from other tasks.
1493        let mut writers = Vec::new();
1494        for w in 0..2 {
1495            let store = std::sync::Arc::clone(&store);
1496            let id = id.clone();
1497            writers.push(tokio::spawn(async move {
1498                for i in 0..50u32 {
1499                    store.append_chunk(&id, i, Stream::Stdout, b"noise").await.unwrap();
1500                    store
1501                        .append_event(
1502                            &id,
1503                            i,
1504                            Level::Info,
1505                            "t",
1506                            &format!("w{w}-{i}"),
1507                            &serde_json::json!({}),
1508                            None,
1509                            &EventSource::Shim {
1510                                lib: "test".into(),
1511                                version: "0.1.0".into(),
1512                            },
1513                        )
1514                        .await
1515                        .unwrap();
1516                }
1517            }));
1518        }
1519
1520        let status = RunStatus::Done { exit_code: 0, ended_at: 2_000_000 };
1521        store.update_status(&id, &status).await.unwrap();
1522
1523        for w in writers.drain(..) {
1524            w.await.unwrap();
1525        }
1526        let got = store.get_run(&id).await.unwrap().unwrap();
1527        assert!(
1528            matches!(got.status, RunStatus::Done { exit_code: 0, .. }),
1529            "terminal status was lost under contention: {:?}",
1530            got.status
1531        );
1532    }
1533
1534    #[tokio::test]
1535    async fn insert_and_get_run() {
1536        let store = TaskStore::open_in_memory().await.unwrap();
1537        let id = TaskRunId::new();
1538        let meta = make_meta(id.clone(), "cargo check");
1539        store.insert_run(&meta).await.unwrap();
1540        let got = store.get_run(&id).await.unwrap().unwrap();
1541        assert_eq!(got.command, "cargo check");
1542        assert!(matches!(got.status, RunStatus::Running));
1543    }
1544
1545    #[tokio::test]
1546    async fn append_chunks_monotonic_seq() {
1547        let store = TaskStore::open_in_memory().await.unwrap();
1548        let id = TaskRunId::new();
1549        store.insert_run(&make_meta(id.clone(), "echo hi")).await.unwrap();
1550
1551        let s0 = store.append_chunk(&id, 0, Stream::Stdout, b"hello\n").await.unwrap();
1552        let s1 = store.append_chunk(&id, 5, Stream::Stdout, b"world\n").await.unwrap();
1553        let s2 = store.append_chunk(&id, 10, Stream::Stderr, b"err\n").await.unwrap();
1554
1555        assert_eq!(s0, 0);
1556        assert_eq!(s1, 1);
1557        assert_eq!(s2, 2);
1558        assert_eq!(store.chunk_count(&id).await.unwrap(), 3);
1559    }
1560
1561    #[tokio::test]
1562    async fn seq_resumes_after_reopen() {
1563        let dir = tempfile::tempdir().unwrap();
1564        let path = dir.path().join("tr.turso");
1565        let id = TaskRunId::new();
1566
1567        {
1568            let store = TaskStore::open(&path).await.unwrap();
1569            store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1570            store.append_chunk(&id, 0, Stream::Stdout, b"a").await.unwrap();
1571            store.append_chunk(&id, 1, Stream::Stdout, b"b").await.unwrap();
1572        }
1573
1574        let store = TaskStore::open(&path).await.unwrap();
1575        let seq = store.append_chunk(&id, 2, Stream::Stdout, b"c").await.unwrap();
1576        assert_eq!(seq, 2);
1577        assert_eq!(store.chunk_count(&id).await.unwrap(), 3);
1578    }
1579
1580    #[tokio::test]
1581    async fn chunk_filter_by_stream() {
1582        let store = TaskStore::open_in_memory().await.unwrap();
1583        let id = TaskRunId::new();
1584        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1585        store.append_chunk(&id, 0, Stream::Stdout, b"out").await.unwrap();
1586        store.append_chunk(&id, 1, Stream::Stderr, b"err").await.unwrap();
1587        store.append_chunk(&id, 2, Stream::Stdout, b"out2").await.unwrap();
1588
1589        let chunks = store
1590            .get_chunks(
1591                &id,
1592                &ChunkFilter {
1593                    stream: Some(Stream::Stdout),
1594                    ..Default::default()
1595                },
1596            )
1597            .await
1598            .unwrap();
1599        assert_eq!(chunks.len(), 2);
1600        assert!(chunks.iter().all(|c| c.stream == Stream::Stdout));
1601    }
1602
1603    #[tokio::test]
1604    async fn update_status_to_done() {
1605        let store = TaskStore::open_in_memory().await.unwrap();
1606        let id = TaskRunId::new();
1607        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1608        store
1609            .update_status(
1610                &id,
1611                &RunStatus::Done {
1612                    exit_code: 0,
1613                    ended_at: 2_000_000,
1614                },
1615            )
1616            .await
1617            .unwrap();
1618        let got = store.get_run(&id).await.unwrap().unwrap();
1619        match got.status {
1620            RunStatus::Done { exit_code, .. } => assert_eq!(exit_code, 0),
1621            other => panic!("unexpected status: {other:?}"),
1622        }
1623    }
1624
1625    #[tokio::test]
1626    async fn list_runs_filter_status() {
1627        let store = TaskStore::open_in_memory().await.unwrap();
1628
1629        for cmd in ["a", "b", "c"] {
1630            let id = TaskRunId::new();
1631            store.insert_run(&make_meta(id.clone(), cmd)).await.unwrap();
1632            if cmd == "b" {
1633                store
1634                    .update_status(
1635                        &id,
1636                        &RunStatus::Done {
1637                            exit_code: 0,
1638                            ended_at: 1_000_001,
1639                        },
1640                    )
1641                    .await
1642                    .unwrap();
1643            }
1644        }
1645
1646        let running = store
1647            .list_runs(&RunFilter {
1648                status: Some("running".to_string()),
1649                ..Default::default()
1650            })
1651            .await
1652            .unwrap();
1653        assert_eq!(running.len(), 2);
1654        let done = store
1655            .list_runs(&RunFilter {
1656                status: Some("done".to_string()),
1657                ..Default::default()
1658            })
1659            .await
1660            .unwrap();
1661        assert_eq!(done.len(), 1);
1662    }
1663
1664    // ─── Event tests ──────────────────────────────────────────────────────────
1665
1666    async fn append_evt(
1667        store: &TaskStore,
1668        run_id: &TaskRunId,
1669        level: Level,
1670        target: &str,
1671        msg: &str,
1672        fields: serde_json::Value,
1673    ) -> u32 {
1674        store
1675            .append_event(
1676                run_id,
1677                0,
1678                level,
1679                target,
1680                msg,
1681                &fields,
1682                None,
1683                &EventSource::Beholder {
1684                    name: "cargo".to_string(),
1685                    version: "1.78".to_string(),
1686                },
1687            )
1688            .await
1689            .unwrap()
1690    }
1691
1692    #[tokio::test]
1693    async fn events_schema_has_structural_indexes() {
1694        let store = TaskStore::open_in_memory().await.unwrap();
1695        let n = store
1696            .count_query(
1697                "SELECT COUNT(*) FROM sqlite_master WHERE type='index' \
1698                 AND name IN ('events_by_target', 'events_by_offset')",
1699                vec![],
1700            )
1701            .await
1702            .unwrap();
1703        assert_eq!(n, 2, "both structural events indexes must exist after open");
1704    }
1705
1706    #[tokio::test]
1707    async fn append_and_query_events() {
1708        let store = TaskStore::open_in_memory().await.unwrap();
1709        let id = TaskRunId::new();
1710        store.insert_run(&make_meta(id.clone(), "cargo check")).await.unwrap();
1711
1712        append_evt(&store, &id, Level::Info, "cargo::rustc", "compiling", serde_json::json!({})).await;
1713        append_evt(
1714            &store,
1715            &id,
1716            Level::Warn,
1717            "cargo::rustc",
1718            "unused import",
1719            serde_json::json!({"file": {"path": "src/lib.rs", "line": 10}}),
1720        )
1721        .await;
1722        append_evt(
1723            &store,
1724            &id,
1725            Level::Error,
1726            "cargo::rustc",
1727            "type mismatch",
1728            serde_json::json!({"error": {"code": "E0308"}, "file": {"path": "src/main.rs", "line": 42}}),
1729        )
1730        .await;
1731
1732        let all = store.query_events(&id, &EventFilter::default()).await.unwrap();
1733        assert_eq!(all.len(), 3);
1734
1735        let errors = store
1736            .query_events(
1737                &id,
1738                &EventFilter {
1739                    min_level: Some(Level::Error),
1740                    ..EventFilter::default()
1741                },
1742            )
1743            .await
1744            .unwrap();
1745        assert_eq!(errors.len(), 1);
1746        assert_eq!(errors[0].msg, "type mismatch");
1747
1748        assert_eq!(store.event_count(&id, None).await.unwrap(), 3);
1749        assert_eq!(store.event_count(&id, Some(Level::Warn)).await.unwrap(), 2);
1750        assert_eq!(store.event_count(&id, Some(Level::Error)).await.unwrap(), 1);
1751    }
1752
1753    #[tokio::test]
1754    async fn query_diagnostics_maps_reserved_fields() {
1755        let store = TaskStore::open_in_memory().await.unwrap();
1756        let id = TaskRunId::new();
1757        store.insert_run(&make_meta(id.clone(), "cargo check")).await.unwrap();
1758
1759        append_evt(
1760            &store,
1761            &id,
1762            Level::Error,
1763            "cargo::rustc",
1764            "type mismatch",
1765            serde_json::json!({
1766                "error": {"code": "E0308"},
1767                "file": {"path": "src/main.rs", "line": 42, "col": 5}
1768            }),
1769        )
1770        .await;
1771
1772        let diags = store.query_diagnostics(&id).await.unwrap();
1773        assert_eq!(diags.len(), 1);
1774        assert_eq!(diags[0].file.as_deref(), Some("src/main.rs"));
1775        assert_eq!(diags[0].line, Some(42));
1776        assert_eq!(diags[0].col, Some(5));
1777        assert_eq!(diags[0].code.as_deref(), Some("E0308"));
1778        assert_eq!(diags[0].severity, Level::Error);
1779    }
1780
1781    #[test]
1782    fn field_path_validation_rules() {
1783        assert!(validate_field_path("$.error.code").is_ok());
1784        assert!(validate_field_path("$.file.path").is_ok());
1785        assert!(validate_field_path("$.test_name").is_ok());
1786        assert!(validate_field_path("$.items[0]").is_ok());
1787
1788        assert!(validate_field_path("").is_err());
1789        assert!(validate_field_path("error.code").is_err());
1790        assert!(validate_field_path("$.'injection").is_err());
1791        assert!(validate_field_path("$.error code").is_err());
1792        assert!(validate_field_path("$.a;b").is_err());
1793    }
1794
1795    #[tokio::test]
1796    async fn ensure_field_index_idempotent() {
1797        let store = TaskStore::open_in_memory().await.unwrap();
1798
1799        store.ensure_field_index("$.error.code").await.unwrap();
1800        store.ensure_field_index("$.error.code").await.unwrap();
1801
1802        let indexes = store.list_field_indexes().await.unwrap();
1803        assert_eq!(indexes.len(), 1);
1804        assert_eq!(indexes[0].field_path, "$.error.code");
1805        assert!(indexes[0].index_name.starts_with("events_field_"));
1806    }
1807
1808    #[tokio::test]
1809    async fn field_index_accelerates_field_filter_query() {
1810        let store = TaskStore::open_in_memory().await.unwrap();
1811        let id = TaskRunId::new();
1812        store.insert_run(&make_meta(id.clone(), "cargo check")).await.unwrap();
1813
1814        for i in 0u32..20 {
1815            append_evt(
1816                &store,
1817                &id,
1818                Level::Error,
1819                "t",
1820                "e",
1821                serde_json::json!({"error": {"code": if i % 3 == 0 { "E0308" } else { "E0001" }}}),
1822            )
1823            .await;
1824        }
1825
1826        store.ensure_field_index("$.error.code").await.unwrap();
1827
1828        let results = store
1829            .query_events(
1830                &id,
1831                &EventFilter {
1832                    field_filter: Some(FieldFilter {
1833                        path: "$.error.code".to_string(),
1834                        value: serde_json::json!("E0308"),
1835                    }),
1836                    ..EventFilter::default()
1837                },
1838            )
1839            .await
1840            .unwrap();
1841        assert_eq!(results.len(), 7);
1842    }
1843
1844    #[tokio::test]
1845    async fn timeline_interleaves_chunks_and_events_by_offset() {
1846        let store = TaskStore::open_in_memory().await.unwrap();
1847        let id = TaskRunId::new();
1848        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1849
1850        store.append_chunk(&id, 10, Stream::Stdout, b"first\n").await.unwrap();
1851        append_evt(&store, &id, Level::Info, "t", "ev1", serde_json::json!({})).await;
1852        store
1853            .append_event(
1854                &id,
1855                20,
1856                Level::Warn,
1857                "t",
1858                "ev2",
1859                &serde_json::json!({}),
1860                None,
1861                &EventSource::Beholder {
1862                    name: "test".into(),
1863                    version: "0".into(),
1864                },
1865            )
1866            .await
1867            .unwrap();
1868        store.append_chunk(&id, 30, Stream::Stdout, b"last\n").await.unwrap();
1869
1870        let ticks = store.timeline_ticks(&id, None, None).await.unwrap();
1871        assert_eq!(ticks.len(), 4);
1872        assert!(matches!(ticks[0], TimelineTick::Event(_)));
1873        assert!(matches!(ticks[1], TimelineTick::Chunk(_)));
1874        assert!(matches!(ticks[2], TimelineTick::Event(_)));
1875        assert!(matches!(ticks[3], TimelineTick::Chunk(_)));
1876    }
1877
1878    #[tokio::test]
1879    async fn timeline_since_offset_filters_correctly() {
1880        let store = TaskStore::open_in_memory().await.unwrap();
1881        let id = TaskRunId::new();
1882        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1883        store.append_chunk(&id, 5, Stream::Stdout, b"a").await.unwrap();
1884        store.append_chunk(&id, 15, Stream::Stdout, b"b").await.unwrap();
1885
1886        let ticks = store.timeline_ticks(&id, Some(10), None).await.unwrap();
1887        assert_eq!(ticks.len(), 1);
1888        if let TimelineTick::Chunk(c) = &ticks[0] {
1889            assert_eq!(c.offset_ms, 15);
1890        } else {
1891            panic!("expected Chunk");
1892        }
1893    }
1894
1895    #[tokio::test]
1896    async fn aggregate_events_group_by_target() {
1897        let store = TaskStore::open_in_memory().await.unwrap();
1898        let id = TaskRunId::new();
1899        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1900
1901        for _ in 0..3 {
1902            append_evt(&store, &id, Level::Warn, "cargo::rustc", "w", serde_json::json!({})).await;
1903        }
1904        for _ in 0..2 {
1905            append_evt(&store, &id, Level::Error, "tsc", "e", serde_json::json!({})).await;
1906        }
1907
1908        let buckets = store
1909            .aggregate_events(&AggregateFilter {
1910                group_by: Some(AggregateGroupBy::Target),
1911                ..Default::default()
1912            })
1913            .await
1914            .unwrap();
1915
1916        assert_eq!(buckets.len(), 2);
1917        assert_eq!(buckets[0].key, "cargo::rustc");
1918        assert_eq!(buckets[0].count, 3);
1919        assert_eq!(buckets[1].key, "tsc");
1920        assert_eq!(buckets[1].count, 2);
1921    }
1922
1923    #[tokio::test]
1924    async fn aggregate_events_group_by_error_code() {
1925        let store = TaskStore::open_in_memory().await.unwrap();
1926        let id = TaskRunId::new();
1927        store.insert_run(&make_meta(id.clone(), "cargo check")).await.unwrap();
1928
1929        append_evt(
1930            &store,
1931            &id,
1932            Level::Error,
1933            "cargo::rustc",
1934            "e0308",
1935            serde_json::json!({"error": {"code": "E0308"}}),
1936        )
1937        .await;
1938        append_evt(
1939            &store,
1940            &id,
1941            Level::Error,
1942            "cargo::rustc",
1943            "e0308 again",
1944            serde_json::json!({"error": {"code": "E0308"}}),
1945        )
1946        .await;
1947        append_evt(
1948            &store,
1949            &id,
1950            Level::Error,
1951            "cargo::rustc",
1952            "e0001",
1953            serde_json::json!({"error": {"code": "E0001"}}),
1954        )
1955        .await;
1956        append_evt(&store, &id, Level::Info, "cargo", "no code", serde_json::json!({})).await;
1957
1958        let buckets = store
1959            .aggregate_events(&AggregateFilter {
1960                group_by: Some(AggregateGroupBy::ErrorCode),
1961                ..Default::default()
1962            })
1963            .await
1964            .unwrap();
1965
1966        assert_eq!(buckets.len(), 2, "only events with error.code are counted");
1967        assert_eq!(buckets[0].key, "E0308");
1968        assert_eq!(buckets[0].count, 2);
1969        assert_eq!(buckets[1].key, "E0001");
1970        assert_eq!(buckets[1].count, 1);
1971    }
1972
1973    #[tokio::test]
1974    async fn event_filter_offset_range() {
1975        let store = TaskStore::open_in_memory().await.unwrap();
1976        let id = TaskRunId::new();
1977        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
1978
1979        store
1980            .append_event(
1981                &id,
1982                5,
1983                Level::Info,
1984                "t",
1985                "early",
1986                &serde_json::json!({}),
1987                None,
1988                &EventSource::Beholder {
1989                    name: "x".into(),
1990                    version: "0".into(),
1991                },
1992            )
1993            .await
1994            .unwrap();
1995        store
1996            .append_event(
1997                &id,
1998                50,
1999                Level::Warn,
2000                "t",
2001                "mid",
2002                &serde_json::json!({}),
2003                None,
2004                &EventSource::Beholder {
2005                    name: "x".into(),
2006                    version: "0".into(),
2007                },
2008            )
2009            .await
2010            .unwrap();
2011        store
2012            .append_event(
2013                &id,
2014                100,
2015                Level::Error,
2016                "t",
2017                "late",
2018                &serde_json::json!({}),
2019                None,
2020                &EventSource::Beholder {
2021                    name: "x".into(),
2022                    version: "0".into(),
2023                },
2024            )
2025            .await
2026            .unwrap();
2027
2028        let events = store
2029            .query_events(
2030                &id,
2031                &EventFilter {
2032                    offset_range: Some((10, 60)),
2033                    ..Default::default()
2034                },
2035            )
2036            .await
2037            .unwrap();
2038
2039        assert_eq!(events.len(), 1);
2040        assert_eq!(events[0].msg, "mid");
2041    }
2042
2043    #[tokio::test]
2044    async fn upsert_and_get_triage_round_trips() {
2045        use crate::types::{KeepRange, SeqRange, Triage};
2046        let store = TaskStore::open_in_memory().await.unwrap();
2047        let id = TaskRunId::new();
2048        store.insert_run(&make_meta(id.clone(), "cargo check")).await.unwrap();
2049
2050        let triage = Triage {
2051            run_id: id.clone(),
2052            synopsis: "Build failed: missing semicolon on line 42.".into(),
2053            keep: vec![KeepRange {
2054                range: SeqRange { lo: 10, hi: 15 },
2055                reason: "primary error".into(),
2056            }],
2057            primary: SeqRange { lo: 10, hi: 15 },
2058            model: "claude-haiku-4-5".into(),
2059            prompt_version: 1,
2060            cached_at: 1_700_000_000,
2061            partial: false,
2062        };
2063
2064        store.upsert_triage(&triage).await.unwrap();
2065        let got = store.get_triage(&id).await.unwrap().expect("triage should exist");
2066
2067        assert_eq!(got.run_id.to_string(), id.to_string());
2068        assert_eq!(got.synopsis, triage.synopsis);
2069        assert_eq!(got.keep.len(), 1);
2070        assert_eq!(got.keep[0].range.lo, 10);
2071        assert_eq!(got.keep[0].range.hi, 15);
2072        assert_eq!(got.keep[0].reason, "primary error");
2073        assert_eq!(got.primary.lo, 10);
2074        assert_eq!(got.primary.hi, 15);
2075        assert_eq!(got.model, "claude-haiku-4-5");
2076        assert_eq!(got.prompt_version, 1);
2077        assert_eq!(got.cached_at, 1_700_000_000);
2078        assert!(!got.partial);
2079    }
2080
2081    #[tokio::test]
2082    async fn get_triage_returns_none_when_absent() {
2083        let store = TaskStore::open_in_memory().await.unwrap();
2084        let id = TaskRunId::new();
2085        assert!(store.get_triage(&id).await.unwrap().is_none());
2086    }
2087
2088    #[tokio::test]
2089    async fn upsert_triage_replace_on_conflict() {
2090        use crate::types::{KeepRange, SeqRange, Triage};
2091        let store = TaskStore::open_in_memory().await.unwrap();
2092        let id = TaskRunId::new();
2093        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
2094
2095        let t1 = Triage {
2096            run_id: id.clone(),
2097            synopsis: "first".into(),
2098            keep: vec![],
2099            primary: SeqRange { lo: 0, hi: 0 },
2100            model: "haiku".into(),
2101            prompt_version: 1,
2102            cached_at: 100,
2103            partial: false,
2104        };
2105        let t2 = Triage {
2106            run_id: id.clone(),
2107            synopsis: "second".into(),
2108            keep: vec![KeepRange {
2109                range: SeqRange { lo: 5, hi: 9 },
2110                reason: "r".into(),
2111            }],
2112            primary: SeqRange { lo: 5, hi: 9 },
2113            model: "ollama:qwen".into(),
2114            prompt_version: 2,
2115            cached_at: 200,
2116            partial: true,
2117        };
2118        store.upsert_triage(&t1).await.unwrap();
2119        store.upsert_triage(&t2).await.unwrap();
2120
2121        let got = store.get_triage(&id).await.unwrap().unwrap();
2122        assert_eq!(got.synopsis, "second");
2123        assert_eq!(got.prompt_version, 2);
2124        assert!(got.partial);
2125    }
2126
2127    #[tokio::test]
2128    async fn event_seq_resumes_after_reopen() {
2129        let dir = tempfile::tempdir().unwrap();
2130        let path = dir.path().join("tr.turso");
2131        let id = TaskRunId::new();
2132
2133        {
2134            let store = TaskStore::open(&path).await.unwrap();
2135            store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
2136            append_evt(&store, &id, Level::Info, "t", "a", serde_json::json!({})).await;
2137            append_evt(&store, &id, Level::Info, "t", "b", serde_json::json!({})).await;
2138        }
2139
2140        let store = TaskStore::open(&path).await.unwrap();
2141        let seq = append_evt(&store, &id, Level::Info, "t", "c", serde_json::json!({})).await;
2142        assert_eq!(seq, 2);
2143        assert_eq!(store.event_count(&id, None).await.unwrap(), 3);
2144    }
2145
2146    // ─── GC sweep tests ───────────────────────────────────────────────────────
2147
2148    fn make_meta_with_age(id: TaskRunId, cmd: &str, age_secs: u64) -> TaskRunMeta {
2149        let now = unix_now();
2150        TaskRunMeta {
2151            id,
2152            command: cmd.to_string(),
2153            cwd: PathBuf::from("/tmp"),
2154            env: vec![],
2155            started_at: now.saturating_sub(age_secs),
2156            status: RunStatus::Running,
2157            label: None,
2158            initiator: Initiator::Human {
2159                camp: "local".to_string(),
2160            },
2161            beholder_status: None,
2162            pinned: false,
2163            origin: None,
2164            host_pid: None,
2165        }
2166    }
2167
2168    #[tokio::test]
2169    async fn gc_sweep_drops_output_for_archived_runs() {
2170        let store = TaskStore::open_in_memory().await.unwrap();
2171        let id = TaskRunId::new();
2172        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
2173        store.append_chunk(&id, 0, Stream::Stdout, b"data").await.unwrap();
2174        store.archive_run(&id).await.unwrap();
2175
2176        let result = store.gc_sweep(&GcConfig::default()).await.unwrap();
2177        assert_eq!(result.archived_runs_cleaned, 1);
2178        assert_eq!(result.chunks_deleted, 1);
2179
2180        assert!(store.get_run(&id).await.unwrap().is_some());
2181        assert_eq!(store.chunk_count(&id).await.unwrap(), 0);
2182    }
2183
2184    #[tokio::test]
2185    async fn gc_sweep_warm_rolloff_drops_old_output() {
2186        let store = TaskStore::open_in_memory().await.unwrap();
2187
2188        let old_id = TaskRunId::new();
2189        let old_meta = make_meta_with_age(old_id.clone(), "cmd", 31 * 24 * 3600);
2190        store.insert_run(&old_meta).await.unwrap();
2191        store.append_chunk(&old_id, 0, Stream::Stdout, b"old").await.unwrap();
2192
2193        let new_id = TaskRunId::new();
2194        let new_meta = make_meta_with_age(new_id.clone(), "cmd", 60);
2195        store.insert_run(&new_meta).await.unwrap();
2196        store.append_chunk(&new_id, 0, Stream::Stdout, b"new").await.unwrap();
2197
2198        let result = store.gc_sweep(&GcConfig::default()).await.unwrap();
2199        assert_eq!(result.warm_rolloff_runs, 1);
2200        assert_eq!(result.chunks_deleted, 1);
2201
2202        assert!(store.get_run(&old_id).await.unwrap().is_some());
2203        assert_eq!(store.chunk_count(&old_id).await.unwrap(), 0);
2204        assert_eq!(store.chunk_count(&new_id).await.unwrap(), 1);
2205    }
2206
2207    #[tokio::test]
2208    async fn gc_sweep_pin_exemption() {
2209        let store = TaskStore::open_in_memory().await.unwrap();
2210
2211        let pinned_id = TaskRunId::new();
2212        let mut pinned_meta = make_meta_with_age(pinned_id.clone(), "cmd", 31 * 24 * 3600);
2213        pinned_meta.pinned = true;
2214        store.insert_run(&pinned_meta).await.unwrap();
2215        store.append_chunk(&pinned_id, 0, Stream::Stdout, b"pinned").await.unwrap();
2216
2217        let result = store.gc_sweep(&GcConfig::default()).await.unwrap();
2218        assert_eq!(result.warm_rolloff_runs, 0, "pinned run must not count toward warm rolloff");
2219        assert_eq!(result.chunks_deleted, 0, "pinned run output must survive gc_sweep");
2220        assert_eq!(store.chunk_count(&pinned_id).await.unwrap(), 1);
2221    }
2222
2223    #[tokio::test]
2224    async fn pin_run_toggles_pin_flag() {
2225        let store = TaskStore::open_in_memory().await.unwrap();
2226        let id = TaskRunId::new();
2227        store.insert_run(&make_meta(id.clone(), "cmd")).await.unwrap();
2228
2229        store.pin_run(&id, true).await.unwrap();
2230        assert!(store.get_run(&id).await.unwrap().unwrap().pinned);
2231
2232        store.pin_run(&id, false).await.unwrap();
2233        assert!(!store.get_run(&id).await.unwrap().unwrap().pinned);
2234    }
2235}