1use std::fs;
4use std::path::{Path, PathBuf};
5use std::sync::atomic::{AtomicU64, Ordering};
6use std::time::{Duration, SystemTime};
7
8use sva_formula::filter::Shape;
9use sva_samples::Buffer;
10use sva_samples::{AutomationFrame, FilterTrace, Label};
11
12use sva_formula::Hash;
13
14use super::evict;
15use super::{Cache, Entry, Expected, Payload, PayloadKind};
16
17pub const IO_NANOS_PER_BYTE: u64 = 3;
19
20pub const MIN_COST: Duration = Duration::from_millis(1);
22pub const DEFAULT_MAX_BYTES: u64 = 16 << 30;
24
25const MAGIC_TIME: &[u8; 4] = b"RBC6";
27const TAG_SAMPLES: u8 = 0;
28
29static WRITE: AtomicU64 = AtomicU64::new(0);
31
32pub struct DiskCache {
33 dir: PathBuf,
34 max_bytes: u64,
35 store_everything: bool,
36 held: AtomicU64,
37 evicted: AtomicU64,
38 faults: AtomicU64,
39}
40
41impl DiskCache {
42 pub fn discover() -> Option<DiskCache> {
45 let dir = match std::env::var_os("SVA_CACHE") {
46 Some(v) if !v.is_empty() => PathBuf::from(v),
47 _ => match std::env::var_os("XDG_CACHE_HOME").filter(|v| !v.is_empty()) {
48 Some(v) => PathBuf::from(v).join("sva"),
49 None => PathBuf::from(std::env::var_os("HOME")?)
50 .join(".cache")
51 .join("sva"),
52 },
53 };
54 Some(DiskCache::at(dir))
55 }
56
57 pub fn at(dir: impl Into<PathBuf>) -> DiskCache {
58 let max_bytes = std::env::var("SVA_CACHE_MAX_BYTES")
59 .ok()
60 .and_then(|v| v.parse().ok())
61 .unwrap_or(DEFAULT_MAX_BYTES);
62 DiskCache::bounded(dir, max_bytes)
63 }
64
65 pub fn bounded(dir: impl Into<PathBuf>, max_bytes: u64) -> DiskCache {
66 DiskCache {
67 dir: dir.into(),
68 max_bytes,
69 store_everything: false,
70 held: AtomicU64::new(0),
71 evicted: AtomicU64::new(0),
72 faults: AtomicU64::new(0),
73 }
74 }
75
76 pub fn storing_everything(mut self) -> DiskCache {
78 self.store_everything = true;
79 self
80 }
81
82 fn path_of(&self, key: Hash) -> PathBuf {
83 let hex = key.to_string();
84 self.dir.join(&hex[..2]).join(format!("{hex}.rbc"))
85 }
86}
87
88impl Cache for DiskCache {
89 fn dir(&self) -> Option<&Path> {
90 Some(&self.dir)
91 }
92
93 fn max_bytes(&self) -> u64 {
94 self.max_bytes
95 }
96
97 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 decode(&bytes, node, rate, width, samples) {
135 Some(entry) => {
136 touch(&path);
137 Some(entry)
138 }
139 None => {
140 self.faults.fetch_add(1, Ordering::Relaxed);
141 let _ = fs::remove_file(&path);
142 None
143 }
144 }
145 }
146
147 fn worth_storing(&self, cost: Duration, bytes: usize, kind: PayloadKind) -> bool {
149 kind == PayloadKind::Samples
150 && (self.store_everything
151 || (cost >= MIN_COST
152 && cost.as_nanos() > u128::from(bytes as u64) * u128::from(IO_NANOS_PER_BYTE)))
153 }
154
155 fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
157 let Payload::Samples(buffer) = payload else {
158 unreachable!("worth_storing offers this store nothing but samples")
159 };
160 let path = self.path_of(key);
161 let Some(shard) = path.parent() else {
162 return;
163 };
164 if fs::create_dir_all(shard).is_err() {
165 self.faults.fetch_add(1, Ordering::Relaxed);
166 return;
167 }
168 let temp = shard.join(format!(
169 "{key}.{}.{}.tmp",
170 std::process::id(),
171 WRITE.fetch_add(1, Ordering::Relaxed)
172 ));
173 if fs::write(&temp, encode(buffer, traces, label)).is_err() {
174 self.faults.fetch_add(1, Ordering::Relaxed);
175 let _ = fs::remove_file(&temp);
176 return;
177 }
178 if fs::rename(&temp, &path).is_err() {
179 self.faults.fetch_add(1, Ordering::Relaxed);
180 let _ = fs::remove_file(&temp);
181 }
182 }
183
184 fn sweep(&self) {
186 let mut entries: Vec<(SystemTime, u64, PathBuf)> = Vec::new();
187 for shard in fs::read_dir(&self.dir).into_iter().flatten().flatten() {
188 for file in fs::read_dir(shard.path()).into_iter().flatten().flatten() {
189 let Ok(meta) = file.metadata() else { continue };
190 let when = meta.modified().unwrap_or(SystemTime::UNIX_EPOCH);
191 entries.push((when, meta.len(), file.path()));
192 }
193 }
194 let swept = evict::to_cap(entries, self.max_bytes, |path| {
195 fs::remove_file(path).is_ok()
196 });
197 self.held.store(swept.held, Ordering::Relaxed);
198 self.evicted.store(swept.evicted, Ordering::Relaxed);
199 }
200}
201
202fn touch(path: &Path) {
203 if let Ok(file) = fs::OpenOptions::new().write(true).open(path) {
204 let _ = file.set_modified(SystemTime::now());
205 }
206}
207
208pub(super) struct Writer(pub Vec<u8>);
209
210impl Writer {
211 pub(super) fn u32(&mut self, v: u32) {
212 self.0.extend_from_slice(&v.to_le_bytes());
213 }
214 fn u64(&mut self, v: u64) {
215 self.0.extend_from_slice(&v.to_le_bytes());
216 }
217 pub(super) fn f64(&mut self, v: f64) {
218 self.0.extend_from_slice(&v.to_le_bytes());
219 }
220}
221
222fn encode(buffer: &Buffer, traces: &[FilterTrace], label: Option<&Label>) -> Vec<u8> {
223 encode_time(buffer, traces, label)
224}
225
226fn write_traces(w: &mut Writer, traces: &[FilterTrace]) {
227 w.u32(traces.len() as u32);
228 for trace in traces {
229 w.u32(trace.site as u32);
230 w.u32(trace.channel.map_or(u32::MAX, |c| c as u32));
231 w.0.push(u8::from(trace.clamped));
232 w.u32(trace.shape.len() as u32);
233 w.0.extend_from_slice(trace.shape.as_bytes());
234 w.f64(trace.trace_secs);
235 w.u32(trace.frames.len() as u32);
236 for f in &trace.frames {
237 w.f64(f.t_secs);
238 w.f64(f.cutoff);
239 w.f64(f.q);
240 w.f64(f.gain_db);
241 }
242 }
243}
244
245fn encode_time(buffer: &Buffer, traces: &[FilterTrace], label: Option<&Label>) -> Vec<u8> {
246 let mut w = Writer(Vec::with_capacity(buffer.len() * buffer.width * 8 + 64));
247 w.0.extend_from_slice(MAGIC_TIME);
248 w.0.push(TAG_SAMPLES);
249 w.u32(buffer.rate);
250 w.u64(buffer.len() as u64);
251 w.u32(buffer.width as u32);
252 write_traces(&mut w, traces);
253 super::label::write(&mut w, label);
254 for c in 0..buffer.width {
255 for &s in buffer.plane(c) {
256 w.0.extend_from_slice(&s.to_le_bytes());
257 }
258 }
259 w.0
260}
261
262pub(super) struct Reader<'a>(&'a [u8]);
263
264impl<'a> Reader<'a> {
265 pub(super) fn take(&mut self, n: usize) -> Option<&'a [u8]> {
266 let (head, rest) = self.0.split_at_checked(n)?;
267 self.0 = rest;
268 Some(head)
269 }
270 pub(super) fn u32(&mut self) -> Option<u32> {
271 Some(u32::from_le_bytes(self.take(4)?.try_into().ok()?))
272 }
273 fn u64(&mut self) -> Option<u64> {
274 Some(u64::from_le_bytes(self.take(8)?.try_into().ok()?))
275 }
276 pub(super) fn f64(&mut self) -> Option<f64> {
277 Some(f64::from_le_bytes(self.take(8)?.try_into().ok()?))
278 }
279}
280
281fn read_f64s(r: &mut Reader, n: usize) -> Option<Vec<f64>> {
282 Some(
283 r.take(n * 8)?
284 .chunks_exact(8)
285 .map(|c| f64::from_le_bytes(c.try_into().expect("chunks_exact(8)")))
286 .collect(),
287 )
288}
289
290fn read_traces(r: &mut Reader, node: &str) -> Option<Vec<FilterTrace>> {
291 let mut traces = Vec::new();
292 for _ in 0..r.u32()? {
293 let site = r.u32()? as usize;
294 let channel = match r.u32()? {
295 u32::MAX => None,
296 c => Some(c as usize),
297 };
298 let clamped = r.take(1)?[0] == 1;
299 let name_len = r.u32()? as usize;
300 let shape = Shape::from_name(std::str::from_utf8(r.take(name_len)?).ok()?)?.name();
301 let trace_secs = r.f64()?;
302 let frame_count = r.u32()? as usize;
303 let mut frames = Vec::with_capacity(frame_count.min(1 << 20));
304 for _ in 0..frame_count {
305 frames.push(AutomationFrame {
306 t_secs: r.f64()?,
307 cutoff: r.f64()?,
308 q: r.f64()?,
309 gain_db: r.f64()?,
310 });
311 }
312 traces.push(FilterTrace {
313 node: node.to_string(),
314 site,
315 channel,
316 shape,
317 clamped,
318 trace_secs,
319 frames,
320 });
321 }
322 Some(traces)
323}
324
325fn decode(bytes: &[u8], node: &str, rate: u32, width: usize, samples: usize) -> Option<Entry> {
329 let mut r = Reader(bytes);
330 if r.take(4)? != MAGIC_TIME || r.take(1)?[0] != TAG_SAMPLES {
331 return None;
332 }
333 decode_time(&mut r, node, rate, width, samples)
334}
335
336fn decode_time(
337 r: &mut Reader,
338 node: &str,
339 sample_rate: u32,
340 width: usize,
341 samples: usize,
342) -> Option<Entry> {
343 if r.u32()? != sample_rate || r.u64()? != samples as u64 || r.u32()? != width as u32 {
344 return None;
345 }
346 let traces = read_traces(r, node)?;
347 let label = super::label::read(r)?;
348 let planes: Vec<Vec<f64>> = (0..width)
349 .map(|_| read_f64s(r, samples))
350 .collect::<Option<_>>()?;
351 if !r.0.is_empty() {
352 return None;
353 }
354 Some(Entry {
355 payload: Payload::Samples(Box::new(Buffer::of_planes(sample_rate, planes))),
356 traces,
357 label,
358 })
359}