1use 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#[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#[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 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 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}