use super::clock::{ChangeJournal, Clock};
use super::metadata::MetadataCache;
use std::sync::{Arc, OnceLock};
use tokio::sync::Semaphore;
use crate::core::{NormalizedPath, Result};
use crate::hash::ContentHash;
const HASH_WORKERS_ENV: &str = "ZCCACHE_HASH_WORKERS";
fn hash_semaphore() -> &'static Arc<Semaphore> {
static HASH_SEMAPHORE: OnceLock<Arc<Semaphore>> = OnceLock::new();
HASH_SEMAPHORE.get_or_init(|| {
let permits = resolve_hash_workers(std::env::var(HASH_WORKERS_ENV).ok().as_deref());
Arc::new(Semaphore::new(permits))
})
}
const HASH_SEMAPHORE_UNLIMITED: usize = usize::MAX >> 4;
fn resolve_hash_workers(env: Option<&str>) -> usize {
if let Some(raw) = env {
let trimmed = raw.trim();
if trimmed.eq_ignore_ascii_case("unlimited") || trimmed == "0" {
return HASH_SEMAPHORE_UNLIMITED;
}
if let Ok(n) = trimmed.parse::<usize>() {
if n >= 1 {
return n;
}
}
tracing::warn!(
env = HASH_WORKERS_ENV,
value = raw,
"invalid {HASH_WORKERS_ENV}; falling back to default"
);
}
let parallelism = std::thread::available_parallelism()
.map(|n| n.get())
.unwrap_or(8);
parallelism.saturating_mul(16).clamp(128, 1024)
}
#[derive(Debug, Clone)]
pub struct ClockLookup {
pub hash: ContentHash,
pub clock: Clock,
}
#[derive(Debug, Clone)]
pub struct CacheSystem {
metadata: Arc<MetadataCache>,
journal: Arc<ChangeJournal>,
}
impl CacheSystem {
#[must_use]
pub fn new() -> Self {
Self {
metadata: Arc::new(MetadataCache::new()),
journal: Arc::new(ChangeJournal::new()),
}
}
#[must_use]
pub fn with_metadata(metadata: MetadataCache) -> Self {
Self {
metadata: Arc::new(metadata),
journal: Arc::new(ChangeJournal::new()),
}
}
#[must_use]
pub fn current_clock(&self) -> Clock {
self.journal.current_clock()
}
#[must_use]
pub fn metadata(&self) -> &MetadataCache {
&self.metadata
}
#[must_use]
pub fn journal(&self) -> &ChangeJournal {
&self.journal
}
pub fn lookup_since(&self, path: &NormalizedPath, since_clock: Clock) -> Result<ClockLookup> {
let clock = self.journal.current_clock();
if !self.journal.changed_since(path, since_clock) {
if let Some(hash) = self.metadata.get_cached_hash_if_stat_valid(path) {
return Ok(ClockLookup { hash, clock });
}
}
let hash = self.metadata.lookup(path.as_path())?;
Ok(ClockLookup {
hash,
clock: self.journal.current_clock(),
})
}
pub async fn lookup_since_async(
&self,
path: NormalizedPath,
since_clock: Clock,
) -> Result<ClockLookup> {
let cache = self.clone();
let _permit = hash_semaphore()
.clone()
.acquire_owned()
.await
.map_err(|err| crate::core::Error::Cache {
message: format!("hash semaphore closed: {err}"),
})?;
tokio::task::spawn_blocking(move || cache.lookup_since(&path, since_clock))
.await
.map_err(|err| crate::core::Error::Cache {
message: format!("cache lookup worker failed: {err}"),
})?
}
pub fn apply_changes(&self, changed_paths: Vec<NormalizedPath>) -> Clock {
for path in &changed_paths {
self.metadata.downgrade(path);
}
self.journal.advance(changed_paths)
}
pub fn apply_changes_with_removals(
&self,
changed: Vec<NormalizedPath>,
removed: Vec<NormalizedPath>,
) -> Clock {
for path in &changed {
self.metadata.downgrade(path);
}
for path in &removed {
self.metadata.remove(path);
}
let mut all_paths = changed;
all_paths.extend(removed);
self.journal.advance(all_paths)
}
pub fn apply_overflow(&self) -> Clock {
self.metadata.downgrade_all();
self.journal.mark_overflow()
}
pub fn rescan_entries(&self) -> usize {
self.metadata.rescan_all()
}
pub async fn rescan_entries_async(&self) -> Result<usize> {
let cache = self.clone();
tokio::task::spawn_blocking(move || cache.rescan_entries())
.await
.map_err(|err| crate::core::Error::Cache {
message: format!("cache rescan worker failed: {err}"),
})
}
pub fn clear(&self) -> Clock {
self.metadata.clear();
self.journal.clear()
}
pub fn trim(&self, max_age: std::time::Duration) -> (usize, usize) {
let meta_removed = self.metadata.trim(max_age);
let journal_removed = self.cleanup_journal();
(meta_removed, journal_removed)
}
pub fn evict_oldest(&self, count: usize) -> (usize, usize) {
let meta_removed = self.metadata.evict_oldest(count);
let journal_removed = self.cleanup_journal();
(meta_removed, journal_removed)
}
fn cleanup_journal(&self) -> usize {
let live: std::collections::HashSet<NormalizedPath> =
self.metadata.paths().into_iter().collect();
self.journal.retain_paths(&live)
}
pub fn register_tracked(&self, paths: &[NormalizedPath]) {
for path in paths {
self.journal.register(path.clone());
}
}
}
impl Default for CacheSystem {
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod tests {
use super::super::metadata::Confidence;
use super::*;
use std::fs;
use std::thread;
use std::time::Duration;
use tempfile::TempDir;
fn create_file(dir: &TempDir, name: &str, content: &str) -> NormalizedPath {
let path = dir.path().join(name);
fs::write(&path, content).expect("failed to create test file");
path.into()
}
fn sleep_for_mtime() {
thread::sleep(Duration::from_millis(1100));
}
#[test]
fn new_cache_is_at_clock_zero() {
let cache = CacheSystem::new();
assert_eq!(cache.current_clock(), Clock::ZERO);
}
#[test]
fn default_creates_new_cache() {
let cache = CacheSystem::default();
assert_eq!(cache.current_clock(), Clock::ZERO);
assert!(cache.metadata().is_empty());
}
#[test]
fn lookup_since_nonexistent_file_returns_error() {
let cache = CacheSystem::new();
let result = cache.lookup_since(&NormalizedPath::from("/no/such/file.c"), Clock::ZERO);
assert!(result.is_err());
}
#[test]
fn resolve_hash_workers_unset_uses_default() {
let n = resolve_hash_workers(None);
assert!(n >= 128, "default must be at least 128 (the floor)");
assert!(n <= 1024, "default must be at most 1024 (the ceiling)");
}
#[test]
fn resolve_hash_workers_explicit_value() {
assert_eq!(resolve_hash_workers(Some("64")), 64);
assert_eq!(resolve_hash_workers(Some("256")), 256);
assert_eq!(resolve_hash_workers(Some(" 16 ")), 16);
}
#[test]
fn resolve_hash_workers_zero_is_unlimited() {
assert_eq!(resolve_hash_workers(Some("0")), HASH_SEMAPHORE_UNLIMITED);
}
#[test]
fn resolve_hash_workers_unlimited_keyword() {
assert_eq!(
resolve_hash_workers(Some("unlimited")),
HASH_SEMAPHORE_UNLIMITED
);
assert_eq!(
resolve_hash_workers(Some("UNLIMITED")),
HASH_SEMAPHORE_UNLIMITED
);
assert_eq!(
resolve_hash_workers(Some(" Unlimited ")),
HASH_SEMAPHORE_UNLIMITED
);
}
#[test]
fn resolve_hash_workers_invalid_falls_back_to_default() {
let n = resolve_hash_workers(Some("not-a-number"));
assert!((128..=1024).contains(&n));
}
#[test]
fn lookup_since_journal_no_change_but_no_cached_hash() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "nohash.h", "data");
let cache = CacheSystem::new();
let c1 = cache.apply_changes(vec![path.clone()]);
let result = cache.lookup_since(&path, c1).unwrap();
let expected = crate::hash::hash_bytes(b"data");
assert_eq!(result.hash, expected);
}
#[test]
fn clock_advances_on_apply() {
let cache = CacheSystem::new();
assert_eq!(cache.current_clock(), Clock::ZERO);
let c1 = cache.apply_changes(vec![]);
assert_eq!(c1.tick(), 1);
let c2 = cache.apply_changes(vec![]);
assert_eq!(c2.tick(), 2);
}
#[test]
fn apply_overflow_advances_clock() {
let cache = CacheSystem::new();
let c1 = cache.apply_changes(vec![]);
let overflow = cache.apply_overflow();
assert!(overflow > c1);
}
#[test]
fn lookup_since_returns_current_clock() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "clocked.h", "tick");
let cache = CacheSystem::new();
let result = cache.lookup_since(&path, Clock::ZERO).unwrap();
assert!(result.clock >= Clock::ZERO);
}
#[test]
fn lookup_since_zero_always_stats() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "a.c", "hello");
let cache = CacheSystem::new();
let result = cache.lookup_since(&path, Clock::ZERO).unwrap();
let expected = crate::hash::hash_bytes(b"hello");
assert_eq!(result.hash, expected);
}
#[test]
fn lookup_since_skips_stat_when_no_changes() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "stable.h", "#pragma once");
let cache = CacheSystem::new();
let r1 = cache.lookup_since(&path, Clock::ZERO).unwrap();
let c1 = cache.apply_changes(vec![path.clone()]);
let r2 = cache.lookup_since(&path, c1).unwrap();
assert_eq!(r1.hash, r2.hash);
}
#[test]
fn lookup_since_stats_when_changed() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "evolving.c", "v1");
let cache = CacheSystem::new();
let r1 = cache.lookup_since(&path, Clock::ZERO).unwrap();
sleep_for_mtime();
fs::write(&path, "v2").unwrap();
let c2 = cache.apply_changes(vec![path.clone()]);
let r2 = cache.lookup_since(&path, Clock::ZERO).unwrap();
assert_ne!(r1.hash, r2.hash);
assert_eq!(r2.hash, crate::hash::hash_bytes(b"v2"));
let r3 = cache.lookup_since(&path, c2).unwrap();
assert_eq!(r2.hash, r3.hash);
}
#[test]
fn apply_changes_downgrades_confidence() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "watched.h", "content");
let cache = CacheSystem::new();
cache.lookup_since(&path, Clock::ZERO).unwrap();
let entry = cache.metadata().get(&path).unwrap();
assert_eq!(entry.confidence, Confidence::High);
cache.apply_changes(vec![path.clone()]);
let entry = cache.metadata().get(&path).unwrap();
assert_eq!(entry.confidence, Confidence::Low);
}
#[test]
fn apply_overflow_downgrades_all() {
let dir = TempDir::new().unwrap();
let path_a = create_file(&dir, "a.h", "aaa");
let path_b = create_file(&dir, "b.h", "bbb");
let cache = CacheSystem::new();
cache.lookup_since(&path_a, Clock::ZERO).unwrap();
cache.lookup_since(&path_b, Clock::ZERO).unwrap();
let c_before = cache.current_clock();
cache.apply_overflow();
let a = cache.metadata().get(&path_a).unwrap();
let b = cache.metadata().get(&path_b).unwrap();
assert_eq!(a.confidence, Confidence::Low);
assert_eq!(b.confidence, Confidence::Low);
assert!(cache.journal().changed_since(&path_a, c_before));
}
#[test]
fn register_tracked_enables_fast_path() {
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "header.h", "content");
let cache = CacheSystem::new();
let r1 = cache.lookup_since(&path, Clock::ZERO).unwrap();
cache.register_tracked(std::slice::from_ref(&path));
let clock = cache.current_clock();
let r2 = cache.lookup_since(&path, clock).unwrap();
assert_eq!(r1.hash, r2.hash);
assert_eq!(r2.hash, crate::hash::hash_bytes(b"content"));
}
#[test]
fn concurrent_lookups_shared_cache() {
use std::sync::Arc;
let dir = TempDir::new().unwrap();
let mut paths = Vec::new();
for i in 0..10 {
paths.push(create_file(
&dir,
&format!("header_{i}.h"),
&format!("content {i}"),
));
}
let cache = Arc::new(CacheSystem::new());
for path in &paths {
cache.lookup_since(path, Clock::ZERO).unwrap();
}
let c1 = cache.apply_changes(paths.clone());
let mut handles = Vec::new();
for _ in 0..8 {
let cache = Arc::clone(&cache);
let paths = paths.clone();
handles.push(std::thread::spawn(move || {
for path in &paths {
let result = cache.lookup_since(path, c1).unwrap();
assert_eq!(result.hash, crate::hash::hash_file(path).unwrap(),);
}
}));
}
for h in handles {
h.join().unwrap();
}
}
#[test]
fn rescan_entries_after_overflow() {
let dir = TempDir::new().unwrap();
let path_a = create_file(&dir, "a.h", "aaa");
let path_b = create_file(&dir, "b.h", "bbb");
let cache = CacheSystem::new();
cache.lookup_since(&path_a, Clock::ZERO).unwrap();
cache.lookup_since(&path_b, Clock::ZERO).unwrap();
cache.apply_overflow();
assert_eq!(
cache.metadata().get(&path_a).unwrap().confidence,
Confidence::Low
);
let promoted = cache.rescan_entries();
assert_eq!(promoted, 2);
assert_eq!(
cache.metadata().get(&path_a).unwrap().confidence,
Confidence::High
);
assert_eq!(
cache.metadata().get(&path_b).unwrap().confidence,
Confidence::High
);
}
#[test]
fn clear_empties_metadata_and_invalidates_clocks() {
let dir = TempDir::new().unwrap();
let path_a = create_file(&dir, "a.h", "aaa");
let path_b = create_file(&dir, "b.h", "bbb");
let cache = CacheSystem::new();
cache.lookup_since(&path_a, Clock::ZERO).unwrap();
cache.lookup_since(&path_b, Clock::ZERO).unwrap();
assert_eq!(cache.metadata().len(), 2);
let overflow_clock = cache.clear();
assert!(cache.metadata().is_empty());
assert!(overflow_clock > Clock::ZERO);
}
#[test]
fn trim_cascades_to_journal() {
let cache = CacheSystem::new();
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "old.c", "content");
cache.lookup_since(&path, Clock::ZERO).unwrap();
cache.apply_changes(vec![path.clone()]);
assert_eq!(cache.journal().last_change_len(), 1);
let (meta_removed, journal_removed) = cache.trim(Duration::ZERO);
assert!(meta_removed >= 1);
assert_eq!(journal_removed, 1);
assert!(cache.metadata().is_empty());
assert_eq!(cache.journal().last_change_len(), 0);
}
#[test]
fn evict_oldest_cascades() {
let cache = CacheSystem::new();
let dir = TempDir::new().unwrap();
let path = create_file(&dir, "evict.c", "data");
cache.lookup_since(&path, Clock::ZERO).unwrap();
cache.apply_changes(vec![path.clone()]);
let (meta_removed, journal_removed) = cache.evict_oldest(10);
assert_eq!(meta_removed, 1);
assert_eq!(journal_removed, 1);
assert!(cache.metadata().is_empty());
}
#[test]
fn apply_changes_with_removals() {
let dir = TempDir::new().unwrap();
let path_mod = create_file(&dir, "mod.c", "modified");
let path_del = create_file(&dir, "del.c", "deleted");
let cache = CacheSystem::new();
cache.lookup_since(&path_mod, Clock::ZERO).unwrap();
cache.lookup_since(&path_del, Clock::ZERO).unwrap();
assert_eq!(cache.metadata().len(), 2);
fs::remove_file(&path_del).unwrap();
cache.apply_changes_with_removals(vec![path_mod], vec![path_del.clone()]);
assert!(cache.metadata().get(&path_del).is_none());
assert_eq!(cache.metadata().len(), 1);
}
}