Skip to main content

mlua_swarm/store/run/
sqlite.rs

1//! `SqliteRunStore` — SQLite-backed [`RunStore`] using [`rusqlite-isle`].
2//!
3//! The `Connection` is confined to a dedicated OS thread by `AsyncIsle`;
4//! every call is a typed closure dispatched over a bounded channel.
5//! `step_entries`, `degradations`, and `result_ref` are stored as JSON
6//! blobs — the former two are pure trace/observability artifacts (not
7//! queried relationally), the latter is caller-defined payload shape.
8//! `append_step_entry`/`append_degradation` run as a read-modify-write
9//! inside a single transaction so concurrent appenders don't clobber each
10//! other's entries.
11//!
12//! ## Schema
13//!
14//! ```sql
15//! CREATE TABLE IF NOT EXISTS runs (
16//!   id                 TEXT PRIMARY KEY,
17//!   task_id            TEXT NOT NULL,
18//!   status             TEXT NOT NULL,      -- JSON-encoded `RunStatus`
19//!   step_entries_json  TEXT NOT NULL,      -- JSON-encoded `Vec<StepEntry>`
20//!   degradations_json  TEXT NOT NULL DEFAULT '[]', -- JSON-encoded `Vec<DegradationEntry>` (GH #32)
21//!   operator_sid       TEXT,
22//!   result_ref_json    TEXT,               -- JSON-encoded `serde_json::Value`, NULL when unset
23//!   input_json         TEXT,               -- opaque launch-input snapshot for resume, NULL when unset
24//!   created_at         INTEGER NOT NULL,
25//!   updated_at         INTEGER NOT NULL
26//! );
27//! CREATE INDEX IF NOT EXISTS ix_runs_task_id ON runs(task_id, created_at);
28//! ```
29//!
30//! `degradations_json` (GH #32) and `input_json` (the resume launch-input
31//! snapshot) were both added after the initial release; each migration is
32//! applied idempotently on open via a `PRAGMA table_info(runs)` existence
33//! check followed by the matching `ALTER TABLE runs ADD COLUMN …` when
34//! missing, so pre-existing database files pick up the columns without a
35//! manual migration step. `input_json` is a nullable `TEXT` (no default) —
36//! rows written before resume support simply read back `None`.
37
38use 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
63/// Idempotently ensures a nullable column named `column` exists on `runs`,
64/// adding it via `ALTER TABLE … ADD COLUMN <column> <decl>` when a
65/// pre-existing database file was created before the column was introduced.
66/// Fresh databases get every column from [`SCHEMA_SQL`] directly; this only
67/// fires the `ALTER TABLE` on older files missing it.
68fn 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
85/// SQLite-backed persistent [`RunStore`].
86///
87/// Open with [`SqliteRunStore::open`] (file path) or
88/// [`SqliteRunStore::open_in_memory`] (tests). Both return the store plus
89/// an [`AsyncIsleDriver`] the caller must `shutdown().await` when done —
90/// dropping the driver without a shutdown call leaves the SQLite thread
91/// as-is until the process exits.
92pub struct SqliteRunStore {
93    isle: AsyncIsle,
94}
95
96impl SqliteRunStore {
97    /// Open (or create) a SQLite database file and run the schema
98    /// migrations.
99    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            // The trace store (`SqliteRunTraceStore`) shares this file
102            // from its own confined connection; a short busy wait
103            // absorbs its write transactions instead of surfacing
104            // SQLITE_BUSY here.
105            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    /// Open an ephemeral in-memory database (tests, doctests).
116    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
133/// One `runs` SELECT row in column order: id, task_id, status,
134/// step_entries_json, degradations_json, operator_sid, result_ref_json,
135/// input_json, created_at, updated_at.
136type 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(&degradations_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    // Ids were minted by us before landing in the table; a prefix mismatch
179    // here means the row predates the issue #13 prefix reconciliation or
180    // the file was written by something else — fail loud either way.
181    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                // Immediate: the trace store shares this file from its own
228                // connection; RESERVED-up-front keeps the busy wait
229                // effective (a DEFERRED read-then-upgrade racing it gets
230                // an instant SQLITE_BUSY instead).
231                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                // Immediate — see `create`'s comment (shared-file busy-wait
353                // effectiveness).
354                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                // Immediate — see `create`'s comment (shared-file busy-wait
401                // effectiveness).
402                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        // A single conditional UPDATE is the compare-and-set: the `AND
472        // status = ?from` predicate makes the read+set atomic at the SQLite
473        // level, so two concurrent resumes cannot both flip the same row.
474        // `rows_affected == 1` = we won the transition; `0` = the row was
475        // absent or no longer `from` (a racing transition already won).
476        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                // `?1 IS NULL OR …` folds each optional filter into one
584                // statement; `LIMIT -1` is SQLite's "no cap". `rowid`
585                // breaks `created_at` ties newest-insertion-first.
586                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// ──────────────────────────────────────────────────────────────────────────
635// tests
636// ──────────────────────────────────────────────────────────────────────────
637
638#[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}