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 {
21 Self { cache_dir }
22 }
23
24 pub fn new(app_name: &str) -> Result<Self> {
30 #[cfg(test)]
31 isolate_cache();
32 if let Some(dir) = std::env::var_os("DATUI_CACHE_DIR") {
33 return Ok(Self {
34 cache_dir: PathBuf::from(dir),
35 });
36 }
37 if running_as_a_cargo_test() {
43 panic!(
44 "DATUI_CACHE_DIR is not set: a test would write to the real cache. \
45 Call common::isolate_cache() (or take the runtime from \
46 common::test_runtime(), which does) before building an App or a \
47 CacheManager."
48 );
49 }
50
51 let cache_dir = dirs::cache_dir()
52 .ok_or_else(|| color_eyre::eyre::eyre!("Could not determine cache directory"))?
53 .join(app_name);
54
55 Ok(Self { cache_dir })
56 }
57
58 pub fn cache_dir(&self) -> &Path {
60 &self.cache_dir
61 }
62
63 pub fn cache_file(&self, filename: &str) -> PathBuf {
65 self.cache_dir.join(filename)
66 }
67
68 pub fn ensure_cache_dir(&self) -> Result<()> {
70 if !self.cache_dir.exists() {
71 fs::create_dir_all(&self.cache_dir)?;
72 }
73 Ok(())
74 }
75
76 pub fn clear_file(&self, filename: &str) -> Result<()> {
78 let file_path = self.cache_file(filename);
79 if file_path.exists() {
80 fs::remove_file(&file_path)?;
81 }
82 Ok(())
83 }
84
85 pub fn clear_all(&self) -> Result<()> {
92 for dir in [Shapes::DIR, Facts::DIR, CloudListings::DIR] {
93 match fs::remove_dir_all(self.cache_file(dir)) {
94 Err(e) if e.kind() != std::io::ErrorKind::NotFound => {
95 log::warn!(target: "datui", "remove the {dir} cache: {e}");
96 }
97 _ => {}
98 }
99 }
100 let Ok(entries) = fs::read_dir(&self.cache_dir) else {
101 return Ok(());
102 };
103 for entry in entries.flatten() {
104 let path = entry.path();
105 let list = entry.file_name().to_string_lossy().contains(HISTORY_SUFFIX);
107 if list && path.is_file() {
108 fs::remove_file(&path).or_log(&format!("remove {}", path.display()));
109 }
110 }
111 Ok(())
112 }
113
114 pub fn load_history_file(&self, history_id: &str) -> Result<Vec<String>> {
120 let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
121
122 let bytes = match fs::read(&history_file) {
123 Ok(bytes) => bytes,
124 Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
125 Err(e) => return Err(e.into()),
126 };
127 let mut history = Vec::new();
128 for line in bytes.split(|&b| b == b'\n') {
129 let line = line.strip_suffix(b"\r").unwrap_or(line);
130 match std::str::from_utf8(line) {
131 Ok(line) if !line.trim().is_empty() => history.push(line.to_string()),
132 Ok(_) => {}
133 Err(e) => {
134 log::warn!(target: "datui", "{history_id} history: skipped a line: {e}")
135 }
136 }
137 }
138
139 Ok(history)
140 }
141
142 fn load_history_or_log(&self, history_id: &str) -> Vec<String> {
144 self.load_history_file(history_id)
145 .inspect_err(|e| log::warn!(target: "datui", "read {history_id} history: {e:#}"))
146 .unwrap_or_default()
147 }
148
149 pub fn update_history_file<F>(&self, history_id: &str, update: F) -> Result<HistoryUpdate>
169 where
170 F: FnOnce(&mut Vec<String>),
171 {
172 use fs2::FileExt;
173
174 self.ensure_cache_dir()?;
175 let lock_path = self.cache_file(&format!("{}_history.lock", history_id));
176 let Some(lock) = lock_file(&lock_path, LOCK_TIMEOUT)? else {
177 log::info!(target: "datui", "{history_id} history not updated: its lock is busy");
178 return Ok(HistoryUpdate::SkippedBusy);
179 };
180
181 let mut entries = self.load_history_file(history_id)?;
185 update(&mut entries);
186 let result = self.save_history_file(history_id, &entries);
187
188 let _ = FileExt::unlock(&lock);
190 result.map(|()| HistoryUpdate::Written)
191 }
192
193 pub fn save_history_file(&self, history_id: &str, history: &[String]) -> Result<()> {
195 self.ensure_cache_dir()?;
196 let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
197
198 let mut text = String::new();
200 for entry in history {
201 text.push_str(entry);
202 text.push('\n');
203 }
204 atomic_write(&history_file, text.as_bytes())?;
207 Ok(())
208 }
209}
210
211const HISTORY_SUFFIX: &str = "_history.txt";
213
214const TERMINAL_MODES: &str = "terminal_modes";
216
217pub(crate) fn atomic_write(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
226 static SERIAL: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
227 let serial = SERIAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
228 let mut name = path.file_name().unwrap_or_default().to_owned();
229 name.push(format!(".{}.{serial}.tmp", std::process::id()));
230 let temp = path.with_file_name(name);
231 let written = (|| {
232 let mut file = fs::File::create(&temp)?;
233 file.write_all(bytes)?;
234 file.sync_all()?;
235 fs::rename(&temp, path)
236 })();
237 if written.is_err() {
238 let _ = fs::remove_file(&temp);
239 }
240 written
241}
242
243pub(crate) fn lock_file(
247 path: &Path,
248 timeout: std::time::Duration,
249) -> std::io::Result<Option<fs::File>> {
250 take_lock(path, timeout, false)
251}
252
253pub(crate) fn lock_file_shared(
255 path: &Path,
256 timeout: std::time::Duration,
257) -> std::io::Result<Option<fs::File>> {
258 take_lock(path, timeout, true)
259}
260
261fn take_lock(
262 path: &Path,
263 timeout: std::time::Duration,
264 shared: bool,
265) -> std::io::Result<Option<fs::File>> {
266 use fs2::FileExt;
267
268 let lock = fs::OpenOptions::new()
269 .create(true)
270 .write(true)
271 .truncate(false)
272 .open(path)?;
273 let deadline = std::time::Instant::now() + timeout;
274 loop {
275 let taken = if shared {
276 FileExt::try_lock_shared(&lock)
277 } else {
278 FileExt::try_lock_exclusive(&lock)
279 };
280 if taken.is_ok() {
281 return Ok(Some(lock));
282 }
283 if std::time::Instant::now() >= deadline {
284 return Ok(None);
285 }
286 std::thread::sleep(std::time::Duration::from_millis(2));
287 }
288}
289
290const LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
303
304pub const MAX_RECENTS: usize = 50;
307
308#[derive(Debug, Clone, Copy, PartialEq, Eq)]
319pub enum HistoryUpdate {
320 Written,
322 SkippedBusy,
324}
325
326impl CacheManager {
327 pub fn load_recents(&self) -> Vec<std::path::PathBuf> {
332 self.load_recents_with_visits().0
333 }
334
335 pub fn load_recents_with_visits(
338 &self,
339 ) -> (Vec<PathBuf>, std::collections::HashMap<PathBuf, Visits>) {
340 let mut recents = Vec::new();
341 let mut visits = std::collections::HashMap::new();
342 for line in self.load_history_or_log("recents") {
343 let (path, seen) = parse_recent(&line);
344 recents.push(PathBuf::from(path));
345 if seen.count > 0 {
346 visits.insert(PathBuf::from(path), seen);
347 }
348 }
349 (recents, visits)
350 }
351
352 pub fn load_folds(&self) -> std::collections::HashMap<String, bool> {
357 self.load_history_or_log("home_folds")
358 .into_iter()
359 .filter_map(|line| {
360 let (title, state) = line.rsplit_once('\t')?;
361 Some((title.to_string(), state.trim() == "1"))
362 })
363 .collect()
364 }
365
366 pub fn save_folds(&self, folds: &std::collections::HashMap<String, bool>) {
368 let mut lines: Vec<String> = folds
369 .iter()
370 .map(|(title, folded)| format!("{title}\t{}", if *folded { 1 } else { 0 }))
371 .collect();
372 lines.sort();
373 self.save_history_file("home_folds", &lines)
374 .or_log("save home folds");
375 }
376
377 pub fn terminal_mode(&self, terminal: &str) -> Option<crate::config::ThemeMode> {
380 self.load_history_or_log(TERMINAL_MODES)
381 .iter()
382 .find_map(|line| match line.split_once('\t')? {
383 (key, "dark") if key == terminal => Some(crate::config::ThemeMode::Dark),
384 (key, "light") if key == terminal => Some(crate::config::ThemeMode::Light),
385 _ => None,
386 })
387 }
388
389 pub fn remember_terminal_mode(&self, terminal: &str, mode: crate::config::ThemeMode) {
391 if self.terminal_mode(terminal) == Some(mode) {
392 return;
393 }
394 let word = match mode {
395 crate::config::ThemeMode::Light => "light",
396 _ => "dark",
397 };
398 self.update_history_file(TERMINAL_MODES, |lines| {
399 lines.retain(|line| line.split_once('\t').is_none_or(|(key, _)| key != terminal));
400 lines.push(format!("{terminal}\t{word}"));
401 })
402 .or_log("remember the terminal's background");
403 }
404
405 pub fn forget_recent(&self, path: &std::path::Path) {
411 let target = path.to_string_lossy().into_owned();
412 self.update_history_file("recents", |recents| {
413 recents.retain(|line| parse_recent(line).0 != target);
414 })
415 .or_log("forget a recent");
416 }
417
418 pub fn forget_recents(&self, paths: &[std::path::PathBuf]) {
420 let targets: Vec<String> = paths
421 .iter()
422 .map(|p| p.to_string_lossy().into_owned())
423 .collect();
424 self.update_history_file("recents", |recents| {
425 recents.retain(|line| !targets.iter().any(|t| t == parse_recent(line).0));
426 })
427 .or_log("forget recents");
428 }
429
430 pub fn clear_recents(&self) {
432 self.update_history_file("recents", |recents| recents.clear())
433 .or_log("clear recents");
434 }
435
436 fn recent_is_worth_keeping(path: &str, mounts: &crate::locality::Mounts) -> bool {
458 let path = std::path::Path::new(path);
459 if crate::locality::object_scheme(path).is_some() || mounts.is_network(path) {
463 return true;
464 }
465 match path.parent() {
466 Some(parent) if !parent.as_os_str().is_empty() => parent.exists(),
467 _ => true,
468 }
469 }
470
471 pub fn push_recent(&self, path: &std::path::Path) -> HistoryUpdate {
477 let looks_like_url = path.to_string_lossy().contains("://");
480 let stored = if looks_like_url {
481 path.to_path_buf()
482 } else {
483 crate::canonical::canonicalize(path)
486 .or_else(|e| match crate::members::split(path) {
487 Some((db, table)) => crate::canonical::canonicalize(&db)
488 .map(|db| crate::members::place(&db, &table)),
489 None => Err(e),
490 })
491 .unwrap_or_else(|_| path.to_path_buf())
492 };
493 let entry = stored.to_string_lossy().into_owned();
494
495 let mounts = crate::locality::Mounts::current();
498
499 let now = unix_now();
500 self.update_history_file("recents", |recents| {
501 let mut visits = Visits::default();
502 if let Some(at) = recents.iter().position(|l| parse_recent(l).0 == entry) {
503 visits = parse_recent(&recents.remove(at)).1;
504 }
505 visits.count = visits.count.saturating_add(1);
506 visits.last = now;
507 recents.insert(0, format!("{entry}\t{}\t{}", visits.count, visits.last));
508 recents.retain(|line| Self::recent_is_worth_keeping(parse_recent(line).0, &mounts));
509 recents.truncate(MAX_RECENTS);
510 })
511 .inspect_err(|e| log::warn!(target: "datui", "record a recent: {e:#}"))
512 .unwrap_or(HistoryUpdate::SkippedBusy)
513 }
514
515 pub fn load_visits(&self) -> std::collections::HashMap<PathBuf, Visits> {
517 self.load_recents_with_visits().1
518 }
519}
520
521fn parse_recent(line: &str) -> (&str, Visits) {
524 let mut parts = line.rsplitn(3, '\t');
525 if let (Some(last), Some(count), Some(path)) = (parts.next(), parts.next(), parts.next())
526 && let (Ok(last), Ok(count)) = (last.parse(), count.parse())
527 {
528 return (path, Visits { count, last });
529 }
530 (line, Visits::default())
531}
532
533fn unix_now() -> u64 {
534 std::time::SystemTime::now()
535 .duration_since(std::time::UNIX_EPOCH)
536 .map(|d| d.as_secs())
537 .unwrap_or_default()
538}
539
540#[derive(Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize)]
543pub struct Visits {
544 pub count: u32,
545 pub last: u64,
547}
548
549impl Visits {
550 pub fn frecency(&self, now: u64) -> f64 {
553 let age = now.saturating_sub(self.last);
554 let weight = match age {
555 a if a < 3_600 => 4.0,
556 a if a < 86_400 => 2.0,
557 a if a < 604_800 => 0.5,
558 _ => 0.25,
559 };
560 f64::from(self.count) * weight
561 }
562}
563
564pub fn by_frecency(
567 mut recents: Vec<PathBuf>,
568 visits: &std::collections::HashMap<PathBuf, Visits>,
569) -> Vec<PathBuf> {
570 let now = unix_now();
571 let score = |p: &PathBuf| visits.get(p).map_or(0.0, |v| v.frecency(now));
572 recents.sort_by(|a, b| score(b).total_cmp(&score(a)));
573 recents
574}
575
576#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
584pub struct DatasetFacts {
585 pub mtime: u64,
587 pub size: u64,
589 pub rows: Option<usize>,
590 pub cols: Option<usize>,
591 #[serde(default)]
595 pub cols_sampled: bool,
596 #[serde(default)]
599 pub columns: Vec<String>,
600 #[serde(default)]
604 pub kind: Option<crate::discover::EntryKind>,
605 #[serde(default)]
608 pub classified_by: u32,
609 #[serde(default)]
613 pub cost: crate::discover::Cost,
614 #[serde(default, skip_serializing_if = "crate::discover::Holds::is_empty")]
619 pub holds: crate::discover::Holds,
620}
621
622#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
635pub struct DatasetShape {
636 pub fingerprint: String,
640 pub files: Vec<CachedFooter>,
643 pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
658 pub taken_at: u64,
660}
661
662#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
670pub struct CachedFooter {
671 pub schema: Option<usize>,
675 pub row_group_rows: Vec<usize>,
677 pub row_group_bytes: Vec<usize>,
681 #[serde(default, skip_serializing_if = "Vec::is_empty")]
685 pub column_bytes: Vec<usize>,
686}
687
688impl DatasetShape {
689 pub fn fingerprint_of<'a>(
703 files: impl IntoIterator<Item = (&'a str, u64, u64, Option<&'a str>)>,
704 ) -> String {
705 let mut hasher = StableHasher::default();
706 let mut count = 0usize;
707 let mut bytes = 0u64;
708 for (key, size, stamp, etag) in files {
709 hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
710 match etag {
711 Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
712 None => hasher.u64(0),
713 };
714 count += 1;
715 bytes = bytes.saturating_add(size);
716 }
717 format!("{count}-{bytes}-{:016x}", hasher.finish())
718 }
719
720 pub fn intern_schema(
722 schemas: &mut Vec<Vec<(String, polars::prelude::DataType)>>,
723 schema: &polars::prelude::Schema,
724 ) -> usize {
725 let columns: Vec<(String, polars::prelude::DataType)> = schema
726 .iter()
727 .map(|(name, dtype)| (name.to_string(), dtype.clone()))
728 .collect();
729 schemas
730 .iter()
731 .position(|s| *s == columns)
732 .unwrap_or_else(|| {
733 schemas.push(columns);
734 schemas.len() - 1
735 })
736 }
737
738 pub fn schema_at(
740 schemas: &[Vec<(String, polars::prelude::DataType)>],
741 at: usize,
742 ) -> Option<polars::prelude::Schema> {
743 let columns = schemas.get(at)?;
744 let mut schema = polars::prelude::Schema::with_capacity(columns.len());
745 for (name, dtype) in columns {
746 schema.with_column(name.as_str().into(), dtype.clone());
747 }
748 Some(schema)
749 }
750}
751
752#[cfg(test)]
764pub(crate) fn isolate_cache() {
765 static SCRATCH: std::sync::Mutex<Vec<tempfile::TempDir>> = std::sync::Mutex::new(Vec::new());
768 unsafe extern "C" {
769 fn atexit(callback: extern "C" fn()) -> std::ffi::c_int;
770 }
771 extern "C" fn remove_scratch_dirs() {
772 if let Ok(mut held) = SCRATCH.lock() {
773 held.clear();
774 }
775 }
776 let scratch_dir = |prefix: &str| {
777 tempfile::Builder::new()
778 .prefix(prefix)
779 .tempdir()
780 .expect("a scratch directory for the test process")
781 };
782
783 static ISOLATE: std::sync::Once = std::sync::Once::new();
784 ISOLATE.call_once(|| {
785 let dir = scratch_dir("datui-unit-cache-");
786 let config_dir = scratch_dir("datui-unit-config-");
787 unsafe { std::env::set_var("DATUI_CACHE_DIR", dir.path()) };
790 unsafe { std::env::set_var("DATUI_CONFIG_DIR", config_dir.path()) };
791 let mut held = SCRATCH.lock().unwrap_or_else(|e| e.into_inner());
792 held.push(dir);
793 held.push(config_dir);
794 unsafe { atexit(remove_scratch_dirs) };
797 });
798}
799
800pub(crate) fn running_as_a_cargo_test() -> bool {
805 std::env::current_exe()
806 .is_ok_and(|exe| cargo_test_layout(&exe, std::env::var_os("CARGO_TARGET_DIR").as_deref()))
807}
808
809fn cargo_test_layout(exe: &Path, target_dir: Option<&std::ffi::OsStr>) -> bool {
814 let Some(deps) = exe.parent() else {
815 return false;
816 };
817 if deps.file_name().is_none_or(|name| name != "deps") {
818 return false;
819 }
820 if target_dir.is_some_and(|dir| exe.starts_with(dir)) {
821 return true;
822 }
823 deps.ancestors()
824 .skip(1)
825 .take(3)
826 .any(|dir| dir.file_name().is_some_and(|name| name == "target"))
827}
828
829#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
832pub struct CloudListing {
833 pub fingerprint: String,
836 pub buckets: Vec<String>,
837 pub listed_at: u64,
839}
840
841pub(crate) struct Shapes;
847
848impl Kind for Shapes {
849 const DIR: &'static str = "shapes";
850 const EXT: &'static str = "shape";
851 const VERSION: u16 = 1;
852 const BUDGET: u64 = 128 << 20;
853 type Value = DatasetShape;
854
855 fn encode(shape: &DatasetShape) -> Result<Vec<u8>> {
856 encode_shape(shape)
857 }
858
859 fn decode(payload: &[u8]) -> Option<DatasetShape> {
860 decode_shape(payload)
861 }
862}
863
864#[derive(Debug, Clone, Default, PartialEq, Eq)]
872pub struct FileFooters {
873 pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
875 pub files: Vec<(u64, CachedFooter)>,
877}
878
879pub fn file_identity(key: &str, size: u64, stamp: u64, etag: Option<&str>) -> u64 {
882 let mut hasher = StableHasher::default();
883 hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
884 match etag {
885 Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
886 None => hasher.u64(0),
887 };
888 hasher.finish()
889}
890
891pub(crate) struct FileFootersKind;
894
895impl Kind for FileFootersKind {
896 const DIR: &'static str = "file_footers";
897 const EXT: &'static str = "footers";
898 const VERSION: u16 = 1;
899 const BUDGET: u64 = 64 << 20;
900 type Value = FileFooters;
901
902 fn encode(footers: &FileFooters) -> Result<Vec<u8>> {
903 let shape = DatasetShape {
904 fingerprint: String::new(),
905 files: footers.files.iter().map(|(_, f)| f.clone()).collect(),
906 schemas: footers.schemas.clone(),
907 taken_at: 0,
908 };
909 let mut out = encode_shape(&shape)?;
910 for (identity, _) in &footers.files {
911 out.extend_from_slice(&identity.to_le_bytes());
912 }
913 Ok(out)
914 }
915
916 fn decode(payload: &[u8]) -> Option<FileFooters> {
917 let count = payload.len().checked_sub(4)?;
919 let (len, rest) = payload.split_first_chunk::<4>()?;
920 let header = usize::try_from(u32::from_le_bytes(*len)).ok()?;
921 let mut body = rest.get(header..)?;
922 let files = usize::try_from(take_varint(&mut body)?).ok()?;
923 let ids = files.checked_mul(8)?;
924 if ids > count {
925 return None;
926 }
927 let (shape, identities) = payload.split_at(payload.len() - ids);
928 let shape = decode_shape(shape)?;
929 (shape.files.len() == files).then(|| FileFooters {
930 schemas: shape.schemas,
931 files: identities
932 .as_chunks::<8>()
933 .0
934 .iter()
935 .map(|id| u64::from_le_bytes(*id))
936 .zip(shape.files)
937 .collect(),
938 })
939 }
940}
941
942pub(crate) struct Facts;
945
946impl Kind for Facts {
947 const DIR: &'static str = "facts";
948 const EXT: &'static str = "facts";
949 const VERSION: u16 = 1;
950 const BUDGET: u64 = 16 << 20;
951 type Value = DatasetFacts;
952
953 fn encode(facts: &DatasetFacts) -> Result<Vec<u8>> {
954 Ok(serde_json::to_vec(facts)?)
955 }
956
957 fn decode(payload: &[u8]) -> Option<DatasetFacts> {
958 serde_json::from_slice(payload).ok()
959 }
960}
961
962pub(crate) struct CloudListings;
964
965impl Kind for CloudListings {
966 const DIR: &'static str = "cloud_listings";
967 const EXT: &'static str = "listing";
968 const VERSION: u16 = 1;
969 const BUDGET: u64 = 4 << 20;
970 type Value = CloudListing;
971
972 fn encode(listing: &CloudListing) -> Result<Vec<u8>> {
973 Ok(serde_json::to_vec(listing)?)
974 }
975
976 fn decode(payload: &[u8]) -> Option<CloudListing> {
977 serde_json::from_slice(payload).ok()
978 }
979}
980
981#[derive(serde::Serialize, serde::Deserialize)]
985struct ShapeHeader {
986 fingerprint: String,
987 schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
988 taken_at: u64,
989}
990fn put_varint(out: &mut Vec<u8>, mut n: u64) {
991 while n >= 0x80 {
992 out.push((n as u8) | 0x80);
993 n >>= 7;
994 }
995 out.push(n as u8);
996}
997
998fn take_varint(bytes: &mut &[u8]) -> Option<u64> {
999 let mut n = 0u64;
1000 for shift in (0..64).step_by(7) {
1001 let (&b, rest) = bytes.split_first()?;
1002 *bytes = rest;
1003 n |= u64::from(b & 0x7f) << shift;
1004 if b & 0x80 == 0 {
1005 return Some(n);
1006 }
1007 }
1008 None
1009}
1010
1011fn put_list(out: &mut Vec<u8>, values: &[usize]) {
1012 put_varint(out, values.len() as u64);
1013 for &v in values {
1014 put_varint(out, v as u64);
1015 }
1016}
1017
1018fn take_list(bytes: &mut &[u8]) -> Option<Vec<usize>> {
1019 let len = usize::try_from(take_varint(bytes)?).ok()?;
1020 if len > bytes.len() {
1023 return None;
1024 }
1025 (0..len)
1026 .map(|_| take_varint(bytes).and_then(|v| usize::try_from(v).ok()))
1027 .collect()
1028}
1029
1030fn encode_shape(shape: &DatasetShape) -> Result<Vec<u8>> {
1031 let header = serde_json::to_vec(&ShapeHeader {
1032 fingerprint: shape.fingerprint.clone(),
1033 schemas: shape.schemas.clone(),
1034 taken_at: shape.taken_at,
1035 })?;
1036 let mut out = Vec::with_capacity(8 + header.len() + shape.files.len() * 8);
1037 out.extend_from_slice(&u32::try_from(header.len())?.to_le_bytes());
1038 out.extend_from_slice(&header);
1039 put_varint(&mut out, shape.files.len() as u64);
1040 for file in &shape.files {
1041 put_varint(&mut out, file.schema.map_or(0, |s| s as u64 + 1));
1042 put_list(&mut out, &file.row_group_rows);
1043 put_list(&mut out, &file.row_group_bytes);
1044 put_list(&mut out, &file.column_bytes);
1045 }
1046 Ok(out)
1047}
1048
1049fn decode_shape(bytes: &[u8]) -> Option<DatasetShape> {
1051 let (len, rest) = bytes.split_first_chunk::<4>()?;
1052 let len = usize::try_from(u32::from_le_bytes(*len)).ok()?;
1053 let (header, mut body) = (rest.get(..len)?, rest.get(len..)?);
1054 let header: ShapeHeader = serde_json::from_slice(header).ok()?;
1055 let count = usize::try_from(take_varint(&mut body)?).ok()?;
1056 if count > body.len() {
1057 return None;
1058 }
1059 let mut files = Vec::with_capacity(count);
1060 let (mut rows, mut bytes_total) = (0usize, 0usize);
1063 for _ in 0..count {
1064 let schema = match take_varint(&mut body)? {
1065 0 => None,
1066 at => Some(usize::try_from(at - 1).ok()?),
1067 };
1068 let footer = CachedFooter {
1069 schema,
1070 row_group_rows: take_list(&mut body)?,
1071 row_group_bytes: take_list(&mut body)?,
1072 column_bytes: take_list(&mut body)?,
1073 };
1074 if let Some(at) = footer.schema
1075 && at >= header.schemas.len()
1076 {
1077 return None;
1078 }
1079 for &n in &footer.row_group_rows {
1080 rows = rows.checked_add(n)?;
1081 }
1082 for &n in footer.row_group_bytes.iter().chain(&footer.column_bytes) {
1083 bytes_total = bytes_total.checked_add(n)?;
1084 }
1085 files.push(footer);
1086 }
1087 body.is_empty().then_some(DatasetShape {
1088 fingerprint: header.fingerprint,
1089 files,
1090 schemas: header.schemas,
1091 taken_at: header.taken_at,
1092 })
1093}
1094
1095impl CacheManager {
1096 pub fn dataset_shapes_kept(&self) -> usize {
1098 Store::<Shapes>::new(self).len()
1099 }
1100
1101 pub fn dataset_shape(&self, path: &str, fingerprint: &str) -> Option<DatasetShape> {
1109 Store::<Shapes>::new(self).get(path, fingerprint)
1110 }
1111
1112 pub fn has_dataset_shape(&self, path: &str) -> bool {
1115 Store::<Shapes>::new(self).file(path).exists()
1116 }
1117
1118 pub fn save_dataset_shape(&self, path: &str, shape: DatasetShape) {
1120 Store::<Shapes>::new(self).put(path, &shape.fingerprint, &shape);
1121 }
1122
1123 pub fn file_footers(&self, path: &str) -> Option<FileFooters> {
1125 Store::<FileFootersKind>::new(self).get(path, "")
1126 }
1127
1128 pub fn save_file_footers(&self, path: &str, footers: &FileFooters) {
1130 Store::<FileFootersKind>::new(self).put(path, "", footers);
1131 }
1132
1133 pub fn cloud_listing(&self, id: &str, fingerprint: &str) -> Option<CloudListing> {
1135 Store::<CloudListings>::new(self).get(id, fingerprint)
1136 }
1137
1138 pub fn load_cloud_listings(&self) -> std::collections::HashMap<String, CloudListing> {
1140 Store::<CloudListings>::new(self)
1141 .scan()
1142 .into_iter()
1143 .collect()
1144 }
1145
1146 pub fn save_cloud_listing(&self, id: &str, listing: CloudListing) {
1148 Store::<CloudListings>::new(self).put(id, &listing.fingerprint.clone(), &listing);
1149 }
1150
1151 pub fn load_hidden_cloud_sources(&self) -> Vec<String> {
1153 self.load_history_or_log("cloud_hidden")
1154 }
1155
1156 pub fn hide_cloud_source(&self, id: &str) {
1158 let id = id.to_string();
1159 self.update_history_file("cloud_hidden", |hidden| {
1160 if !hidden.contains(&id) {
1161 hidden.push(id.clone());
1162 }
1163 })
1164 .or_log("hide a cloud source");
1165 }
1166
1167 pub fn examples_hidden(&self) -> bool {
1170 !self.load_history_or_log("examples_hidden").is_empty()
1171 }
1172
1173 pub fn hide_examples(&self) {
1176 self.save_history_file("examples_hidden", &["hidden".to_string()])
1177 .or_log("hide the example datasets");
1178 }
1179
1180 pub fn load_remembered_places(&self) -> Vec<PathBuf> {
1183 let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1184 if !file.exists() {
1185 return Vec::new();
1186 }
1187 self.load_history_or_log("home_remembered")
1188 .into_iter()
1189 .map(PathBuf::from)
1190 .collect()
1191 }
1192
1193 pub fn clear_remembered_places(&self) {
1195 let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1196 if let Err(e) = std::fs::remove_file(&file)
1197 && e.kind() != std::io::ErrorKind::NotFound
1198 {
1199 log::warn!(target: "datui", "remove {}: {e}", file.display());
1200 }
1201 }
1202
1203 pub fn save_remembered_places(&self, places: &[PathBuf]) -> Result<()> {
1205 let places: Vec<String> = places
1206 .iter()
1207 .map(|p| p.to_string_lossy().into_owned())
1208 .collect();
1209 self.save_history_file("home_remembered", &places)
1210 }
1211}
1212
1213impl CacheManager {
1214 pub fn load_dataset_facts(&self) -> std::collections::HashMap<PathBuf, DatasetFacts> {
1218 Store::<Facts>::new(self)
1219 .scan()
1220 .into_iter()
1221 .map(|(path, facts)| (PathBuf::from(path), facts))
1222 .collect()
1223 }
1224
1225 pub fn dataset_facts(&self, path: &Path) -> Option<DatasetFacts> {
1227 Store::<Facts>::new(self).get(path.to_str()?, "")
1228 }
1229
1230 pub fn touch_dataset_facts<'a>(&self, paths: impl IntoIterator<Item = &'a Path>) {
1233 let store = Store::<Facts>::new(self);
1234 for path in paths.into_iter().filter_map(Path::to_str) {
1235 store.touch(path);
1236 }
1237 }
1238
1239 pub fn record_dataset_facts(&self, facts: &[(PathBuf, DatasetFacts)]) {
1242 Store::<Facts>::new(self).put_all(
1243 facts
1244 .iter()
1245 .filter_map(|(path, facts)| Some((path.to_str()?, "", facts))),
1246 );
1247 }
1248
1249 fn with_cache_lock<F>(&self, name: &str, work: F) -> Result<()>
1252 where
1253 F: FnOnce() -> Result<()>,
1254 {
1255 use fs2::FileExt;
1256
1257 self.ensure_cache_dir()?;
1258 let Some(lock) = lock_file(&self.cache_file(&format!("{name}.lock")), LOCK_TIMEOUT)? else {
1259 log::info!(target: "datui", "{name} cache not updated: its lock is busy");
1260 return Ok(());
1261 };
1262
1263 let result = work();
1264 let _ = FileExt::unlock(&lock);
1265 result
1266 }
1267}
1268
1269#[cfg(test)]
1270mod harness_tests {
1271 #[test]
1277 fn a_cargo_test_binary_is_recognized() {
1278 assert!(super::running_as_a_cargo_test());
1279 }
1280
1281 #[test]
1285 fn a_unit_test_needs_no_setup_to_isolate() {
1286 let cache = super::CacheManager::new(crate::APP_NAME).unwrap();
1287 let config = crate::config::ConfigManager::new(crate::APP_NAME).unwrap();
1288 let scratch = std::env::temp_dir();
1289 assert!(cache.cache_dir().starts_with(&scratch), "{cache:?}");
1290 let config = config.config_dir();
1291 assert!(config.starts_with(&scratch), "{config:?}");
1292 }
1293
1294 #[test]
1297 fn only_cargos_layout_is_a_test() {
1298 use std::ffi::OsStr;
1299 use std::path::Path;
1300 let layout = |exe: &str| super::cargo_test_layout(Path::new(exe), None);
1301 assert!(layout(
1302 "/home/x/src/datui/target/debug/deps/home_test-1a2b3c"
1303 ));
1304 assert!(layout(
1305 "/home/x/src/datui/target/x86_64-unknown-linux-gnu/release/deps/datui-1a2b"
1306 ));
1307 assert!(!layout("/home/x/src/datui/target/debug/datui"));
1308 assert!(!layout("/opt/deps/bin/datui"));
1309 assert!(!layout("/home/x/deps/datui-0.4/bin/datui"));
1310 assert!(!layout("/home/x/src/datui/target/debug/examples/demo"));
1311 assert!(!layout("/home/x/build/datui/debug/deps/home_test-1a2b3c"));
1313 assert!(super::cargo_test_layout(
1314 Path::new("/home/x/build/datui/debug/deps/home_test-1a2b3c"),
1315 Some(OsStr::new("/home/x/build/datui"))
1316 ));
1317 assert!(!super::cargo_test_layout(
1318 Path::new("/home/x/build/deps/datui"),
1319 Some(OsStr::new("/home/x/other"))
1320 ));
1321 }
1322}
1323
1324#[cfg(test)]
1325mod recents_pruning_tests {
1326 use super::*;
1327
1328 fn cache() -> (CacheManager, tempfile::TempDir) {
1329 let dir = tempfile::tempdir().expect("temp dir");
1330 (CacheManager::with_dir(dir.path().to_path_buf()), dir)
1331 }
1332
1333 #[test]
1335 fn a_terminals_last_answer_is_kept_per_terminal() {
1336 use crate::config::ThemeMode;
1337 let (cache, _keep) = cache();
1338 assert_eq!(cache.terminal_mode("WezTerm"), None);
1339 cache.remember_terminal_mode("WezTerm", ThemeMode::Light);
1340 cache.remember_terminal_mode("tmux", ThemeMode::Dark);
1341 assert_eq!(cache.terminal_mode("WezTerm"), Some(ThemeMode::Light));
1342 assert_eq!(cache.terminal_mode("tmux"), Some(ThemeMode::Dark));
1343 cache.remember_terminal_mode("WezTerm", ThemeMode::Dark);
1344 assert_eq!(cache.terminal_mode("WezTerm"), Some(ThemeMode::Dark));
1345 assert_eq!(cache.terminal_mode("tmux"), Some(ThemeMode::Dark));
1346 assert_eq!(cache.terminal_mode(""), None);
1347 }
1348
1349 #[test]
1350 fn a_recent_whose_directory_is_gone_is_forgotten() {
1351 let (cache, _keep) = cache();
1355 let scratch = tempfile::tempdir().expect("scratch");
1356 let dataset = scratch.path().join("people.csv");
1357 std::fs::write(&dataset, b"a,b\n1,2\n").expect("write");
1358 let dataset = crate::canonical::canonicalize(&dataset).expect("canonicalize");
1361
1362 cache.push_recent(&dataset);
1363 assert!(cache.load_recents().iter().any(|p| p == &dataset));
1364
1365 let elsewhere = tempfile::tempdir().expect("elsewhere");
1368 let survivor = elsewhere.path().join("still-here.csv");
1369 std::fs::write(&survivor, b"a\n1\n").expect("write");
1370
1371 drop(scratch);
1373
1374 cache.push_recent(&survivor);
1377 let recents = cache.load_recents();
1378 assert!(
1379 !recents.iter().any(|p| p == &dataset),
1380 "the dead path should be gone; got {recents:?}"
1381 );
1382 assert!(
1383 recents
1384 .iter()
1385 .any(|p| p.file_name() == survivor.file_name()),
1386 "the live path should remain; got {recents:?}"
1387 );
1388 }
1389
1390 #[test]
1391 fn a_deleted_file_in_a_directory_that_still_exists_is_kept() {
1392 let (cache, _keep) = cache();
1396 let scratch = tempfile::tempdir().expect("scratch");
1397 let dataset = scratch.path().join("nightly.parquet");
1398 std::fs::write(&dataset, b"x").expect("write");
1399 let dataset = crate::canonical::canonicalize(&dataset).expect("canonicalize");
1402 cache.push_recent(&dataset);
1403
1404 std::fs::remove_file(&dataset).expect("remove");
1405 let other = scratch.path().join("other.csv");
1406 std::fs::write(&other, b"a\n1\n").expect("write");
1407 cache.push_recent(&other);
1408
1409 assert!(
1410 cache.load_recents().iter().any(|p| p == &dataset),
1411 "a missing file in a live directory should stay"
1412 );
1413 }
1414
1415 #[test]
1416 fn a_remote_recent_is_never_stated_let_alone_dropped() {
1417 let (cache, _keep) = cache();
1420 let scratch = tempfile::tempdir().expect("scratch");
1421 let local = scratch.path().join("local.csv");
1422 std::fs::write(&local, b"a\n1\n").expect("write");
1423
1424 for url in [
1425 "s3://bucket/warehouse/events.parquet",
1426 "gs://bucket/data.csv",
1427 "https://example.com/data.csv",
1428 ] {
1429 cache.push_recent(std::path::Path::new(url));
1430 }
1431 cache.push_recent(&local);
1432
1433 let recents = cache.load_recents();
1434 for url in [
1435 "s3://bucket/warehouse/events.parquet",
1436 "gs://bucket/data.csv",
1437 "https://example.com/data.csv",
1438 ] {
1439 assert!(
1440 recents.iter().any(|p| p.to_string_lossy() == url),
1441 "{url} should have survived; got {recents:?}"
1442 );
1443 }
1444 }
1445}
1446
1447#[cfg(test)]
1448mod dataset_shape_tests {
1449 use super::*;
1450
1451 fn shape(fingerprint: &str, taken_at: u64) -> DatasetShape {
1452 DatasetShape {
1453 fingerprint: fingerprint.to_string(),
1454 files: vec![
1455 CachedFooter {
1456 schema: Some(0),
1457 row_group_rows: vec![10],
1458 row_group_bytes: vec![1_000],
1459 column_bytes: Vec::new(),
1460 },
1461 CachedFooter {
1462 schema: Some(0),
1463 row_group_rows: vec![20],
1464 row_group_bytes: vec![2_000],
1465 column_bytes: Vec::new(),
1466 },
1467 ],
1468 schemas: vec![vec![("id".into(), polars::prelude::DataType::Int64)]],
1469 taken_at,
1470 }
1471 }
1472
1473 #[test]
1476 fn a_shape_that_does_not_hold_together_is_refused() {
1477 let good = shape("fp", 1);
1478 assert_eq!(
1479 decode_shape(&encode_shape(&good).unwrap()),
1480 Some(good.clone())
1481 );
1482 let mut wild = good.clone();
1483 wild.files[0].schema = Some(5);
1484 assert!(decode_shape(&encode_shape(&wild).unwrap()).is_none());
1485 let mut huge = good.clone();
1486 huge.files[0].row_group_rows = vec![usize::MAX, 1];
1487 assert!(decode_shape(&encode_shape(&huge).unwrap()).is_none());
1488 let bytes = encode_shape(&good).unwrap();
1489 for cut in 0..bytes.len() {
1490 assert!(decode_shape(&bytes[..cut]).is_none(), "cut {cut}");
1491 }
1492 }
1493
1494 #[test]
1501 fn a_shape_comes_back_only_for_the_dataset_it_was_taken_from() {
1502 let dir = tempfile::tempdir().unwrap();
1503 let cache = super::CacheManager::with_dir(dir.path().to_path_buf());
1504 cache.save_dataset_shape("s3://b/events/", shape("2-30-abc", 100));
1505
1506 assert_eq!(
1507 cache.dataset_shape("s3://b/events/", "2-30-abc"),
1508 Some(shape("2-30-abc", 100)),
1509 "the same dataset, unchanged"
1510 );
1511 assert_eq!(
1512 cache.dataset_shape("s3://b/events/", "3-40-def"),
1513 None,
1514 "a file added, removed or rewritten since"
1515 );
1516 assert_eq!(
1517 cache.dataset_shape("s3://b/other/", "2-30-abc"),
1518 None,
1519 "and a different dataset that happens to weigh the same"
1520 );
1521 }
1522
1523 #[test]
1525 fn the_fingerprint_moves_when_the_files_do() {
1526 let base = DatasetShape::fingerprint_of([
1527 ("a.parquet", 100, 7, Some("e1")),
1528 ("b.parquet", 200, 8, Some("e2")),
1529 ]);
1530 assert_eq!(
1531 base,
1532 DatasetShape::fingerprint_of([
1533 ("a.parquet", 100, 7, Some("e1")),
1534 ("b.parquet", 200, 8, Some("e2")),
1535 ]),
1536 "the same listing twice is the same fingerprint"
1537 );
1538 for (changed, why) in [
1539 (vec![("a.parquet", 100, 7, Some("e1"))], "a file removed"),
1540 (
1541 vec![
1542 ("a.parquet", 100, 7, Some("e1")),
1543 ("b.parquet", 200, 8, Some("e2")),
1544 ("c.parquet", 50, 9, Some("e3")),
1545 ],
1546 "a file added",
1547 ),
1548 (
1549 vec![
1550 ("a.parquet", 100, 7, Some("e1")),
1551 ("b.parquet", 201, 8, Some("e2")),
1552 ],
1553 "a file resized",
1554 ),
1555 (
1556 vec![
1557 ("a.parquet", 100, 7, Some("e1")),
1558 ("b.parquet", 200, 9, Some("e2")),
1559 ],
1560 "a file rewritten, which the stamp catches",
1561 ),
1562 (
1563 vec![
1564 ("a.parquet", 100, 7, Some("e1")),
1565 ("b.parquet", 200, 8, Some("e9")),
1566 ],
1567 "rewritten within the same second at the same length: only the tag sees it",
1568 ),
1569 (
1570 vec![
1571 ("a.parquet", 100, 7, Some("e1")),
1572 ("b.parquet", 200, 8, None),
1573 ],
1574 "a tag gone",
1575 ),
1576 (
1577 vec![
1578 ("a.parquet", 100, 7, Some("e1")),
1579 ("renamed.parquet", 200, 8, Some("e2")),
1580 ],
1581 "a file renamed, which reorders the positional join the cache is",
1582 ),
1583 ] {
1584 assert_ne!(base, DatasetShape::fingerprint_of(changed), "{why}");
1585 }
1586 }
1587
1588 #[test]
1591 fn the_stable_hash_is_pinned() {
1592 assert_eq!(stable_hash(b"123456789"), 0x995d_c9bb_df19_39fa);
1593 assert_eq!(
1594 DatasetShape::fingerprint_of([("a.parquet", 100, 7, Some("e1"))]),
1595 "1-100-7bb7c5da2f222965"
1596 );
1597 }
1598
1599 #[test]
1602 fn a_shape_comes_back_as_it_was_stored() {
1603 let dir = tempfile::tempdir().unwrap();
1604 let cache = super::CacheManager::with_dir(dir.path().to_path_buf());
1605 let mut stored = shape("f", 7);
1606 stored.files.push(CachedFooter {
1607 schema: None,
1608 row_group_rows: Vec::new(),
1609 row_group_bytes: Vec::new(),
1610 column_bytes: Vec::new(),
1611 });
1612 stored.files.push(CachedFooter {
1613 schema: Some(1),
1614 row_group_rows: vec![0, 300, u32::MAX as usize + 5],
1615 row_group_bytes: vec![1, 2, 3],
1616 column_bytes: vec![128, 1 << 40],
1617 });
1618 stored.schemas.push(vec![
1619 ("id".into(), polars::prelude::DataType::Int64),
1620 (
1621 "tags".into(),
1622 polars::prelude::DataType::List(Box::new(polars::prelude::DataType::String)),
1623 ),
1624 ]);
1625 cache.save_dataset_shape("s3://b/events/", stored.clone());
1626 assert_eq!(cache.dataset_shape("s3://b/events/", "f"), Some(stored));
1627 }
1628
1629 #[test]
1632 #[ignore = "a timing, not a check"]
1633 fn shape_cache_timings() {
1634 use polars::prelude::DataType;
1635 const FILES: usize = 842_000;
1636 let dir = tempfile::tempdir().unwrap();
1637 let cache = super::CacheManager::with_dir(dir.path().to_path_buf());
1638 let big = DatasetShape {
1639 fingerprint: "big".into(),
1640 files: (0..FILES)
1641 .map(|i| CachedFooter {
1642 schema: Some(0),
1643 row_group_rows: vec![1_000 + (i * 7919) % 90_000],
1644 row_group_bytes: vec![10_000 + (i * 104_729) % 900_000],
1645 column_bytes: Vec::new(),
1646 })
1647 .collect(),
1648 schemas: vec![vec![
1649 ("ID".into(), DataType::String),
1650 ("DATE".into(), DataType::String),
1651 ("ELEMENT".into(), DataType::String),
1652 ("DATA_VALUE".into(), DataType::Int32),
1653 ("M_FLAG".into(), DataType::String),
1654 ("Q_FLAG".into(), DataType::String),
1655 ("S_FLAG".into(), DataType::String),
1656 ("OBS_TIME".into(), DataType::String),
1657 ]],
1658 taken_at: 1,
1659 };
1660 let time = |what: &str, work: &mut dyn FnMut()| {
1661 let began = std::time::Instant::now();
1662 work();
1663 eprintln!("{what}: {:.1?}", began.elapsed());
1664 };
1665 time("store by_station", &mut || {
1666 cache.save_dataset_shape("s3://noaa/by_station/", big.clone())
1667 });
1668 time("store a small one beside it", &mut || {
1669 cache.save_dataset_shape("s3://b/small/", shape("small", 2))
1670 });
1671 time("lookup + touch by_station", &mut || {
1672 assert!(
1673 cache
1674 .dataset_shape("s3://noaa/by_station/", "big")
1675 .is_some()
1676 )
1677 });
1678 time("lookup by_station, fingerprint moved", &mut || {
1679 assert!(
1680 cache
1681 .dataset_shape("s3://noaa/by_station/", "moved")
1682 .is_none()
1683 )
1684 });
1685 time("lookup + touch the small one", &mut || {
1686 assert!(cache.dataset_shape("s3://b/small/", "small").is_some())
1687 });
1688 let on_disk = Store::<Shapes>::new(&cache)
1689 .dir()
1690 .read_dir()
1691 .unwrap()
1692 .flatten()
1693 .map(|e| e.metadata().unwrap().len())
1694 .sum::<u64>();
1695 eprintln!("on disk: {:.1} MB", on_disk as f64 / 1e6);
1696 }
1697}
1698
1699#[cfg(test)]
1701mod store_harness_tests {
1702 use super::*;
1703 use std::time::{Duration, SystemTime, UNIX_EPOCH};
1704
1705 trait Sample: Kind + Send + Sync + 'static {
1708 fn sample(variant: u8) -> Self::Value;
1709 fn big() -> Self::Value;
1710 }
1711
1712 impl Sample for Shapes {
1713 fn sample(variant: u8) -> DatasetShape {
1714 DatasetShape {
1715 fingerprint: "fp".into(),
1716 files: vec![CachedFooter {
1717 schema: Some(0),
1718 row_group_rows: vec![usize::from(variant)],
1719 row_group_bytes: vec![100],
1720 column_bytes: vec![8],
1721 }],
1722 schemas: vec![vec![("id".into(), polars::prelude::DataType::Int64)]],
1723 taken_at: 1,
1724 }
1725 }
1726 fn big() -> DatasetShape {
1727 let mut shape = Self::sample(1);
1728 shape.files = vec![shape.files[0].clone(); 2_000];
1729 shape
1730 }
1731 }
1732
1733 impl Sample for Facts {
1734 fn sample(variant: u8) -> DatasetFacts {
1735 DatasetFacts {
1736 mtime: 1,
1737 size: 2,
1738 rows: Some(usize::from(variant)),
1739 cols: Some(1),
1740 columns: vec!["a".into()],
1741 ..Default::default()
1742 }
1743 }
1744 fn big() -> DatasetFacts {
1745 DatasetFacts {
1746 columns: (0..2_000).map(|i| format!("column_{i}")).collect(),
1747 ..Self::sample(1)
1748 }
1749 }
1750 }
1751
1752 impl Sample for CloudListings {
1753 fn sample(variant: u8) -> CloudListing {
1754 CloudListing {
1755 fingerprint: "fp".into(),
1756 buckets: vec![format!("bucket-{variant}")],
1757 listed_at: 1,
1758 }
1759 }
1760 fn big() -> CloudListing {
1761 CloudListing {
1762 buckets: (0..2_000).map(|i| format!("bucket-{i}")).collect(),
1763 ..Self::sample(1)
1764 }
1765 }
1766 }
1767
1768 fn encoded<K: Kind>(value: &K::Value) -> Vec<u8> {
1769 K::encode(value).unwrap()
1770 }
1771
1772 fn same<K: Kind>(a: Option<K::Value>, b: &K::Value) -> bool {
1773 a.is_some_and(|a| encoded::<K>(&a) == encoded::<K>(b))
1774 }
1775
1776 fn store<K: Kind>() -> (Store<K>, tempfile::TempDir) {
1777 let dir = tempfile::tempdir().unwrap();
1778 (
1779 Store::new(&CacheManager::with_dir(dir.path().to_path_buf())),
1780 dir,
1781 )
1782 }
1783
1784 fn age(file: &Path, secs: u64) {
1785 fs::OpenOptions::new()
1786 .write(true)
1787 .open(file)
1788 .unwrap()
1789 .set_modified(UNIX_EPOCH + Duration::from_secs(secs))
1790 .unwrap();
1791 }
1792
1793 fn modified(file: &Path) -> SystemTime {
1794 fs::metadata(file).unwrap().modified().unwrap()
1795 }
1796
1797 fn round_trip<K: Sample>() {
1798 let (store, _dir) = store::<K>();
1799 let value = K::sample(1);
1800 store.put("s3://b/a/", "fp", &value);
1801 assert!(same::<K>(store.get("s3://b/a/", "fp"), &value), "a hit");
1802 assert!(
1803 store.get("s3://b/a/", "moved").is_none(),
1804 "another fingerprint"
1805 );
1806 assert!(store.get("s3://b/z/", "fp").is_none(), "another key");
1807 let scanned = store.scan();
1808 assert_eq!(scanned.len(), 1);
1809 assert_eq!(scanned[0].0, "s3://b/a/");
1810 let other = K::sample(2);
1812 fs::copy(store.file("s3://b/a/"), store.file("s3://b/z/")).unwrap();
1813 assert!(
1814 store.get("s3://b/z/", "fp").is_none(),
1815 "a collision is a miss"
1816 );
1817 store.put("s3://b/a/", "fp", &other);
1818 assert!(same::<K>(store.get("s3://b/a/", "fp"), &other), "replaced");
1819 }
1820
1821 fn damage_is_a_miss<K: Sample>() {
1822 let (store, _dir) = store::<K>();
1823 store.put("k", "fp", &K::sample(1));
1824 let file = store.file("k");
1825 let good = fs::read(&file).unwrap();
1826 for cut in 0..good.len() {
1827 fs::write(&file, &good[..cut]).unwrap();
1828 assert!(store.get("k", "fp").is_none(), "cut at {cut}");
1829 }
1830 for at in 0..good.len() {
1831 for bit in 0..8 {
1832 let mut bad = good.clone();
1833 bad[at] ^= 1 << bit;
1834 fs::write(&file, &bad).unwrap();
1835 assert!(store.get("k", "fp").is_none(), "bit {bit} of byte {at}");
1836 }
1837 }
1838 assert!(store.scan().is_empty(), "nor does a scan see it");
1839 fs::write(&file, &good).unwrap();
1840 assert!(
1841 store.get("k", "fp").is_some(),
1842 "and the good bytes still read"
1843 );
1844 }
1845
1846 fn evicts_by_bytes_in_lru_order<K: Sample>() {
1847 let (store, _dir) = store::<K>();
1848 for (i, key) in ["k/a", "k/b", "k/c"].into_iter().enumerate() {
1849 store.put(key, "fp", &K::sample(1));
1850 age(&store.file(key), 1_000 + i as u64);
1851 }
1852 let one = fs::metadata(store.file("k/a")).unwrap().len();
1853 assert!(store.get("k/a", "fp").is_some());
1855
1856 let store = store.with_budget(3 * one);
1857 store.put("k/d", "fp", &K::sample(1));
1858 assert_eq!(store.len(), 3, "kept to its budget in bytes");
1859 assert!(store.get("k/b", "fp").is_none(), "b went");
1860 for kept in ["k/a", "k/c", "k/d"] {
1861 assert!(store.get(kept, "fp").is_some(), "{kept} stayed");
1862 }
1863
1864 store.put("k/huge", "fp", &K::big());
1866 assert_eq!(store.len(), 1, "everything else made room");
1867 assert!(store.get("k/huge", "fp").is_some());
1868 }
1869
1870 fn a_hit_only_touches<K: Sample>() {
1871 let (store, _dir) = store::<K>();
1872 store.put("k/a", "fp", &K::sample(1));
1873 store.put("k/b", "fp", &K::sample(1));
1874 let (a, b) = (store.file("k/a"), store.file("k/b"));
1875 age(&a, 1_000);
1876 age(&b, 1_000);
1877 let before = fs::read(&a).unwrap();
1878
1879 assert!(store.get("k/a", "fp").is_some());
1880 assert!(
1881 modified(&a) > UNIX_EPOCH + Duration::from_secs(1_000),
1882 "dated now"
1883 );
1884 assert_eq!(fs::read(&a).unwrap(), before, "and not rewritten");
1885 assert_eq!(
1886 modified(&b),
1887 UNIX_EPOCH + Duration::from_secs(1_000),
1888 "b untouched"
1889 );
1890
1891 age(&a, 1_000);
1892 assert!(store.get("k/a", "moved").is_none());
1893 assert_eq!(
1894 modified(&a),
1895 UNIX_EPOCH + Duration::from_secs(1_000),
1896 "a miss is not use"
1897 );
1898 assert_eq!(store.scan().len(), 2);
1899 assert_eq!(
1900 modified(&a),
1901 UNIX_EPOCH + Duration::from_secs(1_000),
1902 "nor is a scan"
1903 );
1904
1905 let inode = fs::metadata(&a).unwrap();
1907 store.put("k/a", "fp", &K::sample(1));
1908 assert!(modified(&a) > UNIX_EPOCH + Duration::from_secs(1_000));
1909 #[cfg(unix)]
1910 {
1911 use std::os::unix::fs::MetadataExt;
1912 assert_eq!(fs::metadata(&a).unwrap().ino(), inode.ino(), "not replaced");
1913 }
1914 let _ = inode;
1915 }
1916
1917 fn sweeps_stale_temp_files<K: Sample>() {
1918 let (store, _dir) = store::<K>();
1919 store.put("k/a", "fp", &K::sample(1));
1920 let stale = store.dir().join("x.shape.1.0.tmp");
1921 let fresh = store.dir().join("x.shape.1.1.tmp");
1922 fs::write(&stale, b"half").unwrap();
1923 fs::write(&fresh, b"half").unwrap();
1924 age(&stale, 1_000);
1925 store.put("k/b", "fp", &K::sample(1));
1926 assert!(!stale.exists(), "a dead writer's temp file goes");
1927 assert!(fresh.exists(), "a live writer's stays");
1928 }
1929
1930 fn clear_all_removes_it<K: Sample>() {
1931 let dir = tempfile::tempdir().unwrap();
1932 let cache = CacheManager::with_dir(dir.path().to_path_buf());
1933 let store = Store::<K>::new(&cache);
1934 store.put("k", "fp", &K::sample(1));
1935 assert!(store.get("k", "fp").is_some());
1936 cache.clear_all().unwrap();
1937 assert!(store.get("k", "fp").is_none());
1938 assert!(!store.dir().exists(), "the kind's directory is gone");
1939 assert!(
1940 cache.cache_file(&format!("{}.lock", K::DIR)).exists(),
1941 "its lock, which another instance may hold, is not"
1942 );
1943 }
1944
1945 fn two_writers_race<K: Sample>() {
1946 let (store, _dir) = store::<K>();
1947 let store = std::sync::Arc::new(store);
1948 let writers: Vec<_> = (1..=2u8)
1949 .map(|variant| {
1950 let store = store.clone();
1951 std::thread::spawn(move || {
1952 for i in 0..25 {
1953 store.put("shared", "fp", &K::sample(variant));
1954 store.put(&format!("own/{variant}/{i}"), "fp", &K::sample(variant));
1955 assert!(store.get("shared", "fp").is_some(), "never torn");
1956 }
1957 })
1958 })
1959 .collect();
1960 for writer in writers {
1961 writer.join().unwrap();
1962 }
1963 let shared = encoded::<K>(&store.get("shared", "fp").unwrap());
1964 assert!(
1965 shared == encoded::<K>(&K::sample(1)) || shared == encoded::<K>(&K::sample(2)),
1966 "one writer's value, whole"
1967 );
1968 assert_eq!(store.len(), 51, "every entry landed");
1969 let temps = fs::read_dir(store.dir())
1970 .unwrap()
1971 .flatten()
1972 .filter(|e| e.path().extension().is_some_and(|x| x == "tmp"))
1973 .count();
1974 assert_eq!(temps, 0, "and no temp file is left");
1975 }
1976
1977 macro_rules! suite {
1978 ($name:ident, $kind:ty) => {
1979 mod $name {
1980 use super::*;
1981 #[test]
1982 fn round_trip() {
1983 super::round_trip::<$kind>();
1984 }
1985 #[test]
1986 fn damage_is_a_miss() {
1987 super::damage_is_a_miss::<$kind>();
1988 }
1989 #[test]
1990 fn evicts_by_bytes_in_lru_order() {
1991 super::evicts_by_bytes_in_lru_order::<$kind>();
1992 }
1993 #[test]
1994 fn a_hit_only_touches() {
1995 super::a_hit_only_touches::<$kind>();
1996 }
1997 #[test]
1998 fn sweeps_stale_temp_files() {
1999 super::sweeps_stale_temp_files::<$kind>();
2000 }
2001 #[test]
2002 fn clear_all_removes_it() {
2003 super::clear_all_removes_it::<$kind>();
2004 }
2005 #[test]
2006 fn two_writers_race() {
2007 super::two_writers_race::<$kind>();
2008 }
2009 }
2010 };
2011 }
2012
2013 suite!(shapes, Shapes);
2014 suite!(facts, Facts);
2015 suite!(cloud_listings, CloudListings);
2016
2017 #[test]
2019 fn the_old_files_are_removed() {
2020 let dir = tempfile::tempdir().unwrap();
2021 let cache = CacheManager::with_dir(dir.path().to_path_buf());
2022 let old = [
2023 "datasets.json",
2024 "dataset_shapes.json",
2025 "cloud_sources.json",
2026 "visits.json",
2027 ];
2028 for name in old {
2029 fs::write(dir.path().join(name), b"{}").unwrap();
2030 }
2031 fs::create_dir(dir.path().join("dataset_shapes")).unwrap();
2032 fs::write(dir.path().join("dataset_shapes/0.shape"), b"x").unwrap();
2033 cache.record_dataset_facts(&[(PathBuf::from("/d"), Facts::sample(1))]);
2034 for name in old {
2035 assert!(!dir.path().join(name).exists(), "{name}");
2036 }
2037 assert!(!dir.path().join("dataset_shapes").exists());
2038 }
2039}
2040
2041#[cfg(test)]
2042mod facts_compat_tests {
2043 use super::DatasetFacts;
2044 use crate::discover::EntryKind;
2045
2046 #[test]
2054 fn an_unknown_kind_costs_only_its_own_row() {
2055 let json = r#"{"mtime":1,"size":2,"rows":3,"cols":4,"columns":[],"kind":"quicksand"}"#;
2056 let facts: DatasetFacts = serde_json::from_str(json).expect("the record still parses");
2057 assert_eq!(
2058 facts.kind,
2059 Some(EntryKind::Unknown),
2060 "a kind from the future reads as unexamined"
2061 );
2062 assert_eq!(facts.rows, Some(3), "and the measurements survive with it");
2063 assert_eq!(facts.classified_by, 0, "recorded before that existed");
2064
2065 let index = r#"{"/a":{"mtime":1,"size":2,"rows":3,"cols":4,"columns":[],"kind":"quicksand"},
2067 "/b":{"mtime":1,"size":2,"rows":9,"cols":1,"columns":[],"kind":"hive"}}"#;
2068 let map: std::collections::HashMap<std::path::PathBuf, DatasetFacts> =
2069 serde_json::from_str(index).expect("the index still parses");
2070 assert_eq!(map.len(), 2, "both rows, not none of them");
2071 }
2072
2073 #[test]
2076 fn clear_all_removes_the_caches_and_lists_only() {
2077 let dir = tempfile::tempdir().expect("a temp dir");
2078 let cache = super::CacheManager::with_dir(dir.path().to_path_buf());
2079 for name in [
2080 "query_history.txt",
2081 "fuzzy_history.txt",
2082 "recents_history.txt",
2083 "recents_history.txt.12.0.tmp",
2084 "cloud_hidden_history.txt",
2085 ] {
2086 std::fs::write(dir.path().join(name), b"x").expect("write");
2087 }
2088 cache.record_dataset_facts(&[("/d".into(), DatasetFacts::default())]);
2089 let kept = [
2090 "facts.lock",
2091 "recents_history.lock",
2092 "datui.log",
2093 "keep.parquet",
2094 "subdir",
2095 ];
2096 for name in kept {
2097 if name == "subdir" {
2098 std::fs::create_dir(dir.path().join(name)).expect("mkdir");
2099 } else {
2100 std::fs::write(dir.path().join(name), b"x").expect("write");
2101 }
2102 }
2103
2104 cache.clear_all().expect("clear");
2105
2106 let mut left: Vec<String> = std::fs::read_dir(dir.path())
2107 .expect("read dir")
2108 .flatten()
2109 .map(|e| e.file_name().to_string_lossy().into_owned())
2110 .collect();
2111 left.sort();
2112 let mut kept = kept.map(String::from).to_vec();
2113 kept.sort();
2114 assert_eq!(left, kept);
2115 }
2116}