use crate::config::{Admission, CacheConfig, CacheStats, PutOptions};
use crate::error::CacheError;
use crate::index::{self, Manifest, ManifestEntry};
use crate::layout;
use crate::policy::{EvictionContext, EvictionEntry};
use dig_store::{get_capsule_identity, CapsuleIdentity};
use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
#[derive(Clone)]
pub struct Cache {
inner: Arc<Mutex<Inner>>,
}
impl std::fmt::Debug for Cache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let inner = self.lock();
f.debug_struct("Cache")
.field("root", &inner.root)
.field("count", &inner.entries.len())
.field("bytes_used", &inner.bytes_used)
.field("capacity", &inner.config.max_bytes)
.finish()
}
}
#[derive(Debug, Clone)]
pub struct CachedCapsule {
id: CapsuleIdentity,
path: PathBuf,
}
impl CachedCapsule {
pub fn id(&self) -> CapsuleIdentity {
self.id
}
pub fn path(&self) -> &Path {
&self.path
}
}
struct Entry {
key_hex: String,
size: u64,
seq: u64,
pinned: bool,
}
struct Inner {
root: PathBuf,
config: CacheConfig,
entries: HashMap<CapsuleIdentity, Entry>,
bytes_used: u64,
next_seq: u64,
}
impl Cache {
pub fn open(root: &Path, config: CacheConfig) -> Result<Cache, CacheError> {
layout::ensure_dirs(root)?;
layout::clean_tmp(root)?;
let scanned = layout::scan_capsules(root)?;
let manifest = index::load(root);
let inner = Inner::rebuild(root.to_path_buf(), config, scanned, manifest)?;
inner.save_manifest()?;
Ok(Cache {
inner: Arc::new(Mutex::new(inner)),
})
}
pub fn put_file(
&self,
id: CapsuleIdentity,
src: &Path,
opts: PutOptions,
) -> Result<Admission, CacheError> {
let size = std::fs::metadata(src)
.map_err(|e| CacheError::io(src, e))?
.len();
let mut inner = self.lock();
inner.reject_if_too_large(id, size)?;
let (staged, _) = layout::stage_file(&inner.root, src)?;
inner.admit_staged(id, staged, size, opts.pinned)
}
pub fn put_bytes(
&self,
id: CapsuleIdentity,
bytes: &[u8],
opts: PutOptions,
) -> Result<Admission, CacheError> {
if opts.check_identity {
verify_claimed_identity(id, bytes)?;
}
let size = bytes.len() as u64;
let mut inner = self.lock();
inner.reject_if_too_large(id, size)?;
let staged = layout::stage_bytes(&inner.root, bytes)?;
inner.admit_staged(id, staged, size, opts.pinned)
}
pub fn get(&self, id: &CapsuleIdentity) -> Option<CachedCapsule> {
let mut inner = self.lock();
let seq = inner.next_seq;
let key_hex = {
let entry = inner.entries.get_mut(id)?;
entry.seq = seq; entry.key_hex.clone()
};
inner.next_seq += 1;
let path = layout::capsule_path(&inner.root, &key_hex);
Some(CachedCapsule { id: *id, path })
}
pub fn get_bytes(&self, id: &CapsuleIdentity) -> Result<Option<Vec<u8>>, CacheError> {
let Some(cached) = self.get(id) else {
return Ok(None);
};
let bytes = std::fs::read(cached.path()).map_err(|e| CacheError::io(cached.path(), e))?;
Ok(Some(bytes))
}
pub fn contains(&self, id: &CapsuleIdentity) -> bool {
self.lock().entries.contains_key(id)
}
pub fn holdings(&self) -> Vec<CapsuleIdentity> {
self.lock().entries.keys().copied().collect()
}
pub fn remove(&self, id: &CapsuleIdentity) -> Result<bool, CacheError> {
let mut inner = self.lock();
if !inner.entries.contains_key(id) {
return Ok(false);
}
inner.drop_entry(id)?;
inner.save_manifest()?;
Ok(true)
}
pub fn pin(&self, id: &CapsuleIdentity) -> Result<bool, CacheError> {
self.set_pinned(id, true)
}
pub fn unpin(&self, id: &CapsuleIdentity) -> Result<bool, CacheError> {
self.set_pinned(id, false)
}
pub fn stats(&self) -> CacheStats {
let inner = self.lock();
CacheStats {
bytes_used: inner.bytes_used,
count: inner.entries.len(),
capacity: inner.config.max_bytes,
}
}
pub fn set_config(&self, config: CacheConfig) -> Result<Admission, CacheError> {
let mut inner = self.lock();
inner.config = config;
let evicted = inner.evict_to_fit();
inner.save_manifest()?;
Ok(Admission { evicted })
}
fn set_pinned(&self, id: &CapsuleIdentity, pinned: bool) -> Result<bool, CacheError> {
let mut inner = self.lock();
match inner.entries.get_mut(id) {
Some(entry) => {
entry.pinned = pinned;
inner.save_manifest()?;
Ok(true)
}
None => Ok(false),
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, Inner> {
self.inner
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner())
}
}
impl Inner {
fn rebuild(
root: PathBuf,
config: CacheConfig,
scanned: Vec<layout::ScannedFile>,
manifest: Option<Manifest>,
) -> Result<Inner, CacheError> {
let by_key: HashMap<String, ManifestEntry> = manifest
.map(|m| {
m.entries
.into_iter()
.map(|e| (e.retrieval_key.clone(), e))
.collect()
})
.unwrap_or_default();
struct Pending {
id: CapsuleIdentity,
key_hex: String,
size: u64,
pinned: bool,
sort_key: u128,
}
let mut pending = Vec::new();
for file in scanned {
let resolved = by_key.get(&file.key_hex).and_then(|entry| {
entry
.identity()
.map(|id| (id, entry.seq as u128, entry.pinned))
});
let (id, sort_key, pinned) = match resolved {
Some((id, seq, pinned)) => (id, seq, pinned),
None => {
match layout::recover_identity(&file.path) {
Ok(id) if layout::retrieval_key_hex(&id) == file.key_hex => {
(id, file.mtime_nanos, false)
}
_ => continue,
}
}
};
pending.push(Pending {
id,
key_hex: file.key_hex,
size: file.size,
pinned,
sort_key,
});
}
pending.sort_by_key(|p| p.sort_key);
let mut entries = HashMap::new();
let mut bytes_used = 0u64;
for (seq, p) in pending.into_iter().enumerate() {
bytes_used = bytes_used.saturating_add(p.size);
entries.insert(
p.id,
Entry {
key_hex: p.key_hex,
size: p.size,
seq: seq as u64,
pinned: p.pinned,
},
);
}
let next_seq = entries.len() as u64;
Ok(Inner {
root,
config,
entries,
bytes_used,
next_seq,
})
}
fn reject_if_too_large(&self, id: CapsuleIdentity, size: u64) -> Result<(), CacheError> {
if size > self.config.max_bytes {
return Err(CacheError::EntryTooLarge {
id: Box::new(id),
size,
capacity: self.config.max_bytes,
});
}
Ok(())
}
fn admit_staged(
&mut self,
id: CapsuleIdentity,
staged: PathBuf,
size: u64,
pinned: bool,
) -> Result<Admission, CacheError> {
let key_hex = layout::retrieval_key_hex(&id);
let evicted = self.select_evictions(&id, size);
if let Err(e) = self.apply_evictions(&evicted) {
layout::discard_staged(&staged);
return Err(e);
}
if let Err(e) = layout::finalize(&self.root, &key_hex, &staged) {
layout::discard_staged(&staged);
return Err(e);
}
if let Some(previous) = self.entries.get(&id) {
self.bytes_used = self.bytes_used.saturating_sub(previous.size);
}
self.bytes_used = self.bytes_used.saturating_add(size);
let seq = self.next_seq;
self.next_seq += 1;
self.entries.insert(
id,
Entry {
key_hex,
size,
seq,
pinned,
},
);
self.save_manifest()?;
Ok(Admission { evicted })
}
fn select_evictions(
&self,
incoming_id: &CapsuleIdentity,
incoming: u64,
) -> Vec<CapsuleIdentity> {
let existing = self.entries.get(incoming_id).map(|e| e.size).unwrap_or(0);
let current_bytes = self.bytes_used.saturating_sub(existing);
let entries: Vec<EvictionEntry> = self
.entries
.iter()
.filter(|(id, _)| *id != incoming_id)
.map(|(id, e)| EvictionEntry {
id: *id,
size: e.size,
last_access: e.seq,
pinned: e.pinned,
})
.collect();
let ctx = EvictionContext {
entries: &entries,
current_bytes,
capacity: self.config.max_bytes,
incoming_size: incoming,
};
self.config.policy.select_evictions(&ctx)
}
fn evict_to_fit(&mut self) -> Vec<CapsuleIdentity> {
let entries: Vec<EvictionEntry> = self
.entries
.iter()
.map(|(id, e)| EvictionEntry {
id: *id,
size: e.size,
last_access: e.seq,
pinned: e.pinned,
})
.collect();
let ctx = EvictionContext {
entries: &entries,
current_bytes: self.bytes_used,
capacity: self.config.max_bytes,
incoming_size: 0,
};
let evicted = self.config.policy.select_evictions(&ctx);
for id in &evicted {
let _ = self.drop_entry(id);
}
evicted
}
fn apply_evictions(&mut self, evicted: &[CapsuleIdentity]) -> Result<(), CacheError> {
for id in evicted {
self.drop_entry(id)?;
}
Ok(())
}
fn drop_entry(&mut self, id: &CapsuleIdentity) -> Result<(), CacheError> {
if let Some(entry) = self.entries.remove(id) {
layout::remove_capsule(&self.root, &entry.key_hex)?;
self.bytes_used = self.bytes_used.saturating_sub(entry.size);
}
Ok(())
}
fn save_manifest(&self) -> Result<(), CacheError> {
let entries = self
.entries
.iter()
.map(|(id, e)| ManifestEntry {
store_id: index::bytes32_hex(&id.store_id),
root_hash: index::bytes32_hex(&id.root_hash),
retrieval_key: e.key_hex.clone(),
size: e.size,
seq: e.seq,
pinned: e.pinned,
})
.collect();
let manifest = Manifest {
version: index::MANIFEST_VERSION,
next_seq: self.next_seq,
entries,
};
index::save(&self.root, &manifest)
}
}
fn verify_claimed_identity(claimed: CapsuleIdentity, bytes: &[u8]) -> Result<(), CacheError> {
let recovered = get_capsule_identity(bytes).map_err(|e| CacheError::CorruptEntry {
path: PathBuf::from("<in-memory bytes>"),
reason: e.to_string(),
})?;
if recovered != claimed {
return Err(CacheError::IdentityMismatch {
claimed: Box::new(claimed),
recovered: Box::new(recovered),
});
}
Ok(())
}