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