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