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