1use 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
30const STALE_TEMP: Duration = Duration::from_secs(3_600);
33
34pub(crate) trait Kind {
36 const DIR: &'static str;
38 const EXT: &'static str;
39 const VERSION: u16;
41 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
48static CRC64: crc::Crc<u64> = crc::Crc::<u64>::new(&crc::CRC_64_XZ);
51
52pub 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 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
79pub fn stable_hash(bytes: &[u8]) -> u64 {
81 CRC64.checksum(bytes)
82}
83
84pub(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 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 pub(crate) fn touch(&self, key: &str) {
126 let file = self.file(key);
127 if file.exists() {
128 touch(&file);
129 }
130 }
131
132 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 #[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 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 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 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 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 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 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 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
301fn 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
314fn 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
325fn 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}