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    fn degraded(&self) -> bool {
169        false
170    }
171}
172
173#[cfg(test)]
174mod tests {
175    use super::*;
176    use crate::serve::history::RunStatus;
177    use std::collections::BTreeMap;
178
179    fn rec(id: &str, status: RunStatus, submitted: DateTime<Utc>) -> RunRecord {
180        let mut r = RunRecord::queued(id.into(), None, BTreeMap::new(), None, submitted);
181        r.status = status;
182        if status.is_terminal() {
183            r.finished_at = Some(submitted);
184        }
185        r
186    }
187
188    #[tokio::test]
189    async fn upsert_then_get_roundtrips() {
190        let h = MemoryHistory::new(Duration::from_secs(60));
191        let r = rec("a", RunStatus::Queued, Utc::now());
192        h.upsert(&r).await.unwrap();
193        assert_eq!(h.get("a").await.unwrap().unwrap().run_id, "a");
194        assert!(h.get("missing").await.unwrap().is_none());
195    }
196
197    #[tokio::test]
198    async fn idempotency_fresh_replay_conflict() {
199        let h = MemoryHistory::new(Duration::from_secs(60));
200        let w = Duration::from_secs(60);
201        assert_eq!(
202            h.claim_idempotency("k", "fp1", "run1", w).await.unwrap(),
203            Claim::Fresh
204        );
205        // Same key + same fingerprint → replay the first run id.
206        assert_eq!(
207            h.claim_idempotency("k", "fp1", "run2", w).await.unwrap(),
208            Claim::Replay("run1".into())
209        );
210        // Same key + different fingerprint → conflict.
211        assert_eq!(
212            h.claim_idempotency("k", "fp2", "run3", w).await.unwrap(),
213            Claim::Conflict
214        );
215    }
216
217    #[tokio::test]
218    async fn expired_claim_is_reclaimable() {
219        let h = MemoryHistory::new(Duration::from_secs(60));
220        // Zero window → any prior claim is immediately expired.
221        let w = Duration::ZERO;
222        assert_eq!(
223            h.claim_idempotency("k", "fp1", "run1", w).await.unwrap(),
224            Claim::Fresh
225        );
226        assert_eq!(
227            h.claim_idempotency("k", "fp2", "run2", w).await.unwrap(),
228            Claim::Fresh
229        );
230    }
231
232    #[tokio::test]
233    async fn delete_respects_terminal_state() {
234        let h = MemoryHistory::new(Duration::from_secs(60));
235        h.upsert(&rec("run", RunStatus::Running, Utc::now()))
236            .await
237            .unwrap();
238        assert_eq!(h.delete("run").await.unwrap(), DeleteOutcome::StillRunning);
239        assert_eq!(h.delete("nope").await.unwrap(), DeleteOutcome::NotFound);
240        h.upsert(&rec("run", RunStatus::Completed, Utc::now()))
241            .await
242            .unwrap();
243        assert_eq!(h.delete("run").await.unwrap(), DeleteOutcome::Deleted);
244        assert!(h.get("run").await.unwrap().is_none());
245    }
246
247    #[tokio::test]
248    async fn delete_also_removes_matching_idem_claim() {
249        // M8 (#146): deleting a run must drop its idempotency claim, so a later
250        // replay of the key starts a fresh run instead of 404-ing on the
251        // now-missing record until the claim self-expires.
252        let h = MemoryHistory::new(Duration::from_secs(3600));
253        let w = Duration::from_secs(3600);
254        assert_eq!(
255            h.claim_idempotency("k", "fp", "r1", w).await.unwrap(),
256            Claim::Fresh
257        );
258        let mut r = RunRecord::queued(
259            "r1".into(),
260            None,
261            BTreeMap::new(),
262            Some("k".into()),
263            Utc::now(),
264        );
265        r.status = RunStatus::Completed;
266        r.finished_at = Some(Utc::now());
267        h.upsert(&r).await.unwrap();
268
269        assert_eq!(h.delete("r1").await.unwrap(), DeleteOutcome::Deleted);
270        // The key is free again → fresh run, not a replay of the deleted one.
271        assert_eq!(
272            h.claim_idempotency("k", "fp", "r2", w).await.unwrap(),
273            Claim::Fresh
274        );
275    }
276
277    #[tokio::test]
278    async fn delete_keeps_claim_owned_by_a_newer_run() {
279        // Guard: deleting an OLD run must not remove a claim a NEWER run owns.
280        let h = MemoryHistory::new(Duration::from_secs(3600));
281        h.claim_idempotency("k", "fp", "r1", Duration::from_secs(3600))
282            .await
283            .unwrap();
284        // r2 re-claims the key (force the prior claim stale with a zero window).
285        assert_eq!(
286            h.claim_idempotency("k", "fp", "r2", Duration::ZERO)
287                .await
288                .unwrap(),
289            Claim::Fresh
290        );
291        let mut r1 = RunRecord::queued(
292            "r1".into(),
293            None,
294            BTreeMap::new(),
295            Some("k".into()),
296            Utc::now(),
297        );
298        r1.status = RunStatus::Completed;
299        r1.finished_at = Some(Utc::now());
300        h.upsert(&r1).await.unwrap();
301        assert_eq!(h.delete("r1").await.unwrap(), DeleteOutcome::Deleted);
302        // The claim still belongs to r2.
303        assert_eq!(
304            h.claim_idempotency("k", "fp", "r3", Duration::from_secs(3600))
305                .await
306                .unwrap(),
307            Claim::Replay("r2".into())
308        );
309    }
310
311    #[tokio::test]
312    async fn list_orders_desc_and_paginates() {
313        let h = MemoryHistory::new(Duration::from_secs(60));
314        let t0 = Utc::now();
315        for (i, id) in ["a", "b", "c"].iter().enumerate() {
316            h.upsert(&rec(
317                id,
318                RunStatus::Completed,
319                t0 + chrono::Duration::seconds(i as i64),
320            ))
321            .await
322            .unwrap();
323        }
324        // Newest first → c, b, a. Page size 2.
325        let page = h
326            .list(&ListFilter {
327                limit: 2,
328                ..Default::default()
329            })
330            .await
331            .unwrap();
332        assert_eq!(
333            page.runs
334                .iter()
335                .map(|r| r.run_id.clone())
336                .collect::<Vec<_>>(),
337            vec!["c", "b"]
338        );
339        assert_eq!(page.next_cursor.as_deref(), Some("b"));
340        // Next page from the cursor → a.
341        let page2 = h
342            .list(&ListFilter {
343                limit: 2,
344                cursor: Some("b".into()),
345                ..Default::default()
346            })
347            .await
348            .unwrap();
349        assert_eq!(
350            page2
351                .runs
352                .iter()
353                .map(|r| r.run_id.clone())
354                .collect::<Vec<_>>(),
355            vec!["a"]
356        );
357        assert!(page2.next_cursor.is_none());
358    }
359
360    #[tokio::test]
361    async fn list_filters_by_status_and_name() {
362        let h = MemoryHistory::new(Duration::from_secs(60));
363        let mut r = rec("x", RunStatus::Failed, Utc::now());
364        r.name = Some("nightly".into());
365        h.upsert(&r).await.unwrap();
366        h.upsert(&rec("y", RunStatus::Completed, Utc::now()))
367            .await
368            .unwrap();
369        let only_failed = h
370            .list(&ListFilter {
371                status: Some(RunStatus::Failed),
372                limit: 50,
373                ..Default::default()
374            })
375            .await
376            .unwrap();
377        assert_eq!(only_failed.runs.len(), 1);
378        assert_eq!(only_failed.runs[0].run_id, "x");
379        // Name filter also works.
380        let by_name = h
381            .list(&ListFilter {
382                name: Some("nightly".into()),
383                limit: 50,
384                ..Default::default()
385            })
386            .await
387            .unwrap();
388        assert_eq!(by_name.runs.len(), 1);
389        assert_eq!(by_name.runs[0].run_id, "x");
390    }
391
392    #[tokio::test]
393    async fn purge_drops_expired_terminal_runs() {
394        let h = MemoryHistory::new(Duration::from_secs(60));
395        h.upsert(&rec(
396            "old",
397            RunStatus::Completed,
398            Utc::now() - chrono::Duration::seconds(10),
399        ))
400        .await
401        .unwrap();
402        h.upsert(&rec("live", RunStatus::Running, Utc::now()))
403            .await
404            .unwrap();
405        // retain_for = 0 → every terminal record is expired; running is kept.
406        let removed = h.purge_expired(Duration::ZERO).await.unwrap();
407        assert_eq!(removed, 1);
408        assert!(h.get("old").await.unwrap().is_none());
409        assert!(h.get("live").await.unwrap().is_some());
410    }
411}