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 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 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 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 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 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
268fn 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
281fn 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
292fn 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}