Skip to main content

datui_lib/cache/
store.rs

1//! One store for every cache kind: a directory per kind, a file per entry.
2//!
3//! `<cache>/<kind>/<stable-hash(key)>.<ext>`, each file framed as
4//!
5//! | bytes | what |
6//! |---|---|
7//! | 4 | `dtuc` |
8//! | 2 | the frame's version |
9//! | 2 | the kind's payload version |
10//! | 4 + n | the key |
11//! | 4 + n | the fingerprint |
12//! | 8 + n | the payload |
13//! | 8 | CRC-64 of everything above |
14//!
15//! A file that does not decode, is of another version, holds another key (a hash
16//! collision) or another fingerprint is a miss. A hit dates the file, and the
17//! modification time is the LRU clock a sweep evicts by once a kind passes its budget.
18
19use super::{CacheManager, atomic_write};
20use crate::logging::LogFailure;
21use color_eyre::Result;
22use std::fs;
23use std::marker::PhantomData;
24use std::path::{Path, PathBuf};
25use std::time::{Duration, SystemTime};
26
27const MAGIC: &[u8; 4] = b"dtuc";
28const FRAME_VERSION: u16 = 1;
29
30/// A temp file older than this is from a writer that died; a live one renames its
31/// temp file within moments of creating it.
32const STALE_TEMP: Duration = Duration::from_secs(3_600);
33
34/// What the cache files of one kind hold, and how much of the disk they may take.
35pub(crate) trait Kind {
36    /// The kind's directory under the cache.
37    const DIR: &'static str;
38    const EXT: &'static str;
39    /// Bumped when the payload's encoding changes: files of another version are misses.
40    const VERSION: u16;
41    /// Bytes kept, all entries together. The entry just written is always kept.
42    const BUDGET: u64;
43    type Value;
44    fn encode(value: &Self::Value) -> Result<Vec<u8>>;
45    fn decode(payload: &[u8]) -> Option<Self::Value>;
46}
47
48/// CRC-64/XZ: fixed by its spec, unlike `DefaultHasher`, whose output may change
49/// with any Rust release and would orphan every file named by it.
50static CRC64: crc::Crc<u64> = crc::Crc::<u64>::new(&crc::CRC_64_XZ);
51
52/// A hash that is the same in every build, for file names and fingerprints.
53pub struct StableHasher(crc::Digest<'static, u64>);
54
55impl Default for StableHasher {
56    fn default() -> Self {
57        Self(CRC64.digest())
58    }
59}
60
61impl StableHasher {
62    /// Length-prefixed, so `("ab", "c")` and `("a", "bc")` differ.
63    pub fn bytes(&mut self, bytes: &[u8]) -> &mut Self {
64        self.u64(bytes.len() as u64);
65        self.0.update(bytes);
66        self
67    }
68
69    pub fn u64(&mut self, n: u64) -> &mut Self {
70        self.0.update(&n.to_le_bytes());
71        self
72    }
73
74    pub fn finish(self) -> u64 {
75        self.0.finalize()
76    }
77}
78
79/// CRC-64/XZ of one byte string.
80pub fn stable_hash(bytes: &[u8]) -> u64 {
81    CRC64.checksum(bytes)
82}
83
84/// The entries of one kind under a cache directory.
85pub(crate) struct Store<K: Kind> {
86    cache: CacheManager,
87    budget: u64,
88    kind: PhantomData<K>,
89}
90
91impl<K: Kind> Store<K> {
92    pub(crate) fn new(cache: &CacheManager) -> Self {
93        Self {
94            cache: cache.clone(),
95            budget: K::BUDGET,
96            kind: PhantomData,
97        }
98    }
99
100    #[cfg(test)]
101    pub(crate) fn with_budget(mut self, budget: u64) -> Self {
102        self.budget = budget;
103        self
104    }
105
106    pub(crate) fn dir(&self) -> PathBuf {
107        self.cache.cache_file(K::DIR)
108    }
109
110    pub(crate) fn file(&self, key: &str) -> PathBuf {
111        self.dir()
112            .join(format!("{:016x}.{}", stable_hash(key.as_bytes()), K::EXT))
113    }
114
115    /// The value stored under `key` at `fingerprint`, dating it as used.
116    pub(crate) fn get(&self, key: &str, fingerprint: &str) -> Option<K::Value> {
117        let file = self.file(key);
118        let bytes = fs::read(&file).ok()?;
119        let value = self.read(&file, &bytes, key, fingerprint)?;
120        touch(&file);
121        Some(value)
122    }
123
124    /// Date `key`'s entry as used, if there is one.
125    pub(crate) fn touch(&self, key: &str) {
126        let file = self.file(key);
127        if file.exists() {
128            touch(&file);
129        }
130    }
131
132    /// Every entry that reads, by key, without dating any.
133    pub(crate) fn scan(&self) -> Vec<(String, K::Value)> {
134        let Ok(entries) = fs::read_dir(self.dir()) else {
135            return Vec::new();
136        };
137        entries
138            .flatten()
139            .map(|e| e.path())
140            .filter(|p| p.extension().is_some_and(|x| x == K::EXT))
141            .filter_map(|file| {
142                let bytes = fs::read(&file).ok()?;
143                let frame = unframe(&bytes, K::VERSION).or_else(|| {
144                    damaged::<K>(&file);
145                    None
146                })?;
147                let value = K::decode(frame.payload).or_else(|| {
148                    damaged::<K>(&file);
149                    None
150                })?;
151                Some((frame.key.to_string(), value))
152            })
153            .collect()
154    }
155
156    /// Entries on disk.
157    #[cfg(test)]
158    pub(crate) fn len(&self) -> usize {
159        fs::read_dir(self.dir()).map_or(0, |entries| {
160            entries
161                .flatten()
162                .filter(|e| e.path().extension().is_some_and(|x| x == K::EXT))
163                .count()
164        })
165    }
166
167    pub(crate) fn put(&self, key: &str, fingerprint: &str, value: &K::Value) {
168        self.put_all([(key, fingerprint, value)]);
169    }
170
171    /// Store each entry, then sweep once.
172    pub(crate) fn put_all<'a>(
173        &self,
174        entries: impl IntoIterator<Item = (&'a str, &'a str, &'a K::Value)>,
175    ) where
176        K::Value: 'a,
177    {
178        let mut written = Vec::new();
179        let mut bytes = 0u64;
180        for (key, fingerprint, value) in entries {
181            let file = self.file(key);
182            let stored = (|| -> Result<()> {
183                let frame = frame(K::VERSION, key, fingerprint, &K::encode(value)?)?;
184                // An entry that has not changed is only dated: a measured directory is
185                // recorded on every look, and most looks find what the last one did.
186                if fs::metadata(&file).is_ok_and(|m| m.len() == frame.len() as u64)
187                    && fs::read(&file).is_ok_and(|old| old == frame)
188                {
189                    touch(&file);
190                    return Ok(());
191                }
192                fs::create_dir_all(self.dir())?;
193                atomic_write(&file, &frame)?;
194                bytes += frame.len() as u64;
195                Ok(())
196            })();
197            match stored {
198                Ok(()) => written.push(file),
199                Err(e) => log::warn!(target: "datui", "save a {} entry: {e:#}", K::DIR),
200            }
201        }
202        if !written.is_empty() && self.may_pass_budget(bytes) {
203            self.sweep(&written);
204        }
205    }
206
207    /// Whether the kind may now be past its budget: unknown until this session's first
208    /// sweep, then what that sweep found plus everything written since. Sweeping reads
209    /// the whole directory, and measuring a directory of thousands writes a batch at a
210    /// time; a sweep per batch grew with the cache.
211    fn may_pass_budget(&self, written: u64) -> bool {
212        let mut swept = self.cache.swept.lock().unwrap_or_else(|e| e.into_inner());
213        match swept.get_mut(K::DIR) {
214            Some(total) => {
215                *total = total.saturating_add(written);
216                *total > self.budget
217            }
218            None => true,
219        }
220    }
221
222    /// Drop stale temp files and, past the budget, the least recently used entries down
223    /// to three quarters of it, never one of `keep`. Under the kind's lock only so two sweeps do not race;
224    /// writes land by rename and need none.
225    fn sweep(&self, keep: &[PathBuf]) {
226        self.cache
227            .with_cache_lock(K::DIR, || {
228                retire_legacy(&self.cache);
229                let Ok(entries) = fs::read_dir(self.dir()) else {
230                    return Ok(());
231                };
232                let now = SystemTime::now();
233                let mut kept = Vec::new();
234                for entry in entries.flatten() {
235                    let path = entry.path();
236                    let Ok(meta) = entry.metadata() else { continue };
237                    let Ok(modified) = meta.modified() else {
238                        continue;
239                    };
240                    if path.extension().is_some_and(|x| x == "tmp") {
241                        if now
242                            .duration_since(modified)
243                            .is_ok_and(|age| age > STALE_TEMP)
244                        {
245                            fs::remove_file(&path).or_log("remove a stale cache temp file");
246                        }
247                    } else if path.extension().is_some_and(|x| x == K::EXT) {
248                        kept.push((modified, meta.len(), path));
249                    }
250                }
251                let mut total: u64 = kept.iter().map(|(_, len, _)| len).sum();
252                let found = self.cache.swept.clone();
253                let note = |total: u64| {
254                    let mut swept = found.lock().unwrap_or_else(|e| e.into_inner());
255                    swept.insert(K::DIR, total);
256                };
257                if total <= self.budget {
258                    note(total);
259                    return Ok(());
260                }
261                // Down to three quarters, so a cache that has filled its budget, where an LRU
262                // cache settles, does not sweep again on the next write.
263                let target = self.budget / 4 * 3;
264                kept.sort();
265                for (_, len, file) in kept {
266                    if total <= target {
267                        break;
268                    }
269                    if !keep.contains(&file) && fs::remove_file(&file).is_ok() {
270                        total -= len;
271                    }
272                }
273                note(total);
274                Ok(())
275            })
276            .or_log(&format!("sweep the {} cache", K::DIR));
277    }
278
279    fn read(&self, file: &Path, bytes: &[u8], key: &str, fingerprint: &str) -> Option<K::Value> {
280        let Some(frame) = unframe(bytes, K::VERSION) else {
281            damaged::<K>(file);
282            return None;
283        };
284        // Another key is a hash collision and another fingerprint a changed dataset:
285        // ordinary misses, not damage.
286        if frame.key != key || frame.fingerprint != fingerprint {
287            return None;
288        }
289        K::decode(frame.payload).or_else(|| {
290            damaged::<K>(file);
291            None
292        })
293    }
294}
295
296/// Files earlier builds wrote, which nothing reads now.
297fn retire_legacy(cache: &CacheManager) {
298    for name in [
299        "datasets.json",
300        "dataset_shapes.json",
301        "cloud_sources.json",
302        "visits.json",
303    ] {
304        let _ = fs::remove_file(cache.cache_file(name));
305    }
306    let _ = fs::remove_dir_all(cache.cache_file("dataset_shapes"));
307}
308
309/// Log the first damaged file of each kind this process meets; one is news, a
310/// thousand is noise.
311fn damaged<K: Kind>(file: &Path) {
312    static LOGGED: std::sync::Mutex<Vec<&'static str>> = std::sync::Mutex::new(Vec::new());
313    let mut logged = LOGGED.lock().unwrap_or_else(|e| e.into_inner());
314    if !logged.contains(&K::DIR) {
315        logged.push(K::DIR);
316        log::warn!(target: "datui", "{} is damaged or from another build; ignoring it", file.display());
317    }
318}
319
320/// Mark a file as used just now. Best effort: an entry that cannot be re-dated is
321/// still an entry that can be used.
322fn touch(file: &Path) {
323    fs::OpenOptions::new()
324        .write(true)
325        .open(file)
326        .and_then(|f| f.set_modified(SystemTime::now()))
327        .or_log("re-date a cache entry");
328}
329
330pub(crate) struct Frame<'a> {
331    pub(crate) key: &'a str,
332    pub(crate) fingerprint: &'a str,
333    pub(crate) payload: &'a [u8],
334}
335
336pub(crate) fn frame(version: u16, key: &str, fingerprint: &str, payload: &[u8]) -> Result<Vec<u8>> {
337    let mut out = Vec::with_capacity(36 + key.len() + fingerprint.len() + payload.len());
338    out.extend_from_slice(MAGIC);
339    out.extend_from_slice(&FRAME_VERSION.to_le_bytes());
340    out.extend_from_slice(&version.to_le_bytes());
341    out.extend_from_slice(&u32::try_from(key.len())?.to_le_bytes());
342    out.extend_from_slice(key.as_bytes());
343    out.extend_from_slice(&u32::try_from(fingerprint.len())?.to_le_bytes());
344    out.extend_from_slice(fingerprint.as_bytes());
345    out.extend_from_slice(&(payload.len() as u64).to_le_bytes());
346    out.extend_from_slice(payload);
347    out.extend_from_slice(&CRC64.checksum(&out).to_le_bytes());
348    Ok(out)
349}
350
351pub(crate) fn unframe(bytes: &[u8], version: u16) -> Option<Frame<'_>> {
352    let (body, sum) = bytes.split_last_chunk::<8>()?;
353    if CRC64.checksum(body) != u64::from_le_bytes(*sum) {
354        return None;
355    }
356    let rest = body.strip_prefix(MAGIC)?;
357    let (frame_version, rest) = rest.split_first_chunk::<2>()?;
358    let (kind_version, rest) = rest.split_first_chunk::<2>()?;
359    if u16::from_le_bytes(*frame_version) != FRAME_VERSION
360        || u16::from_le_bytes(*kind_version) != version
361    {
362        return None;
363    }
364    let (key, rest) = take_str(rest)?;
365    let (fingerprint, rest) = take_str(rest)?;
366    let (len, payload) = rest.split_first_chunk::<8>()?;
367    (usize::try_from(u64::from_le_bytes(*len)).ok()? == payload.len()).then_some(Frame {
368        key,
369        fingerprint,
370        payload,
371    })
372}
373
374fn take_str(bytes: &[u8]) -> Option<(&str, &[u8])> {
375    let (len, rest) = bytes.split_first_chunk::<4>()?;
376    let len = usize::try_from(u32::from_le_bytes(*len)).ok()?;
377    let (text, rest) = (rest.get(..len)?, rest.get(len..)?);
378    Some((std::str::from_utf8(text).ok()?, rest))
379}