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