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 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 assert_eq!(
219 h.claim_idempotency("k", "fp1", "run2", w).await.unwrap(),
220 Claim::Replay("run1".into())
221 );
222 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 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 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 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 let h = MemoryHistory::new(Duration::from_secs(3600));
293 h.claim_idempotency("k", "fp", "r1", Duration::from_secs(3600))
294 .await
295 .unwrap();
296 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 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 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 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 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 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}