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_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
16/// Recompute wins a tie; `examples/io_price.rs` measures a round trip against it.
17pub const IO_NANOS_PER_BYTE: u64 = 3;
18
19/// Below this, nothing is worth a file: an inode is not free.
20pub const MIN_COST: Duration = Duration::from_millis(1);
21/// A 160 s master's node set is about 6 GB; less than that is no cache at all.
22pub const DEFAULT_MAX_BYTES: u64 = 16 << 30;
23
24/// Two threads can reach one key at once, so a per-process temp name is not enough.
25static 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    /// `$SVA_CACHE`, else the XDG cache home. Never under `target/`: `cargo clean` would throw
39    /// away hours of rendering, and an installed binary has no build directory to write into.
40    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    /// The threshold is policy: a caller measuring reuse itself wants every node stored.
79    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            // 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 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    /// Only samples have an encoding here; neither other payload is offered a file.
149    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    /// Renamed into place: a racing reader sees one whole entry or the other.
157    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}