Skip to main content

datui_lib/
cache.rs

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/// Manages cache directory and cache file operations
12#[derive(Clone, Debug)]
13pub struct CacheManager {
14    pub(crate) cache_dir: PathBuf,
15    /// Each kind's bytes on disk as its last sweep found them, plus what this session
16    /// wrote since; shared by clones. See [`store::Store::put_all`].
17    pub(crate) swept:
18        std::sync::Arc<std::sync::Mutex<std::collections::HashMap<&'static str, u64>>>,
19}
20
21impl CacheManager {
22    /// Create a CacheManager rooted at an explicit directory (primarily for testing).
23    pub fn with_dir(cache_dir: PathBuf) -> Self {
24        Self {
25            cache_dir,
26            swept: Default::default(),
27        }
28    }
29
30    /// Create a CacheManager for `app_name`. `DATUI_CACHE_DIR` overrides the location;
31    /// the test suite sets it so opens never record fixtures in the developer's recents.
32    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        // A test reaching the real cache would write fixtures into the developer's recents:
39        // refuse. Test binaries live under `target/<profile>/deps/`; the real binary never.
40        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    /// Get the cache directory path
57    pub fn cache_dir(&self) -> &Path {
58        &self.cache_dir
59    }
60
61    /// Get path to a specific cache file
62    pub fn cache_file(&self, filename: &str) -> PathBuf {
63        self.cache_dir.join(filename)
64    }
65
66    /// Ensure the cache directory exists
67    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    /// Remove every cache kind's directory and line-file list (recents, histories,
75    /// hidden sources, remembered places). Not lock files (another instance may hold
76    /// them), the log, Data Quality copies (owned by their sessions), or views (config).
77    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            // The list, and any temp file a writer of it left.
92            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    /// Load a history file. An unreadable file is an error, never an empty list (the next
101    /// push would write it back and lose every entry); a non-UTF-8 line is skipped alone.
102    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    /// A history file, empty when it is missing or cannot be read; the latter is logged.
126    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    /// Apply `update` to a history file with the read-modify-write under an exclusive
133    /// lock: atomic writes prevent corruption but not lost updates from instances
134    /// opening datasets at once. The lock is awaited up to `LOCK_TIMEOUT`, far past
135    /// realistic contention, then the update is dropped: history never delays what the
136    /// user asked for.
137    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        // Read, modify and write inside the lock. An unreadable file is left alone:
151        // rewriting it from nothing would lose every entry.
152        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        // Released explicitly, though dropping the file would do it too.
157        let _ = FileExt::unlock(&lock);
158        result.map(|()| HistoryUpdate::Written)
159    }
160
161    /// Save history to a history file
162    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        // Oldest first, but we keep the most recent entries.
167        let mut text = String::new();
168        for entry in history {
169            text.push_str(entry);
170            text.push('\n');
171        }
172        // Not truncate-in-place: readers would see it half-written, and concurrent writers
173        // would interleave.
174        atomic_write(&history_file, text.as_bytes())?;
175        Ok(())
176    }
177}
178
179/// What a line-file list's name ends in: `recents_history.txt`.
180const HISTORY_SUFFIX: &str = "_history.txt";
181
182/// The line-file of each terminal's last answer about its background.
183const TERMINAL_MODES: &str = "terminal_modes";
184
185/// Write `bytes` to `path` via a synced sibling temp file renamed over it, so readers,
186/// other instances and crashes see the old file or the new, never part. Every cache
187/// and view write goes through here. The temp name (pid plus a process-wide counter,
188/// ending `.tmp`) is unique per write; sweeps remove stale ones.
189pub(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
207/// Take an exclusive lock on `path` (created if missing), waiting up to `timeout`;
208/// `None` if it stayed busy. Held until the returned file drops.
209pub(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
216/// [`lock_file`], shared: readers hold it together, and never while a writer does.
217pub(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
253/// How long to wait for another instance rewriting a history file: a bound on a
254/// wedged peer, not something realistic contention reaches. Generous because giving
255/// up drops an entry and Windows locks and renames are slow (250ms failed with
256/// sixteen writers); nothing waits on this write.
257const LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
258
259/// Maximum recent paths kept: a few days of work, never needing pagination.
260pub const MAX_RECENTS: usize = 50;
261
262/// Whether a history update happened. Contended updates are abandoned rather than
263/// awaited, so "no error" and "written" differ; callers and tests can tell exactly.
264#[derive(Debug, Clone, Copy, PartialEq, Eq)]
265pub enum HistoryUpdate {
266    /// The lock was taken and the new contents are on disk.
267    Written,
268    /// The lock stayed busy past the deadline, so nothing was written.
269    SkippedBusy,
270}
271
272impl CacheManager {
273    /// Recently opened dataset paths, most recent first: the only thing datui remembers
274    /// about your data between runs, and deleting it loses only ordering.
275    pub fn load_recents(&self) -> Vec<std::path::PathBuf> {
276        self.load_recents_with_visits().0
277    }
278
279    /// Recents, most recent first, with each one's visits, in one read.
280    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    /// Home sections the user folded (`title<TAB>1`) or opened (`title<TAB>0`), by title;
296    /// unlisted sections take their default.
297    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    /// Remember the fold state. Failing to write it loses nothing but a preference.
308    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    /// The background mode `terminal` last reported, for the next start's first frame
319    /// under `theme.mode = "auto"`. One line per terminal: `key<TAB>dark`.
320    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    /// Remember what `terminal` answered, when it differs from what is remembered.
331    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    /// Forget one recent path: an editable recents list is one people trust.
347    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    /// Forget several recently opened paths at once: every recent under one place.
356    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    /// Forget every recently opened path, leaving other caches alone.
368    pub fn clear_recents(&self) {
369        self.update_history_file("recents", |recents| recents.clear())
370            .or_log("clear recents");
371    }
372
373    /// Whether a recorded path is still worth offering. Judged by its directory: a file
374    /// missing from an existing directory may be regenerated in place, but a deleted
375    /// directory takes its recents with it (else a dead root lingers until fifty opens
376    /// push it off). Remote paths are never checked: stat'ing one can hang, and a down
377    /// share is when its recents matter most.
378    fn recent_is_worth_keeping(path: &str, mounts: &crate::home::locality::Mounts) -> bool {
379        let path = std::path::Path::new(path);
380        // The mount table is passed in: reading it per entry meant fifty reads per open.
381        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    /// Record a path as most recently opened, deduplicated and capped. Failures are
391    /// ignored: a convenience list must never hinder opening data. The result tells a
392    /// write from a contended skip.
393    pub fn push_recent(&self, path: &std::path::Path) -> HistoryUpdate {
394        // A URL is recorded as given: canonicalizing is meaningless and would stat a
395        // nonlocal path.
396        let looks_like_url = path.to_string_lossy().contains("://");
397        let stored = if looks_like_url {
398            path.to_path_buf()
399        } else {
400            // A table inside a file of tables has no file of its own: the file is made absolute,
401            // the name kept.
402            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        // One mount table read for the whole prune; kernel-generated, so it cannot block.
413        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
432/// One recents line: `path<TAB>opens<TAB>last opened`, visits kept with the path so
433/// forgetting it forgets them. A bare path has none.
434fn 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/// How often and how lately a dataset was opened: zoxide-style frecency, ranking
452/// Recent and lifting often-opened matches.
453#[derive(Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize)]
454pub struct Visits {
455    pub count: u32,
456    /// Seconds since the epoch.
457    pub last: u64,
458}
459
460impl Visits {
461    /// Opens weighted by recency: ×4 within the hour, ×2 within the day, ×½ within the
462    /// week, ×¼ after.
463    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
475/// `recents` reordered by frecency; ties and pre-visit-count recents keep their order.
476pub 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/// What datui remembers about a measured dataset: a cache, not a catalog. Every field
487/// is re-derivable, and the recorded size and mtime invalidate a changed dataset.
488#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
489pub struct DatasetFacts {
490    /// Modification time in seconds since the epoch, as a fingerprint.
491    pub mtime: u64,
492    /// Size in bytes, the other half of the fingerprint.
493    pub size: u64,
494    pub rows: Option<usize>,
495    pub cols: Option<usize>,
496    /// Whether `cols` is a floor (a sampled large directory), restored with the count so
497    /// a sample is never shown as a total.
498    #[serde(default)]
499    pub cols_sampled: bool,
500    /// Column names, enabling search by column before anything is read this run.
501    #[serde(default)]
502    pub columns: Vec<String>,
503    /// What the dataset turned out to be: recorded since a remote path cannot be
504    /// classified without reading, and a row should read the same in every section.
505    #[serde(default)]
506    pub kind: Option<crate::home::discover::EntryKind>,
507    /// The build rules `kind` came from (see [`crate::home::discover::CLASSIFIER_VERSION`]); `0`
508    /// in records older than this field.
509    #[serde(default)]
510    pub classified_by: u32,
511    /// What opening costs (compression, layout, partitioning), from a footer read a
512    /// remote dataset may not get twice.
513    #[serde(default)]
514    pub cost: crate::home::discover::Cost,
515    /// What one listing of the directory found (its label), restored beside `kind` under
516    /// the same classifier version.
517    #[serde(
518        default,
519        skip_serializing_if = "crate::home::discover::Holds::is_empty"
520    )]
521    pub holds: crate::home::discover::Holds,
522}
523
524/// What an open learned about a dataset's files, so the next open shows columns and
525/// row count without reading footers (seconds for thousands of objects). The listing,
526/// which happens anyway, decides via the fingerprint whether this still holds. A cache,
527/// not a catalog, like [`DatasetFacts`].
528#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
529pub struct DatasetShape {
530    /// What the files looked like when taken; a listing with a different one describes a
531    /// changed dataset, and the rest is ignored.
532    pub fingerprint: String,
533    /// One entry per file, in listing order, which is scan order.
534    pub files: Vec<CachedFooter>,
535    /// The distinct schemas, by index, since thousands of files usually share one. Types
536    /// are Polars' own serialization (names cannot rebuild nested types); if a Polars
537    /// upgrade changes it, entries stop parsing and the cache is just empty.
538    pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
539    /// Seconds since the Unix epoch, for a human reading the file.
540    pub taken_at: u64,
541}
542
543/// What one file's footer said, as much as a reopen needs. It comes back as a
544/// `FileFooter`, so a cached dataset goes through the same code as a fresh read.
545#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
546pub struct CachedFooter {
547    /// Which [`DatasetShape::schemas`] entry this file has; `None` for an unreadable
548    /// footer, remembered as such so a reopen gives the same dataset.
549    pub schema: Option<usize>,
550    /// The rows in each of its row groups, in order.
551    pub row_group_rows: Vec<usize>,
552    /// Compressed bytes of each row group, kept because a note uses them, and a note
553    /// present only on first open is a worse bug than a slow open.
554    pub row_group_bytes: Vec<usize>,
555    /// Uncompressed bytes of each schema column, in order: local footers carry them (the
556    /// only source for binary widths); cloud footers leave this empty.
557    #[serde(default, skip_serializing_if = "Vec::is_empty")]
558    pub column_bytes: Vec<usize>,
559}
560
561impl DatasetShape {
562    /// A listing's fingerprint: file names, count, sizes, mtimes and store tags, all seen
563    /// without opening anything. Names are included since remembered footers join to a
564    /// fresh listing by position; ETags since size and whole-second mtime miss
565    /// same-length rewrites within a second.
566    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    /// Where `schema` sits in `schemas`, added on the end if it is not there yet.
585    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    /// The schema at `at` in `schemas`, or `None` when the table has no such entry.
603    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/// Point the cache and config at this test process's own scratch directories, once
617/// per process (`CacheManager::new` and `ConfigManager::new` call it). Named at random,
618/// not by pid (reused pids inherited old recents), and removed at exit.
619#[cfg(test)]
620pub(crate) fn isolate_cache() {
621    // Held for the process's life; a static never drops, so an exit handler removes them.
622    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        // SAFETY: test-only. Tests run on parallel threads, so this can race another test
643        // reading the environment; accepted in tests and never done outside them.
644        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        // SAFETY: the C runtime's `atexit`, present on every platform std runs on; the
650        // callback only drops the directories above.
651        unsafe { atexit(remove_scratch_dirs) };
652    });
653}
654
655/// Whether this process is a cargo-built test binary (`target/<profile>/deps/…`, or
656/// `target/<triple>/<profile>/deps/` cross-built); the program itself never is.
657pub(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
662/// The directory is `deps` under a `target` at most three levels up, or under cargo's
663/// configured target dir, so a program installed in some `deps` is not mistaken.
664fn 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/// A cloud source's buckets from an earlier run, shown at once while a fresh listing
681/// is out.
682#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
683pub struct CloudListing {
684    /// What the source pointed at when listed; a mismatch means another server, ignored.
685    pub fingerprint: String,
686    pub buckets: Vec<String>,
687    /// Seconds since the Unix epoch.
688    pub listed_at: u64,
689}
690
691/// Dataset shapes by path, at their listing's fingerprint. Bounded by bytes: shape
692/// size follows file count (842k files is megabytes), so a count bound would let a
693/// few huge datasets evict a small one in use.
694pub(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/// The footers a dataset's count has read, each by [`file_identity`], so a recount
713/// reads only new files and a stopped count keeps its progress. By file, unlike
714/// [`DatasetShape`], so it stays right for a changed dataset.
715#[derive(Debug, Clone, Default, PartialEq, Eq)]
716pub struct FileFooters {
717    /// The distinct schemas, as in [`DatasetShape::schemas`].
718    pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
719    /// Each file's identity, and its footer.
720    pub files: Vec<(u64, CachedFooter)>,
721}
722
723/// A file's identity for [`FileFooters`]: key, size, mtime and store tag, any of which
724/// changes on rewrite.
725pub 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
735/// [`FileFooters`] by dataset path; empty fingerprint, as each file checks itself.
736pub(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        // The identities follow the shape, eight bytes a file.
761        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
785/// What home measured of each dataset, by path; records carry their own size and
786/// mtime, so the store fingerprint is empty.
787pub(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
805/// Each cloud source's last bucket listing, by source ID, at the source's fingerprint.
806pub(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/// A shape's header, read before its footers: fingerprint, schemas, time. JSON for the
825/// Polars-serialized schemas; the bulk per-file footers follow as varints.
826#[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    // Each value is at least a byte: a longer length means a broken file, not an
863    // allocation.
864    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
891/// The shape in `bytes`, or `None` for one that does not hold together.
892fn 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    // Dataset totals must fit, so no later sum overflows on a damaged file.
903    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    /// How many dataset shapes are kept.
938    #[cfg(test)]
939    pub fn dataset_shapes_kept(&self) -> usize {
940        Store::<Shapes>::new(self).len()
941    }
942
943    /// The remembered shape of a dataset if `fingerprint` still matches (a mismatched
944    /// entry has no correct use). A hit counts as use, so unchanging datasets stay.
945    pub fn dataset_shape(&self, path: &str, fingerprint: &str) -> Option<DatasetShape> {
946        Store::<Shapes>::new(self).get(path, fingerprint)
947    }
948
949    /// Whether any shape is kept for `path`: one stat.
950    pub fn has_dataset_shape(&self, path: &str) -> bool {
951        Store::<Shapes>::new(self).file(path).exists()
952    }
953
954    /// Remember one dataset's shape, keeping the others while they fit the budget.
955    pub fn save_dataset_shape(&self, path: &str, shape: DatasetShape) {
956        Store::<Shapes>::new(self).put(path, &shape.fingerprint, &shape);
957    }
958
959    /// The footers counts of the dataset at `path` have read, by file.
960    pub fn file_footers(&self, path: &str) -> Option<FileFooters> {
961        Store::<FileFootersKind>::new(self).get(path, "")
962    }
963
964    /// Remember the footers a count of the dataset at `path` read, by file.
965    pub fn save_file_footers(&self, path: &str, footers: &FileFooters) {
966        Store::<FileFootersKind>::new(self).put(path, "", footers);
967    }
968
969    /// The last listing of cloud source `id`, if it was taken at `fingerprint`.
970    pub fn cloud_listing(&self, id: &str, fingerprint: &str) -> Option<CloudListing> {
971        Store::<CloudListings>::new(self).get(id, fingerprint)
972    }
973
974    /// Record one source's listing, keeping the others.
975    pub fn save_cloud_listing(&self, id: &str, listing: CloudListing) {
976        Store::<CloudListings>::new(self).put(id, &listing.fingerprint.clone(), &listing);
977    }
978
979    /// Source IDs hidden from the home screen with Delete.
980    pub fn load_hidden_cloud_sources(&self) -> Vec<String> {
981        self.load_history_or_log("cloud_hidden")
982    }
983
984    /// Hide a source from the home screen until the cache is cleared.
985    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    /// Whether Delete on its heading hid the bundled Example datasets (a user's
996    /// `examples.toml` still shows).
997    pub fn examples_hidden(&self) -> bool {
998        !self.load_history_or_log("examples_hidden").is_empty()
999    }
1000
1001    /// Hide the bundled Example datasets until the cache is cleared.
1002    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    /// The directories Ctrl+D kept here before 0.4.0 (now in `catalog.toml`), in added
1008    /// order.
1009    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    /// Drop the list [`Self::load_remembered_places`] reads, once it is moved.
1021    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    /// Write `places` where Ctrl+D kept them before 0.4.0, for the migration's tests.
1031    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    /// Every remembered dataset fact, not counted as use. Unreadable records are absent:
1042    /// failing to read a cache must never be worse than lacking it.
1043    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    /// What datui knows about one dataset, counted as use.
1052    pub fn dataset_facts(&self, path: &Path) -> Option<DatasetFacts> {
1053        Store::<Facts>::new(self).get(path.to_str()?, "")
1054    }
1055
1056    /// Count these records as used (shown on home); least recently used go first past
1057    /// the budget.
1058    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    /// Record newly measured datasets, one file each; non-UTF-8 paths are skipped.
1066    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    /// Run `work` under the named cache lock, or skip it if contended past the deadline.
1075    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/// The same checks for every kind the [`Store`] holds.
1103#[cfg(test)]
1104mod store_harness_tests;
1105
1106#[cfg(test)]
1107mod facts_compat_tests;