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