1use 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#[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
46const 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#[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 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
135struct 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 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 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 #[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 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 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 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 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 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 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
1013pub 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#[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#[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
1099const 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#[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 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 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}