1use std::fs;
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::time::{Duration, SystemTime};
7
8use sva_samples::{FilterTrace, Label};
9
10use sva_formula::Hash;
11
12use super::entry_bytes::{self, RawF64, SampleCodec, Unread};
13use super::evict;
14use super::{Cache, Entry, Expected, Payload, PayloadKind};
15
16pub const IO_NANOS_PER_BYTE: u64 = 3;
18
19pub const MIN_COST: Duration = Duration::from_millis(1);
21pub const DEFAULT_MAX_BYTES: u64 = 16 << 30;
23
24static WRITE: AtomicU64 = AtomicU64::new(0);
26
27pub struct DiskCache {
28 dir: PathBuf,
29 max_bytes: u64,
30 store_everything: bool,
31 held: AtomicU64,
32 evicted: AtomicU64,
33 faults: AtomicU64,
34 codec: Box<dyn SampleCodec>,
35}
36
37impl DiskCache {
38 pub fn discover() -> Option<DiskCache> {
41 let dir = match std::env::var_os("SVA_CACHE") {
42 Some(v) if !v.is_empty() => PathBuf::from(v),
43 _ => match std::env::var_os("XDG_CACHE_HOME").filter(|v| !v.is_empty()) {
44 Some(v) => PathBuf::from(v).join("sva"),
45 None => PathBuf::from(std::env::var_os("HOME")?)
46 .join(".cache")
47 .join("sva"),
48 },
49 };
50 Some(DiskCache::at(dir))
51 }
52
53 pub fn at(dir: impl Into<PathBuf>) -> DiskCache {
54 let max_bytes = std::env::var("SVA_CACHE_MAX_BYTES")
55 .ok()
56 .and_then(|v| v.parse().ok())
57 .unwrap_or(DEFAULT_MAX_BYTES);
58 DiskCache::bounded(dir, max_bytes)
59 }
60
61 pub fn bounded(dir: impl Into<PathBuf>, max_bytes: u64) -> DiskCache {
62 DiskCache {
63 dir: dir.into(),
64 max_bytes,
65 store_everything: false,
66 held: AtomicU64::new(0),
67 evicted: AtomicU64::new(0),
68 faults: AtomicU64::new(0),
69 codec: Box::new(RawF64),
70 }
71 }
72
73 pub fn coded(mut self, codec: Box<dyn SampleCodec>) -> DiskCache {
74 self.codec = codec;
75 self
76 }
77
78 pub fn storing_everything(mut self) -> DiskCache {
80 self.store_everything = true;
81 self
82 }
83
84 fn path_of(&self, key: Hash) -> PathBuf {
85 let hex = key.to_string();
86 self.dir.join(&hex[..2]).join(format!("{hex}.rbc"))
87 }
88}
89
90impl Cache for DiskCache {
91 fn dir(&self) -> Option<&Path> {
92 Some(&self.dir)
93 }
94
95 fn max_bytes(&self) -> u64 {
96 self.max_bytes
97 }
98
99 fn held_bytes(&self) -> u64 {
100 self.held.load(Ordering::Relaxed)
101 }
102
103 fn evicted_bytes(&self) -> u64 {
104 self.evicted.load(Ordering::Relaxed)
105 }
106
107 fn faults(&self) -> u64 {
108 self.faults.load(Ordering::Relaxed)
109 }
110
111 fn holds(&self, key: Hash) -> bool {
112 self.path_of(key).is_file()
113 }
114
115 fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
116 let Expected::Samples {
117 rate,
118 width,
119 samples,
120 } = expected
121 else {
122 return None;
123 };
124 let path = self.path_of(key);
125 let bytes = match fs::read(&path) {
126 Ok(bytes) => bytes,
127 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return None,
129 Err(_) => {
130 self.faults.fetch_add(1, Ordering::Relaxed);
131 return None;
132 }
133 };
134 match entry_bytes::decode(&bytes, node, (rate, width, samples), self.codec.as_ref()) {
135 Ok(entry) => {
136 touch(&path);
137 Some(entry)
138 }
139 Err(Unread::OtherCodec) => None,
140 Err(Unread::Damaged) => {
141 self.faults.fetch_add(1, Ordering::Relaxed);
142 let _ = fs::remove_file(&path);
143 None
144 }
145 }
146 }
147
148 fn worth_storing(&self, cost: Duration, bytes: usize, kind: PayloadKind) -> bool {
150 kind == PayloadKind::Samples
151 && (self.store_everything
152 || (cost >= MIN_COST
153 && cost.as_nanos() > u128::from(bytes as u64) * u128::from(IO_NANOS_PER_BYTE)))
154 }
155
156 fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
158 let Payload::Samples(buffer) = payload else {
159 unreachable!("worth_storing offers this store nothing but samples")
160 };
161 let path = self.path_of(key);
162 let Some(shard) = path.parent() else {
163 return;
164 };
165 if fs::create_dir_all(shard).is_err() {
166 self.faults.fetch_add(1, Ordering::Relaxed);
167 return;
168 }
169 let temp = shard.join(format!(
170 "{key}.{}.{}.tmp",
171 std::process::id(),
172 WRITE.fetch_add(1, Ordering::Relaxed)
173 ));
174 if fs::write(
175 &temp,
176 entry_bytes::encode(buffer, traces, label, self.codec.as_ref()),
177 )
178 .is_err()
179 {
180 self.faults.fetch_add(1, Ordering::Relaxed);
181 let _ = fs::remove_file(&temp);
182 return;
183 }
184 if fs::rename(&temp, &path).is_err() {
185 self.faults.fetch_add(1, Ordering::Relaxed);
186 let _ = fs::remove_file(&temp);
187 }
188 }
189
190 fn sweep(&self) {
191 let mut entries: Vec<(SystemTime, u64, PathBuf)> = Vec::new();
192 for shard in fs::read_dir(&self.dir).into_iter().flatten().flatten() {
193 for file in fs::read_dir(shard.path()).into_iter().flatten().flatten() {
194 let Ok(meta) = file.metadata() else { continue };
195 let when = meta.modified().unwrap_or(SystemTime::UNIX_EPOCH);
196 entries.push((when, meta.len(), file.path()));
197 }
198 }
199 let swept = evict::to_cap(entries, self.max_bytes, |path| {
200 fs::remove_file(path).is_ok()
201 });
202 self.held.store(swept.held, Ordering::Relaxed);
203 self.evicted.store(swept.evicted, Ordering::Relaxed);
204 }
205}
206
207fn touch(path: &Path) {
208 if let Ok(file) = fs::OpenOptions::new().write(true).open(path) {
209 let _ = file.set_modified(SystemTime::now());
210 }
211}