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