1use 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
35const BUSY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
37const BUSY_RETRY_BASE: std::time::Duration = std::time::Duration::from_millis(10);
39const BUSY_RETRY_MAX: std::time::Duration = std::time::Duration::from_secs(2);
41
42fn is_busy(e: &turso::Error) -> bool {
45 matches!(e, turso::Error::Busy(_) | turso::Error::BusySnapshot(_))
46}
47
48#[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
66const 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#[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 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
155struct 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 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 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 #[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 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 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 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 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 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 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 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
1065pub 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#[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#[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
1151const 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#[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 #[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 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 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 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}