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 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 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 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 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 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 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
296fn 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
309fn 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
320fn 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}