1use super::{
39 DegradationEntry, RunId, RunListFilter, RunRecord, RunStatus, RunStore, RunStoreError,
40 StepEntry, TaskId,
41};
42use async_trait::async_trait;
43use rusqlite::{params, OptionalExtension};
44use rusqlite_isle::{AsyncIsle, AsyncIsleDriver, IsleError};
45use std::path::Path;
46
47const SCHEMA_SQL: &str = "\
48CREATE TABLE IF NOT EXISTS runs (\
49 id TEXT PRIMARY KEY, \
50 task_id TEXT NOT NULL, \
51 status TEXT NOT NULL, \
52 step_entries_json TEXT NOT NULL, \
53 degradations_json TEXT NOT NULL DEFAULT '[]', \
54 operator_sid TEXT, \
55 result_ref_json TEXT, \
56 input_json TEXT, \
57 created_at INTEGER NOT NULL, \
58 updated_at INTEGER NOT NULL\
59);\
60CREATE INDEX IF NOT EXISTS ix_runs_task_id ON runs(task_id, created_at);\
61";
62
63fn migrate_add_column_if_missing(
69 conn: &rusqlite::Connection,
70 column: &str,
71 decl: &str,
72) -> rusqlite::Result<()> {
73 let mut stmt = conn.prepare("PRAGMA table_info(runs)")?;
74 let has_column = stmt
75 .query_map([], |row| row.get::<_, String>(1))?
76 .collect::<Result<Vec<String>, _>>()?
77 .iter()
78 .any(|name| name == column);
79 if !has_column {
80 conn.execute_batch(&format!("ALTER TABLE runs ADD COLUMN {column} {decl};"))?;
81 }
82 Ok(())
83}
84
85pub struct SqliteRunStore {
93 isle: AsyncIsle,
94}
95
96impl SqliteRunStore {
97 pub async fn open(path: impl AsRef<Path>) -> Result<(Self, AsyncIsleDriver), RunStoreError> {
100 let (isle, driver) = AsyncIsle::spawn(path.as_ref().to_path_buf(), |conn| {
101 conn.busy_timeout(std::time::Duration::from_millis(5_000))?;
106 conn.execute_batch(SCHEMA_SQL)?;
107 migrate_add_column_if_missing(conn, "degradations_json", "TEXT NOT NULL DEFAULT '[]'")?;
108 migrate_add_column_if_missing(conn, "input_json", "TEXT")
109 })
110 .await
111 .map_err(map_isle_err)?;
112 Ok((Self { isle }, driver))
113 }
114
115 pub async fn open_in_memory() -> Result<(Self, AsyncIsleDriver), RunStoreError> {
117 let (isle, driver) = AsyncIsle::open_in_memory(|conn| {
118 conn.busy_timeout(std::time::Duration::from_millis(5_000))?;
119 conn.execute_batch(SCHEMA_SQL)?;
120 migrate_add_column_if_missing(conn, "degradations_json", "TEXT NOT NULL DEFAULT '[]'")?;
121 migrate_add_column_if_missing(conn, "input_json", "TEXT")
122 })
123 .await
124 .map_err(map_isle_err)?;
125 Ok((Self { isle }, driver))
126 }
127}
128
129fn map_isle_err(e: IsleError) -> RunStoreError {
130 RunStoreError::Other(format!("sqlite: {e}"))
131}
132
133type RunRow = (
137 String,
138 String,
139 String,
140 String,
141 String,
142 Option<String>,
143 Option<String>,
144 Option<String>,
145 i64,
146 i64,
147);
148
149const RUN_SELECT_COLUMNS: &str = "id, task_id, status, step_entries_json, degradations_json, \
150 operator_sid, result_ref_json, input_json, created_at, updated_at";
151
152fn row_to_record(row: RunRow) -> Result<RunRecord, RunStoreError> {
153 let (
154 id,
155 task_id,
156 status_json,
157 step_entries_json,
158 degradations_json,
159 operator_sid,
160 result_ref_json,
161 input_json,
162 created_at,
163 updated_at,
164 ) = row;
165 let status: RunStatus = serde_json::from_str(&status_json)
166 .map_err(|e| RunStoreError::Other(format!("decode status: {e}")))?;
167 let step_entries: Vec<StepEntry> = serde_json::from_str(&step_entries_json)
168 .map_err(|e| RunStoreError::Other(format!("decode step_entries: {e}")))?;
169 let degradations: Vec<DegradationEntry> = serde_json::from_str(°radations_json)
170 .map_err(|e| RunStoreError::Other(format!("decode degradations: {e}")))?;
171 let result_ref: Option<serde_json::Value> = match result_ref_json {
172 Some(text) => Some(
173 serde_json::from_str(&text)
174 .map_err(|e| RunStoreError::Other(format!("decode result_ref: {e}")))?,
175 ),
176 None => None,
177 };
178 let id = RunId::parse(id).map_err(|e| RunStoreError::Other(format!("decode id: {e}")))?;
182 let task_id =
183 TaskId::parse(task_id).map_err(|e| RunStoreError::Other(format!("decode task_id: {e}")))?;
184 Ok(RunRecord {
185 id,
186 task_id,
187 status,
188 step_entries,
189 degradations,
190 operator_sid,
191 result_ref,
192 input_json,
193 created_at: created_at as u64,
194 updated_at: updated_at as u64,
195 })
196}
197
198#[async_trait]
199impl RunStore for SqliteRunStore {
200 fn name(&self) -> &str {
201 "sqlite"
202 }
203
204 async fn create(&self, record: RunRecord) -> Result<(), RunStoreError> {
205 let id = record.id.to_string();
206 let id_for_conflict = record.id.clone();
207 let task_id = record.task_id.to_string();
208 let status_json = serde_json::to_string(&record.status)
209 .map_err(|e| RunStoreError::Other(format!("encode status: {e}")))?;
210 let step_entries_json = serde_json::to_string(&record.step_entries)
211 .map_err(|e| RunStoreError::Other(format!("encode step_entries: {e}")))?;
212 let degradations_json = serde_json::to_string(&record.degradations)
213 .map_err(|e| RunStoreError::Other(format!("encode degradations: {e}")))?;
214 let operator_sid = record.operator_sid.clone();
215 let result_ref_json = record
216 .result_ref
217 .as_ref()
218 .map(serde_json::to_string)
219 .transpose()
220 .map_err(|e| RunStoreError::Other(format!("encode result_ref: {e}")))?;
221 let input_json = record.input_json.clone();
222 let created_at = record.created_at as i64;
223 let updated_at = record.updated_at as i64;
224
225 self.isle
226 .call(move |conn| {
227 let tx =
232 conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
233 let exists: i64 = tx.query_row(
234 "SELECT COUNT(*) FROM runs WHERE id = ?1",
235 params![id],
236 |row| row.get(0),
237 )?;
238 if exists > 0 {
239 return Err(rusqlite::Error::SqliteFailure(
240 rusqlite::ffi::Error::new(rusqlite::ffi::SQLITE_CONSTRAINT),
241 Some(format!("__mlua_swarm_duplicate:{id}")),
242 ));
243 }
244 tx.execute(
245 "INSERT INTO runs (id, task_id, status, step_entries_json, \
246 degradations_json, operator_sid, result_ref_json, input_json, \
247 created_at, updated_at) \
248 VALUES (?1, ?2, ?3, ?4, ?5, ?6, ?7, ?8, ?9, ?10)",
249 params![
250 id,
251 task_id,
252 status_json,
253 step_entries_json,
254 degradations_json,
255 operator_sid,
256 result_ref_json,
257 input_json,
258 created_at,
259 updated_at,
260 ],
261 )?;
262 tx.commit()?;
263 Ok(())
264 })
265 .await
266 .map_err(|e| match &e {
267 IsleError::Sqlite(rusqlite::Error::SqliteFailure(_, Some(msg)))
268 if msg.starts_with("__mlua_swarm_duplicate:") =>
269 {
270 RunStoreError::Duplicate(id_for_conflict.clone())
271 }
272 _ => map_isle_err(e),
273 })
274 }
275
276 async fn get(&self, id: &RunId) -> Result<RunRecord, RunStoreError> {
277 let id_str = id.to_string();
278 let id_for_notfound = id.clone();
279 let row = self
280 .isle
281 .call(move |conn| {
282 conn.query_row(
283 &format!("SELECT {RUN_SELECT_COLUMNS} FROM runs WHERE id = ?1"),
284 params![id_str],
285 |row| {
286 Ok((
287 row.get::<_, String>(0)?,
288 row.get::<_, String>(1)?,
289 row.get::<_, String>(2)?,
290 row.get::<_, String>(3)?,
291 row.get::<_, String>(4)?,
292 row.get::<_, Option<String>>(5)?,
293 row.get::<_, Option<String>>(6)?,
294 row.get::<_, Option<String>>(7)?,
295 row.get::<_, i64>(8)?,
296 row.get::<_, i64>(9)?,
297 ))
298 },
299 )
300 .optional()
301 })
302 .await
303 .map_err(map_isle_err)?;
304 match row {
305 Some(row) => row_to_record(row),
306 None => Err(RunStoreError::NotFound(id_for_notfound)),
307 }
308 }
309
310 async fn list_by_task(&self, task_id: &TaskId) -> Result<Vec<RunRecord>, RunStoreError> {
311 let task_id_str = task_id.to_string();
312 let rows = self
313 .isle
314 .call(move |conn| {
315 let mut stmt = conn.prepare(&format!(
316 "SELECT {RUN_SELECT_COLUMNS} FROM runs \
317 WHERE task_id = ?1 ORDER BY created_at ASC"
318 ))?;
319 let iter = stmt.query_map(params![task_id_str], |row| {
320 Ok((
321 row.get::<_, String>(0)?,
322 row.get::<_, String>(1)?,
323 row.get::<_, String>(2)?,
324 row.get::<_, String>(3)?,
325 row.get::<_, String>(4)?,
326 row.get::<_, Option<String>>(5)?,
327 row.get::<_, Option<String>>(6)?,
328 row.get::<_, Option<String>>(7)?,
329 row.get::<_, i64>(8)?,
330 row.get::<_, i64>(9)?,
331 ))
332 })?;
333 let mut out = Vec::new();
334 for r in iter {
335 out.push(r?);
336 }
337 Ok(out)
338 })
339 .await
340 .map_err(map_isle_err)?;
341 rows.into_iter().map(row_to_record).collect()
342 }
343
344 async fn append_step_entry(&self, id: &RunId, entry: StepEntry) -> Result<(), RunStoreError> {
345 let id_str = id.to_string();
346 let id_for_notfound = id.clone();
347 let updated_at = crate::types::now_unix() as i64;
348
349 let updated = self
350 .isle
351 .call(move |conn| {
352 let tx =
355 conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
356 let existing: Option<String> = tx
357 .query_row(
358 "SELECT step_entries_json FROM runs WHERE id = ?1",
359 params![id_str],
360 |row| row.get(0),
361 )
362 .optional()?;
363 let Some(existing_json) = existing else {
364 return Ok(false);
365 };
366 let mut entries: Vec<StepEntry> = serde_json::from_str(&existing_json)
367 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
368 entries.push(entry);
369 let new_json = serde_json::to_string(&entries)
370 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
371 tx.execute(
372 "UPDATE runs SET step_entries_json = ?1, updated_at = ?2 WHERE id = ?3",
373 params![new_json, updated_at, id_str],
374 )?;
375 tx.commit()?;
376 Ok(true)
377 })
378 .await
379 .map_err(map_isle_err)?;
380
381 if updated {
382 Ok(())
383 } else {
384 Err(RunStoreError::NotFound(id_for_notfound))
385 }
386 }
387
388 async fn append_degradation(
389 &self,
390 id: &RunId,
391 entry: DegradationEntry,
392 ) -> Result<(), RunStoreError> {
393 let id_str = id.to_string();
394 let id_for_notfound = id.clone();
395 let updated_at = crate::types::now_unix() as i64;
396
397 let updated = self
398 .isle
399 .call(move |conn| {
400 let tx =
403 conn.transaction_with_behavior(rusqlite::TransactionBehavior::Immediate)?;
404 let existing: Option<String> = tx
405 .query_row(
406 "SELECT degradations_json FROM runs WHERE id = ?1",
407 params![id_str],
408 |row| row.get(0),
409 )
410 .optional()?;
411 let Some(existing_json) = existing else {
412 return Ok(false);
413 };
414 let mut entries: Vec<DegradationEntry> = serde_json::from_str(&existing_json)
415 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
416 entries.push(entry);
417 let new_json = serde_json::to_string(&entries)
418 .map_err(|e| rusqlite::Error::ToSqlConversionFailure(Box::new(e)))?;
419 tx.execute(
420 "UPDATE runs SET degradations_json = ?1, updated_at = ?2 WHERE id = ?3",
421 params![new_json, updated_at, id_str],
422 )?;
423 tx.commit()?;
424 Ok(true)
425 })
426 .await
427 .map_err(map_isle_err)?;
428
429 if updated {
430 Ok(())
431 } else {
432 Err(RunStoreError::NotFound(id_for_notfound))
433 }
434 }
435
436 async fn update_status(&self, id: &RunId, status: RunStatus) -> Result<(), RunStoreError> {
437 let id_str = id.to_string();
438 let id_for_notfound = id.clone();
439 let status_json = serde_json::to_string(&status)
440 .map_err(|e| RunStoreError::Other(format!("encode status: {e}")))?;
441 let updated_at = crate::types::now_unix() as i64;
442 let n = self
443 .isle
444 .call(move |conn| {
445 conn.execute(
446 "UPDATE runs SET status = ?1, updated_at = ?2 WHERE id = ?3",
447 params![status_json, updated_at, id_str],
448 )
449 })
450 .await
451 .map_err(map_isle_err)?;
452 if n == 0 {
453 Err(RunStoreError::NotFound(id_for_notfound))
454 } else {
455 Ok(())
456 }
457 }
458
459 async fn try_transition(
460 &self,
461 id: &RunId,
462 from: RunStatus,
463 to: RunStatus,
464 ) -> Result<bool, RunStoreError> {
465 let id_str = id.to_string();
466 let from_json = serde_json::to_string(&from)
467 .map_err(|e| RunStoreError::Other(format!("encode from status: {e}")))?;
468 let to_json = serde_json::to_string(&to)
469 .map_err(|e| RunStoreError::Other(format!("encode to status: {e}")))?;
470 let updated_at = crate::types::now_unix() as i64;
471 let n = self
477 .isle
478 .call(move |conn| {
479 conn.execute(
480 "UPDATE runs SET status = ?1, updated_at = ?2 WHERE id = ?3 AND status = ?4",
481 params![to_json, updated_at, id_str, from_json],
482 )
483 })
484 .await
485 .map_err(map_isle_err)?;
486 Ok(n == 1)
487 }
488
489 async fn set_result(
490 &self,
491 id: &RunId,
492 result_ref: serde_json::Value,
493 ) -> Result<(), RunStoreError> {
494 let id_str = id.to_string();
495 let id_for_notfound = id.clone();
496 let result_ref_json = serde_json::to_string(&result_ref)
497 .map_err(|e| RunStoreError::Other(format!("encode result_ref: {e}")))?;
498 let updated_at = crate::types::now_unix() as i64;
499 let n = self
500 .isle
501 .call(move |conn| {
502 conn.execute(
503 "UPDATE runs SET result_ref_json = ?1, updated_at = ?2 WHERE id = ?3",
504 params![result_ref_json, updated_at, id_str],
505 )
506 })
507 .await
508 .map_err(map_isle_err)?;
509 if n == 0 {
510 Err(RunStoreError::NotFound(id_for_notfound))
511 } else {
512 Ok(())
513 }
514 }
515
516 async fn set_input_json(&self, id: &RunId, input_json: String) -> Result<(), RunStoreError> {
517 let id_str = id.to_string();
518 let id_for_notfound = id.clone();
519 let updated_at = crate::types::now_unix() as i64;
520 let n = self
521 .isle
522 .call(move |conn| {
523 conn.execute(
524 "UPDATE runs SET input_json = ?1, updated_at = ?2 WHERE id = ?3",
525 params![input_json, updated_at, id_str],
526 )
527 })
528 .await
529 .map_err(map_isle_err)?;
530 if n == 0 {
531 Err(RunStoreError::NotFound(id_for_notfound))
532 } else {
533 Ok(())
534 }
535 }
536
537 async fn list_running(&self) -> Result<Vec<RunRecord>, RunStoreError> {
538 let status_json = serde_json::to_string(&RunStatus::Running)
539 .map_err(|e| RunStoreError::Other(format!("encode status: {e}")))?;
540 let rows = self
541 .isle
542 .call(move |conn| {
543 let mut stmt = conn.prepare(&format!(
544 "SELECT {RUN_SELECT_COLUMNS} FROM runs WHERE status = ?1"
545 ))?;
546 let iter = stmt.query_map(params![status_json], |row| {
547 Ok((
548 row.get::<_, String>(0)?,
549 row.get::<_, String>(1)?,
550 row.get::<_, String>(2)?,
551 row.get::<_, String>(3)?,
552 row.get::<_, String>(4)?,
553 row.get::<_, Option<String>>(5)?,
554 row.get::<_, Option<String>>(6)?,
555 row.get::<_, Option<String>>(7)?,
556 row.get::<_, i64>(8)?,
557 row.get::<_, i64>(9)?,
558 ))
559 })?;
560 let mut out = Vec::new();
561 for r in iter {
562 out.push(r?);
563 }
564 Ok(out)
565 })
566 .await
567 .map_err(map_isle_err)?;
568 rows.into_iter().map(row_to_record).collect()
569 }
570
571 async fn list(&self, filter: &RunListFilter) -> Result<Vec<RunRecord>, RunStoreError> {
572 let task_id = filter.task_id.as_ref().map(|t| t.to_string());
573 let status_json = filter
574 .status
575 .map(|s| serde_json::to_string(&s))
576 .transpose()
577 .map_err(|e| RunStoreError::Other(format!("encode status: {e}")))?;
578 let limit = filter.limit.map(|l| l as i64).unwrap_or(-1);
579 let offset = filter.offset.map(|o| o as i64).unwrap_or(0);
580 let rows = self
581 .isle
582 .call(move |conn| {
583 let mut stmt = conn.prepare(&format!(
587 "SELECT {RUN_SELECT_COLUMNS} FROM runs \
588 WHERE (?1 IS NULL OR task_id = ?1) \
589 AND (?2 IS NULL OR status = ?2) \
590 ORDER BY created_at DESC, rowid DESC \
591 LIMIT ?3 OFFSET ?4"
592 ))?;
593 let iter = stmt.query_map(params![task_id, status_json, limit, offset], |row| {
594 Ok((
595 row.get::<_, String>(0)?,
596 row.get::<_, String>(1)?,
597 row.get::<_, String>(2)?,
598 row.get::<_, String>(3)?,
599 row.get::<_, String>(4)?,
600 row.get::<_, Option<String>>(5)?,
601 row.get::<_, Option<String>>(6)?,
602 row.get::<_, Option<String>>(7)?,
603 row.get::<_, i64>(8)?,
604 row.get::<_, i64>(9)?,
605 ))
606 })?;
607 let mut out = Vec::new();
608 for r in iter {
609 out.push(r?);
610 }
611 Ok(out)
612 })
613 .await
614 .map_err(map_isle_err)?;
615 rows.into_iter().map(row_to_record).collect()
616 }
617
618 async fn delete(&self, id: &RunId) -> Result<(), RunStoreError> {
619 let id_str = id.to_string();
620 let id_for_notfound = id.clone();
621 let n = self
622 .isle
623 .call(move |conn| conn.execute("DELETE FROM runs WHERE id = ?1", params![id_str]))
624 .await
625 .map_err(map_isle_err)?;
626 if n == 0 {
627 Err(RunStoreError::NotFound(id_for_notfound))
628 } else {
629 Ok(())
630 }
631 }
632}
633
634#[cfg(test)]
639mod tests {
640 use super::*;
641 use serde_json::json;
642
643 fn mk(id: &str, task_id: &str, created_at: u64) -> RunRecord {
644 RunRecord {
645 id: RunId::parse(id).unwrap(),
646 task_id: TaskId::parse(task_id).unwrap(),
647 status: RunStatus::Pending,
648 step_entries: vec![],
649 degradations: vec![],
650 operator_sid: None,
651 result_ref: None,
652 input_json: None,
653 created_at,
654 updated_at: created_at,
655 }
656 }
657
658 fn mk_degradation(tool: &str, at: u64) -> DegradationEntry {
659 DegradationEntry {
660 tool: tool.to_string(),
661 error: "boom".to_string(),
662 fallback: "cached-default".to_string(),
663 note: None,
664 step_ref: Some("worker".to_string()),
665 attempt: Some(1),
666 at,
667 }
668 }
669
670 #[tokio::test]
671 async fn create_then_get() {
672 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
673 s.create(mk("R-1", "T-1", 100)).await.unwrap();
674 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
675 assert_eq!(got.task_id, TaskId::parse("T-1").unwrap());
676 assert_eq!(got.status, RunStatus::Pending);
677 assert!(got.step_entries.is_empty());
678 assert_eq!(got.result_ref, None);
679 drop(s);
680 driver.shutdown().await.unwrap();
681 }
682
683 #[tokio::test]
684 async fn duplicate_create_rejected() {
685 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
686 s.create(mk("R-1", "T-1", 100)).await.unwrap();
687 let err = s.create(mk("R-1", "T-1", 200)).await.unwrap_err();
688 assert!(matches!(err, RunStoreError::Duplicate(_)), "got: {err:?}");
689 drop(s);
690 driver.shutdown().await.unwrap();
691 }
692
693 #[tokio::test]
694 async fn get_missing_returns_not_found() {
695 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
696 let err = s.get(&RunId::parse("R-nope").unwrap()).await.unwrap_err();
697 assert!(matches!(err, RunStoreError::NotFound(_)));
698 drop(s);
699 driver.shutdown().await.unwrap();
700 }
701
702 #[tokio::test]
703 async fn list_by_task_filters_and_orders_ascending() {
704 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
705 s.create(mk("R-1", "T-1", 300)).await.unwrap();
706 s.create(mk("R-2", "T-2", 50)).await.unwrap();
707 s.create(mk("R-3", "T-1", 100)).await.unwrap();
708 let list = s
709 .list_by_task(&TaskId::parse("T-1").unwrap())
710 .await
711 .unwrap();
712 let ids: Vec<_> = list.iter().map(|r| r.id.to_string()).collect();
713 assert_eq!(ids, vec!["R-3", "R-1"]);
714 drop(s);
715 driver.shutdown().await.unwrap();
716 }
717
718 #[tokio::test]
719 async fn append_step_entry_accumulates_in_order() {
720 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
721 s.create(mk("R-1", "T-1", 100)).await.unwrap();
722 s.append_step_entry(
723 &RunId::parse("R-1").unwrap(),
724 StepEntry::basic(
725 crate::types::StepId::parse("ST-1").unwrap(),
726 Some("step-a".into()),
727 Some("dispatched".into()),
728 None,
729 101,
730 ),
731 )
732 .await
733 .unwrap();
734 s.append_step_entry(
735 &RunId::parse("R-1").unwrap(),
736 StepEntry::basic(
737 crate::types::StepId::parse("ST-2").unwrap(),
738 Some("step-b".into()),
739 Some("passed".into()),
740 None,
741 102,
742 ),
743 )
744 .await
745 .unwrap();
746 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
747 assert_eq!(got.step_entries.len(), 2);
748 assert_eq!(got.step_entries[0].step_ref, Some("step-a".into()));
749 assert_eq!(got.step_entries[1].step_ref, Some("step-b".into()));
750 drop(s);
751 driver.shutdown().await.unwrap();
752 }
753
754 #[tokio::test]
755 async fn append_step_entry_unknown_run_fails() {
756 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
757 let err = s
758 .append_step_entry(
759 &RunId::parse("R-nope").unwrap(),
760 StepEntry::basic(
761 crate::types::StepId::parse("ST-1").unwrap(),
762 None,
763 None,
764 None,
765 1,
766 ),
767 )
768 .await
769 .unwrap_err();
770 assert!(matches!(err, RunStoreError::NotFound(_)));
771 drop(s);
772 driver.shutdown().await.unwrap();
773 }
774
775 #[tokio::test]
776 async fn append_degradation_accumulates_in_order() {
777 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
778 s.create(mk("R-1", "T-1", 100)).await.unwrap();
779 s.append_degradation(
780 &RunId::parse("R-1").unwrap(),
781 mk_degradation("web_search", 101),
782 )
783 .await
784 .unwrap();
785 s.append_degradation(
786 &RunId::parse("R-1").unwrap(),
787 mk_degradation("code_exec", 102),
788 )
789 .await
790 .unwrap();
791 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
792 assert_eq!(got.degradations.len(), 2);
793 assert_eq!(got.degradations[0].tool, "web_search");
794 assert_eq!(got.degradations[1].tool, "code_exec");
795 drop(s);
796 driver.shutdown().await.unwrap();
797 }
798
799 #[tokio::test]
800 async fn append_degradation_unknown_run_fails() {
801 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
802 let err = s
803 .append_degradation(
804 &RunId::parse("R-nope").unwrap(),
805 mk_degradation("web_search", 1),
806 )
807 .await
808 .unwrap_err();
809 assert!(matches!(err, RunStoreError::NotFound(_)));
810 drop(s);
811 driver.shutdown().await.unwrap();
812 }
813
814 #[tokio::test]
815 async fn append_degradation_bumps_updated_at() {
816 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
817 s.create(mk("R-1", "T-1", 100)).await.unwrap();
818 s.append_degradation(
819 &RunId::parse("R-1").unwrap(),
820 mk_degradation("web_search", 200),
821 )
822 .await
823 .unwrap();
824 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
825 assert!(got.updated_at > 100);
826 drop(s);
827 driver.shutdown().await.unwrap();
828 }
829
830 #[tokio::test]
831 async fn update_status_persists() {
832 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
833 s.create(mk("R-1", "T-1", 100)).await.unwrap();
834 s.update_status(&RunId::parse("R-1").unwrap(), RunStatus::Done)
835 .await
836 .unwrap();
837 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
838 assert_eq!(got.status, RunStatus::Done);
839 drop(s);
840 driver.shutdown().await.unwrap();
841 }
842
843 #[tokio::test]
844 async fn set_result_persists() {
845 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
846 s.create(mk("R-1", "T-1", 100)).await.unwrap();
847 s.set_result(&RunId::parse("R-1").unwrap(), json!({"ok": true}))
848 .await
849 .unwrap();
850 let got = s.get(&RunId::parse("R-1").unwrap()).await.unwrap();
851 assert_eq!(got.result_ref, Some(json!({"ok": true})));
852 drop(s);
853 driver.shutdown().await.unwrap();
854 }
855
856 #[tokio::test]
857 async fn persists_across_reopen() {
858 let dir = tempfile::tempdir().unwrap();
859 let path = dir.path().join("runs.db");
860
861 {
862 let (s, driver) = SqliteRunStore::open(&path).await.unwrap();
863 s.create(mk("R-keep", "T-keep", 42)).await.unwrap();
864 s.append_step_entry(
865 &RunId::parse("R-keep").unwrap(),
866 StepEntry::basic(
867 crate::types::StepId::parse("ST-1").unwrap(),
868 Some("step-a".into()),
869 Some("dispatched".into()),
870 None,
871 43,
872 ),
873 )
874 .await
875 .unwrap();
876 drop(s);
877 driver.shutdown().await.unwrap();
878 }
879
880 let (s, driver) = SqliteRunStore::open(&path).await.unwrap();
881 let got = s.get(&RunId::parse("R-keep").unwrap()).await.unwrap();
882 assert_eq!(got.task_id, TaskId::parse("T-keep").unwrap());
883 assert_eq!(got.step_entries.len(), 1);
884 assert_eq!(got.step_entries[0].step_ref, Some("step-a".into()));
885 drop(s);
886 driver.shutdown().await.unwrap();
887 }
888
889 #[tokio::test]
890 async fn list_running_filters_by_status() {
891 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
892 s.create(mk("R-1", "T-1", 100)).await.unwrap();
893 s.create(mk("R-2", "T-2", 200)).await.unwrap();
894 s.create(mk("R-3", "T-3", 300)).await.unwrap();
895 s.update_status(&RunId::parse("R-2").unwrap(), RunStatus::Running)
896 .await
897 .unwrap();
898 s.update_status(&RunId::parse("R-3").unwrap(), RunStatus::Done)
899 .await
900 .unwrap();
901 let running = s.list_running().await.unwrap();
902 assert_eq!(running.len(), 1);
903 assert_eq!(running[0].id, RunId::parse("R-2").unwrap());
904 assert_eq!(running[0].status, RunStatus::Running);
905 drop(s);
906 driver.shutdown().await.unwrap();
907 }
908
909 #[tokio::test]
910 async fn try_transition_is_atomic_compare_and_set() {
911 let (s, driver) = SqliteRunStore::open_in_memory().await.unwrap();
912 s.create(mk("R-1", "T-1", 100)).await.unwrap();
913 s.update_status(&RunId::parse("R-1").unwrap(), RunStatus::Interrupted)
914 .await
915 .unwrap();
916
917 let first = s
918 .try_transition(
919 &RunId::parse("R-1").unwrap(),
920 RunStatus::Interrupted,
921 RunStatus::Running,
922 )
923 .await
924 .unwrap();
925 assert!(first, "first CAS must flip Interrupted -> Running");
926 assert_eq!(
927 s.get(&RunId::parse("R-1").unwrap()).await.unwrap().status,
928 RunStatus::Running
929 );
930
931 let second = s
932 .try_transition(
933 &RunId::parse("R-1").unwrap(),
934 RunStatus::Interrupted,
935 RunStatus::Running,
936 )
937 .await
938 .unwrap();
939 assert!(
940 !second,
941 "a racing second CAS must not flip a now-Running row"
942 );
943
944 let absent = s
945 .try_transition(
946 &RunId::parse("R-nope").unwrap(),
947 RunStatus::Interrupted,
948 RunStatus::Running,
949 )
950 .await
951 .unwrap();
952 assert!(!absent, "an absent Run must report false, not error");
953 drop(s);
954 driver.shutdown().await.unwrap();
955 }
956
957 #[tokio::test]
958 async fn input_json_roundtrips_across_reopen() {
959 let dir = tempfile::tempdir().unwrap();
960 let path = dir.path().join("runs.db");
961 let snapshot = r#"{"blueprint":"snapshot","init_ctx":{}}"#;
962
963 {
964 let (s, driver) = SqliteRunStore::open(&path).await.unwrap();
965 let mut rec = mk("R-keep", "T-keep", 42);
966 rec.input_json = Some(snapshot.to_string());
967 s.create(rec).await.unwrap();
968 drop(s);
969 driver.shutdown().await.unwrap();
970 }
971
972 let (s, driver) = SqliteRunStore::open(&path).await.unwrap();
973 let got = s.get(&RunId::parse("R-keep").unwrap()).await.unwrap();
974 assert_eq!(got.input_json.as_deref(), Some(snapshot));
975 drop(s);
976 driver.shutdown().await.unwrap();
977 }
978}