use std::path::{Path, PathBuf};
const COMPRESSION_EXTENSIONS: &[&str] = &["gz", "bz2", "xz", "zst", "zstd"];
pub const MAX_ENTRIES_PER_DIR: usize = 5_000;
#[derive(Debug, Clone, Copy, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum EntryKind {
File,
Other,
Hive,
MultiFile,
Delta,
Iceberg,
Hudi,
Directory,
#[serde(other)]
Unknown,
}
pub const CLASSIFIER_VERSION: u32 = 7;
impl EntryKind {
pub fn label(self) -> &'static str {
match self {
EntryKind::File | EntryKind::Other => "",
EntryKind::Hive => "hive",
EntryKind::MultiFile => "multi",
EntryKind::Delta => "delta",
EntryKind::Iceberg => "iceberg",
EntryKind::Hudi => "hudi",
EntryKind::Directory => "dir",
EntryKind::Unknown => "",
}
}
pub fn is_dataset(self) -> bool {
!matches!(self, EntryKind::Directory | EntryKind::Other) && !self.is_lake_table()
}
pub fn is_known_dataset(self) -> bool {
self != EntryKind::Unknown && self.is_dataset()
}
pub fn is_lake_table(self) -> bool {
matches!(
self,
EntryKind::Delta | EntryKind::Iceberg | EntryKind::Hudi
)
}
pub fn lake_name(self) -> Option<&'static str> {
match self {
EntryKind::Delta => Some("Delta"),
EntryKind::Iceberg => Some("Iceberg"),
EntryKind::Hudi => Some("Hudi"),
_ => None,
}
}
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Holds {
#[serde(default)]
pub formats: Vec<(String, usize)>,
#[serde(default, alias = "folders")]
pub directories: usize,
#[serde(default)]
pub partitions: usize,
#[serde(default)]
pub not_read: usize,
#[serde(default)]
pub unnamed: usize,
#[serde(default)]
pub skipped: usize,
#[serde(default)]
pub skipped_names: Vec<String>,
#[serde(default)]
pub truncated: bool,
#[serde(default)]
pub dataset_dict: bool,
}
pub(crate) const SKIPPED_NAMES_SHOWN: usize = 4;
impl Holds {
pub fn data_files(&self) -> usize {
self.formats.iter().map(|(_, n)| n).sum()
}
pub fn one_format(&self) -> Option<&str> {
match self.formats.as_slice() {
[(name, _)] => Some(name),
_ => None,
}
}
pub fn model_weights(&self) -> Option<(&str, usize)> {
if !is_model_directory(counts_names(self)) {
return None;
}
self.formats
.iter()
.find(|(name, _)| is_weights(name))
.map(|(name, count)| (name.as_str(), *count))
}
pub fn label(&self) -> String {
let more = if self.truncated { "+" } else { "" };
if let Some((name, count)) = self.model_weights() {
return format!("{count}{more} {name}");
}
match self.formats.as_slice() {
[] if self.directories == 1 => format!("1 dir{more}"),
[] if self.directories > 1 => format!("{}{more} dirs", self.directories),
[] => format!("dir{more}"),
[(name, count)] => format!("{count}{more} {name}"),
_ => "mixed".to_string(),
}
}
pub fn is_empty(&self) -> bool {
self.formats.is_empty()
&& self.directories == 0
&& self.skipped == 0
&& self.not_read == 0
&& self.unnamed == 0
&& self.partitions == 0
&& self.skipped_names.is_empty()
&& !self.dataset_dict
&& !self.truncated
}
pub fn line(&self, with_partitions: bool) -> Option<String> {
let more = if self.truncated { "+" } else { "" };
let mut parts: Vec<String> = self
.formats
.iter()
.map(|(name, count)| format!("{count}{more} {name}"))
.collect();
let plain = self.directories.saturating_sub(self.partitions);
if plain > 0 {
let word = if plain == 1 {
"directory"
} else {
"directories"
};
parts.push(format!("{plain}{more} {word}"));
}
if with_partitions && self.partitions > 0 {
let word = if self.partitions == 1 {
"partition"
} else {
"partitions"
};
parts.push(format!("{}{more} {word}", self.partitions));
}
(!parts.is_empty()).then(|| parts.join(" · "))
}
}
impl Entry {
pub fn enter_lists_tables(&self) -> bool {
self.kind == EntryKind::File
&& self.cost.tables.is_some_and(|n| n > 1)
&& !self.cost.opens_one
&& self.format_spec.is_none()
}
pub fn hidden_by_default(&self) -> bool {
self.kind == EntryKind::Other || self.table.as_ref().is_some_and(|t| t.internal)
}
pub fn label(&self) -> std::borrow::Cow<'static, str> {
match self.kind {
_ if self.opens_whole_directory => "".into(),
EntryKind::Directory | EntryKind::MultiFile if !self.holds.is_empty() => {
self.holds.label().into()
}
EntryKind::File if self.format_spec.is_some() => {
self.format_spec.clone().unwrap_or_default().into()
}
EntryKind::File if self.cost.tables.is_some() => {
let n = self.cost.tables.unwrap_or_default();
format!("{n} {}", if n == 1 { "table" } else { "tables" }).into()
}
EntryKind::File
if crate::FileFormat::from_path(Path::new(&self.name)).is_none()
&& crate::FileFormat::from_path(&self.path).is_some() =>
{
crate::FileFormat::from_path(&self.path)
.map(crate::FileFormat::name)
.unwrap_or_default()
.into()
}
kind => kind.label().into(),
}
}
}
#[derive(Debug, Clone)]
pub struct Entry {
pub path: PathBuf,
pub kind: EntryKind,
pub name: String,
pub size: Option<u64>,
pub modified: Option<std::time::SystemTime>,
pub rows: Option<usize>,
pub cols: Option<usize>,
pub cols_sampled: bool,
pub columns: Vec<String>,
pub cost: Cost,
pub holds: Holds,
pub opens_whole_directory: bool,
pub format_spec: Option<String>,
pub table: Option<TableOf>,
pub measured: bool,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct TableOf {
pub format: Option<crate::FileFormat>,
pub kind: String,
pub internal: bool,
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Cost {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub source: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub uncompressed: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub codec: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub row_groups: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub partitions: Option<Partitions>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub tables: Option<usize>,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub opens_one: bool,
#[serde(default, skip_serializing_if = "std::ops::Not::not")]
pub ipc_stream: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct HowRead {
pub mode: crate::ReadMode,
pub download: bool,
}
pub fn how_read(entry: &Entry) -> Option<HowRead> {
use crate::Stored;
if entry.kind != EntryKind::File {
return None;
}
let stored = if crate::CompressionFormat::from_extension(&entry.path).is_some() {
Stored::Compressed { in_memory: false }
} else if entry.cost.ipc_stream {
Stored::Stream
} else {
Stored::Plain
};
let choice = match &entry.format_spec {
Some(name) => crate::cli::FormatChoice::Spec(name.clone()),
None if entry.table.is_some() => {
crate::cli::FormatChoice::Builtin(entry.table.as_ref().and_then(|t| t.format)?)
}
None => crate::cli::FormatChoice::Builtin(data_format(&entry.path)?),
};
let mode = choice.read_mode(stored)?;
let download = match crate::cloud::source::input_source(&entry.path) {
crate::cloud::source::InputSource::Local(_) => false,
crate::cloud::source::InputSource::Http(_) => {
choice.http_file() == crate::RemoteRead::Downloaded
}
_ => choice.bucket_object(stored) == crate::RemoteRead::Downloaded,
};
Some(HowRead { mode, download })
}
#[derive(Debug, Clone, Default, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub struct Partitions {
pub keys: Vec<String>,
pub first_key_values: Vec<String>,
pub count: usize,
pub more: bool,
}
impl Entry {
pub fn directory(path: &Path) -> Self {
Self::new(path.to_path_buf(), EntryKind::Directory)
}
pub fn for_test(path: &Path, name: &str) -> Self {
Self::new(path.to_path_buf(), EntryKind::File).with_name(name)
}
pub(crate) fn with_name(mut self, name: impl Into<String>) -> Self {
self.name = name.into();
self
}
pub(crate) fn new(path: PathBuf, kind: EntryKind) -> Self {
let name = path
.file_name()
.map(|n| n.to_string_lossy().into_owned())
.unwrap_or_else(|| path.to_string_lossy().into_owned());
Self {
path,
kind,
name,
size: None,
modified: None,
rows: None,
cols: None,
cols_sampled: false,
columns: Vec::new(),
cost: Cost::default(),
holds: Default::default(),
opens_whole_directory: false,
measured: false,
format_spec: None,
table: None,
}
}
pub(crate) fn with_fs_metadata(mut self, meta: &std::fs::Metadata) -> Self {
if meta.is_file() {
self.size = Some(meta.len());
}
self.modified = meta.modified().ok();
self
}
}
pub fn is_parquet_key(key: &str) -> bool {
let key = key.trim_end_matches('/');
let (directory, name) = match key.rsplit_once('/') {
Some((directory, name)) => (directory, name),
None => ("", key),
};
if is_bookkeeping(name) {
return false;
}
if name.to_ascii_lowercase().ends_with(".parquet") {
return true;
}
let directory_name = directory.rsplit('/').next().unwrap_or(directory);
!name.contains('.') && directory_name.to_ascii_lowercase().ends_with(".parquet")
}
#[cfg(test)]
mod parquet_key_tests {
use super::is_parquet_key;
#[test]
fn parquet_without_an_extension_is_known_by_its_directory() {
assert!(is_parquet_key(
"occurrence/2026-09-01/occurrence.parquet/000001"
));
assert!(is_parquet_key("data/part-0.parquet"));
assert!(is_parquet_key("DATA/PART-0.PARQUET"));
assert!(!is_parquet_key(
"occurrence/2026-09-01/occurrence.parquet/_SUCCESS"
));
assert!(!is_parquet_key("occurrence.parquet/.part-0.crc"));
assert!(!is_parquet_key("occurrence/2026-09-01/citation.txt"));
assert!(!is_parquet_key("notes/000001"));
}
}
pub fn sniff_format(path: &Path) -> Option<crate::FileFormat> {
crate::formats::readers::sniff_file(path, crate::formats::readers::Asked::Listing)
}
#[derive(Debug, Clone)]
pub enum Sniffed {
Format,
Spec(std::sync::Arc<crate::formats::Spec>),
}
pub fn sniff_listed(path: &Path, formats: &crate::formats::Registry) -> Option<Sniffed> {
use crate::formats::readers::{Asked, HEAD, head_of, sniff};
let head = head_of(path)?;
if sniff(&head, Some(path), Asked::Listing, |_| true).is_some() {
return Some(Sniffed::Format);
}
formats
.listed(path, &head, head.len() < HEAD)
.map(Sniffed::Spec)
}
pub fn name_spec_file(entry: &mut Entry, spec: &crate::formats::Spec) {
entry.kind = EntryKind::File;
entry.format_spec = Some(spec.name.clone());
if spec.lists_variants() {
entry.cost.tables = Some(spec.records.variants.len());
}
}
pub fn name_unlisted_file(entry: &mut Entry, formats: &crate::formats::Registry) {
if formats.is_empty()
|| entry.kind != EntryKind::File
|| entry.table.is_some()
|| is_data_file(&entry.path)
{
return;
}
let spec = match formats.by_glob(&entry.path, false).into_iter().next() {
Some(spec) => Some(spec),
None if worth_sniffing(&entry.path) && is_regular_file(&entry.path) => {
match sniff_listed(&entry.path, formats) {
Some(Sniffed::Spec(spec)) => Some(spec),
_ => None,
}
}
None => None,
};
if let Some(spec) = spec {
name_spec_file(entry, &spec);
}
}
pub(crate) const MAX_SNIFFS_PER_DIR: usize = 256;
pub fn has_no_extension(path: &Path) -> bool {
path.extension().is_none()
}
pub fn worth_sniffing(path: &Path) -> bool {
path.extension().is_none_or(|e| {
e.eq_ignore_ascii_case("bin") || data_format(path).is_some_and(crate::FileFormat::is_lines)
})
}
pub fn has_parquet_magic(path: &Path) -> bool {
use std::io::{Read, Seek, SeekFrom};
let Ok(mut file) = std::fs::File::open(path) else {
return false;
};
let mut head = [0u8; 4];
let mut tail = [0u8; 4];
file.read_exact(&mut head).is_ok()
&& file.seek(SeekFrom::End(-4)).is_ok()
&& file.read_exact(&mut tail).is_ok()
&& &head == b"PAR1"
&& &tail == b"PAR1"
}
pub fn is_parquet_path(path: &Path) -> bool {
path.extension()
.and_then(|e| e.to_str())
.is_some_and(|e| e.eq_ignore_ascii_case("parquet"))
|| is_parquet_key(&directory_and_name(path))
}
pub fn is_data_file(path: &Path) -> bool {
data_extension(path).is_some() || is_parquet_key(&directory_and_name(path))
}
pub fn data_extension(path: &Path) -> Option<String> {
let name = path.file_name().and_then(|n| n.to_str())?;
let lower = name.to_ascii_lowercase();
let mut parts: Vec<&str> = lower.rsplit('.').collect();
parts.reverse();
if parts.len() < 2 {
return None;
}
let mut idx = parts.len() - 1;
if COMPRESSION_EXTENSIONS.contains(&parts[idx]) && idx > 1 {
idx -= 1;
}
crate::FileFormat::from_extension(parts[idx]).map(|_| parts[idx].to_string())
}
pub const NO_READER: &str = "datui has no reader for this file";
pub fn unreadable_by_name(path: &Path) -> bool {
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
return false;
};
let name = name.to_ascii_lowercase();
let parts: Vec<&str> = name.rsplit('.').collect();
let compressed = |last: &str| COMPRESSION_EXTENSIONS.contains(&last);
let ext = match parts[..] {
[last, inner, _, ..] if compressed(last) => inner,
[last, _] if compressed(last) => return false,
[last, _, ..] => last,
_ => return false,
};
crate::FileFormat::from_extension(ext).is_none()
}
pub fn data_format(path: &Path) -> Option<crate::FileFormat> {
crate::FileFormat::from_name_ending(path)
.or_else(|| crate::FileFormat::from_extension(&data_extension(path)?))
}
pub(crate) fn directory_and_name(path: &Path) -> String {
let name = path.file_name().unwrap_or_default().to_string_lossy();
match path.parent().and_then(|p| p.file_name()) {
Some(directory) => format!("{}/{name}", directory.to_string_lossy()),
None => name.into_owned(),
}
}
pub(crate) fn rank_formats(a: (&str, usize), b: (&str, usize)) -> std::cmp::Ordering {
let text = crate::FileFormat::Text.name();
(a.0 == text)
.cmp(&(b.0 == text))
.then_with(|| b.1.cmp(&a.1))
.then_with(|| (a.0 != "parquet").cmp(&(b.0 != "parquet")))
.then_with(|| a.0.cmp(b.0))
}
pub(crate) fn is_hugging_face_metadata(name: &str) -> bool {
matches!(name, "dataset_info.json" | "state.json")
}
fn order_formats(counts: &mut [(crate::FileFormat, usize)]) {
counts.sort_by(|a, b| rank_formats((a.0.name(), a.1), (b.0.name(), b.1)));
}
pub fn is_bookkeeping(name: &str) -> bool {
if name.ends_with("_$folder$") {
return true;
}
if is_partition_name(name) {
return false;
}
name.starts_with(['_', '.'])
}
fn is_weights(name: &str) -> bool {
name == crate::FileFormat::Safetensors.name() || name == crate::FileFormat::Gguf.name()
}
pub(crate) fn is_model_directory<'a>(names: impl IntoIterator<Item = &'a str>) -> bool {
let mut weights = None;
for name in names {
if is_weights(name) {
if weights.is_some_and(|w| w != name) {
return false;
}
weights = Some(name);
} else if name != crate::FileFormat::Json.name() {
return false;
}
}
weights.is_some()
}
fn counts_names(holds: &Holds) -> impl Iterator<Item = &str> {
holds.formats.iter().map(|(name, _)| name.as_str())
}
pub fn is_partition_name(name: &str) -> bool {
matches!(name.find('='), Some(i) if i > 0)
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub enum DirectoryFormat {
One(crate::FileFormat, Vec<PathBuf>),
Mixed {
format: crate::FileFormat,
files: Vec<PathBuf>,
passed_over: Vec<(crate::FileFormat, usize)>,
},
Deeper,
}
pub fn directory_format(dir: &Path) -> DirectoryFormat {
let Ok(iter) = std::fs::read_dir(dir) else {
return DirectoryFormat::Deeper;
};
let mut by_format: Vec<(crate::FileFormat, Vec<PathBuf>)> = Vec::new();
let mut nameless: Vec<PathBuf> = Vec::new();
let mut partitioned = false;
for entry in iter.flatten() {
let path = entry.path();
let Some(name) = path.file_name().and_then(|n| n.to_str()) else {
continue;
};
if is_bookkeeping(name) {
continue;
}
let is_file = match entry.file_type() {
Ok(kind) if kind.is_symlink() => is_regular_file(&path),
Ok(kind) => kind.is_file(),
Err(_) => is_regular_file(&path),
};
if !is_file {
partitioned |= is_partition_dir(&path);
continue;
}
let Some(found) = data_format(&path) else {
if path.extension().is_none() {
nameless.push(path);
}
continue;
};
match by_format.iter_mut().find(|(f, _)| *f == found) {
Some((_, of_that_format)) => of_that_format.push(path),
None => by_format.push((found, vec![path])),
}
}
if by_format.is_empty() && !nameless.is_empty() {
nameless.sort();
let picks = spread(nameless.len());
let sniffed: Vec<crate::FileFormat> = picks
.iter()
.filter_map(|i| sniff_format(&nameless[*i]))
.collect();
if sniffed.len() == picks.len()
&& let Some(found) = sniffed.first().copied()
&& sniffed.iter().all(|f| *f == found)
{
by_format.push((found, nameless));
}
}
if by_format.iter().any(|(f, _)| !f.is_lines()) {
by_format.retain(|(f, _)| !f.is_lines());
}
if by_format
.iter()
.any(|(f, _)| *f == crate::FileFormat::Arrow)
{
for (format, files) in &mut by_format {
if *format == crate::FileFormat::Json {
files.retain(|f| {
!f.file_name()
.and_then(|n| n.to_str())
.is_some_and(is_hugging_face_metadata)
});
}
}
by_format.retain(|(_, files)| !files.is_empty());
}
if partitioned {
return DirectoryFormat::Deeper;
}
by_format.sort_by(|a, b| rank_formats((a.0.name(), a.1.len()), (b.0.name(), b.1.len())));
if is_model_directory(by_format.iter().map(|(f, _)| f.name()))
&& let Some(at) = by_format.iter().position(|(f, _)| is_weights(f.name()))
{
let weights = by_format.remove(at);
by_format.insert(0, weights);
}
let mut by_format = by_format.into_iter();
let Some((format, mut files)) = by_format.next() else {
return DirectoryFormat::Deeper;
};
files.sort();
let passed_over: Vec<(crate::FileFormat, usize)> =
by_format.map(|(f, of_that)| (f, of_that.len())).collect();
if passed_over.is_empty() {
DirectoryFormat::One(format, files)
} else {
DirectoryFormat::Mixed {
format,
files,
passed_over,
}
}
}
const MAX_HIVE_DEPTH: usize = 16;
pub fn hive_leaf_format(dir: &Path) -> DirectoryFormat {
let mut at = dir.to_path_buf();
for _ in 0..MAX_HIVE_DEPTH {
match directory_format(&at) {
DirectoryFormat::Deeper => match first_partition(&at) {
Some(next) => at = next,
None => return DirectoryFormat::Deeper,
},
settled => return settled,
}
}
DirectoryFormat::Deeper
}
fn first_partition(dir: &Path) -> Option<PathBuf> {
let iter = std::fs::read_dir(dir).ok()?;
iter.flatten()
.take(MAX_ENTRIES_PER_DIR)
.map(|entry| entry.path())
.filter(|path| is_partition_dir(path) && path.is_dir())
.min()
}
fn is_partition_dir(path: &Path) -> bool {
path.file_name()
.and_then(|n| n.to_str())
.is_some_and(is_partition_name)
}
pub fn is_empty_marker(name: &str, size: u64) -> bool {
size == 0 && !name.contains('.')
}
pub fn classify_directory(path: &Path) -> EntryKind {
look_at_directory(path).0
}
pub fn look_at_directory(path: &Path) -> (EntryKind, Holds) {
if let Some(lake) = lake_table(path) {
return (lake, Holds::default());
}
let Ok(iter) = std::fs::read_dir(path) else {
return (EntryKind::Directory, Holds::default());
};
let directory = path.file_name().unwrap_or_default().to_string_lossy();
let rules = Rules {
directory: &directory,
sniff: (!crate::home::is_remote_path(path)).then_some(path),
in_bucket: false,
};
let mut truncated = false;
let seen = iter
.flatten()
.take(MAX_ENTRIES_PER_DIR + 1)
.enumerate()
.map_while(|(at, entry)| {
truncated = at == MAX_ENTRIES_PER_DIR;
(!truncated).then(|| seen_on_disk(&entry))
});
let (kind, mut holds) = classify(seen, &rules);
holds.truncated = truncated;
(kind, holds)
}
fn seen_on_disk(entry: &std::fs::DirEntry) -> Seen {
let (is_dir, is_file) = match entry.file_type() {
Ok(kind) if !kind.is_symlink() => (kind.is_dir(), kind.is_file()),
_ => std::fs::metadata(entry.path()).map_or((false, false), |m| (m.is_dir(), m.is_file())),
};
Seen {
name: entry.file_name().to_string_lossy().into_owned(),
is_dir,
is_file,
size: None,
}
}
#[derive(Debug, Clone)]
pub struct Seen {
pub name: String,
pub is_dir: bool,
pub is_file: bool,
pub size: Option<u64>,
}
pub struct Rules<'a> {
pub directory: &'a str,
pub sniff: Option<&'a Path>,
pub in_bucket: bool,
}
pub fn classify(seen: impl Iterator<Item = Seen>, rules: &Rules) -> (EntryKind, Holds) {
use crate::FileFormat;
let mut holds = Holds::default();
let mut counts: Vec<(FileFormat, usize)> = Vec::new();
let mut skipped: Vec<String> = Vec::new();
let mut hugging_face: Vec<String> = Vec::new();
let mut dict_file = false;
let mut lake: Vec<&'static str> = Vec::new();
let (mut present, mut data_files, mut parquet) = (0usize, 0usize, 0usize);
let mut sniffs_left = rules.sniff.map_or(0, |_| MAX_SNIFFS_PER_DIR);
for s in seen {
if s.is_dir
&& rules.in_bucket
&& let Some(marker) = ["_delta_log", ".hoodie", "metadata", "data"]
.into_iter()
.find(|m| *m == s.name)
{
lake.push(marker);
}
if is_bookkeeping(&s.name) {
skipped.push(s.name);
continue;
}
present += 1;
if s.is_dir {
holds.directories += 1;
holds.partitions += usize::from(is_partition_name(&s.name));
continue;
}
if s.size.is_some_and(|size| is_empty_marker(&s.name, size)) {
holds.not_read += 1;
continue;
}
let name = Path::new(&s.name);
let key = format!("{}/{}", rules.directory, s.name);
let named = data_format(name);
let found = named
.filter(|f| !f.is_lines())
.map(
|f| match crate::formats::model_files::is_safetensors_index(name) {
true => FileFormat::Json,
false => f,
},
)
.or_else(|| is_parquet_key(&key).then_some(FileFormat::Parquet))
.or_else(|| {
let dir = rules.sniff?;
spend_sniff(&mut sniffs_left, s.is_file, name)
.then(|| sniff_format(&dir.join(name)))?
})
.or(named)
.filter(|_| s.is_file);
let Some(found) = found else {
match s.is_file && has_no_extension(name) {
true => holds.unnamed += 1,
false => holds.not_read += 1,
}
continue;
};
if found == FileFormat::Json && is_hugging_face_metadata(&s.name) {
hugging_face.push(s.name.clone());
}
dict_file |= found == FileFormat::Json && s.name == crate::formats::hf_splits::DATASET_DICT;
data_files += 1;
parquet += usize::from(is_parquet_key(&key));
match counts.iter_mut().find(|(f, _)| *f == found) {
Some((_, n)) => *n += 1,
None => counts.push((found, 1)),
}
}
if !counts.iter().any(|(f, _)| *f == FileFormat::Arrow) {
hugging_face.clear();
}
holds.dataset_dict = rules.in_bucket && dict_file && holds.directories > 0;
if holds.dataset_dict {
hugging_face.push(crate::formats::hf_splits::DATASET_DICT.to_string());
}
if let Some((_, n)) = counts.iter_mut().find(|(f, _)| *f == FileFormat::Json) {
*n -= hugging_face.len();
data_files -= hugging_face.len();
present -= hugging_face.len();
skipped.append(&mut hugging_face);
}
if counts.iter().any(|(f, _)| !f.is_lines()) {
for (_, n) in counts.iter_mut().filter(|(f, _)| f.is_lines()) {
data_files -= *n;
holds.not_read += std::mem::take(n);
}
}
counts.retain(|(_, n)| *n > 0);
order_formats(&mut counts);
let one_readable = matches!(counts.as_slice(), [(f, _)] if f.reads_many_files());
holds.formats = counts
.into_iter()
.map(|(f, n)| (f.name().to_string(), n))
.collect();
skipped.sort();
skipped.dedup();
holds.skipped = skipped.len();
skipped.truncate(SKIPPED_NAMES_SHOWN);
holds.skipped_names = skipped;
let marked = |m| lake.contains(&m);
if marked("_delta_log") {
return (EntryKind::Delta, holds);
}
if marked(".hoodie") {
return (EntryKind::Hudi, holds);
}
if marked("metadata") && marked("data") && parquet == 0 {
return (EntryKind::Iceberg, holds);
}
if holds.partitions > 0 && holds.partitions >= data_files {
return (EntryKind::Hive, holds);
}
let one_table = if rules.in_bucket {
parquet > 1 && parquet == data_files && parquet * 2 >= present
} else {
(data_files > 1 && one_readable && data_files * 2 >= present)
|| (holds.partitions == 0 && is_model_directory(counts_names(&holds)))
};
let kind = match one_table {
true => EntryKind::MultiFile,
false => EntryKind::Directory,
};
(kind, holds)
}
fn spend_sniff(left: &mut usize, is_file: bool, name: &Path) -> bool {
let spend = is_file && *left > 0 && worth_sniffing(name);
*left -= usize::from(spend);
spend
}
const ICEBERG_METADATA_PROBE: usize = 64;
fn lake_table(path: &Path) -> Option<EntryKind> {
if path.join("_delta_log").is_dir() {
return Some(EntryKind::Delta);
}
if path.join(".hoodie").is_dir() {
return Some(EntryKind::Hudi);
}
let metadata = path.join("metadata");
if path.join("data").is_dir()
&& metadata.is_dir()
&& std::fs::read_dir(&metadata).is_ok_and(|entries| {
entries
.flatten()
.take(ICEBERG_METADATA_PROBE)
.any(|e| e.file_name().to_string_lossy().ends_with(".metadata.json"))
})
{
return Some(EntryKind::Iceberg);
}
None
}
#[derive(Debug, Clone, Default)]
pub struct Scan {
pub entries: Vec<Entry>,
pub truncated: bool,
}
pub fn scan_dir_specs(dir: &Path, formats: &crate::formats::Registry) -> Scan {
scan_dir_with(dir, formats, |_| {})
}
const LISTING_PROGRESS_EVERY: std::time::Duration = std::time::Duration::from_millis(250);
pub fn scan_dir_progressive(dir: &Path, progress: impl FnMut(&[Entry])) -> Scan {
scan_dir_with(dir, &crate::formats::Registry::default(), progress)
}
fn scan_dir_with(
dir: &Path,
formats: &crate::formats::Registry,
mut progress: impl FnMut(&[Entry]),
) -> Scan {
let Ok(iter) = std::fs::read_dir(dir) else {
return Scan::default();
};
let mut shown = std::time::Instant::now();
let mut entries = Vec::new();
let mut sent = 0usize;
let mut seen = 0usize;
let mut truncated = false;
let mut sniffs_left = if crate::home::is_remote_path(dir) {
0
} else {
MAX_SNIFFS_PER_DIR
};
for dir_entry in iter.flatten().take(MAX_ENTRIES_PER_DIR + 1) {
seen += 1;
if seen > MAX_ENTRIES_PER_DIR {
truncated = true;
break;
}
let path = dir_entry.path();
let name = dir_entry.file_name();
if name.to_string_lossy().starts_with('.') {
continue;
}
let Ok(meta) = dir_entry.metadata() else {
continue;
};
let mut spec = None;
let kind = if meta.is_dir() {
EntryKind::Unknown
} else if meta.is_file() && is_data_file(&path) {
EntryKind::File
} else if spend_sniff(&mut sniffs_left, meta.is_file(), &path) {
match sniff_listed(&path, formats) {
Some(Sniffed::Format) => EntryKind::File,
Some(Sniffed::Spec(found)) => {
spec = Some(found);
EntryKind::File
}
None => EntryKind::Other,
}
} else if meta.is_file() {
EntryKind::Other
} else {
continue;
};
let mut entry = Entry::new(path, kind).with_fs_metadata(&meta);
if let Some(spec) = spec {
name_spec_file(&mut entry, &spec);
}
entries.push(entry);
if shown.elapsed() >= LISTING_PROGRESS_EVERY {
progress(&entries[sent..]);
sent = entries.len();
shown = std::time::Instant::now();
}
}
sort_entries(&mut entries);
Scan { entries, truncated }
}
pub(crate) fn sort_entries(entries: &mut [Entry]) {
entries.sort_by(|a, b| {
let group = |k: EntryKind| match k {
k if k.is_known_dataset() => 0,
EntryKind::Other => 2,
_ => 1,
};
group(a.kind).cmp(&group(b.kind)).then_with(|| {
a.name
.to_ascii_lowercase()
.cmp(&b.name.to_ascii_lowercase())
})
});
}
const MAX_WALK_DEPTH: u8 = 4;
const MAX_FOOTERS_PER_DATASET: usize = 64;
pub fn enrich(entry: &mut Entry) {
enrich_as(entry, &crate::formats::schema_union::ReadAs::default())
}
pub fn enrich_as(entry: &mut Entry, as_read: &crate::formats::schema_union::ReadAs) {
enrich_with(entry, as_read, None)
}
pub fn enrich_with(
entry: &mut Entry,
as_read: &crate::formats::schema_union::ReadAs,
remembered: Option<&crate::cache::CacheManager>,
) {
match entry.kind {
EntryKind::File => {
enrich_parquet(entry);
enrich_tables(entry);
enrich_arrow(entry);
}
EntryKind::Hive | EntryKind::MultiFile => enrich_dataset(entry, as_read, remembered),
EntryKind::Directory | EntryKind::Unknown | EntryKind::Other => {}
EntryKind::Delta | EntryKind::Iceberg | EntryKind::Hudi => {}
}
}
fn enrich_dataset(
entry: &mut Entry,
as_read: &crate::formats::schema_union::ReadAs,
remembered: Option<&crate::cache::CacheManager>,
) {
if entry.kind == EntryKind::Hive {
entry.cost.partitions = partition_layout(&entry.path);
}
let reads_as_parquet = entry.kind == EntryKind::Hive
|| match entry.holds.one_format() {
None => true,
Some(name) => crate::FileFormat::from_name(name) == Some(crate::FileFormat::Parquet),
};
if !reads_as_parquet {
entry.size = None;
judge_by_names(entry, as_read);
return;
}
entry.size = None;
let files = parquet_files_under(&entry.path);
if files.len() > MAX_FOOTERS_PER_DATASET
&& let Some((listed, footers)) = remembered
.and_then(|cache| crate::formats::dataset_files::remembered_footers(&entry.path, cache))
{
measure_from_footers(entry, &listed, &footers);
return;
}
if files.is_empty() || files.len() > MAX_FOOTERS_PER_DATASET {
let sampled = sample_footers(&files);
let names: Vec<Vec<String>> = sampled.iter().map(column_names).collect();
let tops: Vec<Vec<String>> = names
.iter()
.map(|n| crate::formats::schema_union::top_level_columns(n))
.collect();
if entry.kind == EntryKind::MultiFile && one_table_from(&tops) == Some(false) {
let own_files = direct_children(&files, &entry.path);
let own = sample_footers(&own_files);
entry.columns = union_of(&names);
entry.cols_sampled = own.len() < own_files.len();
let top = union_of(&own.iter().map(top_level_names).collect::<Vec<_>>());
downgrade_to_directory(entry, (!top.is_empty()).then_some(top.len()));
return;
}
if let Some(meta) = sampled.first() {
entry.columns = union_of(&names);
let top = union_of(&sampled.iter().map(top_level_names).collect::<Vec<_>>());
entry.cols = Some(top.len() + partition_columns_beyond(entry, &top));
entry.cols_sampled = true;
physical_facts(meta, &mut entry.cost);
entry.cost.uncompressed = None;
}
return;
}
let mut rows = 0usize;
let mut bytes = 0u64;
let mut columns: Vec<String> = Vec::new();
let mut seen_columns = std::collections::HashSet::new();
let mut top_level: Vec<String> = Vec::new();
let mut seen_top_level = std::collections::HashSet::new();
let mut per_file: Vec<Vec<String>> = Vec::with_capacity(files.len());
let mut own_bytes = 0u64;
let mut own_top_level: Vec<String> = Vec::new();
let mut own_seen_top = std::collections::HashSet::new();
let mut cost = Cost::default();
let mut uncompressed = 0u64;
let mut row_groups = 0usize;
for file in &files {
let Some(meta) = crate::formats::parquet_footer::read_parquet_metadata(file) else {
return; };
rows += meta.num_rows;
let names = column_names(&meta);
for name in &names {
if seen_columns.insert(name.clone()) {
columns.push(name.clone());
}
}
for name in top_level_names(&meta) {
if seen_top_level.insert(name.clone()) {
top_level.push(name);
}
}
per_file.push(crate::formats::schema_union::top_level_columns(&names));
let file_bytes = std::fs::metadata(file).map(|m| m.len()).unwrap_or(0);
bytes += file_bytes;
if file.parent() == Some(entry.path.as_path()) {
own_bytes += file_bytes;
for name in top_level_names(&meta) {
if own_seen_top.insert(name.clone()) {
own_top_level.push(name);
}
}
}
let mut per_file = Cost::default();
physical_facts(&meta, &mut per_file);
uncompressed += per_file.uncompressed.unwrap_or(0);
row_groups += per_file.row_groups.unwrap_or(0);
if cost.codec.is_none() {
cost.codec = per_file.codec;
}
}
if entry.kind == EntryKind::MultiFile && !crate::formats::schema_union::is_nested(&per_file) {
entry.size = Some(own_bytes);
entry.columns = columns;
entry.cols_sampled = false;
downgrade_to_directory(
entry,
(!own_top_level.is_empty()).then_some(own_top_level.len()),
);
return;
}
entry.rows = Some(rows);
entry.cols = Some(top_level.len() + partition_columns_beyond(entry, &top_level));
entry.size = Some(bytes);
entry.columns = columns;
cost.uncompressed = (uncompressed > 0).then_some(uncompressed);
cost.row_groups = (row_groups > 0).then_some(row_groups);
cost.partitions = entry.cost.partitions.take();
entry.cost = cost;
}
fn partition_columns_beyond(entry: &Entry, top_level: &[String]) -> usize {
entry.cost.partitions.as_ref().map_or(0, |layout| {
layout
.keys
.iter()
.filter(|key| !top_level.contains(key))
.count()
})
}
fn direct_children(files: &[PathBuf], dir: &Path) -> Vec<PathBuf> {
files
.iter()
.filter(|f| f.parent() == Some(dir))
.cloned()
.collect()
}
pub(crate) fn spread(files: usize) -> Vec<usize> {
let mut picks = match files {
0 => Vec::new(),
n => vec![0, n / 2, n - 1],
};
picks.dedup();
picks
}
pub(crate) fn one_table_from(footers: &[Vec<String>]) -> Option<bool> {
(footers.len() >= 2).then(|| crate::formats::schema_union::is_nested(footers))
}
fn sample_footers(files: &[PathBuf]) -> Vec<crate::formats::parquet_footer::Footer> {
spread(files.len())
.into_iter()
.filter_map(|i| crate::formats::parquet_footer::read_parquet_metadata(&files[i]))
.collect()
}
fn judge_by_names(entry: &mut Entry, as_read: &crate::formats::schema_union::ReadAs) {
if entry.kind != EntryKind::MultiFile {
return;
}
let Some(format) = entry
.holds
.one_format()
.and_then(crate::FileFormat::from_name)
else {
return;
};
let DirectoryFormat::One(_, files) = directory_format(&entry.path) else {
return;
};
let sampled = crate::formats::schema_union::sample_files(&files, format, as_read);
if sampled.nests == Some(false) {
let cols = (!sampled.columns.is_empty()).then_some(sampled.columns.len());
entry.cols_sampled = sampled.read < files.len();
entry.columns = sampled.columns;
downgrade_to_directory(entry, cols);
}
}
fn downgrade_to_directory(entry: &mut Entry, cols: Option<usize>) {
entry.kind = EntryKind::Directory;
entry.rows = None;
entry.cols = cols;
entry.cost = Cost {
partitions: entry.cost.partitions.take(),
..Cost::default()
};
}
fn top_level_names(meta: &crate::formats::parquet_footer::Footer) -> Vec<String> {
meta.schema_descr
.fields()
.iter()
.map(|field| field.name().to_string())
.collect()
}
fn union_of(per_file: &[Vec<String>]) -> Vec<String> {
let mut seen = std::collections::HashSet::new();
per_file
.iter()
.flatten()
.filter(|name| seen.insert(name.as_str()))
.cloned()
.collect()
}
fn parquet_files_under(dir: &Path) -> Vec<PathBuf> {
let mut files = crate::formats::dataset_files::LocalFiles::new(dir)
.first_files(MAX_WALK_DEPTH as usize + 1, MAX_FOOTERS_PER_DATASET);
files.retain(|p| is_regular_file(p));
files
}
fn measure_from_footers(
entry: &mut Entry,
files: &[crate::formats::dataset_files::DatasetFile],
footers: &[Option<crate::formats::schema_union::FileFooter>],
) {
let per_file: Vec<Vec<String>> = footers
.iter()
.flatten()
.map(|f| f.schema.iter_names().map(|n| n.to_string()).collect())
.collect();
let columns = union_of(&per_file);
if entry.kind == EntryKind::MultiFile && !crate::formats::schema_union::is_nested(&per_file) {
let own: Vec<usize> = files
.iter()
.enumerate()
.filter(|(_, f)| Path::new(&f.key).parent() == Some(entry.path.as_path()))
.map(|(i, _)| i)
.collect();
entry.size = Some(own.iter().map(|&i| files[i].size).sum());
let own_columns = union_of(
&own.iter()
.filter_map(|&i| footers[i].as_ref())
.map(|f| f.schema.iter_names().map(|n| n.to_string()).collect())
.collect::<Vec<_>>(),
);
entry.columns = columns;
entry.cols_sampled = false;
downgrade_to_directory(
entry,
(!own_columns.is_empty()).then_some(own_columns.len()),
);
return;
}
let footers: Vec<&crate::formats::schema_union::FileFooter> =
footers.iter().flatten().collect();
let uncompressed: u64 = footers
.iter()
.flat_map(|f| &f.column_bytes)
.map(|(_, bytes)| *bytes as u64)
.sum();
let row_groups: usize = footers.iter().map(|f| f.row_group_rows.len()).sum();
entry.rows = Some(footers.iter().map(|f| f.rows()).sum());
entry.cols = Some(columns.len() + partition_columns_beyond(entry, &columns));
entry.size = Some(files.iter().map(|f| f.size).sum());
entry.columns = columns;
entry.cols_sampled = false;
entry.cost = Cost {
uncompressed: (uncompressed > 0).then_some(uncompressed),
row_groups: (row_groups > 0).then_some(row_groups),
partitions: entry.cost.partitions.take(),
..Cost::default()
};
}
pub fn enrich_parquet(entry: &mut Entry) {
if entry.kind != EntryKind::File {
return;
}
if !is_parquet_path(&entry.path) {
return;
}
if !is_regular_file(&entry.path) {
return;
}
if let Some(meta) = crate::formats::parquet_footer::read_parquet_metadata(&entry.path) {
entry.rows = Some(meta.num_rows);
entry.columns = column_names(&meta);
entry.cols = Some(top_level_names(&meta).len());
physical_facts(&meta, &mut entry.cost);
}
}
pub fn enrich_tables(entry: &mut Entry) {
if entry.kind != EntryKind::File || entry.table.is_some() {
return;
}
let named = data_format(&entry.path);
if !is_regular_file(&entry.path) {
return;
}
let Some(format) = crate::formats::members::holder(&entry.path) else {
if named.is_some_and(|f| f.holds_tables() && crate::formats::readers::of(f).bytes_decide) {
entry.kind = EntryKind::Other;
}
return;
};
let Ok(tables) = crate::formats::members::tables(&entry.path, format) else {
return;
};
let own: Vec<&crate::formats::sqlite::Table> = tables.iter().filter(|t| !t.internal).collect();
entry.cost.tables = Some(own.len());
entry.cost.opens_one = format.opens_one_table();
if let [one] = own.as_slice()
&& !one.columns.is_empty()
{
entry.columns = one.columns.iter().map(|(name, _)| name.clone()).collect();
entry.cols = Some(entry.columns.len());
}
}
pub fn database_rows(file: &Path) -> Vec<Entry> {
let Some(format) = crate::formats::members::holder(file) else {
return Vec::new();
};
let Ok(mut tables) = crate::formats::members::tables(file, format) else {
return Vec::new();
};
if format
.descriptor()
.tables
.as_ref()
.is_some_and(|t| t.by_name)
{
tables.sort_by_cached_key(|t| t.name.to_lowercase());
}
let modified = std::fs::metadata(file).and_then(|m| m.modified()).ok();
tables
.into_iter()
.map(|table| table_entry(file, format, table, modified))
.collect()
}
pub fn split_rows(dir: &Path) -> Vec<Entry> {
let splits = crate::formats::hf_splits::cache_splits(dir);
if splits.len() < 2 {
return Vec::new();
}
splits
.into_iter()
.map(|split| split_entry(dir, split))
.collect()
}
pub fn split_row(path: &Path) -> Option<Entry> {
let (dir, split) = crate::formats::hf_splits::split_place(path)?;
Some(split_entry(&dir, split))
}
fn split_entry(dir: &Path, split: String) -> Entry {
let mut entry = Entry::new(dir.join(&split), EntryKind::File).with_name(split);
entry.table = Some(TableOf {
format: Some(crate::FileFormat::Arrow),
kind: "split".to_string(),
internal: false,
});
entry
}
pub fn variant_rows(file: &Path, formats: &crate::formats::Registry) -> Vec<Entry> {
let Some((spec, tables)) = crate::formats::members::variants(file, formats) else {
return Vec::new();
};
let modified = std::fs::metadata(file).and_then(|m| m.modified()).ok();
tables
.into_iter()
.map(|table| variant_entry(file, &spec, table, modified))
.collect()
}
pub fn variant_row(path: &Path, formats: &crate::formats::Registry) -> Option<Entry> {
let (file, name) = crate::formats::members::split_variant(path, formats)?;
let (spec, tables) = crate::formats::members::variants(&file, formats)?;
let table = tables.into_iter().find(|t| t.name == name)?;
let modified = std::fs::metadata(&file).and_then(|m| m.modified()).ok();
let mut entry = variant_entry(&file, &spec, table, modified);
entry.path = path.to_path_buf();
Some(entry)
}
fn variant_entry(
file: &Path,
spec: &str,
table: crate::formats::sqlite::Table,
modified: Option<std::time::SystemTime>,
) -> Entry {
let mut entry = Entry::new(
crate::formats::members::place(file, &table.name),
EntryKind::File,
)
.with_name(table.name);
entry.modified = modified;
entry.columns = table.columns.into_iter().map(|(name, _)| name).collect();
entry.cols = (!entry.columns.is_empty()).then_some(entry.columns.len());
entry.format_spec = Some(spec.to_string());
entry.table = Some(TableOf {
format: None,
kind: table.kind,
internal: false,
});
entry
}
pub fn table_row(path: &Path) -> Option<Entry> {
let (file, name) = crate::formats::members::split(path)?;
let format = crate::formats::members::holder(&file)?;
let table = crate::formats::members::tables(&file, format)
.ok()?
.into_iter()
.find(|t| t.name == name)?;
let modified = std::fs::metadata(&file).and_then(|m| m.modified()).ok();
let mut entry = table_entry(&file, format, table, modified);
entry.path = path.to_path_buf();
Some(entry)
}
fn table_entry(
file: &Path,
format: crate::FileFormat,
table: crate::formats::sqlite::Table,
modified: Option<std::time::SystemTime>,
) -> Entry {
let mut entry = Entry::new(
crate::formats::members::place(file, &table.name),
EntryKind::File,
)
.with_name(table.name);
entry.modified = modified;
entry.columns = table.columns.into_iter().map(|(name, _)| name).collect();
entry.cols = (!entry.columns.is_empty()).then_some(entry.columns.len());
entry.table = Some(TableOf {
format: Some(format),
kind: table.kind,
internal: table.internal,
});
entry
}
fn read_head<'a>(path: &Path, buf: &'a mut [u8]) -> Option<&'a [u8]> {
use std::io::Read;
let mut file = std::fs::File::open(path).ok()?;
let mut filled = 0;
loop {
match file.read(&mut buf[filled..]) {
Ok(0) => break,
Ok(n) => filled += n,
Err(_) => return None,
}
if filled == buf.len() {
break;
}
}
Some(&buf[..filled])
}
fn enrich_arrow(entry: &mut Entry) {
if entry.kind != EntryKind::File
|| data_format(&entry.path) != Some(crate::FileFormat::Arrow)
|| crate::CompressionFormat::from_extension(&entry.path).is_some()
|| !is_regular_file(&entry.path)
{
return;
}
let mut head = [0u8; 8];
if let Some(head) = read_head(&entry.path, &mut head) {
entry.cost.ipc_stream = !head.starts_with(b"ARROW1");
}
}
pub fn physical_facts(meta: &crate::formats::parquet_footer::Footer, cost: &mut Cost) {
if meta.row_groups.is_empty() {
return;
}
cost.row_groups = Some(meta.row_groups.len());
let mut uncompressed: u64 = 0;
let mut codecs: Vec<String> = Vec::new();
for rg in &meta.row_groups {
uncompressed = uncompressed.saturating_add(rg.total_byte_size() as u64);
for cc in rg.parquet_columns() {
let codec = format!("{:?}", cc.compression()).to_lowercase();
if !codecs.contains(&codec) {
codecs.push(codec);
}
}
}
if uncompressed > 0 {
cost.uncompressed = Some(uncompressed);
}
cost.codec = match codecs.len() {
0 => None,
1 => Some(codecs.remove(0)),
n => Some(format!("mixed ({n})")),
};
}
const MAX_PARTITION_DIRS: usize = 512;
pub fn partition_layout(dir: &Path) -> Option<Partitions> {
let iter = std::fs::read_dir(dir).ok()?;
let mut values: Vec<String> = Vec::new();
let mut keys: Vec<String> = Vec::new();
let mut count = 0usize;
let mut more = false;
for entry in iter.flatten() {
if count >= MAX_PARTITION_DIRS {
more = true;
break;
}
let name = entry.file_name().to_string_lossy().into_owned();
let Some((key, value)) = name.split_once('=') else {
continue;
};
if !entry.path().is_dir() {
continue;
}
if keys.is_empty() {
keys.push(key.to_string());
keys.extend(nested_keys(&entry.path()));
}
values.push(value.to_string());
count += 1;
}
if keys.is_empty() {
return None;
}
values.sort();
values.dedup();
Some(Partitions {
keys,
first_key_values: values,
count,
more,
})
}
fn nested_keys(dir: &Path) -> Vec<String> {
let mut keys = Vec::new();
let mut current = dir.to_path_buf();
for _ in 0..6 {
let Ok(iter) = std::fs::read_dir(¤t) else {
break;
};
let Some(child) = iter
.flatten()
.find(|e| e.file_name().to_string_lossy().contains('=') && e.path().is_dir())
else {
break;
};
let name = child.file_name().to_string_lossy().into_owned();
let Some((key, _)) = name.split_once('=') else {
break;
};
keys.push(key.to_string());
current = child.path();
}
keys
}
pub fn format_rows(rows: usize) -> String {
let r = rows as f64;
if rows >= 1_000_000_000 {
format!("{:.1}B", r / 1e9)
} else if rows >= 1_000_000 {
format!("{:.1}M", r / 1e6)
} else if rows >= 10_000 {
format!("{:.0}k", r / 1e3)
} else if rows >= 1_000 {
let mut out = String::new();
let digits = rows.to_string();
for (i, c) in digits.chars().enumerate() {
if i > 0 && (digits.len() - i).is_multiple_of(3) {
out.push(',');
}
out.push(c);
}
out
} else {
rows.to_string()
}
}
pub fn format_age(t: std::time::SystemTime) -> String {
let Ok(elapsed) = t.elapsed() else {
return String::new();
};
let secs = elapsed.as_secs();
if secs < 60 {
"now".to_string()
} else if secs < 3600 {
format!("{}m", secs / 60)
} else if secs < 86_400 {
format!("{}h", secs / 3600)
} else if secs < 86_400 * 365 {
format!("{}d", secs / 86_400)
} else {
format!("{}y", secs / (86_400 * 365))
}
}
pub type SchemaPreview = Vec<(String, polars::prelude::DataType)>;
fn table_preview(entry: &Entry) -> Option<Option<SchemaPreview>> {
let (file, format, name) = match &entry.table {
Some(table) => match (crate::formats::members::split(&entry.path), table.format) {
(Some((file, _)), Some(format)) => (file, format, Some(entry.name.as_str())),
_ => return Some(None),
},
None if is_regular_file(&entry.path) => {
let format = crate::formats::members::holder(&entry.path).or_else(|| {
data_format(&entry.path).filter(|f| {
f.holds_tables() && crate::formats::readers::of(*f).table_schema.is_some()
})
})?;
(entry.path.clone(), format, None)
}
None => return None,
};
Some(
crate::formats::readers::of(format)
.table_schema
.and_then(|schema| schema(&file, name)),
)
}
fn first_parquet_under(dir: &Path, depth: u8) -> Option<PathBuf> {
if depth > MAX_WALK_DEPTH {
return None;
}
let mut subdirs = Vec::new();
for entry in std::fs::read_dir(dir).ok()?.flatten().take(64) {
let path = entry.path();
if path.is_dir() {
subdirs.push(path);
} else if is_parquet_path(&path) && is_regular_file(&path) {
return Some(path);
}
}
subdirs.sort();
subdirs
.into_iter()
.take(4)
.find_map(|d| first_parquet_under(&d, depth + 1))
}
pub fn column_names(meta: &crate::formats::parquet_footer::Footer) -> Vec<String> {
meta.schema_descr
.columns()
.iter()
.map(|c| c.path_in_schema.join("."))
.collect()
}
fn is_regular_file(path: &Path) -> bool {
std::fs::metadata(path)
.map(|m| m.file_type().is_file())
.unwrap_or(false)
}
pub fn schema_preview(entry: &Entry) -> Option<SchemaPreview> {
use polars::prelude::{ParquetReader, Schema, SchemaExt, SerReader};
let file_path = match entry.kind {
EntryKind::File => {
if let Some(preview) = table_preview(entry) {
return preview;
}
if !is_parquet_path(&entry.path) {
return None;
}
entry.path.clone()
}
EntryKind::Hive | EntryKind::MultiFile => first_parquet_under(&entry.path, 0)?,
EntryKind::Directory | EntryKind::Unknown | EntryKind::Other => return None,
EntryKind::Delta | EntryKind::Iceberg | EntryKind::Hudi => return None,
};
if !is_regular_file(&file_path) {
return None;
}
let file = std::fs::File::open(&file_path).ok()?;
let mut reader = ParquetReader::new(file);
let arrow_schema = reader.schema().ok()?;
let schema = Schema::from_arrow_schema(arrow_schema.as_ref());
let mut preview: SchemaPreview = Vec::new();
if entry.kind == EntryKind::Hive
&& let Ok(below) = file_path.strip_prefix(&entry.path)
{
for part in below.parent().into_iter().flat_map(Path::components) {
let part = part.as_os_str().to_string_lossy();
if let Some((key, value)) = part.split_once('=')
&& !key.is_empty()
&& schema.get(key).is_none()
{
let dtype = if value.is_empty() || value == "__HIVE_DEFAULT_PARTITION__" {
polars::prelude::DataType::String
} else {
polars::io::csv::read::schema_inference::infer_field_schema(value, true, false)
};
preview.push((key.to_string(), dtype));
}
}
}
preview.extend(
schema
.iter()
.map(|(name, dtype)| (name.to_string(), dtype.clone())),
);
Some(preview)
}
#[cfg(test)]
mod classification_tests;