use std::fs;
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::time::{Duration, SystemTime};
use sva_samples::{FilterTrace, Label};
use sva_formula::Hash;
use super::entry_bytes::{self, RawF64, SampleCodec, Unread};
use super::evict;
use super::{Cache, Entry, Expected, Payload, PayloadKind};
pub const IO_NANOS_PER_BYTE: u64 = 3;
pub const MIN_COST: Duration = Duration::from_millis(1);
pub const DEFAULT_MAX_BYTES: u64 = 16 << 30;
static WRITE: AtomicU64 = AtomicU64::new(0);
pub struct DiskCache {
dir: PathBuf,
max_bytes: u64,
store_everything: bool,
held: AtomicU64,
evicted: AtomicU64,
faults: AtomicU64,
codec: Box<dyn SampleCodec>,
}
impl DiskCache {
pub fn discover() -> Option<DiskCache> {
let dir = match std::env::var_os("SVA_CACHE") {
Some(v) if !v.is_empty() => PathBuf::from(v),
_ => match std::env::var_os("XDG_CACHE_HOME").filter(|v| !v.is_empty()) {
Some(v) => PathBuf::from(v).join("sva"),
None => PathBuf::from(std::env::var_os("HOME")?)
.join(".cache")
.join("sva"),
},
};
Some(DiskCache::at(dir))
}
pub fn at(dir: impl Into<PathBuf>) -> DiskCache {
let max_bytes = std::env::var("SVA_CACHE_MAX_BYTES")
.ok()
.and_then(|v| v.parse().ok())
.unwrap_or(DEFAULT_MAX_BYTES);
DiskCache::bounded(dir, max_bytes)
}
pub fn bounded(dir: impl Into<PathBuf>, max_bytes: u64) -> DiskCache {
DiskCache {
dir: dir.into(),
max_bytes,
store_everything: false,
held: AtomicU64::new(0),
evicted: AtomicU64::new(0),
faults: AtomicU64::new(0),
codec: Box::new(RawF64),
}
}
pub fn coded(mut self, codec: Box<dyn SampleCodec>) -> DiskCache {
self.codec = codec;
self
}
pub fn storing_everything(mut self) -> DiskCache {
self.store_everything = true;
self
}
fn path_of(&self, key: Hash) -> PathBuf {
let hex = key.to_string();
self.dir.join(&hex[..2]).join(format!("{hex}.rbc"))
}
}
impl Cache for DiskCache {
fn dir(&self) -> Option<&Path> {
Some(&self.dir)
}
fn max_bytes(&self) -> u64 {
self.max_bytes
}
fn held_bytes(&self) -> u64 {
self.held.load(Ordering::Relaxed)
}
fn evicted_bytes(&self) -> u64 {
self.evicted.load(Ordering::Relaxed)
}
fn faults(&self) -> u64 {
self.faults.load(Ordering::Relaxed)
}
fn holds(&self, key: Hash) -> bool {
self.path_of(key).is_file()
}
fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
let Expected::Samples {
rate,
width,
samples,
} = expected
else {
return None;
};
let path = self.path_of(key);
let bytes = match fs::read(&path) {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return None,
Err(_) => {
self.faults.fetch_add(1, Ordering::Relaxed);
return None;
}
};
match entry_bytes::decode(&bytes, node, (rate, width, samples), self.codec.as_ref()) {
Ok(entry) => {
touch(&path);
Some(entry)
}
Err(Unread::OtherCodec) => None,
Err(Unread::Damaged) => {
self.faults.fetch_add(1, Ordering::Relaxed);
let _ = fs::remove_file(&path);
None
}
}
}
fn worth_storing(&self, cost: Duration, bytes: usize, kind: PayloadKind) -> bool {
kind == PayloadKind::Samples
&& (self.store_everything
|| (cost >= MIN_COST
&& cost.as_nanos() > u128::from(bytes as u64) * u128::from(IO_NANOS_PER_BYTE)))
}
fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
let Payload::Samples(buffer) = payload else {
unreachable!("worth_storing offers this store nothing but samples")
};
let path = self.path_of(key);
let Some(shard) = path.parent() else {
return;
};
if fs::create_dir_all(shard).is_err() {
self.faults.fetch_add(1, Ordering::Relaxed);
return;
}
let temp = shard.join(format!(
"{key}.{}.{}.tmp",
std::process::id(),
WRITE.fetch_add(1, Ordering::Relaxed)
));
if fs::write(
&temp,
entry_bytes::encode(buffer, traces, label, self.codec.as_ref()),
)
.is_err()
{
self.faults.fetch_add(1, Ordering::Relaxed);
let _ = fs::remove_file(&temp);
return;
}
if fs::rename(&temp, &path).is_err() {
self.faults.fetch_add(1, Ordering::Relaxed);
let _ = fs::remove_file(&temp);
}
}
fn sweep(&self) {
let mut entries: Vec<(SystemTime, u64, PathBuf)> = Vec::new();
for shard in fs::read_dir(&self.dir).into_iter().flatten().flatten() {
for file in fs::read_dir(shard.path()).into_iter().flatten().flatten() {
let Ok(meta) = file.metadata() else { continue };
let when = meta.modified().unwrap_or(SystemTime::UNIX_EPOCH);
entries.push((when, meta.len(), file.path()));
}
}
let swept = evict::to_cap(entries, self.max_bytes, |path| {
fs::remove_file(path).is_ok()
});
self.held.store(swept.held, Ordering::Relaxed);
self.evicted.store(swept.evicted, Ordering::Relaxed);
}
}
fn touch(path: &Path) {
if let Ok(file) = fs::OpenOptions::new().write(true).open(path) {
let _ = file.set_modified(SystemTime::now());
}
}