1use md5::{Digest, Md5};
2use std::collections::HashMap;
3use std::sync::OnceLock;
4use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
5use std::time::{Duration, Instant, SystemTime};
6
7use super::tokens::count_tokens;
8
9fn instant_base() -> Instant {
13 static BASE: OnceLock<Instant> = OnceLock::new();
14 *BASE.get_or_init(Instant::now)
15}
16
17fn encode_instant(i: Instant) -> u64 {
18 i.saturating_duration_since(instant_base()).as_millis() as u64
19}
20
21fn decode_instant(ms: u64) -> Instant {
22 instant_base() + Duration::from_millis(ms)
23}
24
25fn normalize_key(path: &str) -> String {
26 crate::core::pathutil::normalize_tool_path(path)
27}
28
29pub(crate) const DEFAULT_CACHE_MAX_TOKENS: usize = 2_000_000;
34
35fn resolve_cache_max_tokens(env: Option<&str>, configured: usize) -> usize {
42 if let Some(raw) = env
43 && let Ok(n) = raw.trim().parse::<usize>()
44 && n > 0
45 {
46 return n;
47 }
48 if configured > 0 {
49 configured
50 } else {
51 DEFAULT_CACHE_MAX_TOKENS
52 }
53}
54
55pub(crate) fn max_cache_tokens() -> usize {
60 resolve_cache_max_tokens(
61 std::env::var("LEAN_CTX_CACHE_MAX_TOKENS").ok().as_deref(),
62 crate::core::config::Config::load().cache_max_tokens,
63 )
64}
65
66#[derive(Debug)]
72pub struct CacheEntry {
73 compressed_content: Vec<u8>,
74 pub hash: String,
75 pub line_count: usize,
76 pub original_tokens: usize,
77 read_count: AtomicU32,
78 pub path: String,
79 last_access: AtomicU64,
80 pub stored_mtime: Option<SystemTime>,
81 pub compressed_outputs: HashMap<String, String>,
83 pub full_content_delivered: bool,
86 pub delivered_conversation: Option<String>,
91 pub last_mode: String,
93}
94
95const ZSTD_LEVEL: i32 = 3;
96
97fn zstd_compress(data: &str) -> Vec<u8> {
98 zstd::encode_all(data.as_bytes(), ZSTD_LEVEL).unwrap_or_else(|_| data.as_bytes().to_vec())
99}
100
101fn zstd_decompress(data: &[u8]) -> Option<String> {
102 zstd::decode_all(data)
103 .ok()
104 .and_then(|v| String::from_utf8(v).ok())
105}
106
107impl CacheEntry {
108 pub fn new(
110 content: &str,
111 hash: String,
112 line_count: usize,
113 original_tokens: usize,
114 path: String,
115 stored_mtime: Option<SystemTime>,
116 ) -> Self {
117 let compressed_content = zstd_compress(content);
118 Self {
119 compressed_content,
120 hash,
121 line_count,
122 original_tokens,
123 read_count: AtomicU32::new(1),
124 path,
125 last_access: AtomicU64::new(encode_instant(Instant::now())),
126 stored_mtime,
127 compressed_outputs: HashMap::new(),
128 full_content_delivered: false,
129 delivered_conversation: None,
130 last_mode: String::new(),
131 }
132 }
133
134 pub fn read_count(&self) -> u32 {
136 self.read_count.load(Ordering::Relaxed)
137 }
138
139 pub fn bump_read_count(&self) -> u32 {
141 self.read_count.fetch_add(1, Ordering::Relaxed) + 1
142 }
143
144 pub fn set_read_count(&self, n: u32) {
146 self.read_count.store(n, Ordering::Relaxed);
147 }
148
149 pub fn last_access(&self) -> Instant {
151 decode_instant(self.last_access.load(Ordering::Relaxed))
152 }
153
154 pub fn touch(&self) {
156 self.last_access
157 .store(encode_instant(Instant::now()), Ordering::Relaxed);
158 }
159
160 pub fn set_last_access(&self, when: Instant) {
162 self.last_access
163 .store(encode_instant(when), Ordering::Relaxed);
164 }
165
166 pub fn content(&self) -> Option<String> {
168 zstd_decompress(&self.compressed_content)
169 }
170
171 pub fn set_content(&mut self, content: &str) {
173 self.compressed_content = zstd_compress(content);
174 }
175
176 pub fn compressed_size(&self) -> usize {
178 self.compressed_content.len()
179 }
180}
181
182#[derive(Debug, Clone)]
184pub struct StoreResult {
185 pub line_count: usize,
186 pub original_tokens: usize,
187 pub read_count: u32,
188 pub was_hit: bool,
189 pub full_content_delivered: bool,
191}
192
193impl CacheEntry {
194 pub fn eviction_score_legacy(&self, now: Instant) -> f64 {
196 let elapsed = now
197 .checked_duration_since(self.last_access())
198 .unwrap_or_default()
199 .as_secs_f64();
200 let recency = 1.0 / (1.0 + elapsed.sqrt());
201 let frequency = (self.read_count() as f64 + 1.0).ln();
202 let size_value = (self.original_tokens as f64 + 1.0).ln();
203 recency * 0.4 + frequency * 0.3 + size_value * 0.3
204 }
205
206 pub fn get_compressed(&self, mode_key: &str) -> Option<&String> {
207 self.compressed_outputs.get(mode_key)
208 }
209
210 pub fn set_compressed(&mut self, mode_key: &str, output: String) {
211 const MAX_COMPRESSED_VARIANTS: usize = 3;
212 if self.compressed_outputs.len() >= MAX_COMPRESSED_VARIANTS
213 && !self.compressed_outputs.contains_key(mode_key)
214 && let Some(oldest_key) = self.compressed_outputs.keys().next().cloned()
215 {
216 self.compressed_outputs.remove(&oldest_key);
217 }
218 self.compressed_outputs.insert(mode_key.to_string(), output);
219 }
220
221 pub fn mark_full_delivered(&mut self, conversation: Option<String>) {
222 self.full_content_delivered = true;
223 self.delivered_conversation = conversation;
224 }
225}
226
227const RRF_K: f64 = 60.0;
228
229const HEBBIAN_PROTECT_WEIGHT: f64 = 0.05;
234const HEBBIAN_ACTIVE_SET: usize = 8;
237
238pub fn eviction_scores_rrf(entries: &[(&String, &CacheEntry)], now: Instant) -> Vec<(String, f64)> {
243 if entries.is_empty() {
244 return Vec::new();
245 }
246
247 let n = entries.len();
248
249 let mut recency_order: Vec<usize> = (0..n).collect();
250 recency_order.sort_by(|&a, &b| {
251 let elapsed_a = now
252 .checked_duration_since(entries[a].1.last_access())
253 .unwrap_or_default()
254 .as_secs_f64();
255 let elapsed_b = now
256 .checked_duration_since(entries[b].1.last_access())
257 .unwrap_or_default()
258 .as_secs_f64();
259 elapsed_a
260 .partial_cmp(&elapsed_b)
261 .unwrap_or(std::cmp::Ordering::Equal)
262 });
263
264 let mut frequency_order: Vec<usize> = (0..n).collect();
265 frequency_order.sort_by(|&a, &b| entries[b].1.read_count().cmp(&entries[a].1.read_count()));
266
267 let mut size_order: Vec<usize> = (0..n).collect();
268 size_order.sort_by(|&a, &b| {
269 entries[b]
270 .1
271 .original_tokens
272 .cmp(&entries[a].1.original_tokens)
273 });
274
275 let mut recency_ranks = vec![0usize; n];
276 let mut frequency_ranks = vec![0usize; n];
277 let mut size_ranks = vec![0usize; n];
278
279 for (rank, &idx) in recency_order.iter().enumerate() {
280 recency_ranks[idx] = rank;
281 }
282 for (rank, &idx) in frequency_order.iter().enumerate() {
283 frequency_ranks[idx] = rank;
284 }
285 for (rank, &idx) in size_order.iter().enumerate() {
286 size_ranks[idx] = rank;
287 }
288
289 entries
290 .iter()
291 .enumerate()
292 .map(|(i, (path, _))| {
293 let score = 1.0 / (RRF_K + recency_ranks[i] as f64)
294 + 1.0 / (RRF_K + frequency_ranks[i] as f64)
295 + 1.0 / (RRF_K + size_ranks[i] as f64);
296 ((*path).clone(), score)
297 })
298 .collect()
299}
300
301fn apply_hebbian_bonus(scores: &mut [(String, f64)], bonus: &HashMap<String, f64>) {
304 if bonus.is_empty() {
305 return;
306 }
307 for s in scores.iter_mut() {
308 if let Some(b) = bonus.get(&s.0) {
309 s.1 += *b;
310 }
311 }
312}
313
314#[derive(Debug, Default)]
319pub struct CacheStats {
320 total_reads: AtomicU64,
321 cache_hits: AtomicU64,
322 total_original_tokens: AtomicU64,
323 total_sent_tokens: AtomicU64,
324 files_tracked: AtomicU64,
325}
326
327impl CacheStats {
328 pub fn total_reads(&self) -> u64 {
330 self.total_reads.load(Ordering::Relaxed)
331 }
332
333 pub fn cache_hits(&self) -> u64 {
335 self.cache_hits.load(Ordering::Relaxed)
336 }
337
338 pub fn total_original_tokens(&self) -> u64 {
340 self.total_original_tokens.load(Ordering::Relaxed)
341 }
342
343 pub fn total_sent_tokens(&self) -> u64 {
345 self.total_sent_tokens.load(Ordering::Relaxed)
346 }
347
348 pub fn files_tracked(&self) -> u64 {
350 self.files_tracked.load(Ordering::Relaxed)
351 }
352
353 pub fn hit_rate(&self) -> f64 {
355 let total = self.total_reads();
356 if total == 0 {
357 return 0.0;
358 }
359 (self.cache_hits() as f64 / total as f64) * 100.0
360 }
361
362 pub fn tokens_saved(&self) -> u64 {
364 self.total_original_tokens()
365 .saturating_sub(self.total_sent_tokens())
366 }
367
368 pub fn savings_percent(&self) -> f64 {
370 let original = self.total_original_tokens();
371 if original == 0 {
372 return 0.0;
373 }
374 (self.tokens_saved() as f64 / original as f64) * 100.0
375 }
376}
377
378#[derive(Clone, Debug)]
380pub struct SharedBlock {
381 pub canonical_path: String,
382 pub canonical_ref: String,
383 pub start_line: usize,
384 pub end_line: usize,
385 pub content: String,
386}
387
388pub struct SessionCache {
391 entries: HashMap<String, CacheEntry>,
392 file_refs: HashMap<String, String>,
393 next_ref: usize,
394 stats: CacheStats,
395 shared_blocks: Vec<SharedBlock>,
396 co_access: crate::core::hebbian_cache::CoAccessMatrix,
400}
401
402impl Default for SessionCache {
403 fn default() -> Self {
404 Self::new()
405 }
406}
407
408impl SessionCache {
409 pub fn new() -> Self {
411 Self {
412 entries: HashMap::new(),
413 file_refs: HashMap::new(),
414 next_ref: 1,
415 shared_blocks: Vec::new(),
416 stats: CacheStats::default(),
417 co_access: crate::core::hebbian_cache::CoAccessMatrix::new(),
418 }
419 }
420
421 pub fn record_co_access(&mut self, path: &str) {
425 let key = normalize_key(path);
426 self.co_access
427 .record_access(crate::core::hebbian_cache::path_hash(&key));
428 }
429
430 pub fn flush_co_access(&mut self) {
433 self.co_access.end_burst();
434 }
435
436 pub(crate) fn hebbian_eviction_bonus(&self) -> HashMap<String, f64> {
442 use crate::core::hebbian_cache::path_hash;
443 if self.entries.is_empty() {
444 return HashMap::new();
445 }
446 let mut by_recency: Vec<(&String, Instant)> = self
447 .entries
448 .iter()
449 .map(|(k, e)| (k, e.last_access()))
450 .collect();
451 by_recency.sort_by_key(|(_, t)| std::cmp::Reverse(*t));
452 let active: Vec<u64> = by_recency
453 .iter()
454 .take(HEBBIAN_ACTIVE_SET)
455 .map(|(k, _)| path_hash(k))
456 .collect();
457
458 let mut out = HashMap::new();
459 for k in self.entries.keys() {
460 let h = path_hash(k);
461 let peers: Vec<u64> = active.iter().copied().filter(|&a| a != h).collect();
463 let strength = self.co_access.association_strength(h, &peers);
464 if strength > 0.0 {
465 out.insert(k.clone(), f64::from(strength) * HEBBIAN_PROTECT_WEIGHT);
466 }
467 }
468 if !out.is_empty() {
469 crate::core::introspect::tick("hebbian_cache");
470 }
471 out
472 }
473
474 pub fn get_file_ref(&mut self, path: &str) -> String {
476 let key = normalize_key(path);
477 if let Some(r) = self.file_refs.get(&key) {
478 return r.clone();
479 }
480 let r = format!("F{}", self.next_ref);
481 self.next_ref += 1;
482 self.file_refs.insert(key, r.clone());
483 r
484 }
485
486 pub fn get_file_ref_readonly(&self, path: &str) -> Option<String> {
488 self.file_refs.get(&normalize_key(path)).cloned()
489 }
490
491 pub fn get(&self, path: &str) -> Option<&CacheEntry> {
493 self.entries.get(&normalize_key(path))
494 }
495
496 pub fn get_mut(&mut self, path: &str) -> Option<&mut CacheEntry> {
498 self.entries.get_mut(&normalize_key(path))
499 }
500
501 pub fn get_full_content(&self, path: &str) -> Option<String> {
504 self.entries
505 .get(&normalize_key(path))
506 .and_then(CacheEntry::content)
507 }
508
509 pub fn current_full_content(&self, path: &str) -> Option<(String, usize)> {
522 let entry = self.entries.get(&normalize_key(path))?;
523 if is_cache_entry_stale_verified(&entry.path, entry.stored_mtime, &entry.hash)
524 && let Ok(fresh) = crate::core::io_boundary::read_file_lossy(&entry.path)
525 {
526 let tokens = count_tokens(&fresh);
531 return Some((fresh, tokens));
532 }
533 Some((entry.content()?, entry.original_tokens))
534 }
535
536 pub fn record_cache_hit(&self, path: &str) -> Option<&CacheEntry> {
542 let key = normalize_key(path);
543 let ref_label = self
544 .file_refs
545 .get(&key)
546 .cloned()
547 .unwrap_or_else(|| "F?".to_string());
548 let entry = self.entries.get(&key)?;
549 let new_count = entry.bump_read_count();
550 entry.touch();
551 self.stats.total_reads.fetch_add(1, Ordering::Relaxed);
552 self.stats.cache_hits.fetch_add(1, Ordering::Relaxed);
553 self.stats
554 .total_original_tokens
555 .fetch_add(entry.original_tokens as u64, Ordering::Relaxed);
556 let hit_msg = format!("{ref_label} cached {new_count}t {}L", entry.line_count);
557 let sent_tokens = count_tokens(&hit_msg) as u64;
558 self.stats
559 .total_sent_tokens
560 .fetch_add(sent_tokens, Ordering::Relaxed);
561 crate::core::events::emit_cache_hit(
562 path,
563 (entry.original_tokens as u64).saturating_sub(sent_tokens),
564 );
565 Some(entry)
566 }
567
568 pub fn store(&mut self, path: &str, content: &str) -> StoreResult {
570 let key = normalize_key(path);
571 self.co_access
574 .record_access(crate::core::hebbian_cache::path_hash(&key));
575 let hash = compute_md5(content);
576 let line_count = content.lines().count();
577 let original_tokens = count_tokens(content);
578 let stored_mtime = std::fs::metadata(path).and_then(|m| m.modified()).ok();
579 let now = Instant::now();
580
581 self.stats.total_reads.fetch_add(1, Ordering::Relaxed);
582 self.stats
583 .total_original_tokens
584 .fetch_add(original_tokens as u64, Ordering::Relaxed);
585
586 if let Some(existing) = self.entries.get_mut(&key) {
587 existing.set_last_access(now);
588 if stored_mtime.is_some() {
589 existing.stored_mtime = stored_mtime;
590 }
591 if existing.hash == hash {
592 let new_count = existing.bump_read_count();
593 self.stats.cache_hits.fetch_add(1, Ordering::Relaxed);
594 let hit_msg = format!(
595 "{} cached {new_count}t {}L",
596 self.file_refs.get(&key).unwrap_or(&"F?".to_string()),
597 existing.line_count,
598 );
599 let sent_tokens = count_tokens(&hit_msg) as u64;
600 self.stats
601 .total_sent_tokens
602 .fetch_add(sent_tokens, Ordering::Relaxed);
603 return StoreResult {
604 line_count: existing.line_count,
605 original_tokens: existing.original_tokens,
606 read_count: new_count,
607 was_hit: true,
608 full_content_delivered: existing.full_content_delivered,
609 };
610 }
611 existing.compressed_outputs.clear();
612 existing.set_content(content);
613 existing.hash = hash;
614 existing.line_count = line_count;
615 existing.original_tokens = original_tokens;
616 let new_count = existing.bump_read_count();
617 existing.full_content_delivered = false;
618 existing.delivered_conversation = None;
619 if stored_mtime.is_some() {
620 existing.stored_mtime = stored_mtime;
621 }
622 self.stats
623 .total_sent_tokens
624 .fetch_add(original_tokens as u64, Ordering::Relaxed);
625 return StoreResult {
626 line_count,
627 original_tokens,
628 read_count: new_count,
629 was_hit: false,
630 full_content_delivered: false,
631 };
632 }
633
634 self.evict_if_needed(original_tokens);
635 self.get_file_ref(&key);
636
637 let entry = CacheEntry::new(
638 content,
639 hash,
640 line_count,
641 original_tokens,
642 key.clone(),
643 stored_mtime,
644 );
645
646 self.entries.insert(key, entry);
647 self.stats.files_tracked.fetch_add(1, Ordering::Relaxed);
648 self.stats
649 .total_sent_tokens
650 .fetch_add(original_tokens as u64, Ordering::Relaxed);
651 StoreResult {
652 line_count,
653 original_tokens,
654 read_count: 1,
655 was_hit: false,
656 full_content_delivered: false,
657 }
658 }
659
660 pub fn total_cached_tokens(&self) -> usize {
662 self.entries.values().map(|e| e.original_tokens).sum()
663 }
664
665 pub fn evict_if_needed(&mut self, incoming_tokens: usize) {
668 let max_tokens = max_cache_tokens();
669 let current = self.total_cached_tokens();
670 if current + incoming_tokens <= max_tokens {
671 return;
672 }
673
674 let now = Instant::now();
675 let all: Vec<(&String, &CacheEntry)> = self.entries.iter().collect();
676 let mut scores = eviction_scores_rrf(&all, now);
677 apply_hebbian_bonus(&mut scores, &self.hebbian_eviction_bonus());
678 scores.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
680
681 let mut freed = 0usize;
682 let mut redelivered = 0u64;
683 let target = (current + incoming_tokens).saturating_sub(max_tokens);
684
685 for (path, _score) in &scores {
686 if freed >= target {
687 break;
688 }
689 if let Some(entry) = self.entries.remove(path) {
690 freed += entry.original_tokens;
691 if entry.full_content_delivered {
692 redelivered += 1;
693 }
694 self.file_refs.remove(path);
695 }
696 }
697 crate::core::cache_telemetry::record_eviction(redelivered);
698 }
699
700 pub fn get_all_entries(&self) -> Vec<(&String, &CacheEntry)> {
702 self.entries.iter().collect()
703 }
704
705 pub fn get_stats(&self) -> &CacheStats {
707 &self.stats
708 }
709
710 pub fn file_ref_map(&self) -> &HashMap<String, String> {
712 &self.file_refs
713 }
714
715 pub fn set_shared_blocks(&mut self, blocks: Vec<SharedBlock>) {
717 self.shared_blocks = blocks;
718 }
719
720 pub fn get_shared_blocks(&self) -> &[SharedBlock] {
722 &self.shared_blocks
723 }
724
725 pub fn apply_dedup(&self, path: &str, content: &str) -> Option<String> {
727 if self.shared_blocks.is_empty() {
728 return None;
729 }
730 let refs: Vec<&SharedBlock> = self
731 .shared_blocks
732 .iter()
733 .filter(|b| b.canonical_path != path && content.contains(&b.content))
734 .collect();
735 if refs.is_empty() {
736 return None;
737 }
738 let mut result = content.to_string();
739 for block in refs {
740 result = result.replacen(
741 &block.content,
742 &format!(
743 "[= {}:{}-{}]",
744 block.canonical_ref, block.start_line, block.end_line
745 ),
746 1,
747 );
748 }
749 Some(result)
750 }
751
752 pub fn invalidate(&mut self, path: &str) -> bool {
754 self.entries.remove(&normalize_key(path)).is_some()
755 }
756
757 pub fn get_compressed(&self, path: &str, mode_key: &str) -> Option<&String> {
760 let key = normalize_key(path);
761 let entry = self.entries.get(&key)?;
762 let result = entry.get_compressed(mode_key)?;
763 entry.bump_read_count();
764 entry.touch();
765 self.stats.total_reads.fetch_add(1, Ordering::Relaxed);
766 self.stats.cache_hits.fetch_add(1, Ordering::Relaxed);
767 self.stats
768 .total_original_tokens
769 .fetch_add(entry.original_tokens as u64, Ordering::Relaxed);
770 let sent = count_tokens(result) as u64;
771 self.stats
772 .total_sent_tokens
773 .fetch_add(sent, Ordering::Relaxed);
774 crate::core::events::emit_cache_hit(
775 path,
776 (entry.original_tokens as u64).saturating_sub(sent),
777 );
778 Some(result)
779 }
780
781 pub fn mark_full_delivered(&mut self, path: &str) {
786 let conversation = crate::core::conversation::current_conversation_id();
787 let key = normalize_key(path);
788 let file_ref = self.file_refs.get(&key).cloned();
789 if let Some(entry) = self.entries.get_mut(&key) {
790 entry.mark_full_delivered(conversation.clone());
791 crate::core::read_stub_index::record(crate::core::read_stub_index::StubRecord::new(
795 key.clone(),
796 entry.hash.clone(),
797 entry.stored_mtime,
798 entry.line_count,
799 file_ref.unwrap_or_default(),
800 conversation,
801 ));
802 }
803 }
804
805 pub fn set_compressed(&mut self, path: &str, mode_key: &str, output: String) {
807 if let Some(entry) = self.entries.get_mut(&normalize_key(path)) {
808 entry.set_compressed(mode_key, output);
809 }
810 }
811
812 pub fn reset_delivery_flags(&mut self) -> usize {
816 let mut count = 0;
817 for entry in self.entries.values_mut() {
818 if entry.full_content_delivered {
819 entry.full_content_delivered = false;
820 count += 1;
821 }
822 }
823 count
824 }
825
826 pub fn is_full_delivered(&self, path: &str) -> bool {
828 self.entries
829 .get(&normalize_key(path))
830 .is_some_and(|e| e.full_content_delivered)
831 }
832
833 pub fn count_full_delivered(&self) -> usize {
837 self.entries
838 .values()
839 .filter(|e| e.full_content_delivered)
840 .count()
841 }
842
843 pub fn trim_compressed_outputs(&mut self) -> usize {
846 let mut trimmed = 0;
847 for entry in self.entries.values_mut() {
848 if !entry.compressed_outputs.is_empty() {
849 entry.compressed_outputs.clear();
850 trimmed += 1;
851 }
852 }
853 trimmed
854 }
855
856 pub fn evict_probationary(&mut self) -> usize {
859 let to_remove: Vec<String> = self
860 .entries
861 .iter()
862 .filter(|(_, e)| e.read_count() <= 1)
863 .map(|(k, _)| k.clone())
864 .collect();
865 let count = to_remove.len();
866 let mut redelivered = 0u64;
867 for key in &to_remove {
868 if self
869 .entries
870 .remove(key)
871 .is_some_and(|e| e.full_content_delivered)
872 {
873 redelivered += 1;
874 }
875 self.file_refs.remove(key);
876 }
877 crate::core::cache_telemetry::record_eviction(redelivered);
878 count
879 }
880
881 pub fn evict_to_budget(&mut self, target_tokens: usize) {
883 let current = self.total_cached_tokens();
884 if current <= target_tokens {
885 return;
886 }
887 let now = Instant::now();
888 let all: Vec<(&String, &CacheEntry)> = self.entries.iter().collect();
889 let mut scores = eviction_scores_rrf(&all, now);
890 apply_hebbian_bonus(&mut scores, &self.hebbian_eviction_bonus());
891 scores.sort_by(|a, b| a.1.partial_cmp(&b.1).unwrap_or(std::cmp::Ordering::Equal));
892
893 let mut freed = 0usize;
894 let mut redelivered = 0u64;
895 let target_free = current.saturating_sub(target_tokens);
896 for (path, _score) in &scores {
897 if freed >= target_free {
898 break;
899 }
900 if let Some(entry) = self.entries.remove(path) {
901 freed += entry.original_tokens;
902 if entry.full_content_delivered {
903 redelivered += 1;
904 }
905 self.file_refs.remove(path);
906 }
907 }
908 crate::core::cache_telemetry::record_eviction(redelivered);
909 }
910
911 pub fn approximate_bytes(&self) -> usize {
913 let entries_bytes: usize = self
914 .entries
915 .values()
916 .map(|e| {
917 e.compressed_content.len()
918 + e.hash.len()
919 + e.path.len()
920 + e.compressed_outputs
921 .iter()
922 .map(|(k, v)| k.len() + v.len())
923 .sum::<usize>()
924 + 128 })
926 .sum();
927 let refs_bytes: usize = self.file_refs.iter().map(|(k, v)| k.len() + v.len()).sum();
928 let blocks_bytes: usize = self
929 .shared_blocks
930 .iter()
931 .map(|b| b.canonical_path.len() + b.canonical_ref.len() + b.content.len() + 32)
932 .sum();
933 entries_bytes + refs_bytes + blocks_bytes
934 }
935
936 const MAX_SHARED_BLOCKS: usize = 100;
937
938 pub fn trim_shared_blocks(&mut self) {
940 if self.shared_blocks.len() > Self::MAX_SHARED_BLOCKS {
941 let excess = self.shared_blocks.len() - Self::MAX_SHARED_BLOCKS;
942 self.shared_blocks.drain(..excess);
943 }
944 }
945
946 pub fn clear(&mut self) -> usize {
948 let count = self.entries.len();
949 self.entries.clear();
950 self.file_refs.clear();
951 self.shared_blocks.clear();
952 self.next_ref = 1;
953 self.stats = CacheStats::default();
954 count
955 }
956}
957
958pub fn file_mtime(path: &str) -> Option<SystemTime> {
959 std::fs::metadata(path).and_then(|m| m.modified()).ok()
960}
961
962pub fn is_cache_entry_stale(path: &str, cached_mtime: Option<SystemTime>) -> bool {
963 let current = file_mtime(path);
964 match (cached_mtime, current) {
965 (None, None) => false,
967 (Some(_), None) | (None, Some(_)) => true,
969 (Some(cached), Some(current)) => current != cached,
972 }
973}
974
975const VERIFY_HASH_CAP_BYTES: u64 = 8 * 1024 * 1024;
978
979fn cache_verify_enabled() -> bool {
980 std::env::var("LEAN_CTX_CACHE_VERIFY").map_or(true, |v| v != "0")
981}
982
983pub fn is_cache_entry_stale_verified(
998 path: &str,
999 cached_mtime: Option<SystemTime>,
1000 cached_hash: &str,
1001) -> bool {
1002 if is_cache_entry_stale(path, cached_mtime) {
1003 return true;
1004 }
1005 if cached_hash.is_empty() || !cache_verify_enabled() {
1006 return false;
1007 }
1008 let Ok(meta) = std::fs::metadata(path) else {
1009 return true;
1011 };
1012 if meta.len() > VERIFY_HASH_CAP_BYTES {
1013 return false;
1014 }
1015 match std::fs::read(path) {
1016 Ok(bytes) => compute_md5(&String::from_utf8_lossy(&bytes)) != cached_hash,
1018 Err(_) => true,
1019 }
1020}
1021
1022fn compute_md5(content: &str) -> String {
1023 let mut hasher = Md5::new();
1024 hasher.update(content.as_bytes());
1025 crate::core::agent_identity::hex_encode(&hasher.finalize())
1026}
1027
1028#[cfg(test)]
1029mod tests {
1030 use super::*;
1031 use std::time::Duration;
1032
1033 #[test]
1034 fn cache_stores_and_retrieves() {
1035 let mut cache = SessionCache::new();
1036 let result = cache.store("/test/file.rs", "fn main() {}");
1037 assert!(!result.was_hit);
1038 assert_eq!(result.line_count, 1);
1039 assert!(cache.get("/test/file.rs").is_some());
1040 }
1041
1042 #[test]
1043 fn cache_hit_on_same_content() {
1044 let mut cache = SessionCache::new();
1045 cache.store("/test/file.rs", "content");
1046 let result = cache.store("/test/file.rs", "content");
1047 assert!(result.was_hit, "same content should be a cache hit");
1048 }
1049
1050 #[test]
1051 fn cache_miss_on_changed_content() {
1052 let mut cache = SessionCache::new();
1053 cache.store("/test/file.rs", "old content");
1054 let result = cache.store("/test/file.rs", "new content");
1055 assert!(!result.was_hit, "changed content should not be a cache hit");
1056 }
1057
1058 #[test]
1059 fn file_refs_are_sequential() {
1060 let mut cache = SessionCache::new();
1061 assert_eq!(cache.get_file_ref("/a.rs"), "F1");
1062 assert_eq!(cache.get_file_ref("/b.rs"), "F2");
1063 assert_eq!(cache.get_file_ref("/a.rs"), "F1"); }
1065
1066 #[test]
1067 fn cache_clear_resets_everything() {
1068 let mut cache = SessionCache::new();
1069 cache.store("/a.rs", "a");
1070 cache.store("/b.rs", "b");
1071 let count = cache.clear();
1072 assert_eq!(count, 2);
1073 assert!(cache.get("/a.rs").is_none());
1074 assert_eq!(cache.get_file_ref("/c.rs"), "F1"); }
1076
1077 #[test]
1078 fn cache_invalidate_removes_entry() {
1079 let mut cache = SessionCache::new();
1080 cache.store("/test.rs", "test");
1081 assert!(cache.invalidate("/test.rs"));
1082 assert!(!cache.invalidate("/nonexistent.rs"));
1083 }
1084
1085 #[test]
1086 fn cache_stats_track_correctly() {
1087 let mut cache = SessionCache::new();
1088 cache.store("/a.rs", "hello");
1089 cache.store("/a.rs", "hello"); let stats = cache.get_stats();
1091 assert_eq!(stats.total_reads(), 2);
1092 assert_eq!(stats.cache_hits(), 1);
1093 assert!(stats.hit_rate() > 0.0);
1094 }
1095
1096 #[test]
1097 fn current_full_content_serves_cached_when_fresh() {
1098 let dir = tempfile::tempdir().unwrap();
1099 let file = dir.path().join("handover.md");
1100 std::fs::write(&file, "HANDOVER V1\n").unwrap();
1101 let path = file.to_str().unwrap();
1102
1103 let mut cache = SessionCache::new();
1104 cache.store(path, "HANDOVER V1\n");
1105
1106 let (content, tokens) = cache.current_full_content(path).unwrap();
1107 assert_eq!(content, "HANDOVER V1\n");
1108 assert!(tokens > 0);
1109 }
1110
1111 #[test]
1112 fn current_full_content_rereads_when_file_changed() {
1113 let dir = tempfile::tempdir().unwrap();
1116 let file = dir.path().join("handover.md");
1117 std::fs::write(&file, "HANDOVER V1\n").unwrap();
1118 let path = file.to_str().unwrap();
1119
1120 let mut cache = SessionCache::new();
1121 cache.store(path, "HANDOVER V1\n");
1122
1123 std::thread::sleep(std::time::Duration::from_millis(10));
1125 std::fs::write(&file, "HANDOVER V2 CHANGED\n").unwrap();
1126
1127 let (content, _) = cache.current_full_content(path).unwrap();
1128 assert_eq!(
1129 content, "HANDOVER V2 CHANGED\n",
1130 "stale cached copy must be re-read from disk, not served as-is"
1131 );
1132 }
1133
1134 #[test]
1135 fn current_full_content_none_without_entry() {
1136 let cache = SessionCache::new();
1137 assert!(cache.current_full_content("/no/such/file.rs").is_none());
1138 }
1139
1140 #[test]
1141 fn current_full_content_falls_back_to_cache_when_file_unreadable() {
1142 let dir = tempfile::tempdir().unwrap();
1147 let canon = dir.path().canonicalize().unwrap();
1148 let file = canon.join("gone.md");
1149 std::fs::write(&file, "ORIGINAL\n").unwrap();
1150 let path = file.to_str().unwrap().to_string();
1151
1152 let mut cache = SessionCache::new();
1153 cache.store(&path, "ORIGINAL\n");
1154 std::fs::remove_file(&file).unwrap();
1155
1156 let (content, _) = cache.current_full_content(&path).unwrap();
1157 assert_eq!(
1158 content, "ORIGINAL\n",
1159 "unreadable file must fall back to last-known cached content"
1160 );
1161 }
1162
1163 #[test]
1164 fn record_cache_hit_works_through_shared_ref() {
1165 let mut cache = SessionCache::new();
1166 cache.store("/x.rs", "hello world");
1167 let shared: &SessionCache = &cache;
1169 assert!(shared.record_cache_hit("/x.rs").is_some());
1170 assert!(shared.record_cache_hit("/x.rs").is_some());
1171 assert_eq!(cache.get("/x.rs").unwrap().read_count(), 3);
1173 assert_eq!(cache.get_stats().cache_hits(), 2);
1174 }
1175
1176 #[test]
1177 fn concurrent_cache_hits_are_lossless() {
1178 use std::sync::Arc;
1179 let mut cache = SessionCache::new();
1180 cache.store("/a.rs", "a");
1181 cache.store("/b.rs", "b");
1182 let cache = Arc::new(cache);
1185 let threads = 8;
1186 let iters = 1_000;
1187 let handles: Vec<_> = (0..threads)
1188 .map(|_| {
1189 let c = Arc::clone(&cache);
1190 std::thread::spawn(move || {
1191 for _ in 0..iters {
1192 c.record_cache_hit("/a.rs");
1193 c.record_cache_hit("/b.rs");
1194 }
1195 })
1196 })
1197 .collect();
1198 for h in handles {
1199 h.join().unwrap();
1200 }
1201 let total = (threads * iters) as u64;
1202 assert_eq!(cache.get_stats().cache_hits(), total * 2);
1203 assert_eq!(cache.get("/a.rs").unwrap().read_count(), 1 + total as u32);
1204 assert_eq!(cache.get("/b.rs").unwrap().read_count(), 1 + total as u32);
1205 }
1206
1207 #[test]
1208 fn hebbian_eviction_bonus_is_wired() {
1209 let _ = count_tokens("warmup");
1218 let mut cache = SessionCache::new();
1219 cache.store("/a.rs", "fn a() {}");
1220 cache.store("/b.rs", "fn b() {}");
1221 cache.flush_co_access(); let bonus = cache.hebbian_eviction_bonus();
1223 assert!(
1224 !bonus.is_empty(),
1225 "co-accessed reads must yield a Hebbian eviction bonus (#3 wired)"
1226 );
1227 }
1228
1229 #[test]
1230 fn md5_is_deterministic() {
1231 let h1 = compute_md5("test content");
1232 let h2 = compute_md5("test content");
1233 assert_eq!(h1, h2);
1234 assert_ne!(h1, compute_md5("different"));
1235 }
1236
1237 #[test]
1238 fn rrf_eviction_prefers_recent() {
1239 let key_a = "a.rs".to_string();
1240 let key_b = "b.rs".to_string();
1241 let recent = CacheEntry::new("a", "h1".to_string(), 1, 10, "/a.rs".to_string(), None);
1244 let old = CacheEntry::new("b", "h2".to_string(), 1, 10, "/b.rs".to_string(), None);
1245 let t_old = Instant::now();
1246 std::thread::sleep(std::time::Duration::from_millis(10));
1247 let t_recent = Instant::now();
1248 old.set_last_access(t_old);
1249 recent.set_last_access(t_recent);
1250 let now = Instant::now();
1251 let entries: Vec<(&String, &CacheEntry)> = vec![(&key_a, &recent), (&key_b, &old)];
1252 let scores = eviction_scores_rrf(&entries, now);
1253 let score_a = scores.iter().find(|(p, _)| p == "a.rs").unwrap().1;
1254 let score_b = scores.iter().find(|(p, _)| p == "b.rs").unwrap().1;
1255 assert!(
1256 score_a > score_b,
1257 "recently accessed entries should score higher via RRF"
1258 );
1259 }
1260
1261 #[test]
1262 fn rrf_eviction_prefers_frequent() {
1263 let now = Instant::now();
1264 let key_a = "a.rs".to_string();
1265 let key_b = "b.rs".to_string();
1266 let frequent = {
1267 let e = CacheEntry::new("a", "h1".to_string(), 1, 10, "/a.rs".to_string(), None);
1268 e.set_read_count(20);
1269 e
1270 };
1271 let rare = CacheEntry::new("b", "h2".to_string(), 1, 10, "/b.rs".to_string(), None);
1272 let entries: Vec<(&String, &CacheEntry)> = vec![(&key_a, &frequent), (&key_b, &rare)];
1273 let scores = eviction_scores_rrf(&entries, now);
1274 let score_a = scores.iter().find(|(p, _)| p == "a.rs").unwrap().1;
1275 let score_b = scores.iter().find(|(p, _)| p == "b.rs").unwrap().1;
1276 assert!(
1277 score_a > score_b,
1278 "frequently accessed entries should score higher via RRF"
1279 );
1280 }
1281
1282 #[test]
1283 fn cache_budget_resolver_precedence() {
1284 assert_eq!(resolve_cache_max_tokens(Some("250000"), 999), 250_000);
1286 assert_eq!(resolve_cache_max_tokens(Some(" 80000 "), 0), 80_000);
1287 assert_eq!(resolve_cache_max_tokens(Some("0"), 123_456), 123_456);
1289 assert_eq!(resolve_cache_max_tokens(Some(""), 123_456), 123_456);
1290 assert_eq!(resolve_cache_max_tokens(Some("lots"), 123_456), 123_456);
1291 assert_eq!(resolve_cache_max_tokens(None, 42_000), 42_000);
1293 assert_eq!(resolve_cache_max_tokens(None, 0), DEFAULT_CACHE_MAX_TOKENS);
1295 assert_eq!(
1296 resolve_cache_max_tokens(Some("0"), 0),
1297 DEFAULT_CACHE_MAX_TOKENS
1298 );
1299 }
1300
1301 #[test]
1302 fn evict_if_needed_removes_lowest_score() {
1303 crate::test_env::set_var("LEAN_CTX_CACHE_MAX_TOKENS", "50");
1304 let mut cache = SessionCache::new();
1305 let big_content = "a]".repeat(30); cache.store("/old.rs", &big_content);
1307 let new_content = "b ".repeat(30); cache.store("/new.rs", &new_content);
1311 assert!(
1316 cache.total_cached_tokens() <= 60,
1317 "eviction should have kicked in"
1318 );
1319 crate::test_env::remove_var("LEAN_CTX_CACHE_MAX_TOKENS");
1320 }
1321
1322 #[test]
1323 fn stale_detection_flags_newer_file() {
1324 let dir = tempfile::tempdir().unwrap();
1325 let path = dir.path().join("stale.txt");
1326 let p = path.to_string_lossy().to_string();
1327
1328 std::fs::write(&path, "one").unwrap();
1329 let mut cache = SessionCache::new();
1330 cache.store(&p, "one");
1331
1332 let entry = cache.get(&p).unwrap();
1333 assert!(!is_cache_entry_stale(&p, entry.stored_mtime));
1334
1335 std::thread::sleep(Duration::from_secs(1));
1337 std::fs::write(&path, "two").unwrap();
1338
1339 let entry = cache.get(&p).unwrap();
1340 assert!(is_cache_entry_stale(&p, entry.stored_mtime));
1341 }
1342
1343 #[test]
1345 fn stale_detection_flags_backward_mtime() {
1346 let dir = tempfile::tempdir().unwrap();
1347 let path = dir.path().join("backward.txt");
1348 let p = path.to_string_lossy().to_string();
1349
1350 std::fs::write(&path, "one").unwrap();
1351 let mut cache = SessionCache::new();
1352 cache.store(&p, "one");
1353 let entry_mtime = cache.get(&p).unwrap().stored_mtime;
1354 assert!(!is_cache_entry_stale(&p, entry_mtime));
1355
1356 std::fs::write(&path, "zero").unwrap();
1358 let f = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
1359 f.set_modified(SystemTime::now() - Duration::from_hours(1))
1360 .unwrap();
1361 drop(f);
1362
1363 assert!(
1364 is_cache_entry_stale(&p, entry_mtime),
1365 "older mtime must read as stale"
1366 );
1367 }
1368
1369 #[test]
1372 fn verified_staleness_catches_same_mtime_content_change() {
1373 let dir = tempfile::tempdir().unwrap();
1374 let path = dir.path().join("sneaky.txt");
1375 let p = path.to_string_lossy().to_string();
1376
1377 std::fs::write(&path, "one").unwrap();
1378 let original_mtime = std::fs::metadata(&path).unwrap().modified().unwrap();
1379 let mut cache = SessionCache::new();
1380 cache.store(&p, "one");
1381 let (mtime, hash) = {
1382 let e = cache.get(&p).unwrap();
1383 (e.stored_mtime, e.hash.clone())
1384 };
1385
1386 assert!(!is_cache_entry_stale_verified(&p, mtime, &hash));
1388
1389 std::fs::write(&path, "two").unwrap();
1391 let f = std::fs::OpenOptions::new().write(true).open(&path).unwrap();
1392 f.set_modified(original_mtime).unwrap();
1393 drop(f);
1394
1395 assert!(
1396 !is_cache_entry_stale(&p, mtime),
1397 "test premise: the mtime check alone is fooled"
1398 );
1399 assert!(
1400 is_cache_entry_stale_verified(&p, mtime, &hash),
1401 "hash verification must catch the change"
1402 );
1403 }
1404
1405 #[test]
1406 fn verified_staleness_flags_unreadable_file() {
1407 let mut cache = SessionCache::new();
1408 cache.store("/nonexistent/file.rs", "content");
1409 let (mtime, hash) = {
1410 let e = cache.get("/nonexistent/file.rs").unwrap();
1411 (e.stored_mtime, e.hash.clone())
1412 };
1413 assert!(is_cache_entry_stale_verified(
1414 "/nonexistent/file.rs",
1415 mtime,
1416 &hash
1417 ));
1418 }
1419
1420 #[test]
1421 fn compressed_outputs_cached_and_retrieved() {
1422 let mut cache = SessionCache::new();
1423 cache.store("/test.rs", "fn main() {}");
1424 cache.set_compressed("/test.rs", "map", "compressed map output".to_string());
1425 assert_eq!(
1426 cache.get_compressed("/test.rs", "map"),
1427 Some(&"compressed map output".to_string())
1428 );
1429 assert_eq!(cache.get_compressed("/test.rs", "signatures"), None);
1430 }
1431
1432 #[test]
1433 fn compressed_outputs_cleared_on_content_change() {
1434 let mut cache = SessionCache::new();
1435 cache.store("/test.rs", "old content");
1436 cache.set_compressed("/test.rs", "map", "old map".to_string());
1437 assert!(cache.get_compressed("/test.rs", "map").is_some());
1438
1439 cache.store("/test.rs", "new content");
1440 assert_eq!(cache.get_compressed("/test.rs", "map"), None);
1441 }
1442
1443 #[test]
1444 fn compressed_outputs_survive_same_content_store() {
1445 let mut cache = SessionCache::new();
1446 cache.store("/test.rs", "content");
1447 cache.set_compressed("/test.rs", "map", "cached map".to_string());
1448
1449 let result = cache.store("/test.rs", "content");
1450 assert!(result.was_hit);
1451 assert_eq!(
1452 cache.get_compressed("/test.rs", "map"),
1453 Some(&"cached map".to_string())
1454 );
1455 }
1456
1457 #[test]
1458 fn compressed_outputs_cleared_on_invalidate() {
1459 let mut cache = SessionCache::new();
1460 cache.store("/test.rs", "content");
1461 cache.set_compressed("/test.rs", "signatures", "cached sigs".to_string());
1462 cache.invalidate("/test.rs");
1463 assert_eq!(cache.get_compressed("/test.rs", "signatures"), None);
1464 }
1465
1466 #[test]
1467 fn compressed_outputs_cleared_on_clear() {
1468 let mut cache = SessionCache::new();
1469 cache.store("/a.rs", "a");
1470 cache.set_compressed("/a.rs", "map", "map_a".to_string());
1471 cache.clear();
1472 assert_eq!(cache.get_compressed("/a.rs", "map"), None);
1473 }
1474}