Skip to main content

faucet_cli/serve/history/
memory.rs

1//! `DashMap`-backed run history (default backend). Lost on restart; that is the
2//! documented memory-backend trade-off. Idempotency claims live in a second map
3//! and are pruned both lazily (on re-claim) and by `purge_expired`.
4
5use super::{Claim, DeleteOutcome, HistoryError, ListFilter, ListPage, RunHistory, RunRecord};
6use async_trait::async_trait;
7use chrono::{DateTime, Utc};
8use dashmap::DashMap;
9use std::time::Duration;
10
11struct IdemEntry {
12    run_id: String,
13    fingerprint: String,
14    claimed_at: DateTime<Utc>,
15}
16
17pub struct MemoryHistory {
18    runs: DashMap<String, RunRecord>,
19    idem: DashMap<String, IdemEntry>,
20    /// Retention window for idempotency claims (separate from run retention).
21    idem_retention: Duration,
22}
23
24impl MemoryHistory {
25    pub fn new(idem_retention: Duration) -> Self {
26        Self {
27            runs: DashMap::new(),
28            idem: DashMap::new(),
29            idem_retention,
30        }
31    }
32}
33
34/// True when `claimed_at` is older than `window` relative to `now`. A claim
35/// timestamped in the future (clock skew) is treated as *not* expired.
36fn is_expired(claimed_at: DateTime<Utc>, now: DateTime<Utc>, window: Duration) -> bool {
37    now.signed_duration_since(claimed_at)
38        .to_std()
39        .map(|age| age >= window)
40        .unwrap_or(false)
41}
42
43#[async_trait]
44impl RunHistory for MemoryHistory {
45    async fn claim_idempotency(
46        &self,
47        key: &str,
48        fingerprint: &str,
49        run_id: &str,
50        window: Duration,
51    ) -> Result<Claim, HistoryError> {
52        use dashmap::mapref::entry::Entry;
53        let now = Utc::now();
54        // Holding the entry locks the shard, so claim is atomic under contention.
55        match self.idem.entry(key.to_string()) {
56            Entry::Occupied(mut e) => {
57                let expired = is_expired(e.get().claimed_at, now, window);
58                if expired {
59                    e.insert(IdemEntry {
60                        run_id: run_id.to_string(),
61                        fingerprint: fingerprint.to_string(),
62                        claimed_at: now,
63                    });
64                    Ok(Claim::Fresh)
65                } else if e.get().fingerprint == fingerprint {
66                    Ok(Claim::Replay(e.get().run_id.clone()))
67                } else {
68                    Ok(Claim::Conflict)
69                }
70            }
71            Entry::Vacant(v) => {
72                v.insert(IdemEntry {
73                    run_id: run_id.to_string(),
74                    fingerprint: fingerprint.to_string(),
75                    claimed_at: now,
76                });
77                Ok(Claim::Fresh)
78            }
79        }
80    }
81
82    async fn upsert(&self, rec: &RunRecord) -> Result<(), HistoryError> {
83        self.runs.insert(rec.run_id.clone(), rec.clone());
84        Ok(())
85    }
86
87    async fn get(&self, id: &str) -> Result<Option<RunRecord>, HistoryError> {
88        Ok(self.runs.get(id).map(|r| r.clone()))
89    }
90
91    async fn list(&self, filter: &ListFilter) -> Result<ListPage, HistoryError> {
92        let mut rows: Vec<RunRecord> = self
93            .runs
94            .iter()
95            .map(|r| r.clone())
96            .filter(|r| filter.status.is_none_or(|s| r.status == s))
97            .filter(|r| {
98                filter
99                    .name
100                    .as_deref()
101                    .is_none_or(|n| r.name.as_deref() == Some(n))
102            })
103            .filter(|r| filter.since.is_none_or(|t| r.submitted_at >= t))
104            .filter(|r| filter.until.is_none_or(|t| r.submitted_at <= t))
105            .collect();
106        // (submitted_at DESC, run_id DESC)
107        rows.sort_by(|a, b| {
108            b.submitted_at
109                .cmp(&a.submitted_at)
110                .then_with(|| b.run_id.cmp(&a.run_id))
111        });
112        // Cursor = last run_id seen on the previous page; skip past it.
113        if let Some(cursor) = &filter.cursor
114            && let Some(pos) = rows.iter().position(|r| &r.run_id == cursor)
115        {
116            rows.drain(..=pos);
117        }
118        let limit = filter.limit.max(1);
119        let next_cursor = if rows.len() > limit {
120            Some(rows[limit - 1].run_id.clone())
121        } else {
122            None
123        };
124        rows.truncate(limit);
125        Ok(ListPage {
126            runs: rows,
127            next_cursor,
128        })
129    }
130
131    async fn delete(&self, id: &str) -> Result<DeleteOutcome, HistoryError> {
132        let Some(rec) = self.runs.get(id).map(|r| r.clone()) else {
133            return Ok(DeleteOutcome::NotFound);
134        };
135        if !rec.status.is_terminal() {
136            return Ok(DeleteOutcome::StillRunning);
137        }
138        self.runs.remove(id);
139        // Also drop this run's idempotency claim so a replay of the key starts a
140        // fresh run instead of 404-ing on the now-deleted record until the claim
141        // self-expires (#146 M8). Only remove it if the claim still points at
142        // THIS run — a newer run may have re-claimed the key after expiry.
143        if let Some(key) = rec.idempotency_key.as_deref() {
144            self.idem.remove_if(key, |_, e| e.run_id == id);
145        }
146        Ok(DeleteOutcome::Deleted)
147    }
148
149    async fn purge_expired(&self, retain_for: Duration) -> Result<usize, HistoryError> {
150        let now = Utc::now();
151        let before = self.runs.len();
152        self.runs.retain(|_, r| {
153            !r.status.is_terminal()
154                || r.finished_at
155                    .map(|f| !is_expired(f, now, retain_for))
156                    .unwrap_or(true)
157        });
158        // Also drop stale idempotency claims so the map stays bounded.
159        self.idem
160            .retain(|_, e| !is_expired(e.claimed_at, now, self.idem_retention));
161        Ok(before.saturating_sub(self.runs.len()))
162    }
163
164    async fn recover_orphans(&self) -> Result<usize, HistoryError> {
165        Ok(0)
166    }
167
168    async fn cancel_pending(&self, run_id: &str) -> Result<bool, HistoryError> {
169        use crate::serve::history::RunStatus;
170        if let Some(mut r) = self.runs.get_mut(run_id)
171            && r.status == RunStatus::Pending
172        {
173            r.status = RunStatus::Cancelled;
174            r.finished_at = Some(Utc::now());
175            return Ok(true);
176        }
177        Ok(false)
178    }
179
180    fn degraded(&self) -> bool {
181        false
182    }
183}
184
185#[cfg(test)]
186mod tests {
187    use super::*;
188    use crate::serve::history::RunStatus;
189    use std::collections::BTreeMap;
190
191    fn rec(id: &str, status: RunStatus, submitted: DateTime<Utc>) -> RunRecord {
192        let mut r = RunRecord::queued(id.into(), None, BTreeMap::new(), None, submitted);
193        r.status = status;
194        if status.is_terminal() {
195            r.finished_at = Some(submitted);
196        }
197        r
198    }
199
200    #[tokio::test]
201    async fn upsert_then_get_roundtrips() {
202        let h = MemoryHistory::new(Duration::from_secs(60));
203        let r = rec("a", RunStatus::Queued, Utc::now());
204        h.upsert(&r).await.unwrap();
205        assert_eq!(h.get("a").await.unwrap().unwrap().run_id, "a");
206        assert!(h.get("missing").await.unwrap().is_none());
207    }
208
209    #[tokio::test]
210    async fn idempotency_fresh_replay_conflict() {
211        let h = MemoryHistory::new(Duration::from_secs(60));
212        let w = Duration::from_secs(60);
213        assert_eq!(
214            h.claim_idempotency("k", "fp1", "run1", w).await.unwrap(),
215            Claim::Fresh
216        );
217        // Same key + same fingerprint → replay the first run id.
218        assert_eq!(
219            h.claim_idempotency("k", "fp1", "run2", w).await.unwrap(),
220            Claim::Replay("run1".into())
221        );
222        // Same key + different fingerprint → conflict.
223        assert_eq!(
224            h.claim_idempotency("k", "fp2", "run3", w).await.unwrap(),
225            Claim::Conflict
226        );
227    }
228
229    #[tokio::test]
230    async fn expired_claim_is_reclaimable() {
231        let h = MemoryHistory::new(Duration::from_secs(60));
232        // Zero window → any prior claim is immediately expired.
233        let w = Duration::ZERO;
234        assert_eq!(
235            h.claim_idempotency("k", "fp1", "run1", w).await.unwrap(),
236            Claim::Fresh
237        );
238        assert_eq!(
239            h.claim_idempotency("k", "fp2", "run2", w).await.unwrap(),
240            Claim::Fresh
241        );
242    }
243
244    #[tokio::test]
245    async fn delete_respects_terminal_state() {
246        let h = MemoryHistory::new(Duration::from_secs(60));
247        h.upsert(&rec("run", RunStatus::Running, Utc::now()))
248            .await
249            .unwrap();
250        assert_eq!(h.delete("run").await.unwrap(), DeleteOutcome::StillRunning);
251        assert_eq!(h.delete("nope").await.unwrap(), DeleteOutcome::NotFound);
252        h.upsert(&rec("run", RunStatus::Completed, Utc::now()))
253            .await
254            .unwrap();
255        assert_eq!(h.delete("run").await.unwrap(), DeleteOutcome::Deleted);
256        assert!(h.get("run").await.unwrap().is_none());
257    }
258
259    #[tokio::test]
260    async fn delete_also_removes_matching_idem_claim() {
261        // M8 (#146): deleting a run must drop its idempotency claim, so a later
262        // replay of the key starts a fresh run instead of 404-ing on the
263        // now-missing record until the claim self-expires.
264        let h = MemoryHistory::new(Duration::from_secs(3600));
265        let w = Duration::from_secs(3600);
266        assert_eq!(
267            h.claim_idempotency("k", "fp", "r1", w).await.unwrap(),
268            Claim::Fresh
269        );
270        let mut r = RunRecord::queued(
271            "r1".into(),
272            None,
273            BTreeMap::new(),
274            Some("k".into()),
275            Utc::now(),
276        );
277        r.status = RunStatus::Completed;
278        r.finished_at = Some(Utc::now());
279        h.upsert(&r).await.unwrap();
280
281        assert_eq!(h.delete("r1").await.unwrap(), DeleteOutcome::Deleted);
282        // The key is free again → fresh run, not a replay of the deleted one.
283        assert_eq!(
284            h.claim_idempotency("k", "fp", "r2", w).await.unwrap(),
285            Claim::Fresh
286        );
287    }
288
289    #[tokio::test]
290    async fn delete_keeps_claim_owned_by_a_newer_run() {
291        // Guard: deleting an OLD run must not remove a claim a NEWER run owns.
292        let h = MemoryHistory::new(Duration::from_secs(3600));
293        h.claim_idempotency("k", "fp", "r1", Duration::from_secs(3600))
294            .await
295            .unwrap();
296        // r2 re-claims the key (force the prior claim stale with a zero window).
297        assert_eq!(
298            h.claim_idempotency("k", "fp", "r2", Duration::ZERO)
299                .await
300                .unwrap(),
301            Claim::Fresh
302        );
303        let mut r1 = RunRecord::queued(
304            "r1".into(),
305            None,
306            BTreeMap::new(),
307            Some("k".into()),
308            Utc::now(),
309        );
310        r1.status = RunStatus::Completed;
311        r1.finished_at = Some(Utc::now());
312        h.upsert(&r1).await.unwrap();
313        assert_eq!(h.delete("r1").await.unwrap(), DeleteOutcome::Deleted);
314        // The claim still belongs to r2.
315        assert_eq!(
316            h.claim_idempotency("k", "fp", "r3", Duration::from_secs(3600))
317                .await
318                .unwrap(),
319            Claim::Replay("r2".into())
320        );
321    }
322
323    #[tokio::test]
324    async fn list_orders_desc_and_paginates() {
325        let h = MemoryHistory::new(Duration::from_secs(60));
326        let t0 = Utc::now();
327        for (i, id) in ["a", "b", "c"].iter().enumerate() {
328            h.upsert(&rec(
329                id,
330                RunStatus::Completed,
331                t0 + chrono::Duration::seconds(i as i64),
332            ))
333            .await
334            .unwrap();
335        }
336        // Newest first → c, b, a. Page size 2.
337        let page = h
338            .list(&ListFilter {
339                limit: 2,
340                ..Default::default()
341            })
342            .await
343            .unwrap();
344        assert_eq!(
345            page.runs
346                .iter()
347                .map(|r| r.run_id.clone())
348                .collect::<Vec<_>>(),
349            vec!["c", "b"]
350        );
351        assert_eq!(page.next_cursor.as_deref(), Some("b"));
352        // Next page from the cursor → a.
353        let page2 = h
354            .list(&ListFilter {
355                limit: 2,
356                cursor: Some("b".into()),
357                ..Default::default()
358            })
359            .await
360            .unwrap();
361        assert_eq!(
362            page2
363                .runs
364                .iter()
365                .map(|r| r.run_id.clone())
366                .collect::<Vec<_>>(),
367            vec!["a"]
368        );
369        assert!(page2.next_cursor.is_none());
370    }
371
372    #[tokio::test]
373    async fn list_filters_by_status_and_name() {
374        let h = MemoryHistory::new(Duration::from_secs(60));
375        let mut r = rec("x", RunStatus::Failed, Utc::now());
376        r.name = Some("nightly".into());
377        h.upsert(&r).await.unwrap();
378        h.upsert(&rec("y", RunStatus::Completed, Utc::now()))
379            .await
380            .unwrap();
381        let only_failed = h
382            .list(&ListFilter {
383                status: Some(RunStatus::Failed),
384                limit: 50,
385                ..Default::default()
386            })
387            .await
388            .unwrap();
389        assert_eq!(only_failed.runs.len(), 1);
390        assert_eq!(only_failed.runs[0].run_id, "x");
391        // Name filter also works.
392        let by_name = h
393            .list(&ListFilter {
394                name: Some("nightly".into()),
395                limit: 50,
396                ..Default::default()
397            })
398            .await
399            .unwrap();
400        assert_eq!(by_name.runs.len(), 1);
401        assert_eq!(by_name.runs[0].run_id, "x");
402    }
403
404    #[tokio::test]
405    async fn purge_drops_expired_terminal_runs() {
406        let h = MemoryHistory::new(Duration::from_secs(60));
407        h.upsert(&rec(
408            "old",
409            RunStatus::Completed,
410            Utc::now() - chrono::Duration::seconds(10),
411        ))
412        .await
413        .unwrap();
414        h.upsert(&rec("live", RunStatus::Running, Utc::now()))
415            .await
416            .unwrap();
417        // retain_for = 0 → every terminal record is expired; running is kept.
418        let removed = h.purge_expired(Duration::ZERO).await.unwrap();
419        assert_eq!(removed, 1);
420        assert!(h.get("old").await.unwrap().is_none());
421        assert!(h.get("live").await.unwrap().is_some());
422    }
423}