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}
16
17impl CacheManager {
18    /// Create a CacheManager rooted at an explicit directory (primarily for testing).
19    pub fn with_dir(cache_dir: PathBuf) -> Self {
20        Self { cache_dir }
21    }
22
23    /// Create a CacheManager for `app_name`. `DATUI_CACHE_DIR` overrides the location;
24    /// the test suite sets it so opens never record fixtures in the developer's recents.
25    pub fn new(app_name: &str) -> Result<Self> {
26        #[cfg(test)]
27        isolate_cache();
28        if let Some(dir) = std::env::var_os("DATUI_CACHE_DIR") {
29            return Ok(Self {
30                cache_dir: PathBuf::from(dir),
31            });
32        }
33        // A test reaching the real cache would write fixtures into the developer's recents:
34        // refuse. Test binaries live under `target/<profile>/deps/`; the real binary never.
35        if running_as_a_cargo_test() {
36            panic!(
37                "DATUI_CACHE_DIR is not set: a test would write to the real cache. \
38                 Call common::isolate_cache() (or take the runtime from \
39                 common::test_runtime(), which does) before building an App or a \
40                 CacheManager."
41            );
42        }
43
44        let cache_dir = dirs::cache_dir()
45            .ok_or_else(|| color_eyre::eyre::eyre!("Could not determine cache directory"))?
46            .join(app_name);
47
48        Ok(Self { cache_dir })
49    }
50
51    /// Get the cache directory path
52    pub fn cache_dir(&self) -> &Path {
53        &self.cache_dir
54    }
55
56    /// Get path to a specific cache file
57    pub fn cache_file(&self, filename: &str) -> PathBuf {
58        self.cache_dir.join(filename)
59    }
60
61    /// Ensure the cache directory exists
62    pub fn ensure_cache_dir(&self) -> Result<()> {
63        if !self.cache_dir.exists() {
64            fs::create_dir_all(&self.cache_dir)?;
65        }
66        Ok(())
67    }
68
69    /// Remove every cache kind's directory and line-file list (recents, histories,
70    /// hidden sources, remembered places). Not lock files (another instance may hold
71    /// them), the log, Data Quality copies (owned by their sessions), or views (config).
72    pub fn clear_all(&self) -> Result<()> {
73        for dir in [Shapes::DIR, Facts::DIR, CloudListings::DIR] {
74            match fs::remove_dir_all(self.cache_file(dir)) {
75                Err(e) if e.kind() != std::io::ErrorKind::NotFound => {
76                    log::warn!(target: "datui", "remove the {dir} cache: {e}");
77                }
78                _ => {}
79            }
80        }
81        let Ok(entries) = fs::read_dir(&self.cache_dir) else {
82            return Ok(());
83        };
84        for entry in entries.flatten() {
85            let path = entry.path();
86            // The list, and any temp file a writer of it left.
87            let list = entry.file_name().to_string_lossy().contains(HISTORY_SUFFIX);
88            if list && path.is_file() {
89                fs::remove_file(&path).or_log(&format!("remove {}", path.display()));
90            }
91        }
92        Ok(())
93    }
94
95    /// Load a history file. An unreadable file is an error, never an empty list (the next
96    /// push would write it back and lose every entry); a non-UTF-8 line is skipped alone.
97    pub fn load_history_file(&self, history_id: &str) -> Result<Vec<String>> {
98        let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
99
100        let bytes = match fs::read(&history_file) {
101            Ok(bytes) => bytes,
102            Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
103            Err(e) => return Err(e.into()),
104        };
105        let mut history = Vec::new();
106        for line in bytes.split(|&b| b == b'\n') {
107            let line = line.strip_suffix(b"\r").unwrap_or(line);
108            match std::str::from_utf8(line) {
109                Ok(line) if !line.trim().is_empty() => history.push(line.to_string()),
110                Ok(_) => {}
111                Err(e) => {
112                    log::warn!(target: "datui", "{history_id} history: skipped a line: {e}")
113                }
114            }
115        }
116
117        Ok(history)
118    }
119
120    /// A history file, empty when it is missing or cannot be read; the latter is logged.
121    fn load_history_or_log(&self, history_id: &str) -> Vec<String> {
122        self.load_history_file(history_id)
123            .inspect_err(|e| log::warn!(target: "datui", "read {history_id} history: {e:#}"))
124            .unwrap_or_default()
125    }
126
127    /// Apply `update` to a history file with the read-modify-write under an exclusive
128    /// lock: atomic writes prevent corruption but not lost updates from instances
129    /// opening datasets at once. The lock is awaited up to `LOCK_TIMEOUT`, far past
130    /// realistic contention, then the update is dropped: history never delays what the
131    /// user asked for.
132    pub fn update_history_file<F>(&self, history_id: &str, update: F) -> Result<HistoryUpdate>
133    where
134        F: FnOnce(&mut Vec<String>),
135    {
136        use fs2::FileExt;
137
138        self.ensure_cache_dir()?;
139        let lock_path = self.cache_file(&format!("{}_history.lock", history_id));
140        let Some(lock) = lock_file(&lock_path, LOCK_TIMEOUT)? else {
141            log::info!(target: "datui", "{history_id} history not updated: its lock is busy");
142            return Ok(HistoryUpdate::SkippedBusy);
143        };
144
145        // Read, modify and write inside the lock. An unreadable file is left alone:
146        // rewriting it from nothing would lose every entry.
147        let mut entries = self.load_history_file(history_id)?;
148        update(&mut entries);
149        let result = self.save_history_file(history_id, &entries);
150
151        // Released explicitly, though dropping the file would do it too.
152        let _ = FileExt::unlock(&lock);
153        result.map(|()| HistoryUpdate::Written)
154    }
155
156    /// Save history to a history file
157    pub fn save_history_file(&self, history_id: &str, history: &[String]) -> Result<()> {
158        self.ensure_cache_dir()?;
159        let history_file = self.cache_file(&format!("{history_id}{HISTORY_SUFFIX}"));
160
161        // Oldest first, but we keep the most recent entries.
162        let mut text = String::new();
163        for entry in history {
164            text.push_str(entry);
165            text.push('\n');
166        }
167        // Not truncate-in-place: readers would see it half-written, and concurrent writers
168        // would interleave.
169        atomic_write(&history_file, text.as_bytes())?;
170        Ok(())
171    }
172}
173
174/// What a line-file list's name ends in: `recents_history.txt`.
175const HISTORY_SUFFIX: &str = "_history.txt";
176
177/// The line-file of each terminal's last answer about its background.
178const TERMINAL_MODES: &str = "terminal_modes";
179
180/// Write `bytes` to `path` via a synced sibling temp file renamed over it, so readers,
181/// other instances and crashes see the old file or the new, never part. Every cache
182/// and view write goes through here. The temp name (pid plus a process-wide counter,
183/// ending `.tmp`) is unique per write; sweeps remove stale ones.
184pub(crate) fn atomic_write(path: &Path, bytes: &[u8]) -> std::io::Result<()> {
185    static SERIAL: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
186    let serial = SERIAL.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
187    let mut name = path.file_name().unwrap_or_default().to_owned();
188    name.push(format!(".{}.{serial}.tmp", std::process::id()));
189    let temp = path.with_file_name(name);
190    let written = (|| {
191        let mut file = fs::File::create(&temp)?;
192        file.write_all(bytes)?;
193        file.sync_all()?;
194        fs::rename(&temp, path)
195    })();
196    if written.is_err() {
197        let _ = fs::remove_file(&temp);
198    }
199    written
200}
201
202/// Take an exclusive lock on `path` (created if missing), waiting up to `timeout`;
203/// `None` if it stayed busy. Held until the returned file drops.
204pub(crate) fn lock_file(
205    path: &Path,
206    timeout: std::time::Duration,
207) -> std::io::Result<Option<fs::File>> {
208    take_lock(path, timeout, false)
209}
210
211/// [`lock_file`], shared: readers hold it together, and never while a writer does.
212pub(crate) fn lock_file_shared(
213    path: &Path,
214    timeout: std::time::Duration,
215) -> std::io::Result<Option<fs::File>> {
216    take_lock(path, timeout, true)
217}
218
219fn take_lock(
220    path: &Path,
221    timeout: std::time::Duration,
222    shared: bool,
223) -> std::io::Result<Option<fs::File>> {
224    use fs2::FileExt;
225
226    let lock = fs::OpenOptions::new()
227        .create(true)
228        .write(true)
229        .truncate(false)
230        .open(path)?;
231    let deadline = std::time::Instant::now() + timeout;
232    loop {
233        let taken = if shared {
234            FileExt::try_lock_shared(&lock)
235        } else {
236            FileExt::try_lock_exclusive(&lock)
237        };
238        if taken.is_ok() {
239            return Ok(Some(lock));
240        }
241        if std::time::Instant::now() >= deadline {
242            return Ok(None);
243        }
244        std::thread::sleep(std::time::Duration::from_millis(2));
245    }
246}
247
248/// How long to wait for another instance rewriting a history file: a bound on a
249/// wedged peer, not something realistic contention reaches. Generous because giving
250/// up drops an entry and Windows locks and renames are slow (250ms failed with
251/// sixteen writers); nothing waits on this write.
252const LOCK_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(2);
253
254/// Maximum recent paths kept: a few days of work, never needing pagination.
255pub const MAX_RECENTS: usize = 50;
256
257/// Whether a history update happened. Contended updates are abandoned rather than
258/// awaited, so "no error" and "written" differ; callers and tests can tell exactly.
259#[derive(Debug, Clone, Copy, PartialEq, Eq)]
260pub enum HistoryUpdate {
261    /// The lock was taken and the new contents are on disk.
262    Written,
263    /// The lock stayed busy past the deadline, so nothing was written.
264    SkippedBusy,
265}
266
267impl CacheManager {
268    /// Recently opened dataset paths, most recent first: the only thing datui remembers
269    /// about your data between runs, and deleting it loses only ordering.
270    pub fn load_recents(&self) -> Vec<std::path::PathBuf> {
271        self.load_recents_with_visits().0
272    }
273
274    /// Recents, most recent first, with each one's visits, in one read.
275    pub fn load_recents_with_visits(
276        &self,
277    ) -> (Vec<PathBuf>, std::collections::HashMap<PathBuf, Visits>) {
278        let mut recents = Vec::new();
279        let mut visits = std::collections::HashMap::new();
280        for line in self.load_history_or_log("recents") {
281            let (path, seen) = parse_recent(&line);
282            recents.push(PathBuf::from(path));
283            if seen.count > 0 {
284                visits.insert(PathBuf::from(path), seen);
285            }
286        }
287        (recents, visits)
288    }
289
290    /// Home sections the user folded (`title<TAB>1`) or opened (`title<TAB>0`), by title;
291    /// unlisted sections take their default.
292    pub fn load_folds(&self) -> std::collections::HashMap<String, bool> {
293        self.load_history_or_log("home_folds")
294            .into_iter()
295            .filter_map(|line| {
296                let (title, state) = line.rsplit_once('\t')?;
297                Some((title.to_string(), state.trim() == "1"))
298            })
299            .collect()
300    }
301
302    /// Remember the fold state. Failing to write it loses nothing but a preference.
303    pub fn save_folds(&self, folds: &std::collections::HashMap<String, bool>) {
304        let mut lines: Vec<String> = folds
305            .iter()
306            .map(|(title, folded)| format!("{title}\t{}", if *folded { 1 } else { 0 }))
307            .collect();
308        lines.sort();
309        self.save_history_file("home_folds", &lines)
310            .or_log("save home folds");
311    }
312
313    /// The background mode `terminal` last reported, for the next start's first frame
314    /// under `theme.mode = "auto"`. One line per terminal: `key<TAB>dark`.
315    pub fn terminal_mode(&self, terminal: &str) -> Option<crate::config::ThemeMode> {
316        self.load_history_or_log(TERMINAL_MODES)
317            .iter()
318            .find_map(|line| match line.split_once('\t')? {
319                (key, "dark") if key == terminal => Some(crate::config::ThemeMode::Dark),
320                (key, "light") if key == terminal => Some(crate::config::ThemeMode::Light),
321                _ => None,
322            })
323    }
324
325    /// Remember what `terminal` answered, when it differs from what is remembered.
326    pub fn remember_terminal_mode(&self, terminal: &str, mode: crate::config::ThemeMode) {
327        if self.terminal_mode(terminal) == Some(mode) {
328            return;
329        }
330        let word = match mode {
331            crate::config::ThemeMode::Light => "light",
332            _ => "dark",
333        };
334        self.update_history_file(TERMINAL_MODES, |lines| {
335            lines.retain(|line| line.split_once('\t').is_none_or(|(key, _)| key != terminal));
336            lines.push(format!("{terminal}\t{word}"));
337        })
338        .or_log("remember the terminal's background");
339    }
340
341    /// Forget one recent path: an editable recents list is one people trust.
342    pub fn forget_recent(&self, path: &std::path::Path) {
343        let target = path.to_string_lossy().into_owned();
344        self.update_history_file("recents", |recents| {
345            recents.retain(|line| parse_recent(line).0 != target);
346        })
347        .or_log("forget a recent");
348    }
349
350    /// Forget several recently opened paths at once: every recent under one place.
351    pub fn forget_recents(&self, paths: &[std::path::PathBuf]) {
352        let targets: Vec<String> = paths
353            .iter()
354            .map(|p| p.to_string_lossy().into_owned())
355            .collect();
356        self.update_history_file("recents", |recents| {
357            recents.retain(|line| !targets.iter().any(|t| t == parse_recent(line).0));
358        })
359        .or_log("forget recents");
360    }
361
362    /// Forget every recently opened path, leaving other caches alone.
363    pub fn clear_recents(&self) {
364        self.update_history_file("recents", |recents| recents.clear())
365            .or_log("clear recents");
366    }
367
368    /// Whether a recorded path is still worth offering. Judged by its directory: a file
369    /// missing from an existing directory may be regenerated in place, but a deleted
370    /// directory takes its recents with it (else a dead root lingers until fifty opens
371    /// push it off). Remote paths are never checked: stat'ing one can hang, and a down
372    /// share is when its recents matter most.
373    fn recent_is_worth_keeping(path: &str, mounts: &crate::home::locality::Mounts) -> bool {
374        let path = std::path::Path::new(path);
375        // The mount table is passed in: reading it per entry meant fifty reads per open.
376        if crate::home::locality::object_scheme(path).is_some() || mounts.is_network(path) {
377            return true;
378        }
379        match path.parent() {
380            Some(parent) if !parent.as_os_str().is_empty() => parent.exists(),
381            _ => true,
382        }
383    }
384
385    /// Record a path as most recently opened, deduplicated and capped. Failures are
386    /// ignored: a convenience list must never hinder opening data. The result tells a
387    /// write from a contended skip.
388    pub fn push_recent(&self, path: &std::path::Path) -> HistoryUpdate {
389        // A URL is recorded as given: canonicalizing is meaningless and would stat a
390        // nonlocal path.
391        let looks_like_url = path.to_string_lossy().contains("://");
392        let stored = if looks_like_url {
393            path.to_path_buf()
394        } else {
395            // A table inside a file of tables has no file of its own: the file is made absolute,
396            // the name kept.
397            crate::canonical::canonicalize(path)
398                .or_else(|e| match crate::formats::members::split(path) {
399                    Some((db, table)) => crate::canonical::canonicalize(&db)
400                        .map(|db| crate::formats::members::place(&db, &table)),
401                    None => Err(e),
402                })
403                .unwrap_or_else(|_| path.to_path_buf())
404        };
405        let entry = stored.to_string_lossy().into_owned();
406
407        // One mount table read for the whole prune; kernel-generated, so it cannot block.
408        let mounts = crate::home::locality::Mounts::current();
409
410        let now = unix_now();
411        self.update_history_file("recents", |recents| {
412            let mut visits = Visits::default();
413            if let Some(at) = recents.iter().position(|l| parse_recent(l).0 == entry) {
414                visits = parse_recent(&recents.remove(at)).1;
415            }
416            visits.count = visits.count.saturating_add(1);
417            visits.last = now;
418            recents.insert(0, format!("{entry}\t{}\t{}", visits.count, visits.last));
419            recents.retain(|line| Self::recent_is_worth_keeping(parse_recent(line).0, &mounts));
420            recents.truncate(MAX_RECENTS);
421        })
422        .inspect_err(|e| log::warn!(target: "datui", "record a recent: {e:#}"))
423        .unwrap_or(HistoryUpdate::SkippedBusy)
424    }
425}
426
427/// One recents line: `path<TAB>opens<TAB>last opened`, visits kept with the path so
428/// forgetting it forgets them. A bare path has none.
429fn parse_recent(line: &str) -> (&str, Visits) {
430    let mut parts = line.rsplitn(3, '\t');
431    if let (Some(last), Some(count), Some(path)) = (parts.next(), parts.next(), parts.next())
432        && let (Ok(last), Ok(count)) = (last.parse(), count.parse())
433    {
434        return (path, Visits { count, last });
435    }
436    (line, Visits::default())
437}
438
439fn unix_now() -> u64 {
440    std::time::SystemTime::now()
441        .duration_since(std::time::UNIX_EPOCH)
442        .map(|d| d.as_secs())
443        .unwrap_or_default()
444}
445
446/// How often and how lately a dataset was opened: zoxide-style frecency, ranking
447/// Recent and lifting often-opened matches.
448#[derive(Debug, Clone, Copy, Default, PartialEq, serde::Serialize, serde::Deserialize)]
449pub struct Visits {
450    pub count: u32,
451    /// Seconds since the epoch.
452    pub last: u64,
453}
454
455impl Visits {
456    /// Opens weighted by recency: ×4 within the hour, ×2 within the day, ×½ within the
457    /// week, ×¼ after.
458    pub fn frecency(&self, now: u64) -> f64 {
459        let age = now.saturating_sub(self.last);
460        let weight = match age {
461            a if a < 3_600 => 4.0,
462            a if a < 86_400 => 2.0,
463            a if a < 604_800 => 0.5,
464            _ => 0.25,
465        };
466        f64::from(self.count) * weight
467    }
468}
469
470/// `recents` reordered by frecency; ties and pre-visit-count recents keep their order.
471pub fn by_frecency(
472    mut recents: Vec<PathBuf>,
473    visits: &std::collections::HashMap<PathBuf, Visits>,
474) -> Vec<PathBuf> {
475    let now = unix_now();
476    let score = |p: &PathBuf| visits.get(p).map_or(0.0, |v| v.frecency(now));
477    recents.sort_by(|a, b| score(b).total_cmp(&score(a)));
478    recents
479}
480
481/// What datui remembers about a measured dataset: a cache, not a catalog. Every field
482/// is re-derivable, and the recorded size and mtime invalidate a changed dataset.
483#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
484pub struct DatasetFacts {
485    /// Modification time in seconds since the epoch, as a fingerprint.
486    pub mtime: u64,
487    /// Size in bytes, the other half of the fingerprint.
488    pub size: u64,
489    pub rows: Option<usize>,
490    pub cols: Option<usize>,
491    /// Whether `cols` is a floor (a sampled large directory), restored with the count so
492    /// a sample is never shown as a total.
493    #[serde(default)]
494    pub cols_sampled: bool,
495    /// Column names, enabling search by column before anything is read this run.
496    #[serde(default)]
497    pub columns: Vec<String>,
498    /// What the dataset turned out to be: recorded since a remote path cannot be
499    /// classified without reading, and a row should read the same in every section.
500    #[serde(default)]
501    pub kind: Option<crate::home::discover::EntryKind>,
502    /// The build rules `kind` came from (see [`crate::home::discover::CLASSIFIER_VERSION`]); `0`
503    /// in records older than this field.
504    #[serde(default)]
505    pub classified_by: u32,
506    /// What opening costs (compression, layout, partitioning), from a footer read a
507    /// remote dataset may not get twice.
508    #[serde(default)]
509    pub cost: crate::home::discover::Cost,
510    /// What one listing of the directory found (its label), restored beside `kind` under
511    /// the same classifier version.
512    #[serde(
513        default,
514        skip_serializing_if = "crate::home::discover::Holds::is_empty"
515    )]
516    pub holds: crate::home::discover::Holds,
517}
518
519/// What an open learned about a dataset's files, so the next open shows columns and
520/// row count without reading footers (seconds for thousands of objects). The listing,
521/// which happens anyway, decides via the fingerprint whether this still holds. A cache,
522/// not a catalog, like [`DatasetFacts`].
523#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
524pub struct DatasetShape {
525    /// What the files looked like when taken; a listing with a different one describes a
526    /// changed dataset, and the rest is ignored.
527    pub fingerprint: String,
528    /// One entry per file, in listing order, which is scan order.
529    pub files: Vec<CachedFooter>,
530    /// The distinct schemas, by index, since thousands of files usually share one. Types
531    /// are Polars' own serialization (names cannot rebuild nested types); if a Polars
532    /// upgrade changes it, entries stop parsing and the cache is just empty.
533    pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
534    /// Seconds since the Unix epoch, for a human reading the file.
535    pub taken_at: u64,
536}
537
538/// What one file's footer said, as much as a reopen needs. It comes back as a
539/// `FileFooter`, so a cached dataset goes through the same code as a fresh read.
540#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
541pub struct CachedFooter {
542    /// Which [`DatasetShape::schemas`] entry this file has; `None` for an unreadable
543    /// footer, remembered as such so a reopen gives the same dataset.
544    pub schema: Option<usize>,
545    /// The rows in each of its row groups, in order.
546    pub row_group_rows: Vec<usize>,
547    /// Compressed bytes of each row group, kept because a note uses them, and a note
548    /// present only on first open is a worse bug than a slow open.
549    pub row_group_bytes: Vec<usize>,
550    /// Uncompressed bytes of each schema column, in order: local footers carry them (the
551    /// only source for binary widths); cloud footers leave this empty.
552    #[serde(default, skip_serializing_if = "Vec::is_empty")]
553    pub column_bytes: Vec<usize>,
554}
555
556impl DatasetShape {
557    /// A listing's fingerprint: file names, count, sizes, mtimes and store tags, all seen
558    /// without opening anything. Names are included since remembered footers join to a
559    /// fresh listing by position; ETags since size and whole-second mtime miss
560    /// same-length rewrites within a second.
561    pub fn fingerprint_of<'a>(
562        files: impl IntoIterator<Item = (&'a str, u64, u64, Option<&'a str>)>,
563    ) -> String {
564        let mut hasher = StableHasher::default();
565        let mut count = 0usize;
566        let mut bytes = 0u64;
567        for (key, size, stamp, etag) in files {
568            hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
569            match etag {
570                Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
571                None => hasher.u64(0),
572            };
573            count += 1;
574            bytes = bytes.saturating_add(size);
575        }
576        format!("{count}-{bytes}-{:016x}", hasher.finish())
577    }
578
579    /// Where `schema` sits in `schemas`, added on the end if it is not there yet.
580    pub fn intern_schema(
581        schemas: &mut Vec<Vec<(String, polars::prelude::DataType)>>,
582        schema: &polars::prelude::Schema,
583    ) -> usize {
584        let columns: Vec<(String, polars::prelude::DataType)> = schema
585            .iter()
586            .map(|(name, dtype)| (name.to_string(), dtype.clone()))
587            .collect();
588        schemas
589            .iter()
590            .position(|s| *s == columns)
591            .unwrap_or_else(|| {
592                schemas.push(columns);
593                schemas.len() - 1
594            })
595    }
596
597    /// The schema at `at` in `schemas`, or `None` when the table has no such entry.
598    pub fn schema_at(
599        schemas: &[Vec<(String, polars::prelude::DataType)>],
600        at: usize,
601    ) -> Option<polars::prelude::Schema> {
602        let columns = schemas.get(at)?;
603        let mut schema = polars::prelude::Schema::with_capacity(columns.len());
604        for (name, dtype) in columns {
605            schema.with_column(name.as_str().into(), dtype.clone());
606        }
607        Some(schema)
608    }
609}
610
611/// Point the cache and config at this test process's own scratch directories, once
612/// per process (`CacheManager::new` and `ConfigManager::new` call it). Named at random,
613/// not by pid (reused pids inherited old recents), and removed at exit.
614#[cfg(test)]
615pub(crate) fn isolate_cache() {
616    // Held for the process's life; a static never drops, so an exit handler removes them.
617    static SCRATCH: std::sync::Mutex<Vec<tempfile::TempDir>> = std::sync::Mutex::new(Vec::new());
618    unsafe extern "C" {
619        fn atexit(callback: extern "C" fn()) -> std::ffi::c_int;
620    }
621    extern "C" fn remove_scratch_dirs() {
622        if let Ok(mut held) = SCRATCH.lock() {
623            held.clear();
624        }
625    }
626    let scratch_dir = |prefix: &str| {
627        tempfile::Builder::new()
628            .prefix(prefix)
629            .tempdir()
630            .expect("a scratch directory for the test process")
631    };
632
633    static ISOLATE: std::sync::Once = std::sync::Once::new();
634    ISOLATE.call_once(|| {
635        let dir = scratch_dir("datui-unit-cache-");
636        let config_dir = scratch_dir("datui-unit-config-");
637        // SAFETY: test-only. Tests run on parallel threads, so this can race another test
638        // reading the environment; accepted in tests and never done outside them.
639        unsafe { std::env::set_var("DATUI_CACHE_DIR", dir.path()) };
640        unsafe { std::env::set_var("DATUI_CONFIG_DIR", config_dir.path()) };
641        let mut held = SCRATCH.lock().unwrap_or_else(|e| e.into_inner());
642        held.push(dir);
643        held.push(config_dir);
644        // SAFETY: the C runtime's `atexit`, present on every platform std runs on; the
645        // callback only drops the directories above.
646        unsafe { atexit(remove_scratch_dirs) };
647    });
648}
649
650/// Whether this process is a cargo-built test binary (`target/<profile>/deps/…`, or
651/// `target/<triple>/<profile>/deps/` cross-built); the program itself never is.
652pub(crate) fn running_as_a_cargo_test() -> bool {
653    std::env::current_exe()
654        .is_ok_and(|exe| cargo_test_layout(&exe, std::env::var_os("CARGO_TARGET_DIR").as_deref()))
655}
656
657/// The directory is `deps` under a `target` at most three levels up, or under cargo's
658/// configured target dir, so a program installed in some `deps` is not mistaken.
659fn cargo_test_layout(exe: &Path, target_dir: Option<&std::ffi::OsStr>) -> bool {
660    let Some(deps) = exe.parent() else {
661        return false;
662    };
663    if deps.file_name().is_none_or(|name| name != "deps") {
664        return false;
665    }
666    if target_dir.is_some_and(|dir| exe.starts_with(dir)) {
667        return true;
668    }
669    deps.ancestors()
670        .skip(1)
671        .take(3)
672        .any(|dir| dir.file_name().is_some_and(|name| name == "target"))
673}
674
675/// A cloud source's buckets from an earlier run, shown at once while a fresh listing
676/// is out.
677#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
678pub struct CloudListing {
679    /// What the source pointed at when listed; a mismatch means another server, ignored.
680    pub fingerprint: String,
681    pub buckets: Vec<String>,
682    /// Seconds since the Unix epoch.
683    pub listed_at: u64,
684}
685
686/// Dataset shapes by path, at their listing's fingerprint. Bounded by bytes: shape
687/// size follows file count (842k files is megabytes), so a count bound would let a
688/// few huge datasets evict a small one in use.
689pub(crate) struct Shapes;
690
691impl Kind for Shapes {
692    const DIR: &'static str = "shapes";
693    const EXT: &'static str = "shape";
694    const VERSION: u16 = 1;
695    const BUDGET: u64 = 128 << 20;
696    type Value = DatasetShape;
697
698    fn encode(shape: &DatasetShape) -> Result<Vec<u8>> {
699        encode_shape(shape)
700    }
701
702    fn decode(payload: &[u8]) -> Option<DatasetShape> {
703        decode_shape(payload)
704    }
705}
706
707/// The footers a dataset's count has read, each by [`file_identity`], so a recount
708/// reads only new files and a stopped count keeps its progress. By file, unlike
709/// [`DatasetShape`], so it stays right for a changed dataset.
710#[derive(Debug, Clone, Default, PartialEq, Eq)]
711pub struct FileFooters {
712    /// The distinct schemas, as in [`DatasetShape::schemas`].
713    pub schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
714    /// Each file's identity, and its footer.
715    pub files: Vec<(u64, CachedFooter)>,
716}
717
718/// A file's identity for [`FileFooters`]: key, size, mtime and store tag, any of which
719/// changes on rewrite.
720pub fn file_identity(key: &str, size: u64, stamp: u64, etag: Option<&str>) -> u64 {
721    let mut hasher = StableHasher::default();
722    hasher.bytes(key.as_bytes()).u64(size).u64(stamp);
723    match etag {
724        Some(etag) => hasher.u64(1).bytes(etag.as_bytes()),
725        None => hasher.u64(0),
726    };
727    hasher.finish()
728}
729
730/// [`FileFooters`] by dataset path; empty fingerprint, as each file checks itself.
731pub(crate) struct FileFootersKind;
732
733impl Kind for FileFootersKind {
734    const DIR: &'static str = "file_footers";
735    const EXT: &'static str = "footers";
736    const VERSION: u16 = 1;
737    const BUDGET: u64 = 64 << 20;
738    type Value = FileFooters;
739
740    fn encode(footers: &FileFooters) -> Result<Vec<u8>> {
741        let shape = DatasetShape {
742            fingerprint: String::new(),
743            files: footers.files.iter().map(|(_, f)| f.clone()).collect(),
744            schemas: footers.schemas.clone(),
745            taken_at: 0,
746        };
747        let mut out = encode_shape(&shape)?;
748        for (identity, _) in &footers.files {
749            out.extend_from_slice(&identity.to_le_bytes());
750        }
751        Ok(out)
752    }
753
754    fn decode(payload: &[u8]) -> Option<FileFooters> {
755        // The identities follow the shape, eight bytes a file.
756        let count = payload.len().checked_sub(4)?;
757        let (len, rest) = payload.split_first_chunk::<4>()?;
758        let header = usize::try_from(u32::from_le_bytes(*len)).ok()?;
759        let mut body = rest.get(header..)?;
760        let files = usize::try_from(take_varint(&mut body)?).ok()?;
761        let ids = files.checked_mul(8)?;
762        if ids > count {
763            return None;
764        }
765        let (shape, identities) = payload.split_at(payload.len() - ids);
766        let shape = decode_shape(shape)?;
767        (shape.files.len() == files).then(|| FileFooters {
768            schemas: shape.schemas,
769            files: identities
770                .as_chunks::<8>()
771                .0
772                .iter()
773                .map(|id| u64::from_le_bytes(*id))
774                .zip(shape.files)
775                .collect(),
776        })
777    }
778}
779
780/// What home measured of each dataset, by path; records carry their own size and
781/// mtime, so the store fingerprint is empty.
782pub(crate) struct Facts;
783
784impl Kind for Facts {
785    const DIR: &'static str = "facts";
786    const EXT: &'static str = "facts";
787    const VERSION: u16 = 1;
788    const BUDGET: u64 = 16 << 20;
789    type Value = DatasetFacts;
790
791    fn encode(facts: &DatasetFacts) -> Result<Vec<u8>> {
792        Ok(serde_json::to_vec(facts)?)
793    }
794
795    fn decode(payload: &[u8]) -> Option<DatasetFacts> {
796        serde_json::from_slice(payload).ok()
797    }
798}
799
800/// Each cloud source's last bucket listing, by source ID, at the source's fingerprint.
801pub(crate) struct CloudListings;
802
803impl Kind for CloudListings {
804    const DIR: &'static str = "cloud_listings";
805    const EXT: &'static str = "listing";
806    const VERSION: u16 = 1;
807    const BUDGET: u64 = 4 << 20;
808    type Value = CloudListing;
809
810    fn encode(listing: &CloudListing) -> Result<Vec<u8>> {
811        Ok(serde_json::to_vec(listing)?)
812    }
813
814    fn decode(payload: &[u8]) -> Option<CloudListing> {
815        serde_json::from_slice(payload).ok()
816    }
817}
818
819/// A shape's header, read before its footers: fingerprint, schemas, time. JSON for the
820/// Polars-serialized schemas; the bulk per-file footers follow as varints.
821#[derive(serde::Serialize, serde::Deserialize)]
822struct ShapeHeader {
823    fingerprint: String,
824    schemas: Vec<Vec<(String, polars::prelude::DataType)>>,
825    taken_at: u64,
826}
827fn put_varint(out: &mut Vec<u8>, mut n: u64) {
828    while n >= 0x80 {
829        out.push((n as u8) | 0x80);
830        n >>= 7;
831    }
832    out.push(n as u8);
833}
834
835fn take_varint(bytes: &mut &[u8]) -> Option<u64> {
836    let mut n = 0u64;
837    for shift in (0..64).step_by(7) {
838        let (&b, rest) = bytes.split_first()?;
839        *bytes = rest;
840        n |= u64::from(b & 0x7f) << shift;
841        if b & 0x80 == 0 {
842            return Some(n);
843        }
844    }
845    None
846}
847
848fn put_list(out: &mut Vec<u8>, values: &[usize]) {
849    put_varint(out, values.len() as u64);
850    for &v in values {
851        put_varint(out, v as u64);
852    }
853}
854
855fn take_list(bytes: &mut &[u8]) -> Option<Vec<usize>> {
856    let len = usize::try_from(take_varint(bytes)?).ok()?;
857    // Each value is at least a byte: a longer length means a broken file, not an
858    // allocation.
859    if len > bytes.len() {
860        return None;
861    }
862    (0..len)
863        .map(|_| take_varint(bytes).and_then(|v| usize::try_from(v).ok()))
864        .collect()
865}
866
867fn encode_shape(shape: &DatasetShape) -> Result<Vec<u8>> {
868    let header = serde_json::to_vec(&ShapeHeader {
869        fingerprint: shape.fingerprint.clone(),
870        schemas: shape.schemas.clone(),
871        taken_at: shape.taken_at,
872    })?;
873    let mut out = Vec::with_capacity(8 + header.len() + shape.files.len() * 8);
874    out.extend_from_slice(&u32::try_from(header.len())?.to_le_bytes());
875    out.extend_from_slice(&header);
876    put_varint(&mut out, shape.files.len() as u64);
877    for file in &shape.files {
878        put_varint(&mut out, file.schema.map_or(0, |s| s as u64 + 1));
879        put_list(&mut out, &file.row_group_rows);
880        put_list(&mut out, &file.row_group_bytes);
881        put_list(&mut out, &file.column_bytes);
882    }
883    Ok(out)
884}
885
886/// The shape in `bytes`, or `None` for one that does not hold together.
887fn decode_shape(bytes: &[u8]) -> Option<DatasetShape> {
888    let (len, rest) = bytes.split_first_chunk::<4>()?;
889    let len = usize::try_from(u32::from_le_bytes(*len)).ok()?;
890    let (header, mut body) = (rest.get(..len)?, rest.get(len..)?);
891    let header: ShapeHeader = serde_json::from_slice(header).ok()?;
892    let count = usize::try_from(take_varint(&mut body)?).ok()?;
893    if count > body.len() {
894        return None;
895    }
896    let mut files = Vec::with_capacity(count);
897    // Dataset totals must fit, so no later sum overflows on a damaged file.
898    let (mut rows, mut bytes_total) = (0usize, 0usize);
899    for _ in 0..count {
900        let schema = match take_varint(&mut body)? {
901            0 => None,
902            at => Some(usize::try_from(at - 1).ok()?),
903        };
904        let footer = CachedFooter {
905            schema,
906            row_group_rows: take_list(&mut body)?,
907            row_group_bytes: take_list(&mut body)?,
908            column_bytes: take_list(&mut body)?,
909        };
910        if let Some(at) = footer.schema
911            && at >= header.schemas.len()
912        {
913            return None;
914        }
915        for &n in &footer.row_group_rows {
916            rows = rows.checked_add(n)?;
917        }
918        for &n in footer.row_group_bytes.iter().chain(&footer.column_bytes) {
919            bytes_total = bytes_total.checked_add(n)?;
920        }
921        files.push(footer);
922    }
923    body.is_empty().then_some(DatasetShape {
924        fingerprint: header.fingerprint,
925        files,
926        schemas: header.schemas,
927        taken_at: header.taken_at,
928    })
929}
930
931impl CacheManager {
932    /// How many dataset shapes are kept.
933    #[cfg(test)]
934    pub fn dataset_shapes_kept(&self) -> usize {
935        Store::<Shapes>::new(self).len()
936    }
937
938    /// The remembered shape of a dataset if `fingerprint` still matches (a mismatched
939    /// entry has no correct use). A hit counts as use, so unchanging datasets stay.
940    pub fn dataset_shape(&self, path: &str, fingerprint: &str) -> Option<DatasetShape> {
941        Store::<Shapes>::new(self).get(path, fingerprint)
942    }
943
944    /// Whether any shape is kept for `path`: one stat.
945    pub fn has_dataset_shape(&self, path: &str) -> bool {
946        Store::<Shapes>::new(self).file(path).exists()
947    }
948
949    /// Remember one dataset's shape, keeping the others while they fit the budget.
950    pub fn save_dataset_shape(&self, path: &str, shape: DatasetShape) {
951        Store::<Shapes>::new(self).put(path, &shape.fingerprint, &shape);
952    }
953
954    /// The footers counts of the dataset at `path` have read, by file.
955    pub fn file_footers(&self, path: &str) -> Option<FileFooters> {
956        Store::<FileFootersKind>::new(self).get(path, "")
957    }
958
959    /// Remember the footers a count of the dataset at `path` read, by file.
960    pub fn save_file_footers(&self, path: &str, footers: &FileFooters) {
961        Store::<FileFootersKind>::new(self).put(path, "", footers);
962    }
963
964    /// The last listing of cloud source `id`, if it was taken at `fingerprint`.
965    pub fn cloud_listing(&self, id: &str, fingerprint: &str) -> Option<CloudListing> {
966        Store::<CloudListings>::new(self).get(id, fingerprint)
967    }
968
969    /// Record one source's listing, keeping the others.
970    pub fn save_cloud_listing(&self, id: &str, listing: CloudListing) {
971        Store::<CloudListings>::new(self).put(id, &listing.fingerprint.clone(), &listing);
972    }
973
974    /// Source IDs hidden from the home screen with Delete.
975    pub fn load_hidden_cloud_sources(&self) -> Vec<String> {
976        self.load_history_or_log("cloud_hidden")
977    }
978
979    /// Hide a source from the home screen until the cache is cleared.
980    pub fn hide_cloud_source(&self, id: &str) {
981        let id = id.to_string();
982        self.update_history_file("cloud_hidden", |hidden| {
983            if !hidden.contains(&id) {
984                hidden.push(id.clone());
985            }
986        })
987        .or_log("hide a cloud source");
988    }
989
990    /// Whether Delete on its heading hid the bundled Example datasets (a user's
991    /// `examples.toml` still shows).
992    pub fn examples_hidden(&self) -> bool {
993        !self.load_history_or_log("examples_hidden").is_empty()
994    }
995
996    /// Hide the bundled Example datasets until the cache is cleared.
997    pub fn hide_examples(&self) {
998        self.save_history_file("examples_hidden", &["hidden".to_string()])
999            .or_log("hide the example datasets");
1000    }
1001
1002    /// The directories Ctrl+D kept here before 0.4.0 (now in `catalog.toml`), in added
1003    /// order.
1004    pub fn load_remembered_places(&self) -> Vec<PathBuf> {
1005        let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1006        if !file.exists() {
1007            return Vec::new();
1008        }
1009        self.load_history_or_log("home_remembered")
1010            .into_iter()
1011            .map(PathBuf::from)
1012            .collect()
1013    }
1014
1015    /// Drop the list [`Self::load_remembered_places`] reads, once it is moved.
1016    pub fn clear_remembered_places(&self) {
1017        let file = self.cache_file(&format!("home_remembered{HISTORY_SUFFIX}"));
1018        if let Err(e) = std::fs::remove_file(&file)
1019            && e.kind() != std::io::ErrorKind::NotFound
1020        {
1021            log::warn!(target: "datui", "remove {}: {e}", file.display());
1022        }
1023    }
1024
1025    /// Write `places` where Ctrl+D kept them before 0.4.0, for the migration's tests.
1026    pub fn save_remembered_places(&self, places: &[PathBuf]) -> Result<()> {
1027        let places: Vec<String> = places
1028            .iter()
1029            .map(|p| p.to_string_lossy().into_owned())
1030            .collect();
1031        self.save_history_file("home_remembered", &places)
1032    }
1033}
1034
1035impl CacheManager {
1036    /// Every remembered dataset fact, not counted as use. Unreadable records are absent:
1037    /// failing to read a cache must never be worse than lacking it.
1038    pub fn load_dataset_facts(&self) -> std::collections::HashMap<PathBuf, DatasetFacts> {
1039        Store::<Facts>::new(self)
1040            .scan()
1041            .into_iter()
1042            .map(|(path, facts)| (PathBuf::from(path), facts))
1043            .collect()
1044    }
1045
1046    /// What datui knows about one dataset, counted as use.
1047    pub fn dataset_facts(&self, path: &Path) -> Option<DatasetFacts> {
1048        Store::<Facts>::new(self).get(path.to_str()?, "")
1049    }
1050
1051    /// Count these records as used (shown on home); least recently used go first past
1052    /// the budget.
1053    pub fn touch_dataset_facts<'a>(&self, paths: impl IntoIterator<Item = &'a Path>) {
1054        let store = Store::<Facts>::new(self);
1055        for path in paths.into_iter().filter_map(Path::to_str) {
1056            store.touch(path);
1057        }
1058    }
1059
1060    /// Record newly measured datasets, one file each; non-UTF-8 paths are skipped.
1061    pub fn record_dataset_facts(&self, facts: &[(PathBuf, DatasetFacts)]) {
1062        Store::<Facts>::new(self).put_all(
1063            facts
1064                .iter()
1065                .filter_map(|(path, facts)| Some((path.to_str()?, "", facts))),
1066        );
1067    }
1068
1069    /// Run `work` under the named cache lock, or skip it if contended past the deadline.
1070    fn with_cache_lock<F>(&self, name: &str, work: F) -> Result<()>
1071    where
1072        F: FnOnce() -> Result<()>,
1073    {
1074        use fs2::FileExt;
1075
1076        self.ensure_cache_dir()?;
1077        let Some(lock) = lock_file(&self.cache_file(&format!("{name}.lock")), LOCK_TIMEOUT)? else {
1078            log::info!(target: "datui", "{name} cache not updated: its lock is busy");
1079            return Ok(());
1080        };
1081
1082        let result = work();
1083        let _ = FileExt::unlock(&lock);
1084        result
1085    }
1086}
1087
1088#[cfg(test)]
1089mod harness_tests;
1090
1091#[cfg(test)]
1092mod recents_pruning_tests;
1093
1094#[cfg(test)]
1095mod dataset_shape_tests;
1096
1097/// The same checks for every kind the [`Store`] holds.
1098#[cfg(test)]
1099mod store_harness_tests;
1100
1101#[cfg(test)]
1102mod facts_compat_tests;