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 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 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 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 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
269fn 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
282fn 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
293fn 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}