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        // Made once a batch, not asked again per entry.
181        let mut made_dir = false;
182        for (key, fingerprint, value) in entries {
183            let file = self.file(key);
184            let stored = (|| -> Result<()> {
185                let frame = frame(K::VERSION, key, fingerprint, &K::encode(value)?)?;
186                // An entry that has not changed is only dated: a measured directory is
187                // recorded on every look, and most looks find what the last one did.
188                if fs::metadata(&file).is_ok_and(|m| m.len() == frame.len() as u64)
189                    && fs::read(&file).is_ok_and(|old| old == frame)
190                {
191                    touch(&file);
192                    return Ok(());
193                }
194                if !made_dir {
195                    fs::create_dir_all(self.dir())?;
196                    made_dir = true;
197                }
198                atomic_write(&file, &frame)?;
199                bytes += frame.len() as u64;
200                Ok(())
201            })();
202            match stored {
203                Ok(()) => written.push(file),
204                Err(e) => log::warn!(target: "datui", "save a {} entry: {e:#}", K::DIR),
205            }
206        }
207        if !written.is_empty() && self.may_pass_budget(bytes) {
208            self.sweep(&written);
209        }
210    }
211
212    /// Whether the kind may now be past its budget: unknown until this session's first
213    /// sweep, then what that sweep found plus everything written since. Sweeping reads
214    /// the whole directory, and measuring a directory of thousands writes a batch at a
215    /// time; a sweep per batch grew with the cache.
216    fn may_pass_budget(&self, written: u64) -> bool {
217        let mut swept = self.cache.swept.lock().unwrap_or_else(|e| e.into_inner());
218        match swept.get_mut(K::DIR) {
219            Some(total) => {
220                *total = total.saturating_add(written);
221                *total > self.budget
222            }
223            None => true,
224        }
225    }
226
227    /// Drop stale temp files and, past the budget, the least recently used entries down
228    /// to three quarters of it, never one of `keep`. Under the kind's lock only so two sweeps do not race;
229    /// writes land by rename and need none.
230    fn sweep(&self, keep: &[PathBuf]) {
231        self.cache
232            .with_cache_lock(K::DIR, || {
233                retire_legacy(&self.cache);
234                let Ok(entries) = fs::read_dir(self.dir()) else {
235                    return Ok(());
236                };
237                let now = SystemTime::now();
238                let mut kept = Vec::new();
239                for entry in entries.flatten() {
240                    let path = entry.path();
241                    let Ok(meta) = entry.metadata() else { continue };
242                    let Ok(modified) = meta.modified() else {
243                        continue;
244                    };
245                    if path.extension().is_some_and(|x| x == "tmp") {
246                        if now
247                            .duration_since(modified)
248                            .is_ok_and(|age| age > STALE_TEMP)
249                        {
250                            fs::remove_file(&path).or_log("remove a stale cache temp file");
251                        }
252                    } else if path.extension().is_some_and(|x| x == K::EXT) {
253                        kept.push((modified, meta.len(), path));
254                    }
255                }
256                let mut total: u64 = kept.iter().map(|(_, len, _)| len).sum();
257                let found = self.cache.swept.clone();
258                let note = |total: u64| {
259                    let mut swept = found.lock().unwrap_or_else(|e| e.into_inner());
260                    swept.insert(K::DIR, total);
261                };
262                if total <= self.budget {
263                    note(total);
264                    return Ok(());
265                }
266                // Down to three quarters, so a cache that has filled its budget, where an LRU
267                // cache settles, does not sweep again on the next write.
268                let target = self.budget / 4 * 3;
269                kept.sort();
270                for (_, len, file) in kept {
271                    if total <= target {
272                        break;
273                    }
274                    if !keep.contains(&file) && fs::remove_file(&file).is_ok() {
275                        total -= len;
276                    }
277                }
278                note(total);
279                Ok(())
280            })
281            .or_log(&format!("sweep the {} cache", K::DIR));
282    }
283
284    fn read(&self, file: &Path, bytes: &[u8], key: &str, fingerprint: &str) -> Option<K::Value> {
285        let Some(frame) = unframe(bytes, K::VERSION) else {
286            damaged::<K>(file);
287            return None;
288        };
289        // Another key is a hash collision and another fingerprint a changed dataset:
290        // ordinary misses, not damage.
291        if frame.key != key || frame.fingerprint != fingerprint {
292            return None;
293        }
294        K::decode(frame.payload).or_else(|| {
295            damaged::<K>(file);
296            None
297        })
298    }
299}
300
301/// Files earlier builds wrote, which nothing reads now.
302fn retire_legacy(cache: &CacheManager) {
303    for name in [
304        "datasets.json",
305        "dataset_shapes.json",
306        "cloud_sources.json",
307        "visits.json",
308    ] {
309        let _ = fs::remove_file(cache.cache_file(name));
310    }
311    let _ = fs::remove_dir_all(cache.cache_file("dataset_shapes"));
312}
313
314/// Log the first damaged file of each kind this process meets; one is news, a
315/// thousand is noise.
316fn damaged<K: Kind>(file: &Path) {
317    static LOGGED: std::sync::Mutex<Vec<&'static str>> = std::sync::Mutex::new(Vec::new());
318    let mut logged = LOGGED.lock().unwrap_or_else(|e| e.into_inner());
319    if !logged.contains(&K::DIR) {
320        logged.push(K::DIR);
321        log::warn!(target: "datui", "{} is damaged or from another build; ignoring it", file.display());
322    }
323}
324
325/// Mark a file as used just now. Best effort: an entry that cannot be re-dated is
326/// still an entry that can be used.
327fn touch(file: &Path) {
328    fs::OpenOptions::new()
329        .write(true)
330        .open(file)
331        .and_then(|f| f.set_modified(SystemTime::now()))
332        .or_log("re-date a cache entry");
333}
334
335pub(crate) struct Frame<'a> {
336    pub(crate) key: &'a str,
337    pub(crate) fingerprint: &'a str,
338    pub(crate) payload: &'a [u8],
339}
340
341pub(crate) fn frame(version: u16, key: &str, fingerprint: &str, payload: &[u8]) -> Result<Vec<u8>> {
342    let mut out = Vec::with_capacity(36 + key.len() + fingerprint.len() + payload.len());
343    out.extend_from_slice(MAGIC);
344    out.extend_from_slice(&FRAME_VERSION.to_le_bytes());
345    out.extend_from_slice(&version.to_le_bytes());
346    out.extend_from_slice(&u32::try_from(key.len())?.to_le_bytes());
347    out.extend_from_slice(key.as_bytes());
348    out.extend_from_slice(&u32::try_from(fingerprint.len())?.to_le_bytes());
349    out.extend_from_slice(fingerprint.as_bytes());
350    out.extend_from_slice(&(payload.len() as u64).to_le_bytes());
351    out.extend_from_slice(payload);
352    out.extend_from_slice(&CRC64.checksum(&out).to_le_bytes());
353    Ok(out)
354}
355
356pub(crate) fn unframe(bytes: &[u8], version: u16) -> Option<Frame<'_>> {
357    let (body, sum) = bytes.split_last_chunk::<8>()?;
358    if CRC64.checksum(body) != u64::from_le_bytes(*sum) {
359        return None;
360    }
361    let rest = body.strip_prefix(MAGIC)?;
362    let (frame_version, rest) = rest.split_first_chunk::<2>()?;
363    let (kind_version, rest) = rest.split_first_chunk::<2>()?;
364    if u16::from_le_bytes(*frame_version) != FRAME_VERSION
365        || u16::from_le_bytes(*kind_version) != version
366    {
367        return None;
368    }
369    let (key, rest) = take_str(rest)?;
370    let (fingerprint, rest) = take_str(rest)?;
371    let (len, payload) = rest.split_first_chunk::<8>()?;
372    (usize::try_from(u64::from_le_bytes(*len)).ok()? == payload.len()).then_some(Frame {
373        key,
374        fingerprint,
375        payload,
376    })
377}
378
379fn take_str(bytes: &[u8]) -> Option<(&str, &[u8])> {
380    let (len, rest) = bytes.split_first_chunk::<4>()?;
381    let len = usize::try_from(u32::from_le_bytes(*len)).ok()?;
382    let (text, rest) = (rest.get(..len)?, rest.get(len..)?);
383    Some((std::str::from_utf8(text).ok()?, rest))
384}