1use crate::logging::LogFailure;
2use color_eyre::Result;
3use std::fs;
4use std::io::Write;
5use std::path::{Path, PathBuf};
6
7mod store;
8pub(crate) use store::{Kind, Store};
9pub use store::{StableHasher, stable_hash};
10
11#[derive(Clone, Debug)]
13pub struct CacheManager {
14 pub(crate) cache_dir: PathBuf,
15 pub(crate) swept:
18 std::sync::Arc<std::sync::Mutex<std::collections::HashMap<&'static str, u64>>>,
19}
20
21impl CacheManager {
22 pub fn with_dir(cache_dir: PathBuf) -> Self {
24 Self {
25 cache_dir,
26 swept: Default::default(),
27 }
28 }
29
30 pub fn new(app_name: &str) -> Result<Self> {
33 #[cfg(test)]
34 isolate_cache();
35 if let Some(dir) = std::env::var_os("DATUI_CACHE_DIR") {
36 return Ok(Self::with_dir(PathBuf::from(dir)));
37 }
38 if running_as_a_cargo_test() {
41 panic!(
42 "DATUI_CACHE_DIR is not set: a test would write to the real cache. \
43 Call common::isolate_cache() (or take the runtime from \
44 common::test_runtime(), which does) before building an App or a \
45 CacheManager."
46 );
47 }
48
49 let cache_dir = dirs::cache_dir()
50 .ok_or_else(|| color_eyre::eyre::eyre!("Could not determine cache directory"))?
51 .join(app_name);
52
53 Ok(Self::with_dir(cache_dir))
54 }
55
56 pub fn cache_dir(&self) -> &Path {
58 &self.cache_dir
59 }
60
61 pub fn cache_file(&self, filename: &str) -> PathBuf {
63 self.cache_dir.join(filename)
64 }
65
66 pub fn ensure_cache_dir(&self) -> Result<()> {
68 if !self.cache_dir.exists() {
69 fs::create_dir_all(&self.cache_dir)?;
70 }
71 Ok(())
72 }
73
74 pub fn clear_all(&self) -> Result<()> {
78 for dir in [Shapes::DIR, Facts::DIR, CloudListings::DIR] {
79 match fs::remove_dir_all(self.cache_file(dir)) {
80 Err(e) if e.kind() != std::io::ErrorKind::NotFound => {
81 log::warn!(target: "datui", "remove the {dir} cache: {e}");
82 }
83 _ => {}
84 }
85 }
86 let Ok(entries) = fs::read_dir(&self.cache_dir) else {
87 return Ok(());
88 };
89 for entry in entries.flatten() {
90 let path = entry.path();
91 let list = entry.file_name().to_string_lossy().contains(HISTORY_SUFFIX);
93 if list && path.is_file() {
94 fs::remove_file(&path).or_log(&format!("remove {}", path.display()));
95 }
96 }
97 Ok(())
98 }
99
100 pub fn load_history_file(&self, history_id: &str) -> Result<Vec<String>> {
103 let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
104
105 let bytes = match fs::read(&history_file) {
106 Ok(bytes) => bytes,
107 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
108 Err(e) => return Err(e.into()),
109 };
110 let mut history = Vec::new();
111 for line in bytes.split(|&b| b == b'\n') {
112 let line = line.strip_suffix(b"\r").unwrap_or(line);
113 match std::str::from_utf8(line) {
114 Ok(line) if !line.trim().is_empty() => history.push(line.to_string()),
115 Ok(_) => {}
116 Err(e) => {
117 log::warn!(target: "datui", "{history_id} history: skipped a line: {e}")
118 }
119 }
120 }
121
122 Ok(history)
123 }
124
125 fn load_history_or_log(&self, history_id: &str) -> Vec<String> {
127 self.load_history_file(history_id)
128 .inspect_err(|e| log::warn!(target: "datui", "read {history_id} history: {e:#}"))
129 .unwrap_or_default()
130 }
131
132 pub fn update_history_file<F>(&self, history_id: &str, update: F) -> Result<HistoryUpdate>
138 where
139 F: FnOnce(&mut Vec<String>),
140 {
141 use fs2::FileExt;
142
143 self.ensure_cache_dir()?;
144 let lock_path = self.cache_file(&format!("{}_history.lock", history_id));
145 let Some(lock) = lock_file(&lock_path, LOCK_TIMEOUT)? else {
146 log::info!(target: "datui", "{history_id} history not updated: its lock is busy");
147 return Ok(HistoryUpdate::SkippedBusy);
148 };
149
150 let mut entries = self.load_history_file(history_id)?;
153 update(&mut entries);
154 let result = self.save_history_file(history_id, &entries);
155
156 let _ = FileExt::unlock(&lock);
158 result.map(|()| HistoryUpdate::Written)
159 }
160
161 pub fn save_history_file(&self, history_id: &str, history: &[String]) -> Result<()> {
163 self.ensure_cache_dir()?;
164 let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
165
166 let mut text = String::new();
168 for entry in history {
169 text.push_str(entry);
170 text.push('\n');
171 }
172 atomic_write(&history_file, text.as_bytes())?;
175 Ok(())
176 }
177}
178
179const HISTORY_SUFFIX: &str = "_history.txt";
181
182const TERMINAL_MODES: &str = "terminal_modes";
184
185pub(crate) fn atomic_write(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
190 static SERIAL: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
191 let serial = SERIAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
192 let mut name = path.file_name().unwrap_or_default().to_owned();
193 name.push(format!(".{}.{serial}.tmp", std::process::id()));
194 let temp = path.with_file_name(name);
195 let written = (|| {
196 let mut file = fs::File::create(&temp)?;
197 file.write_all(bytes)?;
198 file.sync_all()?;
199 fs::rename(&temp, path)
200 })();
201 if written.is_err() {
202 let _ = fs::remove_file(&temp);
203 }
204 written
205}
206
207pub(crate) fn lock_file(
210 path: &Path,
211 timeout: std::time::Duration,
212) -> std::io::Result<Option<fs::File>> {
213 take_lock(path, timeout, false)
214}
215
216pub(crate) fn lock_file_shared(
218 path: &Path,
219 timeout: std::time::Duration,
220) -> std::io::Result<Option<fs::File>> {
221 take_lock(path, timeout, true)
222}
223
224fn take_lock(
225 path: &Path,
226 timeout: std::time::Duration,
227 shared: bool,
228) -> std::io::Result<Option<fs::File>> {
229 use fs2::FileExt;
230
231 let lock = fs::OpenOptions::new()
232 .create(true)
233 .write(true)
234 .truncate(false)
235 .open(path)?;
236 let deadline = std::time::Instant::now() + timeout;
237 loop {
238 let taken = if shared {
239 FileExt::try_lock_shared(&lock)
240 } else {
241 FileExt::try_lock_exclusive(&lock)
242 };
243 if taken.is_ok() {
244 return Ok(Some(lock));
245 }
246 if std::time::Instant::now() >= deadline {
247 return Ok(None);
248 }
249 std::thread::sleep(std::time::Duration::from_millis(2));
250 }
251}
252
253const LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
258
259pub const MAX_RECENTS: usize = 50;
261
262#[derive(Debug, Clone, Copy, PartialEq, Eq)]
265pub enum HistoryUpdate {
266 Written,
268 SkippedBusy,
270}
271
272impl CacheManager {
273 pub fn load_recents(&self) -> Vec<std::path::PathBuf> {
276 self.load_recents_with_visits().0
277 }
278
279 pub fn load_recents_with_visits(
281 &self,
282 ) -> (Vec<PathBuf>, std::collections::HashMap<PathBuf, Visits>) {
283 let mut recents = Vec::new();
284 let mut visits = std::collections::HashMap::new();
285 for line in self.load_history_or_log("recents") {
286 let (path, seen) = parse_recent(&line);
287 recents.push(PathBuf::from(path));
288 if seen.count > 0 {
289 visits.insert(PathBuf::from(path), seen);
290 }
291 }
292 (recents, visits)
293 }
294
295 pub fn load_folds(&self) -> std::collections::HashMap<String, bool> {
298 self.load_history_or_log("home_folds")
299 .into_iter()
300 .filter_map(|line| {
301 let (title, state) = line.rsplit_once('\t')?;
302 Some((title.to_string(), state.trim() == "1"))
303 })
304 .collect()
305 }
306
307 pub fn save_folds(&self, folds: &std::collections::HashMap<String, bool>) {
309 let mut lines: Vec<String> = folds
310 .iter()
311 .map(|(title, folded)| format!("{title}\t{}", if *folded { 1 } else { 0 }))
312 .collect();
313 lines.sort();
314 self.save_history_file("home_folds", &lines)
315 .or_log("save home folds");
316 }
317
318 pub fn terminal_mode(&self, terminal: &str) -> Option<crate::config::ThemeMode> {
321 self.load_history_or_log(TERMINAL_MODES)
322 .iter()
323 .find_map(|line| match line.split_once('\t')? {
324 (key, "dark") if key == terminal => Some(crate::config::ThemeMode::Dark),
325 (key, "light") if key == terminal => Some(crate::config::ThemeMode::Light),
326 _ => None,
327 })
328 }
329
330 pub fn remember_terminal_mode(&self, terminal: &str, mode: crate::config::ThemeMode) {
332 if self.terminal_mode(terminal) == Some(mode) {
333 return;
334 }
335 let word = match mode {
336 crate::config::ThemeMode::Light => "light",
337 _ => "dark",
338 };
339 self.update_history_file(TERMINAL_MODES, |lines| {
340 lines.retain(|line| line.split_once('\t').is_none_or(|(key, _)| key != terminal));
341 lines.push(format!("{terminal}\t{word}"));
342 })
343 .or_log("remember the terminal's background");
344 }
345
346 pub fn forget_recent(&self, path: &std::path::Path) {
348 let target = path.to_string_lossy().into_owned();
349 self.update_history_file("recents", |recents| {
350 recents.retain(|line| parse_recent(line).0 != target);
351 })
352 .or_log("forget a recent");
353 }
354
355 pub fn forget_recents(&self, paths: &[std::path::PathBuf]) {
357 let targets: Vec<String> = paths
358 .iter()
359 .map(|p| p.to_string_lossy().into_owned())
360 .collect();
361 self.update_history_file("recents", |recents| {
362 recents.retain(|line| !targets.iter().any(|t| t == parse_recent(line).0));
363 })
364 .or_log("forget recents");
365 }
366
367 pub fn clear_recents(&self) {
369 self.update_history_file("recents", |recents| recents.clear())
370 .or_log("clear recents");
371 }
372
373 fn recent_is_worth_keeping(path: &str, mounts: &crate::home::locality::Mounts) -> bool {
379 let path = std::path::Path::new(path);
380 if crate::home::locality::object_scheme(path).is_some() || mounts.is_network(path) {
382 return true;
383 }
384 match path.parent() {
385 Some(parent) if !parent.as_os_str().is_empty() => parent.exists(),
386 _ => true,
387 }
388 }
389
390 pub fn push_recent(&self, path: &std::path::Path) -> HistoryUpdate {
394 let looks_like_url = path.to_string_lossy().contains("://");
397 let stored = if looks_like_url {
398 path.to_path_buf()
399 } else {
400 crate::canonical::canonicalize(path)
403 .or_else(|e| match crate::formats::members::split(path) {
404 Some((db, table)) => crate::canonical::canonicalize(&db)
405 .map(|db| crate::formats::members::place(&db, &table)),
406 None => Err(e),
407 })
408 .unwrap_or_else(|_| path.to_path_buf())
409 };
410 let entry = stored.to_string_lossy().into_owned();
411
412 let mounts = crate::home::locality::Mounts::current();
414
415 let now = unix_now();
416 self.update_history_file("recents", |recents| {
417 let mut visits = Visits::default();
418 if let Some(at) = recents.iter().position(|l| parse_recent(l).0 == entry) {
419 visits = parse_recent(&recents.remove(at)).1;
420 }
421 visits.count = visits.count.saturating_add(1);
422 visits.last = now;
423 recents.insert(0, format!("{entry}\t{}\t{}", visits.count, visits.last));
424 recents.retain(|line| Self::recent_is_worth_keeping(parse_recent(line).0, &mounts));
425 recents.truncate(MAX_RECENTS);
426 })
427 .inspect_err(|e| log::warn!(target: "datui", "record a recent: {e:#}"))
428 .unwrap_or(HistoryUpdate::SkippedBusy)
429 }
430}
431
432fn parse_recent(line: &str) -> (&str, Visits) {
435 let mut parts = line.rsplitn(3, '\t');
436 if let (Some(last), Some(count), Some(path)) = (parts.next(), parts.next(), parts.next())
437 && let (Ok(last), Ok(count)) = (last.parse(), count.parse())
438 {
439 return (path, Visits { count, last });
440 }
441 (line, Visits::default())
442}
443
444fn unix_now() -> u64 {
445 std::time::SystemTime::now()
446 .duration_since(std::time::UNIX_EPOCH)
447 .map(|d| d.as_secs())
448 .unwrap_or_default()
449}
450
451#[derive(Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize)]
454pub struct Visits {
455 pub count: u32,
456 pub last: u64,
458}
459
460impl Visits {
461 pub fn frecency(&self, now: u64) -> f64 {
464 let age = now.saturating_sub(self.last);
465 let weight = match age {
466 a if a < 3_600 => 4.0,
467 a if a < 86_400 => 2.0,
468 a if a < 604_800 => 0.5,
469 _ => 0.25,
470 };
471 f64::from(self.count) * weight
472 }
473}
474
475pub fn by_frecency(
477 mut recents: Vec<PathBuf>,
478 visits: &std::collections::HashMap<PathBuf, Visits>,
479) -> Vec<PathBuf> {
480 let now = unix_now();
481 let score = |p: &PathBuf| visits.get(p).map_or(0.0, |v| v.frecency(now));
482 recents.sort_by(|a, b| score(b).total_cmp(&score(a)));
483 recents
484}
485
486#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
489pub struct DatasetFacts {
490 pub mtime: u64,
492 pub size: u64,
494 pub rows: Option<usize>,
495 pub cols: Option<usize>,
496 #[serde(default)]
499 pub cols_sampled: bool,
500 #[serde(default)]
502 pub columns: Vec<String>,
503 #[serde(default)]
506 pub kind: Option<crate::home::discover::EntryKind>,
507 #[serde(default)]
510 pub classified_by: u32,
511 #[serde(default)]
514 pub cost: crate::home::discover::Cost,
515 #[serde(
518 default,
519 skip_serializing_if = "crate::home::discover::Holds::is_empty"
520 )]
521 pub holds: crate::home::discover::Holds,
522}
523
524#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
529pub struct DatasetShape {
530 pub fingerprint: String,
533 pub files: Vec<CachedFooter>,
535 pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
539 pub taken_at: u64,
541}
542
543#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
546pub struct CachedFooter {
547 pub schema: Option<usize>,
550 pub row_group_rows: Vec<usize>,
552 pub row_group_bytes: Vec<usize>,
555 #[serde(default, skip_serializing_if = "Vec::is_empty")]
558 pub column_bytes: Vec<usize>,
559}
560
561impl DatasetShape {
562 pub fn fingerprint_of<'a>(
567 files: impl IntoIterator<Item = (&'a str, u64, u64, Option<&'a str>)>,
568 ) -> String {
569 let mut hasher = StableHasher::default();
570 let mut count = 0usize;
571 let mut bytes = 0u64;
572 for (key, size, stamp, etag) in files {
573 hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
574 match etag {
575 Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
576 None => hasher.u64(0),
577 };
578 count += 1;
579 bytes = bytes.saturating_add(size);
580 }
581 format!("{count}-{bytes}-{:016x}", hasher.finish())
582 }
583
584 pub fn intern_schema(
586 schemas: &mut Vec<Vec<(String, polars::prelude::DataType)>>,
587 schema: &polars::prelude::Schema,
588 ) -> usize {
589 let columns: Vec<(String, polars::prelude::DataType)> = schema
590 .iter()
591 .map(|(name, dtype)| (name.to_string(), dtype.clone()))
592 .collect();
593 schemas
594 .iter()
595 .position(|s| *s == columns)
596 .unwrap_or_else(|| {
597 schemas.push(columns);
598 schemas.len() - 1
599 })
600 }
601
602 pub fn schema_at(
604 schemas: &[Vec<(String, polars::prelude::DataType)>],
605 at: usize,
606 ) -> Option<polars::prelude::Schema> {
607 let columns = schemas.get(at)?;
608 let mut schema = polars::prelude::Schema::with_capacity(columns.len());
609 for (name, dtype) in columns {
610 schema.with_column(name.as_str().into(), dtype.clone());
611 }
612 Some(schema)
613 }
614}
615
616#[cfg(test)]
620pub(crate) fn isolate_cache() {
621 static SCRATCH: std::sync::Mutex<Vec<tempfile::TempDir>> = std::sync::Mutex::new(Vec::new());
623 unsafe extern "C" {
624 fn atexit(callback: extern "C" fn()) -> std::ffi::c_int;
625 }
626 extern "C" fn remove_scratch_dirs() {
627 if let Ok(mut held) = SCRATCH.lock() {
628 held.clear();
629 }
630 }
631 let scratch_dir = |prefix: &str| {
632 tempfile::Builder::new()
633 .prefix(prefix)
634 .tempdir()
635 .expect("a scratch directory for the test process")
636 };
637
638 static ISOLATE: std::sync::Once = std::sync::Once::new();
639 ISOLATE.call_once(|| {
640 let dir = scratch_dir("datui-unit-cache-");
641 let config_dir = scratch_dir("datui-unit-config-");
642 unsafe { std::env::set_var("DATUI_CACHE_DIR", dir.path()) };
645 unsafe { std::env::set_var("DATUI_CONFIG_DIR", config_dir.path()) };
646 let mut held = SCRATCH.lock().unwrap_or_else(|e| e.into_inner());
647 held.push(dir);
648 held.push(config_dir);
649 unsafe { atexit(remove_scratch_dirs) };
652 });
653}
654
655pub(crate) fn running_as_a_cargo_test() -> bool {
658 std::env::current_exe()
659 .is_ok_and(|exe| cargo_test_layout(&exe, std::env::var_os("CARGO_TARGET_DIR").as_deref()))
660}
661
662fn cargo_test_layout(exe: &Path, target_dir: Option<&std::ffi::OsStr>) -> bool {
665 let Some(deps) = exe.parent() else {
666 return false;
667 };
668 if deps.file_name().is_none_or(|name| name != "deps") {
669 return false;
670 }
671 if target_dir.is_some_and(|dir| exe.starts_with(dir)) {
672 return true;
673 }
674 deps.ancestors()
675 .skip(1)
676 .take(3)
677 .any(|dir| dir.file_name().is_some_and(|name| name == "target"))
678}
679
680#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
683pub struct CloudListing {
684 pub fingerprint: String,
686 pub buckets: Vec<String>,
687 pub listed_at: u64,
689}
690
691pub(crate) struct Shapes;
695
696impl Kind for Shapes {
697 const DIR: &'static str = "shapes";
698 const EXT: &'static str = "shape";
699 const VERSION: u16 = 1;
700 const BUDGET: u64 = 128 << 20;
701 type Value = DatasetShape;
702
703 fn encode(shape: &DatasetShape) -> Result<Vec<u8>> {
704 encode_shape(shape)
705 }
706
707 fn decode(payload: &[u8]) -> Option<DatasetShape> {
708 decode_shape(payload)
709 }
710}
711
712#[derive(Debug, Clone, Default, PartialEq, Eq)]
716pub struct FileFooters {
717 pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
719 pub files: Vec<(u64, CachedFooter)>,
721}
722
723pub fn file_identity(key: &str, size: u64, stamp: u64, etag: Option<&str>) -> u64 {
726 let mut hasher = StableHasher::default();
727 hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
728 match etag {
729 Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
730 None => hasher.u64(0),
731 };
732 hasher.finish()
733}
734
735pub(crate) struct FileFootersKind;
737
738impl Kind for FileFootersKind {
739 const DIR: &'static str = "file_footers";
740 const EXT: &'static str = "footers";
741 const VERSION: u16 = 1;
742 const BUDGET: u64 = 64 << 20;
743 type Value = FileFooters;
744
745 fn encode(footers: &FileFooters) -> Result<Vec<u8>> {
746 let shape = DatasetShape {
747 fingerprint: String::new(),
748 files: footers.files.iter().map(|(_, f)| f.clone()).collect(),
749 schemas: footers.schemas.clone(),
750 taken_at: 0,
751 };
752 let mut out = encode_shape(&shape)?;
753 for (identity, _) in &footers.files {
754 out.extend_from_slice(&identity.to_le_bytes());
755 }
756 Ok(out)
757 }
758
759 fn decode(payload: &[u8]) -> Option<FileFooters> {
760 let count = payload.len().checked_sub(4)?;
762 let (len, rest) = payload.split_first_chunk::<4>()?;
763 let header = usize::try_from(u32::from_le_bytes(*len)).ok()?;
764 let mut body = rest.get(header..)?;
765 let files = usize::try_from(take_varint(&mut body)?).ok()?;
766 let ids = files.checked_mul(8)?;
767 if ids > count {
768 return None;
769 }
770 let (shape, identities) = payload.split_at(payload.len() - ids);
771 let shape = decode_shape(shape)?;
772 (shape.files.len() == files).then(|| FileFooters {
773 schemas: shape.schemas,
774 files: identities
775 .as_chunks::<8>()
776 .0
777 .iter()
778 .map(|id| u64::from_le_bytes(*id))
779 .zip(shape.files)
780 .collect(),
781 })
782 }
783}
784
785pub(crate) struct Facts;
788
789impl Kind for Facts {
790 const DIR: &'static str = "facts";
791 const EXT: &'static str = "facts";
792 const VERSION: u16 = 1;
793 const BUDGET: u64 = 16 << 20;
794 type Value = DatasetFacts;
795
796 fn encode(facts: &DatasetFacts) -> Result<Vec<u8>> {
797 Ok(serde_json::to_vec(facts)?)
798 }
799
800 fn decode(payload: &[u8]) -> Option<DatasetFacts> {
801 serde_json::from_slice(payload).ok()
802 }
803}
804
805pub(crate) struct CloudListings;
807
808impl Kind for CloudListings {
809 const DIR: &'static str = "cloud_listings";
810 const EXT: &'static str = "listing";
811 const VERSION: u16 = 1;
812 const BUDGET: u64 = 4 << 20;
813 type Value = CloudListing;
814
815 fn encode(listing: &CloudListing) -> Result<Vec<u8>> {
816 Ok(serde_json::to_vec(listing)?)
817 }
818
819 fn decode(payload: &[u8]) -> Option<CloudListing> {
820 serde_json::from_slice(payload).ok()
821 }
822}
823
824#[derive(serde::Serialize, serde::Deserialize)]
827struct ShapeHeader {
828 fingerprint: String,
829 schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
830 taken_at: u64,
831}
832fn put_varint(out: &mut Vec<u8>, mut n: u64) {
833 while n >= 0x80 {
834 out.push((n as u8) | 0x80);
835 n >>= 7;
836 }
837 out.push(n as u8);
838}
839
840fn take_varint(bytes: &mut &[u8]) -> Option<u64> {
841 let mut n = 0u64;
842 for shift in (0..64).step_by(7) {
843 let (&b, rest) = bytes.split_first()?;
844 *bytes = rest;
845 n |= u64::from(b & 0x7f) << shift;
846 if b & 0x80 == 0 {
847 return Some(n);
848 }
849 }
850 None
851}
852
853fn put_list(out: &mut Vec<u8>, values: &[usize]) {
854 put_varint(out, values.len() as u64);
855 for &v in values {
856 put_varint(out, v as u64);
857 }
858}
859
860fn take_list(bytes: &mut &[u8]) -> Option<Vec<usize>> {
861 let len = usize::try_from(take_varint(bytes)?).ok()?;
862 if len > bytes.len() {
865 return None;
866 }
867 (0..len)
868 .map(|_| take_varint(bytes).and_then(|v| usize::try_from(v).ok()))
869 .collect()
870}
871
872fn encode_shape(shape: &DatasetShape) -> Result<Vec<u8>> {
873 let header = serde_json::to_vec(&ShapeHeader {
874 fingerprint: shape.fingerprint.clone(),
875 schemas: shape.schemas.clone(),
876 taken_at: shape.taken_at,
877 })?;
878 let mut out = Vec::with_capacity(8 + header.len() + shape.files.len() * 8);
879 out.extend_from_slice(&u32::try_from(header.len())?.to_le_bytes());
880 out.extend_from_slice(&header);
881 put_varint(&mut out, shape.files.len() as u64);
882 for file in &shape.files {
883 put_varint(&mut out, file.schema.map_or(0, |s| s as u64 + 1));
884 put_list(&mut out, &file.row_group_rows);
885 put_list(&mut out, &file.row_group_bytes);
886 put_list(&mut out, &file.column_bytes);
887 }
888 Ok(out)
889}
890
891fn decode_shape(bytes: &[u8]) -> Option<DatasetShape> {
893 let (len, rest) = bytes.split_first_chunk::<4>()?;
894 let len = usize::try_from(u32::from_le_bytes(*len)).ok()?;
895 let (header, mut body) = (rest.get(..len)?, rest.get(len..)?);
896 let header: ShapeHeader = serde_json::from_slice(header).ok()?;
897 let count = usize::try_from(take_varint(&mut body)?).ok()?;
898 if count > body.len() {
899 return None;
900 }
901 let mut files = Vec::with_capacity(count);
902 let (mut rows, mut bytes_total) = (0usize, 0usize);
904 for _ in 0..count {
905 let schema = match take_varint(&mut body)? {
906 0 => None,
907 at => Some(usize::try_from(at - 1).ok()?),
908 };
909 let footer = CachedFooter {
910 schema,
911 row_group_rows: take_list(&mut body)?,
912 row_group_bytes: take_list(&mut body)?,
913 column_bytes: take_list(&mut body)?,
914 };
915 if let Some(at) = footer.schema
916 && at >= header.schemas.len()
917 {
918 return None;
919 }
920 for &n in &footer.row_group_rows {
921 rows = rows.checked_add(n)?;
922 }
923 for &n in footer.row_group_bytes.iter().chain(&footer.column_bytes) {
924 bytes_total = bytes_total.checked_add(n)?;
925 }
926 files.push(footer);
927 }
928 body.is_empty().then_some(DatasetShape {
929 fingerprint: header.fingerprint,
930 files,
931 schemas: header.schemas,
932 taken_at: header.taken_at,
933 })
934}
935
936impl CacheManager {
937 #[cfg(test)]
939 pub fn dataset_shapes_kept(&self) -> usize {
940 Store::<Shapes>::new(self).len()
941 }
942
943 pub fn dataset_shape(&self, path: &str, fingerprint: &str) -> Option<DatasetShape> {
946 Store::<Shapes>::new(self).get(path, fingerprint)
947 }
948
949 pub fn has_dataset_shape(&self, path: &str) -> bool {
951 Store::<Shapes>::new(self).file(path).exists()
952 }
953
954 pub fn save_dataset_shape(&self, path: &str, shape: DatasetShape) {
956 Store::<Shapes>::new(self).put(path, &shape.fingerprint, &shape);
957 }
958
959 pub fn file_footers(&self, path: &str) -> Option<FileFooters> {
961 Store::<FileFootersKind>::new(self).get(path, "")
962 }
963
964 pub fn save_file_footers(&self, path: &str, footers: &FileFooters) {
966 Store::<FileFootersKind>::new(self).put(path, "", footers);
967 }
968
969 pub fn cloud_listing(&self, id: &str, fingerprint: &str) -> Option<CloudListing> {
971 Store::<CloudListings>::new(self).get(id, fingerprint)
972 }
973
974 pub fn save_cloud_listing(&self, id: &str, listing: CloudListing) {
976 Store::<CloudListings>::new(self).put(id, &listing.fingerprint.clone(), &listing);
977 }
978
979 pub fn load_hidden_cloud_sources(&self) -> Vec<String> {
981 self.load_history_or_log("cloud_hidden")
982 }
983
984 pub fn hide_cloud_source(&self, id: &str) {
986 let id = id.to_string();
987 self.update_history_file("cloud_hidden", |hidden| {
988 if !hidden.contains(&id) {
989 hidden.push(id.clone());
990 }
991 })
992 .or_log("hide a cloud source");
993 }
994
995 pub fn examples_hidden(&self) -> bool {
998 !self.load_history_or_log("examples_hidden").is_empty()
999 }
1000
1001 pub fn hide_examples(&self) {
1003 self.save_history_file("examples_hidden", &["hidden".to_string()])
1004 .or_log("hide the example datasets");
1005 }
1006
1007 pub fn load_remembered_places(&self) -> Vec<PathBuf> {
1010 let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1011 if !file.exists() {
1012 return Vec::new();
1013 }
1014 self.load_history_or_log("home_remembered")
1015 .into_iter()
1016 .map(PathBuf::from)
1017 .collect()
1018 }
1019
1020 pub fn clear_remembered_places(&self) {
1022 let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1023 if let Err(e) = std::fs::remove_file(&file)
1024 && e.kind() != std::io::ErrorKind::NotFound
1025 {
1026 log::warn!(target: "datui", "remove {}: {e}", file.display());
1027 }
1028 }
1029
1030 pub fn save_remembered_places(&self, places: &[PathBuf]) -> Result<()> {
1032 let places: Vec<String> = places
1033 .iter()
1034 .map(|p| p.to_string_lossy().into_owned())
1035 .collect();
1036 self.save_history_file("home_remembered", &places)
1037 }
1038}
1039
1040impl CacheManager {
1041 pub fn load_dataset_facts(&self) -> std::collections::HashMap<PathBuf, DatasetFacts> {
1044 Store::<Facts>::new(self)
1045 .scan()
1046 .into_iter()
1047 .map(|(path, facts)| (PathBuf::from(path), facts))
1048 .collect()
1049 }
1050
1051 pub fn dataset_facts(&self, path: &Path) -> Option<DatasetFacts> {
1053 Store::<Facts>::new(self).get(path.to_str()?, "")
1054 }
1055
1056 pub fn touch_dataset_facts<'a>(&self, paths: impl IntoIterator<Item = &'a Path>) {
1059 let store = Store::<Facts>::new(self);
1060 for path in paths.into_iter().filter_map(Path::to_str) {
1061 store.touch(path);
1062 }
1063 }
1064
1065 pub fn record_dataset_facts(&self, facts: &[(PathBuf, DatasetFacts)]) {
1067 Store::<Facts>::new(self).put_all(
1068 facts
1069 .iter()
1070 .filter_map(|(path, facts)| Some((path.to_str()?, "", facts))),
1071 );
1072 }
1073
1074 fn with_cache_lock<F>(&self, name: &str, work: F) -> Result<()>
1076 where
1077 F: FnOnce() -> Result<()>,
1078 {
1079 use fs2::FileExt;
1080
1081 self.ensure_cache_dir()?;
1082 let Some(lock) = lock_file(&self.cache_file(&format!("{name}.lock")), LOCK_TIMEOUT)? else {
1083 log::info!(target: "datui", "{name} cache not updated: its lock is busy");
1084 return Ok(());
1085 };
1086
1087 let result = work();
1088 let _ = FileExt::unlock(&lock);
1089 result
1090 }
1091}
1092
1093#[cfg(test)]
1094mod harness_tests;
1095
1096#[cfg(test)]
1097mod recents_pruning_tests;
1098
1099#[cfg(test)]
1100mod dataset_shape_tests;
1101
1102#[cfg(test)]
1104mod store_harness_tests;
1105
1106#[cfg(test)]
1107mod facts_compat_tests;