use std::collections::HashMap;
use std::sync::OnceLock;
use std::sync::atomic::{AtomicU32, AtomicU64, Ordering};
use std::time::{Duration, Instant, SystemTime};
fn instant_base() -> Instant {
static BASE: OnceLock<Instant> = OnceLock::new();
*BASE.get_or_init(Instant::now)
}
fn encode_instant(i: Instant) -> u64 {
i.saturating_duration_since(instant_base()).as_millis() as u64
}
fn decode_instant(ms: u64) -> Instant {
instant_base() + Duration::from_millis(ms)
}
pub(super) fn normalize_key(path: &str) -> String {
crate::core::pathutil::normalize_tool_path(path)
}
pub(crate) const DEFAULT_CACHE_MAX_TOKENS: usize = 2_000_000;
pub(super) fn resolve_cache_max_tokens(env: Option<&str>, configured: usize) -> usize {
if let Some(raw) = env
&& let Ok(n) = raw.trim().parse::<usize>()
&& n > 0
{
return n;
}
if configured > 0 {
configured
} else {
DEFAULT_CACHE_MAX_TOKENS
}
}
pub(crate) fn max_cache_tokens() -> usize {
resolve_cache_max_tokens(
std::env::var("LEAN_CTX_CACHE_MAX_TOKENS").ok().as_deref(),
crate::core::config::Config::load().cache_max_tokens,
)
}
#[derive(Debug)]
pub struct CacheEntry {
pub(super) compressed_content: Vec<u8>,
pub hash: String,
pub line_count: usize,
pub original_tokens: usize,
read_count: AtomicU32,
reread_since_full_delivery: AtomicU32,
pub path: String,
last_access: AtomicU64,
pub stored_mtime: Option<SystemTime>,
pub compressed_outputs: HashMap<String, String>,
pub full_content_delivered: bool,
pub delivered_conversation: Option<String>,
pub last_mode: String,
}
const ZSTD_LEVEL: i32 = 3;
fn zstd_compress(data: &str) -> Vec<u8> {
zstd::encode_all(data.as_bytes(), ZSTD_LEVEL).unwrap_or_else(|_| data.as_bytes().to_vec())
}
fn zstd_decompress(data: &[u8]) -> Option<String> {
zstd::decode_all(data)
.ok()
.and_then(|v| String::from_utf8(v).ok())
}
impl CacheEntry {
pub fn new(
content: &str,
hash: String,
line_count: usize,
original_tokens: usize,
path: String,
stored_mtime: Option<SystemTime>,
) -> Self {
let compressed_content = zstd_compress(content);
Self {
compressed_content,
hash,
line_count,
original_tokens,
read_count: AtomicU32::new(1),
reread_since_full_delivery: AtomicU32::new(0),
path,
last_access: AtomicU64::new(encode_instant(Instant::now())),
stored_mtime,
compressed_outputs: HashMap::new(),
full_content_delivered: false,
delivered_conversation: None,
last_mode: String::new(),
}
}
pub fn read_count(&self) -> u32 {
self.read_count.load(Ordering::Relaxed)
}
pub fn bump_read_count(&self) -> u32 {
self.read_count.fetch_add(1, Ordering::Relaxed) + 1
}
pub fn bump_reread(&self) -> u32 {
self.reread_since_full_delivery
.fetch_add(1, Ordering::Relaxed)
+ 1
}
pub fn reset_reread_count(&self) {
self.reread_since_full_delivery.store(0, Ordering::Relaxed);
}
pub fn set_read_count(&self, n: u32) {
self.read_count.store(n, Ordering::Relaxed);
}
pub fn last_access(&self) -> Instant {
decode_instant(self.last_access.load(Ordering::Relaxed))
}
pub fn touch(&self) {
self.last_access
.store(encode_instant(Instant::now()), Ordering::Relaxed);
}
pub fn set_last_access(&self, when: Instant) {
self.last_access
.store(encode_instant(when), Ordering::Relaxed);
}
pub fn content(&self) -> Option<String> {
zstd_decompress(&self.compressed_content)
}
pub fn set_content(&mut self, content: &str) {
self.compressed_content = zstd_compress(content);
}
pub fn compressed_size(&self) -> usize {
self.compressed_content.len()
}
}
#[derive(Debug, Clone)]
pub struct StoreResult {
pub line_count: usize,
pub original_tokens: usize,
pub read_count: u32,
pub was_hit: bool,
pub full_content_delivered: bool,
}
impl CacheEntry {
pub fn eviction_score_legacy(&self, now: Instant) -> f64 {
let elapsed = now
.checked_duration_since(self.last_access())
.unwrap_or_default()
.as_secs_f64();
let recency = 1.0 / (1.0 + elapsed.sqrt());
let frequency = (self.read_count() as f64 + 1.0).ln();
let size_value = (self.original_tokens as f64 + 1.0).ln();
recency * 0.4 + frequency * 0.3 + size_value * 0.3
}
pub fn get_compressed(&self, mode_key: &str) -> Option<&String> {
self.compressed_outputs.get(mode_key)
}
pub fn set_compressed(&mut self, mode_key: &str, output: String) {
const MAX_COMPRESSED_VARIANTS: usize = 3;
if self.compressed_outputs.len() >= MAX_COMPRESSED_VARIANTS
&& !self.compressed_outputs.contains_key(mode_key)
&& let Some(oldest_key) = self.compressed_outputs.keys().next().cloned()
{
self.compressed_outputs.remove(&oldest_key);
}
self.compressed_outputs.insert(mode_key.to_string(), output);
}
pub fn mark_full_delivered(&mut self, conversation: Option<String>) {
self.full_content_delivered = true;
self.reset_reread_count();
self.delivered_conversation = conversation;
}
}
const RRF_K: f64 = 60.0;
pub(super) const HEBBIAN_PROTECT_WEIGHT: f64 = 0.05;
pub(super) const HEBBIAN_ACTIVE_SET: usize = 8;
pub fn eviction_scores_rrf(entries: &[(&String, &CacheEntry)], now: Instant) -> Vec<(String, f64)> {
if entries.is_empty() {
return Vec::new();
}
let n = entries.len();
let mut recency_order: Vec<usize> = (0..n).collect();
recency_order.sort_by(|&a, &b| {
let elapsed_a = now
.checked_duration_since(entries[a].1.last_access())
.unwrap_or_default()
.as_secs_f64();
let elapsed_b = now
.checked_duration_since(entries[b].1.last_access())
.unwrap_or_default()
.as_secs_f64();
elapsed_a
.partial_cmp(&elapsed_b)
.unwrap_or(std::cmp::Ordering::Equal)
});
let mut frequency_order: Vec<usize> = (0..n).collect();
frequency_order.sort_by(|&a, &b| entries[b].1.read_count().cmp(&entries[a].1.read_count()));
let mut size_order: Vec<usize> = (0..n).collect();
size_order.sort_by(|&a, &b| {
entries[b]
.1
.original_tokens
.cmp(&entries[a].1.original_tokens)
});
let mut recency_ranks = vec![0usize; n];
let mut frequency_ranks = vec![0usize; n];
let mut size_ranks = vec![0usize; n];
for (rank, &idx) in recency_order.iter().enumerate() {
recency_ranks[idx] = rank;
}
for (rank, &idx) in frequency_order.iter().enumerate() {
frequency_ranks[idx] = rank;
}
for (rank, &idx) in size_order.iter().enumerate() {
size_ranks[idx] = rank;
}
entries
.iter()
.enumerate()
.map(|(i, (path, _))| {
let score = 1.0 / (RRF_K + recency_ranks[i] as f64)
+ 1.0 / (RRF_K + frequency_ranks[i] as f64)
+ 1.0 / (RRF_K + size_ranks[i] as f64);
((*path).clone(), score)
})
.collect()
}
pub(super) fn apply_hebbian_bonus(scores: &mut [(String, f64)], bonus: &HashMap<String, f64>) {
if bonus.is_empty() {
return;
}
for s in scores.iter_mut() {
if let Some(b) = bonus.get(&s.0) {
s.1 += *b;
}
}
}
#[derive(Debug, Default)]
pub struct CacheStats {
pub(super) total_reads: AtomicU64,
pub(super) cache_hits: AtomicU64,
pub(super) total_original_tokens: AtomicU64,
pub(super) total_sent_tokens: AtomicU64,
pub(super) files_tracked: AtomicU64,
}
impl CacheStats {
pub fn total_reads(&self) -> u64 {
self.total_reads.load(Ordering::Relaxed)
}
pub fn cache_hits(&self) -> u64 {
self.cache_hits.load(Ordering::Relaxed)
}
pub fn total_original_tokens(&self) -> u64 {
self.total_original_tokens.load(Ordering::Relaxed)
}
pub fn total_sent_tokens(&self) -> u64 {
self.total_sent_tokens.load(Ordering::Relaxed)
}
pub fn files_tracked(&self) -> u64 {
self.files_tracked.load(Ordering::Relaxed)
}
pub fn hit_rate(&self) -> f64 {
let total = self.total_reads();
if total == 0 {
return 0.0;
}
(self.cache_hits() as f64 / total as f64) * 100.0
}
pub fn tokens_saved(&self) -> u64 {
self.total_original_tokens()
.saturating_sub(self.total_sent_tokens())
}
pub fn savings_percent(&self) -> f64 {
let original = self.total_original_tokens();
if original == 0 {
return 0.0;
}
(self.tokens_saved() as f64 / original as f64) * 100.0
}
}
#[derive(Clone, Debug)]
pub struct SharedBlock {
pub canonical_path: String,
pub canonical_ref: String,
pub start_line: usize,
pub end_line: usize,
pub content: String,
}