durable_workflow/
sticky_workflow_cache.rs

1//! Bounded durable wire history. Replay always gets a fresh decoded snapshot.
2
3use serde_json::Value;
4use std::{
5    collections::VecDeque,
6    time::{Duration, Instant},
7};
8
9#[derive(Clone, Debug, PartialEq, Eq)]
10pub(crate) struct CacheKey {
11    pub workflow_id: String,
12    pub run_id: String,
13    pub build_id: String,
14}
15
16#[derive(Clone, Debug, PartialEq, Eq)]
17pub(crate) struct ResumeCursor {
18    pub token: String,
19    pub offset: usize,
20}
21
22/// Replay counters and retained encoded history. Decoding/replay memory is additional.
23#[derive(Clone, Copy, Debug, Default, PartialEq, Eq, serde::Serialize)]
24pub struct StickyCacheMetrics {
25    pub hit: u64,
26    pub miss: u64,
27    pub eviction: u64,
28    pub forced_cold_replay: u64,
29    pub entries: usize,
30    pub history_bytes: usize,
31}
32
33#[derive(Debug)]
34struct Entry {
35    key: CacheKey,
36    encoded: Box<[u8]>,
37    expires_at: Instant,
38    resume: Option<ResumeCursor>,
39}
40
41#[derive(Debug)]
42pub(crate) struct StickyWorkflowCache {
43    capacity: usize,
44    max_bytes: usize,
45    ttl: Duration,
46    entries: VecDeque<Entry>,
47    history_bytes: usize,
48    metrics: StickyCacheMetrics,
49}
50
51/// Opt-in bounds for a worker's durable-history cache.
52#[derive(Clone, Debug)]
53pub struct StickyCacheOptions {
54    pub(crate) capacity: usize,
55    pub(crate) max_bytes: usize,
56    pub(crate) ttl: Duration,
57}
58
59impl StickyCacheOptions {
60    /// Bound the retained run count. Zero disables caching. Defaults to 16 MiB and 300 seconds.
61    pub fn new(capacity: usize) -> Self {
62        Self {
63            capacity,
64            max_bytes: 16 * 1024 * 1024,
65            ttl: Duration::from_secs(300),
66        }
67    }
68    pub fn max_history_bytes(mut self, bytes: usize) -> Self {
69        self.max_bytes = bytes;
70        self
71    }
72    /// Whole seconds, from 1 to 3600. Reads do not extend the original expiry.
73    pub fn ttl(mut self, ttl: Duration) -> Self {
74        self.ttl = ttl;
75        self
76    }
77}
78
79pub(crate) fn complete_history(history: &[Value]) -> bool {
80    let Some(first) = history.first() else {
81        return false;
82    };
83    let starts_workflow = first["event_type"] == "WorkflowStarted"
84        || (first["event_type"] == "StartAccepted"
85            && history
86                .get(1)
87                .is_some_and(|event| event["event_type"] == "WorkflowStarted"));
88    starts_workflow
89        && history.iter().enumerate().all(|(index, event)| {
90            event.is_object() && event["sequence"].as_u64() == u64::try_from(index + 1).ok()
91        })
92}
93
94impl StickyWorkflowCache {
95    pub fn new(capacity: usize, max_bytes: usize, ttl: Duration) -> Result<Self, &'static str> {
96        if max_bytes == 0 {
97            return Err("sticky cache byte limit must be positive");
98        }
99        if ttl < Duration::from_secs(1) || ttl > Duration::from_secs(3600) {
100            return Err("sticky cache TTL must be between 1 and 3600 seconds");
101        }
102        Ok(Self {
103            capacity,
104            max_bytes,
105            ttl,
106            entries: VecDeque::new(),
107            history_bytes: 0,
108            metrics: StickyCacheMetrics::default(),
109        })
110    }
111
112    pub fn enabled(&self) -> bool {
113        self.capacity > 0
114    }
115
116    pub fn empty(&self) -> Self {
117        Self::new(self.capacity, self.max_bytes, self.ttl).expect("validated cache options")
118    }
119    pub fn ttl_seconds(&self) -> u64 {
120        self.ttl.as_secs()
121    }
122
123    fn expire(&mut self, now: Instant) {
124        self.entries.retain(|entry| entry.expires_at > now);
125        self.history_bytes = self.entries.iter().map(|entry| entry.encoded.len()).sum();
126    }
127
128    pub fn discard(&mut self, key: &CacheKey) {
129        if let Some(index) = self.entries.iter().position(|entry| &entry.key == key) {
130            let entry = self.entries.remove(index).expect("known cache entry");
131            self.history_bytes -= entry.encoded.len();
132        }
133    }
134
135    pub fn lookup(
136        &mut self,
137        key: &CacheKey,
138        now: Instant,
139    ) -> Option<(Vec<Value>, Option<ResumeCursor>)> {
140        self.expire(now);
141        let index = self.entries.iter().position(|entry| &entry.key == key)?;
142        let entry = self.entries.remove(index)?;
143        let history =
144            serde_json::from_slice(&entry.encoded).expect("cache encodes its own history");
145        let resume = entry.resume.clone();
146        self.entries.push_back(entry);
147        Some((history, resume))
148    }
149
150    pub fn remember(
151        &mut self,
152        key: CacheKey,
153        history: &[Value],
154        resume: Option<ResumeCursor>,
155        now: Instant,
156    ) -> bool {
157        self.expire(now);
158        self.discard(&key);
159        if !self.enabled()
160            || key.workflow_id.is_empty()
161            || key.run_id.is_empty()
162            || key.build_id.is_empty()
163            || !complete_history(history)
164            || resume
165                .as_ref()
166                .is_some_and(|cursor| cursor.token.is_empty() || cursor.offset > history.len())
167        {
168            return false;
169        }
170        let Ok(encoded) = serde_json::to_vec(history) else {
171            return false;
172        };
173        if encoded.len() > self.max_bytes {
174            return false;
175        }
176        while self.entries.len() >= self.capacity
177            || self.history_bytes > self.max_bytes - encoded.len()
178        {
179            let Some(oldest) = self.entries.pop_front() else {
180                return false;
181            };
182            self.history_bytes -= oldest.encoded.len();
183            self.metrics.eviction = self.metrics.eviction.saturating_add(1);
184        }
185        self.history_bytes += encoded.len();
186        self.entries.push_back(Entry {
187            key,
188            encoded: encoded.into_boxed_slice(),
189            expires_at: now + self.ttl,
190            resume,
191        });
192        true
193    }
194
195    pub fn record_replay(&mut self, hit: bool, forced: bool) {
196        if hit {
197            self.metrics.hit = self.metrics.hit.saturating_add(1);
198        } else {
199            self.metrics.miss = self.metrics.miss.saturating_add(1);
200        }
201        if forced {
202            self.metrics.forced_cold_replay = self.metrics.forced_cold_replay.saturating_add(1);
203        }
204    }
205
206    pub fn metrics(&mut self, now: Instant) -> StickyCacheMetrics {
207        self.expire(now);
208        StickyCacheMetrics {
209            entries: self.entries.len(),
210            history_bytes: self.history_bytes,
211            ..self.metrics
212        }
213    }
214
215    pub fn clear(&mut self) {
216        self.entries.clear();
217        self.history_bytes = 0;
218    }
219}
220
221#[cfg(test)]
222mod tests {
223    use super::*;
224    use serde_json::json;
225
226    fn key(run: &str) -> CacheKey {
227        CacheKey {
228            workflow_id: "workflow".into(),
229            run_id: run.into(),
230            build_id: "build".into(),
231        }
232    }
233    fn history() -> Vec<Value> {
234        vec![
235            json!({"event_type":"StartAccepted", "sequence":1}),
236            json!({"event_type":"WorkflowStarted", "sequence":2, "payload":{"value":[1,2]}}),
237        ]
238    }
239    fn cache(capacity: usize, max_bytes: usize) -> StickyWorkflowCache {
240        StickyWorkflowCache::new(capacity, max_bytes, Duration::from_secs(1)).unwrap()
241    }
242
243    #[test]
244    fn canonical_start_and_contiguous_sequences_are_required() {
245        assert!(complete_history(&history()));
246        assert!(complete_history(&[
247            json!({"event_type":"WorkflowStarted","sequence":1})
248        ]));
249        assert!(!complete_history(&[]));
250        assert!(!complete_history(&history()[..1]));
251        for bad in [json!(true), json!(-2), json!(2.0), json!(3), Value::Null] {
252            let mut events = history();
253            events[1]["sequence"] = bad;
254            assert!(!complete_history(&events));
255        }
256        let mut events = history();
257        events[1]["event_type"] = json!("ActivityCompleted");
258        assert!(!complete_history(&events));
259        assert!(!complete_history(&[json!(false)]));
260    }
261
262    #[test]
263    fn replay_mutations_cannot_change_retained_history_or_cursor() {
264        let now = Instant::now();
265        let mut cache = cache(2, 10_000);
266        let cursor = Some(ResumeCursor {
267            token: "opaque/lease-token".into(),
268            offset: 1,
269        });
270        assert!(cache.remember(key("run"), &history(), cursor.clone(), now));
271        let (mut replay, mut replay_cursor) = cache.lookup(&key("run"), now).unwrap();
272        replay[1]["payload"]["value"][0] = json!(999);
273        replay_cursor.as_mut().unwrap().token.clear();
274        assert_eq!(cache.lookup(&key("run"), now), Some((history(), cursor)));
275    }
276
277    #[test]
278    fn lru_lookup_changes_entry_eviction_order() {
279        let now = Instant::now();
280        let mut cache = cache(2, 10_000);
281        for run in ["a", "b"] {
282            assert!(cache.remember(key(run), &history(), None, now));
283        }
284        assert!(cache.lookup(&key("a"), now).is_some());
285        assert!(cache.remember(key("c"), &history(), None, now));
286        assert!(cache.lookup(&key("b"), now).is_none());
287        assert!(cache.lookup(&key("a"), now).is_some());
288        assert!(cache.lookup(&key("c"), now).is_some());
289        assert_eq!(cache.metrics(now).eviction, 1);
290    }
291
292    #[test]
293    fn encoded_byte_budget_evicts_even_below_entry_capacity() {
294        let now = Instant::now();
295        let bytes = serde_json::to_vec(&history()).unwrap().len();
296        let mut cache = cache(10, bytes);
297        assert!(cache.remember(key("a"), &history(), None, now));
298        assert!(cache.remember(key("b"), &history(), None, now));
299        assert!(cache.lookup(&key("a"), now).is_none());
300        assert_eq!(cache.metrics(now).history_bytes, bytes);
301        assert_eq!(cache.metrics(now).entries, 1);
302        assert_eq!(cache.metrics(now).eviction, 1);
303    }
304
305    #[test]
306    fn oversized_or_invalid_replacement_discards_only_its_own_old_entry() {
307        let now = Instant::now();
308        let mut cache = cache(2, 1000);
309        for run in ["a", "b"] {
310            assert!(cache.remember(key(run), &history(), None, now));
311        }
312        let mut oversized = history();
313        oversized[1]["payload"] = json!("x".repeat(1000));
314        assert!(!cache.remember(key("a"), &oversized, None, now));
315        assert!(cache.lookup(&key("a"), now).is_none());
316        assert!(cache.lookup(&key("b"), now).is_some());
317        assert!(!cache.remember(key("b"), &[], None, now));
318        assert_eq!(cache.metrics(now).history_bytes, 0);
319        assert_eq!(cache.metrics(now).eviction, 0);
320    }
321
322    #[test]
323    fn expiry_uses_original_admission_time_and_releases_bytes() {
324        let now = Instant::now();
325        let mut cache = cache(1, 10_000);
326        assert!(cache.remember(key("a"), &history(), None, now));
327        assert!(cache
328            .lookup(&key("a"), now + Duration::from_millis(999))
329            .is_some());
330        assert!(cache
331            .lookup(&key("a"), now + Duration::from_secs(1))
332            .is_none());
333        assert_eq!(cache.metrics(now + Duration::from_secs(1)).history_bytes, 0);
334        assert_eq!(cache.metrics(now + Duration::from_secs(1)).eviction, 0);
335    }
336
337    #[test]
338    fn workflow_run_and_build_identity_are_all_part_of_the_key() {
339        let now = Instant::now();
340        let mut cache = cache(1, 10_000);
341        assert!(cache.remember(key("a"), &history(), None, now));
342        for different in [
343            CacheKey {
344                workflow_id: "other".into(),
345                ..key("a")
346            },
347            key("other"),
348            CacheKey {
349                build_id: "other".into(),
350                ..key("a")
351            },
352        ] {
353            assert!(cache.lookup(&different, now).is_none());
354        }
355        assert!(cache.lookup(&key("a"), now).is_some());
356    }
357
358    #[test]
359    fn disabled_profile_never_retains_or_claims_history() {
360        let now = Instant::now();
361        let mut cache = cache(0, 10_000);
362        assert!(!cache.enabled());
363        assert!(!cache.remember(key("a"), &history(), None, now));
364        assert!(cache.lookup(&key("a"), now).is_none());
365        assert_eq!(cache.metrics(now).entries, 0);
366    }
367
368    #[test]
369    fn invalid_limits_identity_and_cursors_are_refused() {
370        assert!(StickyWorkflowCache::new(1, 0, Duration::from_secs(1)).is_err());
371        assert!(StickyWorkflowCache::new(1, 1, Duration::from_millis(999)).is_err());
372        assert!(StickyWorkflowCache::new(1, 1, Duration::from_secs(3601)).is_err());
373        let now = Instant::now();
374        let mut cache = cache(1, 10_000);
375        for invalid in [
376            CacheKey {
377                workflow_id: "".into(),
378                ..key("a")
379            },
380            CacheKey {
381                run_id: "".into(),
382                ..key("a")
383            },
384            CacheKey {
385                build_id: "".into(),
386                ..key("a")
387            },
388        ] {
389            assert!(!cache.remember(invalid, &history(), None, now));
390        }
391        for cursor in [
392            ResumeCursor {
393                token: "".into(),
394                offset: 0,
395            },
396            ResumeCursor {
397                token: "opaque".into(),
398                offset: 3,
399            },
400        ] {
401            assert!(!cache.remember(key("a"), &history(), Some(cursor), now));
402        }
403    }
404
405    #[test]
406    fn terminal_discard_and_shutdown_clear_preserve_replay_counters() {
407        let now = Instant::now();
408        let mut cache = cache(2, 10_000);
409        for run in ["a", "b"] {
410            assert!(cache.remember(key(run), &history(), None, now));
411        }
412        cache.record_replay(true, false);
413        cache.record_replay(false, true);
414        cache.discard(&key("a"));
415        cache.discard(&key("missing"));
416        assert_eq!(cache.metrics(now).entries, 1);
417        cache.clear();
418        assert_eq!(
419            cache.metrics(now),
420            StickyCacheMetrics {
421                hit: 1,
422                miss: 1,
423                forced_cold_replay: 1,
424                ..StickyCacheMetrics::default()
425            }
426        );
427    }
428}