faucet_cli/serve/history/
memory.rs1use 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 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
34fn 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 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 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 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 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 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 assert_eq!(
207 h.claim_idempotency("k", "fp1", "run2", w).await.unwrap(),
208 Claim::Replay("run1".into())
209 );
210 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 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 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 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 let h = MemoryHistory::new(Duration::from_secs(3600));
281 h.claim_idempotency("k", "fp", "r1", Duration::from_secs(3600))
282 .await
283 .unwrap();
284 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 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 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 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 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 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}