dragonfly_client_rs/
reuse_cache.rs

1//! Bounded process-local results. A cache belongs to one loaded rules snapshot.
2use serde::{de::DeserializeOwned, Deserialize, Serialize};
3use std::{
4    collections::{HashMap, HashSet, VecDeque},
5    fs,
6    path::Path,
7    sync::{
8        atomic::{AtomicBool, AtomicU64, Ordering},
9        Mutex,
10    },
11    time::Instant,
12};
13
14#[derive(Clone, Copy, Debug, Default, Deserialize, Serialize, PartialEq, Eq)]
15#[serde(rename_all = "snake_case")]
16pub enum CacheMode {
17    #[default]
18    Off,
19    Observe,
20    Reuse,
21}
22
23struct Entry {
24    content: Vec<u8>,
25    result: Vec<u8>,
26}
27
28#[derive(Default)]
29struct State {
30    entries: HashMap<String, Entry>,
31    order: VecDeque<String>,
32    bytes: usize,
33    quarantined: HashSet<String>,
34}
35
36pub(crate) struct ReuseCache {
37    pub mode: CacheMode,
38    remote: Option<Mutex<crate::durable_cache::DurableCache>>,
39    max_entries: usize,
40    max_bytes: usize,
41    state: Mutex<State>,
42    hits: AtomicU64,
43    disabled: AtomicBool,
44}
45
46impl ReuseCache {
47    pub(crate) fn new(mode: CacheMode, max_entries: usize, max_bytes: usize) -> Self {
48        Self {
49            mode,
50            remote: None,
51            max_entries,
52            max_bytes,
53            state: Mutex::new(State::default()),
54            hits: AtomicU64::new(0),
55            disabled: AtomicBool::new(false),
56        }
57    }
58
59    pub(crate) fn set_database(&mut self, remote: crate::durable_cache::DurableCache) {
60        self.remote = Some(Mutex::new(remote));
61    }
62
63    pub(crate) fn uses_database(&self) -> bool {
64        self.remote.is_some()
65    }
66
67    pub(crate) fn begin_job(&self, job: &crate::client::Job) {
68        if let Some(remote) = &self.remote {
69            if let Ok(mut remote) = remote.lock() {
70                remote.begin_job(crate::durable_cache::Lease {
71                    name: job.name.clone(),
72                    version: job.version.clone(),
73                    assignment_id: job.assignment_id.clone(),
74                    attempt: job.attempt,
75                });
76            }
77        }
78    }
79
80    pub(crate) fn prefetch(
81        &self,
82        keys: &[(String, crate::durable_cache::Key)],
83        stats: &mut CacheStats,
84    ) {
85        self.remote_operation(stats, |remote| remote.prefetch(keys));
86    }
87
88    pub(crate) fn flush(&self, stats: &mut CacheStats) {
89        self.remote_operation(stats, crate::durable_cache::DurableCache::flush);
90    }
91
92    fn remote_operation(
93        &self,
94        stats: &mut CacheStats,
95        operation: impl FnOnce(&mut crate::durable_cache::DurableCache) -> color_eyre::Result<()>,
96    ) {
97        if self.mode == CacheMode::Off || self.disabled.load(Ordering::Acquire) {
98            return;
99        }
100        let Some(remote) = &self.remote else {
101            return;
102        };
103        let started = Instant::now();
104        let result = (|| {
105            let mut remote = remote
106                .lock()
107                .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?;
108            let result = operation(&mut remote);
109            if remote.revoked {
110                self.disabled.store(true, Ordering::Release);
111            }
112            result
113        })();
114        stats.overhead_us += started.elapsed().as_micros();
115        if let Err(error) = result {
116            stats.errors += 1;
117            tracing::warn!(event="scan_reuse_error", %error, "Database cache unavailable; scanning normally");
118        }
119    }
120
121    pub(crate) fn should_reuse(&self) -> bool {
122        self.mode == CacheMode::Reuse
123            && !self.disabled.load(Ordering::Relaxed)
124            && !(self.hits.fetch_add(1, Ordering::Relaxed) + 1).is_multiple_of(100)
125    }
126
127    pub(crate) fn is_disabled(&self) -> bool {
128        self.disabled.load(Ordering::Acquire)
129    }
130
131    pub(crate) fn quarantine(&self, key: &str, stats: &mut CacheStats) {
132        let Ok(mut storage) = self.state.lock() else {
133            self.disabled.store(true, Ordering::Release);
134            stats.errors += 1;
135            tracing::error!(
136                event = "scan_cache_quarantine_failed",
137                "Cache lock poisoned; disabling local reuse"
138            );
139            return;
140        };
141        // Quarantine failures remain blocked locally and retry on the next job.
142        if storage.quarantined.len() >= 4096 {
143            self.disabled.store(true, Ordering::Release);
144            stats.errors += 1;
145            tracing::error!(
146                event = "scan_cache_quarantine_failed",
147                "Local quarantine capacity reached; disabling local reuse"
148            );
149            return;
150        }
151        storage.quarantined.insert(key.to_owned());
152        if let Some(entry) = storage.entries.remove(key) {
153            storage.bytes -= entry.content.len() + entry.result.len();
154        }
155        storage.order.retain(|entry| entry != key);
156        drop(storage);
157        self.remote_operation(stats, |remote| remote.quarantine(key));
158    }
159
160    fn is_quarantined(&self, key: &str) -> color_eyre::Result<bool> {
161        Ok(self
162            .state
163            .lock()
164            .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?
165            .quarantined
166            .contains(key))
167    }
168
169    #[cfg(test)]
170    pub(crate) fn disable(&self) {
171        self.disabled.store(true, Ordering::Release);
172        if let Some(remote) = &self.remote {
173            if let Ok(mut remote) = remote.lock() {
174                if let Err(error) = remote.revoke() {
175                    tracing::error!(event="scan_cache_revocation_failed", %error, "Failed to persist cache revocation");
176                }
177            }
178        }
179    }
180
181    pub(crate) fn clear(&mut self) {
182        self.state = Mutex::new(State::default());
183        self.disabled.store(false, Ordering::Relaxed);
184        self.hits.store(0, Ordering::Relaxed);
185    }
186
187    pub(crate) fn lookup<T: DeserializeOwned>(
188        &self,
189        key: &str,
190        path: &Path,
191        stats: &mut CacheStats,
192    ) -> Option<T> {
193        if self.mode == CacheMode::Off || self.disabled.load(Ordering::Relaxed) {
194            return None;
195        }
196        let started = Instant::now();
197        stats.lookups += 1;
198        let result = self.read(key, path);
199        stats.overhead_us += started.elapsed().as_micros();
200        match result {
201            Ok(Some(value)) => {
202                stats.candidate_files += 1;
203                Some(value)
204            }
205            Ok(None) => None,
206            Err(error) => {
207                stats.errors += 1;
208                tracing::warn!(event = "scan_reuse_error", %error, "Cache lookup failed; scanning normally");
209                None
210            }
211        }
212    }
213
214    fn read<T: DeserializeOwned>(&self, key: &str, path: &Path) -> color_eyre::Result<Option<T>> {
215        if self.is_quarantined(key)? {
216            return Ok(None);
217        }
218        if let Some(remote) = &self.remote {
219            return remote
220                .lock()
221                .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?
222                .lookup(key);
223        }
224        let state = self
225            .state
226            .lock()
227            .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?;
228        let Some(entry) = state.entries.get(key) else {
229            return Ok(None);
230        };
231        // Hashes locate candidates; exact bytes are required before trusting results.
232        if fs::read(path)? != entry.content {
233            return Ok(None);
234        }
235        Ok(Some(serde_json::from_slice(&entry.result)?))
236    }
237
238    pub(crate) fn insert<T: Serialize>(
239        &self,
240        key: String,
241        path: &Path,
242        value: &T,
243        stats: &mut CacheStats,
244    ) {
245        if self.mode == CacheMode::Off
246            || self.max_entries == 0
247            || self.disabled.load(Ordering::Relaxed)
248        {
249            return;
250        }
251        let started = Instant::now();
252        if let Err(error) = self.write(key, path, value, stats) {
253            stats.errors += 1;
254            tracing::warn!(event = "scan_reuse_error", %error, "Cache write failed; preserving scan result");
255        }
256        stats.overhead_us += started.elapsed().as_micros();
257    }
258
259    fn write<T: Serialize>(
260        &self,
261        key: String,
262        path: &Path,
263        value: &T,
264        stats: &mut CacheStats,
265    ) -> color_eyre::Result<()> {
266        if self.is_quarantined(&key)? {
267            return Ok(());
268        }
269        if let Some(remote) = &self.remote {
270            return remote
271                .lock()
272                .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?
273                .insert(&key, value);
274        }
275        if path.metadata()?.len() > u64::try_from(self.max_bytes)? {
276            return Ok(());
277        }
278        let result = serde_json::to_vec(value)?;
279        let content = fs::read(path)?;
280        let size = content.len().saturating_add(result.len());
281        if size > self.max_bytes {
282            return Ok(());
283        }
284        let mut storage = self
285            .state
286            .lock()
287            .map_err(|_| color_eyre::eyre::eyre!("cache lock poisoned"))?;
288        if storage.entries.contains_key(&key) {
289            return Ok(());
290        }
291        while storage.entries.len() >= self.max_entries
292            || storage.bytes.saturating_add(size) > self.max_bytes
293        {
294            let Some(oldest) = storage.order.pop_front() else {
295                break;
296            };
297            if let Some(entry) = storage.entries.remove(&oldest) {
298                storage.bytes -= entry.content.len() + entry.result.len();
299                stats.evicted_files += 1;
300            }
301        }
302        storage.bytes += size;
303        storage.order.push_back(key.clone());
304        storage.entries.insert(key, Entry { content, result });
305        stats.inserted_files += 1;
306        Ok(())
307    }
308}
309
310/// Per-job work counters, emitted even when a job exits early with an error.
311#[derive(Clone, Serialize)]
312pub(crate) struct CacheStats {
313    scanner: &'static str,
314    pub mode: CacheMode,
315    pub lookups: u64,
316    pub candidate_files: u64,
317    pub reused_files: u64,
318    pub reused_bytes: u64,
319    pub inserted_files: u64,
320    pub evicted_files: u64,
321    pub errors: u64,
322    pub validated_files: u64,
323    pub mismatched_files: u64,
324    pub overhead_us: u128,
325    pub engine_us: u128,
326    pub engine_files: u64,
327    pub engine_bytes: u64,
328}
329
330impl CacheStats {
331    pub(crate) fn new(scanner: &'static str, mode: CacheMode) -> Self {
332        Self {
333            scanner,
334            mode,
335            lookups: 0,
336            candidate_files: 0,
337            reused_files: 0,
338            reused_bytes: 0,
339            inserted_files: 0,
340            evicted_files: 0,
341            errors: 0,
342            validated_files: 0,
343            mismatched_files: 0,
344            overhead_us: 0,
345            engine_us: 0,
346            engine_files: 0,
347            engine_bytes: 0,
348        }
349    }
350}
351
352impl CacheStats {
353    pub(crate) fn emit(&self) {
354        tracing::info!(event = "scan_reuse", scanner = self.scanner, mode = ?self.mode,
355            lookups = self.lookups, candidate_files = self.candidate_files,
356            reused_files = self.reused_files, reused_bytes = self.reused_bytes,
357            inserted_files = self.inserted_files, evicted_files = self.evicted_files,
358            cache_errors = self.errors, validated_files = self.validated_files, mismatched_files = self.mismatched_files, overhead_us = self.overhead_us,
359            engine_us = self.engine_us, engine_files = self.engine_files, engine_bytes = self.engine_bytes,
360            "Cross-package scan reuse statistics");
361    }
362}
363
364#[cfg(test)]
365mod tests {
366    use super::*;
367
368    #[test]
369    fn reuse_audits_every_hundredth_hit_and_can_be_disabled() {
370        let mut cache = ReuseCache::new(CacheMode::Reuse, 10, 1024);
371        for _ in 0..99 {
372            assert!(cache.should_reuse());
373        }
374        assert!(!cache.should_reuse());
375        assert!(cache.should_reuse());
376        cache.disable();
377        assert!(!cache.should_reuse());
378        cache.clear();
379        assert!(cache.should_reuse());
380    }
381
382    #[test]
383    fn exact_bytes_bounds_and_rule_reset() {
384        let dir = tempfile::tempdir().unwrap();
385        let path = dir.path().join("target");
386        fs::write(&path, b"clean").unwrap();
387        let mut cache = ReuseCache::new(CacheMode::Reuse, 1, 64);
388        let mut stats = CacheStats::new("test", CacheMode::Reuse);
389        cache.insert("key".into(), &path, &vec!["finding"], &mut stats);
390        assert_eq!(
391            cache
392                .lookup::<Vec<String>>("key", &path, &mut stats)
393                .unwrap(),
394            vec!["finding"]
395        );
396        fs::write(&path, b"evil!").unwrap();
397        assert!(cache
398            .lookup::<Vec<String>>("key", &path, &mut stats)
399            .is_none());
400        cache.insert("new".into(), &path, &Vec::<String>::new(), &mut stats);
401        assert_eq!(stats.evicted_files, 1);
402        assert!(cache
403            .lookup::<Vec<String>>("new", &path, &mut stats)
404            .unwrap()
405            .is_empty());
406        cache.clear();
407        assert!(cache
408            .lookup::<Vec<String>>("new", &path, &mut stats)
409            .is_none());
410        fs::write(&path, [0_u8; 65]).unwrap();
411        cache.insert("large".into(), &path, &0, &mut stats);
412        assert!(cache.state.lock().unwrap().entries.is_empty());
413    }
414
415    #[test]
416    fn disabled_cache_has_no_io_and_errors_fall_back() {
417        let missing = Path::new("/nonexistent/scan-reuse-test");
418        let cache = ReuseCache::new(CacheMode::Off, 1, 10);
419        let mut stats = CacheStats::new("test", CacheMode::Off);
420        cache.insert("key".into(), missing, &0, &mut stats);
421        assert!(cache.lookup::<u32>("key", missing, &mut stats).is_none());
422        assert_eq!((stats.lookups, stats.errors), (0, 0));
423        let cache = ReuseCache::new(CacheMode::Observe, 1, 10);
424        cache.insert("key".into(), missing, &0, &mut stats);
425        assert_eq!(stats.errors, 1);
426    }
427}