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