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 const ENGINE_DIR_PREFIX: &str = "render-";
28
29fn engine_dir(name: &str) -> bool {
30    name.strip_prefix(ENGINE_DIR_PREFIX)
31        .is_some_and(|hex| hex.len() == 16 && hex.bytes().all(|b| b.is_ascii_hexdigit()))
32}
33
34/// A shard every engine once shared.
35fn shared_shard(name: &str) -> bool {
36    name.len() == 2 && name.bytes().all(|b| b.is_ascii_hexdigit())
37}
38
39pub struct DiskCache {
40    dir: PathBuf,
41    max_bytes: u64,
42    store_everything: bool,
43    held: AtomicU64,
44    evicted: AtomicU64,
45    faults: AtomicU64,
46    codec: Box<dyn SampleCodec>,
47}
48
49impl DiskCache {
50    /// `$SVA_CACHE`, else the XDG cache home. Never under `target/`: `cargo clean` would throw
51    /// away hours of rendering, and an installed binary has no build directory to write into.
52    pub fn discover() -> Option<DiskCache> {
53        let root = match std::env::var_os("SVA_CACHE") {
54            Some(v) if !v.is_empty() => PathBuf::from(v),
55            _ => match std::env::var_os("XDG_CACHE_HOME").filter(|v| !v.is_empty()) {
56                Some(v) => PathBuf::from(v).join("sva"),
57                None => PathBuf::from(std::env::var_os("HOME")?)
58                    .join(".cache")
59                    .join("sva"),
60            },
61        };
62        Some(DiskCache::under(root))
63    }
64
65    /// A key hashes what a node says, not the engine rendering it, so each engine keeps its own
66    /// directory under `root`; no other one's can answer it, so they go.
67    pub fn under(root: impl Into<PathBuf>) -> DiskCache {
68        let root = root.into();
69        let current = format!("{ENGINE_DIR_PREFIX}{:016x}", crate::RENDER_FINGERPRINT);
70        let mut stuck = 0;
71        for entry in fs::read_dir(&root).into_iter().flatten().flatten() {
72            let name = entry.file_name();
73            let Some(name) = name.to_str() else { continue };
74            let is_dir = entry.file_type().is_ok_and(|t| t.is_dir());
75            if is_dir && name != current && (engine_dir(name) || shared_shard(name)) {
76                stuck += u64::from(fs::remove_dir_all(entry.path()).is_err());
77            }
78        }
79        let cache = DiskCache::at(root.join(current));
80        cache.faults.store(stuck, Ordering::Relaxed);
81        cache
82    }
83
84    pub fn at(dir: impl Into<PathBuf>) -> DiskCache {
85        let max_bytes = std::env::var("SVA_CACHE_MAX_BYTES")
86            .ok()
87            .and_then(|v| v.parse().ok())
88            .unwrap_or(DEFAULT_MAX_BYTES);
89        DiskCache::bounded(dir, max_bytes)
90    }
91
92    pub fn bounded(dir: impl Into<PathBuf>, max_bytes: u64) -> DiskCache {
93        DiskCache {
94            dir: dir.into(),
95            max_bytes,
96            store_everything: false,
97            held: AtomicU64::new(0),
98            evicted: AtomicU64::new(0),
99            faults: AtomicU64::new(0),
100            codec: Box::new(RawF64),
101        }
102    }
103
104    pub fn coded(mut self, codec: Box<dyn SampleCodec>) -> DiskCache {
105        self.codec = codec;
106        self
107    }
108
109    /// The threshold is policy: a caller measuring reuse itself wants every node stored.
110    pub fn storing_everything(mut self) -> DiskCache {
111        self.store_everything = true;
112        self
113    }
114
115    fn path_of(&self, key: Hash) -> PathBuf {
116        let hex = key.to_string();
117        self.dir.join(&hex[..2]).join(format!("{hex}.rbc"))
118    }
119}
120
121impl Cache for DiskCache {
122    fn dir(&self) -> Option<&Path> {
123        Some(&self.dir)
124    }
125
126    fn max_bytes(&self) -> u64 {
127        self.max_bytes
128    }
129
130    fn held_bytes(&self) -> u64 {
131        self.held.load(Ordering::Relaxed)
132    }
133
134    fn evicted_bytes(&self) -> u64 {
135        self.evicted.load(Ordering::Relaxed)
136    }
137
138    fn faults(&self) -> u64 {
139        self.faults.load(Ordering::Relaxed)
140    }
141
142    fn holds(&self, key: Hash) -> bool {
143        self.path_of(key).is_file()
144    }
145
146    fn load(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
147        let Expected::Samples {
148            rate,
149            width,
150            samples,
151        } = expected
152        else {
153            return None;
154        };
155        let path = self.path_of(key);
156        let bytes = match fs::read(&path) {
157            Ok(bytes) => bytes,
158            // Absent is a cold miss; anything else is a store that could not answer.
159            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return None,
160            Err(_) => {
161                self.faults.fetch_add(1, Ordering::Relaxed);
162                return None;
163            }
164        };
165        match entry_bytes::decode(&bytes, node, (rate, width, samples), self.codec.as_ref()) {
166            Ok(entry) => {
167                touch(&path);
168                Some(entry)
169            }
170            Err(Unread::OtherCodec) => None,
171            Err(Unread::Damaged) => {
172                self.faults.fetch_add(1, Ordering::Relaxed);
173                let _ = fs::remove_file(&path);
174                None
175            }
176        }
177    }
178
179    fn peek(&self, key: Hash, node: &str, expected: Expected) -> Option<Entry> {
180        let Expected::Samples {
181            rate,
182            width,
183            samples,
184        } = expected
185        else {
186            return None;
187        };
188        let bytes = fs::read(self.path_of(key)).ok()?;
189        entry_bytes::decode(&bytes, node, (rate, width, samples), self.codec.as_ref()).ok()
190    }
191
192    /// Only samples have an encoding here; neither other payload is offered a file.
193    fn worth_storing(&self, cost: Duration, bytes: usize, kind: PayloadKind) -> bool {
194        kind == PayloadKind::Samples
195            && (self.store_everything
196                || (cost >= MIN_COST
197                    && cost.as_nanos() > u128::from(bytes as u64) * u128::from(IO_NANOS_PER_BYTE)))
198    }
199
200    /// Renamed into place: a racing reader sees one whole entry or the other.
201    fn store(&self, key: Hash, payload: &Payload, traces: &[FilterTrace], label: Option<&Label>) {
202        let Payload::Samples(buffer) = payload else {
203            unreachable!("worth_storing offers this store nothing but samples")
204        };
205        let path = self.path_of(key);
206        let Some(shard) = path.parent() else {
207            return;
208        };
209        if fs::create_dir_all(shard).is_err() {
210            self.faults.fetch_add(1, Ordering::Relaxed);
211            return;
212        }
213        let temp = shard.join(format!(
214            "{key}.{}.{}.tmp",
215            std::process::id(),
216            WRITE.fetch_add(1, Ordering::Relaxed)
217        ));
218        if fs::write(
219            &temp,
220            entry_bytes::encode(buffer, traces, label, self.codec.as_ref()),
221        )
222        .is_err()
223        {
224            self.faults.fetch_add(1, Ordering::Relaxed);
225            let _ = fs::remove_file(&temp);
226            return;
227        }
228        if fs::rename(&temp, &path).is_err() {
229            self.faults.fetch_add(1, Ordering::Relaxed);
230            let _ = fs::remove_file(&temp);
231        }
232    }
233
234    fn sweep(&self) {
235        let mut entries: Vec<(SystemTime, u64, PathBuf)> = Vec::new();
236        for shard in fs::read_dir(&self.dir).into_iter().flatten().flatten() {
237            for file in fs::read_dir(shard.path()).into_iter().flatten().flatten() {
238                let Ok(meta) = file.metadata() else { continue };
239                let when = meta.modified().unwrap_or(SystemTime::UNIX_EPOCH);
240                entries.push((when, meta.len(), file.path()));
241            }
242        }
243        let swept = evict::to_cap(entries, self.max_bytes, |path| {
244            fs::remove_file(path).is_ok()
245        });
246        self.held.store(swept.held, Ordering::Relaxed);
247        self.evicted.store(swept.evicted, Ordering::Relaxed);
248    }
249}
250
251fn touch(path: &Path) {
252    if let Ok(file) = fs::OpenOptions::new().write(true).open(path) {
253        let _ = file.set_modified(SystemTime::now());
254    }
255}