Skip to main content

sva_engine/cache/
disk.rs

1// Concern: keeps rendered buffers on disk under their content hash, across processes | Non-concern: what a store is for (cache.rs), computing a hash | IO: (Hash) -> a buffer + traces
2
3use 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
17/// Recompute wins a tie; `examples/io_price.rs` measures a round trip against it.
18pub const IO_NANOS_PER_BYTE: u64 = 3;
19
20/// Below this, nothing is worth a file: an inode is not free.
21pub const MIN_COST: Duration = Duration::from_millis(1);
22/// A 160 s master's node set is about 6 GB; less than that is no cache at all.
23pub const DEFAULT_MAX_BYTES: u64 = 16 << 30;
24
25/// A tag after the magic makes an older file a miss rather than a misread.
26const MAGIC_TIME: &[u8; 4] = b"RBC6";
27const TAG_SAMPLES: u8 = 0;
28
29/// Two threads can reach one key at once, so a per-process temp name is not enough.
30static 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    /// `$SVA_CACHE`, else the XDG cache home. Never under `target/`: `cargo clean` would throw
43    /// away hours of rendering, and an installed binary has no build directory to write into.
44    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    /// The threshold is policy: a caller measuring reuse itself wants every node stored.
77    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    /// A render that reuses nothing and evicted a lot has outgrown its budget, and nothing
98    /// else in the report says so.
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            // Absent is a cold miss; anything else is a store that could not answer.
128            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    /// Only samples have an encoding here; neither other payload is offered a file.
148    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    /// Renamed into place: a racing reader sees one whole entry or the other.
156    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    /// The read time is the mtime `load` refreshes.
185    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
325/// Every field is checked against what the caller asked for, not merely parsed: the hash
326/// already rules out a mismatch, so one here means the hash's domain is wrong and the only
327/// safe answer is to re-render.
328fn 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}