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