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