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;
pub(crate) fn shorten(name: &str, width: usize) -> String {
let chars: Vec<char> = name.chars().collect();
if chars.len() <= width {
return name.to_string();
}
let ellipsis = crate::glyphs::get().ellipsis;
let room = width.saturating_sub(ellipsis.chars().count());
let head = room.div_ceil(2);
let tail = room - head;
format!(
"{}{ellipsis}{}",
chars[..head].iter().collect::<String>(),
chars[chars.len() - tail..].iter().collect::<String>()
)
}
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>,
}
#[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::source::input_source(&entry.path) {
crate::source::InputSource::Local(_) => false,
crate::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 {
let mut entry = Self::new(path.to_path_buf(), EntryKind::File);
entry.name = name.to_string();
entry
}
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,
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::readers::sniff_file(path, crate::readers::Asked::Listing)
}
#[derive(Debug, Clone)]
pub enum Sniffed {
Format(crate::FileFormat),
Spec(std::sync::Arc<crate::formats::Spec>),
}
pub fn sniff_listed(path: &Path, formats: &crate::formats::Registry) -> Option<Sniffed> {
use crate::readers::{Asked, HEAD, head_of, sniff};
let head = head_of(path)?;
if let Some(format) = sniff(&head, Some(path), Asked::Listing, |_| true) {
return Some(Sniffed::Format(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 mut picks = vec![0, nameless.len() / 2, nameless.len() - 1];
picks.dedup();
let sniffed: Vec<crate::FileFormat> = picks
.iter()
.filter_map(|i| nameless.get(*i))
.filter_map(|f| sniff_format(f))
.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 classify_directory(path: &Path) -> EntryKind {
look_at_directory(path).0
}
pub fn look_at_directory(path: &Path) -> (EntryKind, Holds) {
let mut holds = Holds::default();
if let Some(lake) = lake_table(path) {
return (lake, holds);
}
let Ok(iter) = std::fs::read_dir(path) else {
return (EntryKind::Directory, holds);
};
let mut partitions = 0usize;
let mut data_files = 0usize;
let mut seen = 0usize;
let mut sniffs_left = if crate::home::is_remote_path(path) {
0
} else {
MAX_SNIFFS_PER_DIR
};
let mut counts: Vec<(crate::FileFormat, usize)> = Vec::new();
let mut format: Option<crate::FileFormat> = None;
let mut mixed_formats = false;
let mut hugging_face: Vec<String> = Vec::new();
for (entries, entry) in iter.flatten().take(MAX_ENTRIES_PER_DIR + 1).enumerate() {
if entries >= MAX_ENTRIES_PER_DIR {
holds.truncated = true;
break;
}
let entry_path = entry.path();
let name = entry.file_name();
let name = name.to_string_lossy().into_owned();
if is_bookkeeping(&name) {
holds.skipped += 1;
if holds.skipped_names.last().is_none_or(|last| &name < last)
|| holds.skipped_names.len() < SKIPPED_NAMES_SHOWN
{
holds.skipped_names.push(name);
holds.skipped_names.sort();
holds.skipped_names.truncate(SKIPPED_NAMES_SHOWN);
}
continue;
}
let followed = |path: &Path| {
std::fs::metadata(path).map_or((false, false), |m| (m.is_dir(), m.is_file()))
};
let (is_dir, is_file) = match entry.file_type() {
Ok(kind) if kind.is_symlink() => followed(&entry_path),
Ok(kind) => (kind.is_dir(), kind.is_file()),
Err(_) => followed(&entry_path),
};
if is_dir {
holds.directories += 1;
if is_partition_dir(&entry_path) {
partitions += 1;
}
} else if let Some(found) = data_format(&entry_path)
.filter(|f| !f.is_lines())
.map(
|found| match crate::model_files::is_safetensors_index(&entry_path) {
true => crate::FileFormat::Json,
false => found,
},
)
.or_else(|| {
is_parquet_key(&directory_and_name(&entry_path))
.then_some(crate::FileFormat::Parquet)
})
.or_else(|| {
(is_file && sniffs_left > 0 && worth_sniffing(&entry_path)).then(|| {
sniffs_left -= 1;
sniff_format(&entry_path)
})?
})
.or_else(|| data_format(&entry_path))
.filter(|_| is_file)
{
if found == crate::FileFormat::Json && is_hugging_face_metadata(&name) {
hugging_face.push(name.clone());
}
data_files += 1;
match counts.iter_mut().find(|(f, _)| *f == found) {
Some((_, n)) => *n += 1,
None => counts.push((found, 1)),
}
match format {
None => format = Some(found),
Some(first) if first != found => mixed_formats = true,
Some(_) => {}
}
} else if is_file && has_no_extension(&entry_path) {
holds.unnamed += 1;
} else {
holds.not_read += 1;
}
seen += 1;
}
if !hugging_face.is_empty() && counts.iter().any(|(f, _)| *f == crate::FileFormat::Arrow) {
let n = hugging_face.len();
for (format, count) in &mut counts {
if *format == crate::FileFormat::Json {
*count -= n;
}
}
counts.retain(|(_, count)| *count > 0);
data_files -= n;
seen -= n;
holds.skipped += n;
holds.skipped_names.extend(hugging_face);
holds.skipped_names.sort();
holds.skipped_names.truncate(SKIPPED_NAMES_SHOWN);
mixed_formats = counts.len() > 1;
format = counts.first().map(|(f, _)| *f);
}
if counts.iter().any(|(f, _)| !f.is_lines())
&& let Some(at) = counts.iter().position(|(f, _)| f.is_lines())
{
let (_, n) = counts.remove(at);
data_files -= n;
holds.not_read += n;
mixed_formats = counts.len() > 1;
format = counts.first().map(|(f, _)| *f);
}
holds.partitions = partitions;
order_formats(&mut counts);
holds.formats = counts
.into_iter()
.map(|(f, n)| (f.name().to_string(), n))
.collect();
if partitions > 0 && partitions >= data_files {
return (EntryKind::Hive, holds);
}
let readable_as_one = format.is_some_and(crate::FileFormat::reads_many_files);
let homogeneous = data_files > 1 && !mixed_formats && readable_as_one;
let mostly_data = data_files * 2 >= seen;
let model = partitions == 0 && is_model_directory(counts_names(&holds));
let kind = if (homogeneous && mostly_data) || model {
EntryKind::MultiFile
} else {
EntryKind::Directory
};
(kind, holds)
}
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
}
pub fn scan_dir(dir: &Path) -> Vec<Entry> {
scan_dir_bounded(dir).entries
}
#[derive(Debug, Clone, Default)]
pub struct Scan {
pub entries: Vec<Entry>,
pub truncated: bool,
}
pub fn scan_dir_bounded(dir: &Path) -> Scan {
scan_dir_progressive(dir, |_| {})
}
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 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 meta.is_file() && sniffs_left > 0 && worth_sniffing(&path) {
sniffs_left -= 1;
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 {
let mut so_far = entries.clone();
sort_entries(&mut so_far);
progress(&so_far);
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::schema_union::ReadAs::default())
}
pub fn enrich_as(entry: &mut Entry, as_read: &crate::schema_union::ReadAs) {
enrich_with(entry, as_read, None)
}
pub fn enrich_with(
entry: &mut Entry,
as_read: &crate::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::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::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();
if entry.kind == EntryKind::MultiFile && !files_nest(&names) {
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::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::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::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()
}
fn sample_footers(files: &[PathBuf]) -> Vec<crate::parquet_footer::Footer> {
if files.is_empty() {
return Vec::new();
}
let mut picks = vec![0, files.len() / 2, files.len() - 1];
picks.dedup();
picks
.iter()
.filter_map(|i| files.get(*i))
.filter_map(|file| crate::parquet_footer::read_parquet_metadata(file))
.collect()
}
fn files_nest(sampled: &[Vec<String>]) -> bool {
let per_file: Vec<Vec<String>> = sampled
.iter()
.map(|names| crate::schema_union::top_level_columns(names))
.collect();
per_file.len() < 2 || crate::schema_union::is_nested(&per_file)
}
fn judge_by_names(entry: &mut Entry, as_read: &crate::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::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::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::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::dataset_files::DatasetFile],
footers: &[Option<crate::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::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::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::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::members::holder(&entry.path) else {
if named.is_some_and(|f| f.holds_tables() && crate::readers::of(f).bytes_decide) {
entry.kind = EntryKind::Other;
}
return;
};
let Ok(tables) = crate::members::tables(&entry.path, format) else {
return;
};
let own: Vec<&crate::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::members::holder(file) else {
return Vec::new();
};
let Ok(mut tables) = crate::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::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::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);
entry.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::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::members::split_variant(path, formats)?;
let (spec, tables) = crate::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::sqlite::Table,
modified: Option<std::time::SystemTime>,
) -> Entry {
let mut entry = Entry::new(crate::members::place(file, &table.name), EntryKind::File);
entry.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::members::split(path)?;
let format = crate::members::holder(&file)?;
let table = crate::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::sqlite::Table,
modified: Option<std::time::SystemTime>,
) -> Entry {
let mut entry = Entry::new(crate::members::place(file, &table.name), EntryKind::File);
entry.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::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_size(bytes: u64) -> String {
const UNITS: &[&str] = &["B", "KB", "MB", "GB", "TB"];
let mut value = bytes as f64;
let mut unit = 0;
while value >= 1024.0 && unit < UNITS.len() - 1 {
value /= 1024.0;
unit += 1;
}
if unit == 0 {
format!("{} {}", bytes, UNITS[0])
} else if value >= 100.0 {
format!("{:.0} {}", value, UNITS[unit])
} else {
format!("{:.1} {}", value, UNITS[unit])
}
}
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::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::members::holder(&entry.path).or_else(|| {
data_format(&entry.path)
.filter(|f| f.holds_tables() && crate::readers::of(*f).table_schema.is_some())
})?;
(entry.path.clone(), format, None)
}
None => return None,
};
Some(
crate::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::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 {
use super::*;
use polars::prelude::*;
#[test]
fn how_a_row_is_read() {
use crate::ReadMode::*;
let how = |path: &str| how_read(&Entry::for_test(Path::new(path), path));
let at = |path: &str, mode, download| {
assert_eq!(how(path), Some(HowRead { mode, download }), "{path}");
};
at("/d/a.parquet", Lazy, false);
at("/d/a.csv", Lazy, false);
at("/d/a.csv.gz", Decompressed, false);
at("/d/a.json", InMemory, false);
at("/d/a.gpx", Converted, false);
at("/d/a.arrow", Lazy, false);
at("s3://b/a.parquet", Lazy, false);
at("s3://b/a.csv", Lazy, true);
at("gs://b/a.json", InMemory, true);
at("https://example.com/a.parquet", Lazy, true);
at("s3://b/m.safetensors", InMemory, false);
at("https://example.com/m.gguf", InMemory, false);
assert_eq!(how("/d/a.parquet.gz"), None, "does not open");
assert_eq!(how("/d/README"), None);
let mut stream = Entry::for_test(Path::new("/d/x.arrow"), "x.arrow");
stream.cost.ipc_stream = true;
assert_eq!(how_read(&stream).map(|h| h.mode), Some(Converted));
let mut spec = Entry::for_test(Path::new("/d/day.l2.zst"), "day.l2.zst");
spec.format_spec = Some("acme.l2feed".into());
assert_eq!(how_read(&spec).map(|h| h.mode), Some(Decompressed));
at("/d/shop.db", Lazy, false);
at("s3://b/shop.sqlite", Lazy, true);
let mut table = Entry::for_test(Path::new("/d/shop.db/orders"), "orders");
table.table = Some(TableOf {
format: Some(crate::FileFormat::Sqlite),
kind: "table".into(),
internal: false,
});
assert_eq!(how_read(&table).map(|h| h.mode), Some(Lazy));
assert_eq!(how_read(&Entry::directory(Path::new("/d/x"))), None);
}
#[test]
fn measuring_an_arrow_file_tells_a_stream() {
let dir = tempfile::tempdir().unwrap();
let file = dir.path().join("file.arrow");
std::fs::write(&file, b"ARROW1\0\0rest").unwrap();
let stream = dir.path().join("stream.arrow");
std::fs::write(&stream, b"\xff\xff\xff\xff\x10\x01\0\0").unwrap();
for (path, is_stream) in [(file, false), (stream, true)] {
let mut entry = Entry::for_test(&path, "x.arrow");
enrich(&mut entry);
assert_eq!(entry.cost.ipc_stream, is_stream, "{}", path.display());
}
}
#[test]
fn what_is_offered_and_what_opens_are_one_list() {
for ext in [
"parquet", "csv", "tsv", "psv", "json", "jsonl", "ndjson", "arrow", "arrows", "ipc",
"feather", "avro", "orc", "xls", "xlsx", "xlsm", "xlsb",
] {
let named = PathBuf::from(format!("sales.{ext}"));
assert!(
data_format(&named).is_some(),
".{ext} opens, so the home screen must offer it"
);
}
assert_eq!(
data_format(Path::new("README.txt")),
Some(crate::FileFormat::Text)
);
assert!(data_format(Path::new("notes")).is_none());
}
#[test]
fn a_format_name_round_trips_only_through_from_name() {
use crate::FileFormat;
for format in FileFormat::ALL {
assert_eq!(
FileFormat::from_name(format.name()),
Some(format),
"{} is a name",
format.name()
);
}
assert_eq!(FileFormat::from_extension("excel"), None);
assert_eq!(FileFormat::from_name("xlsx"), None);
}
#[test]
fn a_compressed_name_reads_as_the_format_under_it() {
assert_eq!(
data_format(Path::new("sales.csv.gz")),
Some(crate::FileFormat::Csv)
);
assert_eq!(
data_format(Path::new("events.json.zst")),
Some(crate::FileFormat::Json)
);
}
#[test]
fn one_format_under_several_names_is_not_a_mixture() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("a.arrow"), b"x").unwrap();
std::fs::write(dir.path().join("b.ipc"), b"x").unwrap();
assert_eq!(classify_directory(dir.path()), EntryKind::MultiFile);
}
#[test]
fn a_file_datui_does_not_read_does_not_disqualify_a_directory() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "a.parquet", &["id"]);
write(dir.path(), "b.parquet", &["id"]);
std::fs::write(dir.path().join("README.txt"), b"notes").unwrap();
#[cfg(feature = "cloud")]
let objects: Vec<(String, u64)> = [
("out/a.parquet", 100u64),
("out/b.parquet", 100),
("out/README.txt", 12),
]
.iter()
.map(|(k, s)| ((*k).to_string(), *s))
.collect();
#[cfg(feature = "cloud")]
assert_eq!(
classify_directory(dir.path()),
crate::cloud_browse::look_at_listing("out/", &[], &objects).0,
"the two routes answer the same directory alike"
);
assert_eq!(classify_directory(dir.path()), EntryKind::MultiFile);
}
#[test]
fn a_label_counts_what_is_there_rather_than_naming_a_decision() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.parquet", "b.parquet", "c.parquet"] {
write(dir.path(), name, &["id"]);
}
std::fs::create_dir_all(dir.path().join("archive")).unwrap();
std::fs::write(dir.path().join("notes.csv"), b"x").unwrap();
std::fs::write(dir.path().join("_SUCCESS"), b"").unwrap();
std::fs::write(dir.path().join(".part.crc"), b"").unwrap();
let entry = measured(dir.path());
assert_eq!(entry.label(), "mixed", "two formats is two formats");
assert_eq!(
entry.holds.line(true).as_deref(),
Some("3 parquet · 1 csv · 1 directory"),
"and the pane says what the label boiled down"
);
assert_eq!(entry.holds.data_files(), 4);
assert_eq!(entry.holds.directories, 1);
}
#[test]
fn a_directory_of_one_format_is_labelled_by_it() {
let dir = tempfile::tempdir().unwrap();
for i in 0..12 {
write(dir.path(), &format!("part-{i:05}.parquet"), &["id", "ts"]);
}
let entry = measured(dir.path());
assert_eq!(entry.label(), "12 parquet");
assert_eq!(entry.holds.line(true).as_deref(), Some("12 parquet"));
let plain = tempfile::tempdir().unwrap();
for i in 0..20 {
std::fs::write(plain.path().join(format!("note{i}.md")), b"x").unwrap();
}
let plain = measured(plain.path());
assert_eq!(plain.label(), "dir");
assert_eq!(plain.holds.not_read, 20);
assert_eq!(plain.holds.line(true), None);
}
#[test]
fn a_kind_that_names_itself_keeps_its_name() {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join("year=2024")).unwrap();
std::fs::create_dir_all(dir.path().join("year=2025")).unwrap();
assert_eq!(measured(dir.path()).label(), "hive");
let lake = tempfile::tempdir().unwrap();
std::fs::create_dir_all(lake.path().join("_delta_log")).unwrap();
assert_eq!(measured(lake.path()).label(), "delta");
let unlooked = Entry::new(PathBuf::from("/nowhere"), EntryKind::Unknown);
assert_eq!(unlooked.label(), "");
}
#[test]
fn a_directory_read_whole_does_not_claim_there_is_more() {
let dir = tempfile::tempdir().unwrap();
for i in 0..MAX_ENTRIES_PER_DIR {
std::fs::write(dir.path().join(format!("f{i:05}.csv")), b"x").unwrap();
}
let holds = look_at_directory(dir.path()).1;
assert!(!holds.truncated, "every entry was read");
assert_eq!(holds.label(), format!("{MAX_ENTRIES_PER_DIR} csv"));
std::fs::write(dir.path().join("one-more.csv"), b"x").unwrap();
let holds = look_at_directory(dir.path()).1;
assert!(holds.truncated, "and now there is more than was read");
assert!(holds.label().contains('+'));
}
#[test]
fn a_lake_table_is_not_counted() {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join("_delta_log")).unwrap();
write(dir.path(), "part-00000.parquet", &["id"]);
write(dir.path(), "part-00001.parquet", &["id"]);
let (kind, holds) = look_at_directory(dir.path());
assert_eq!(kind, EntryKind::Delta);
assert!(holds.is_empty(), "and its label is the format's own name");
let entry = measured(dir.path());
assert_eq!(entry.label(), "delta");
}
#[test]
fn a_cut_short_listing_does_not_claim_a_directory_is_empty() {
let seen = Holds {
skipped: 5000,
truncated: true,
..Default::default()
};
assert_eq!(seen.label(), "dir+");
let whole = Holds {
skipped: 3,
..Default::default()
};
assert_eq!(whole.label(), "dir");
let nothing_yet = Holds {
truncated: true,
..Default::default()
};
assert!(!nothing_yet.is_empty());
assert_eq!(nothing_yet.label(), "dir+");
assert!(Holds::default().is_empty());
let mixed = Holds {
formats: vec![("parquet".to_string(), 3), ("csv".to_string(), 2)],
truncated: true,
..Default::default()
};
assert_eq!(mixed.label(), "mixed", "more files cannot unmake it");
}
#[test]
fn a_long_name_keeps_both_ends() {
let name = ".part-00000-8f3a91c2-7b4d-4e19-a6f0-c1d2e3f4a5b6-c000.snappy.parquet.crc";
let line = shorten(name, 24);
assert!(line.starts_with(".part-00000"), "the head: {line}");
assert!(line.ends_with(".crc"), "and the tail: {line}");
assert!(!line.contains("8f3a91c2"), "the middle goes: {line}");
assert!(line.chars().count() <= 24, "{line}");
}
#[test]
fn what_a_directory_holds_reads_the_same_twice() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "a.parquet", &["id"]);
for marker in [
"_SUCCESS",
"_committed_9",
"_committed_1",
".crc",
"_started_4",
] {
std::fs::write(dir.path().join(marker), b"").unwrap();
}
let first = look_at_directory(dir.path()).1;
for _ in 0..8 {
assert_eq!(look_at_directory(dir.path()).1, first);
}
assert_eq!(
first.skipped_names,
vec![".crc", "_SUCCESS", "_committed_1", "_committed_9"],
"the first four by name, of five"
);
assert_eq!(first.skipped, 5);
}
#[test]
fn a_directories_numbers_are_what_opening_it_gives() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.parquet", "b.parquet", "c.parquet"] {
write(dir.path(), name, &["id", "legacy"]);
}
let archive = dir.path().join("archive");
std::fs::create_dir_all(&archive).unwrap();
for i in 0..20 {
write(&archive, &format!("old-{i}.parquet"), &["id", "legacy"]);
}
let entry = measured(dir.path());
assert_eq!(
entry.label(),
"3 parquet",
"three files are directly inside"
);
assert_eq!(
entry.holds.line(true).as_deref(),
Some("3 parquet · 1 directory")
);
assert_eq!(entry.rows, Some(23), "and opening it reads all of them");
}
#[test]
fn a_width_is_a_floor_only_when_a_footer_went_unread() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.parquet", "b.parquet", "c.parquet"] {
write(dir.path(), name, &["id", "ts"]);
}
let archive = dir.path().join("archive");
std::fs::create_dir_all(&archive).unwrap();
for i in 0..MAX_FOOTERS_PER_DATASET + 6 {
write(
&archive,
&format!("old-{i:03}.parquet"),
&["wholly", "different"],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Directory, "not one table");
assert_eq!(entry.label(), "3 parquet");
assert_eq!(entry.cols, Some(2), "id and ts");
assert!(
!entry.cols_sampled,
"all three of its own footers were read"
);
}
#[test]
fn a_big_directory_is_still_found_by_a_column_one_level_down() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.parquet", "b.parquet", "c.parquet"] {
write(dir.path(), name, &["id", "ts"]);
}
let archive = dir.path().join("archive");
std::fs::create_dir_all(&archive).unwrap();
for i in 0..MAX_FOOTERS_PER_DATASET + 6 {
write(
&archive,
&format!("old-{i:03}.parquet"),
&["wholly", "different"],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Directory);
assert_eq!(entry.label(), "3 parquet");
assert_eq!(entry.cols, Some(2), "id and ts");
assert!(
entry.columns.contains(&"wholly".to_string()),
"{:?}",
entry.columns
);
assert!(entry.columns.contains(&"id".to_string()));
}
#[test]
fn a_width_over_a_directories_own_files_is_a_floor_when_there_are_too_many() {
let dir = tempfile::tempdir().unwrap();
for i in 0..MAX_FOOTERS_PER_DATASET + 6 {
write(
dir.path(),
&format!("f-{i:03}.parquet"),
&[&format!("c{i}")],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Directory, "not one table");
assert_eq!(entry.label(), "70 parquet");
assert!(
entry.cols_sampled,
"three of seventy footers were read, so the width is a floor"
);
}
#[test]
fn a_directory_read_as_one_table_is_sized_by_everything_under_it() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "a.parquet", &["id", "ts"]);
write(dir.path(), "b.parquet", &["id", "ts"]);
let more = dir.path().join("more");
std::fs::create_dir_all(&more).unwrap();
write(&more, "c.parquet", &["id", "ts"]);
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile, "one table");
let all: u64 = [
dir.path().join("a.parquet"),
dir.path().join("b.parquet"),
more.join("c.parquet"),
]
.iter()
.map(|p| std::fs::metadata(p).unwrap().len())
.sum();
assert_eq!(entry.size, Some(all));
assert_eq!(entry.rows, Some(3));
}
#[test]
fn nothing_counted_is_the_only_thing_holds_calls_empty() {
assert!(Holds::default().is_empty());
let one = |f: fn(&mut Holds)| {
let mut h = Holds::default();
f(&mut h);
h
};
for (what, holds) in [
("a data file", one(|h| h.formats.push(("csv".into(), 1)))),
("a directory", one(|h| h.directories = 1)),
("a partition", one(|h| h.partitions = 1)),
("a file it cannot read", one(|h| h.not_read = 1)),
("a writer's own file", one(|h| h.skipped = 1)),
(
"the name of one",
one(|h| h.skipped_names.push("_SUCCESS".into())),
),
("a listing cut short", one(|h| h.truncated = true)),
] {
assert!(!holds.is_empty(), "{what} is something to say");
}
}
#[test]
fn formats_that_tie_are_ordered_by_name_whatever_order_they_arrived_in() {
use crate::FileFormat;
let mut counts = vec![
(FileFormat::Json, 2),
(FileFormat::Csv, 2),
(FileFormat::Parquet, 5),
];
order_formats(&mut counts);
assert_eq!(
counts,
vec![
(FileFormat::Parquet, 5),
(FileFormat::Csv, 2),
(FileFormat::Json, 2)
]
);
}
#[test]
fn the_label_and_the_read_pick_the_same_format_on_a_tie() {
use crate::FileFormat;
let tmp = tempfile::TempDir::new().unwrap();
for name in ["a.csv", "b.csv", "c.parquet", "d.parquet"] {
std::fs::write(tmp.path().join(name), b"x").unwrap();
}
let (_, holds) = look_at_directory(tmp.path());
assert_eq!(
holds.formats.first().map(|(f, n)| (f.as_str(), *n)),
Some(("parquet", 2)),
"the label names Parquet first: {:?}",
holds.formats
);
match directory_format(tmp.path()) {
DirectoryFormat::Mixed { format, .. } => assert_eq!(
format,
FileFormat::Parquet,
"and so does the reader the open picks"
),
other => panic!("a directory of two formats is mixed, got {other:?}"),
}
}
#[test]
fn a_hugging_face_dataset_is_its_shards() {
use crate::FileFormat;
let tmp = tempfile::TempDir::new().unwrap();
for name in [
"data-00000-of-00002.arrow",
"data-00001-of-00002.arrow",
"dataset_info.json",
"state.json",
] {
std::fs::write(tmp.path().join(name), b"x").unwrap();
}
let (kind, holds) = look_at_directory(tmp.path());
assert_eq!(kind, EntryKind::MultiFile);
assert_eq!(holds.formats, [("arrow".to_string(), 2)]);
assert_eq!(holds.skipped, 2);
assert_eq!(holds.skipped_names, ["dataset_info.json", "state.json"]);
match directory_format(tmp.path()) {
DirectoryFormat::One(FileFormat::Arrow, files) => assert_eq!(files.len(), 2),
other => panic!("the shards are the dataset, got {other:?}"),
}
std::fs::remove_file(tmp.path().join("data-00001-of-00002.arrow")).unwrap();
assert!(matches!(
directory_format(tmp.path()),
DirectoryFormat::One(FileFormat::Arrow, _)
));
let json = tempfile::TempDir::new().unwrap();
for name in ["state.json", "other.json"] {
std::fs::write(json.path().join(name), b"{}").unwrap();
}
assert_eq!(
look_at_directory(json.path()).1.formats,
[("json".to_string(), 2)]
);
}
#[test]
fn partitions_carry_a_directory_only_while_they_are_the_most_of_it() {
let laid_out = |strays: usize| {
let dir = tempfile::tempdir().unwrap();
for year in ["year=2024", "year=2025"] {
let part = dir.path().join(year);
std::fs::create_dir_all(&part).unwrap();
write(&part, "data.parquet", &["id"]);
}
for i in 0..strays {
write(dir.path(), &format!("stray-{i}.parquet"), &["id"]);
}
classify_directory(dir.path())
};
assert_eq!(
laid_out(2),
EntryKind::Hive,
"two partitions against two files beside them"
);
assert_ne!(
laid_out(3),
EntryKind::Hive,
"one more file than partitions is a directory that holds a key=value"
);
}
#[test]
fn a_folder_marker_is_bookkeeping_even_beside_a_partition() {
assert!(is_bookkeeping("year=2024_$folder$"));
assert!(is_bookkeeping("alpha_$folder$"));
assert!(!is_bookkeeping("year=2024"), "the partition itself is data");
assert!(
!is_bookkeeping("_date=2024-01-01"),
"Spark partitions on internal columns"
);
}
#[test]
fn a_table_hidden_under_a_directory_still_downgrades_it() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.parquet", "b.parquet"] {
write(dir.path(), name, &["id", "ts"]);
}
let archive = dir.path().join("archive");
std::fs::create_dir_all(&archive).unwrap();
write(
&archive,
"other.parquet",
&["wholly", "different", "columns"],
);
let entry = measured(dir.path());
assert_eq!(
entry.kind,
EntryKind::Directory,
"a union over these is not one table"
);
assert_eq!(entry.rows, None);
assert_eq!(entry.label(), "2 parquet");
assert_eq!(entry.cols, Some(2), "id and ts");
assert!(entry.columns.contains(&"wholly".to_string()));
assert!(entry.columns.contains(&"id".to_string()));
let own: u64 = ["a.parquet", "b.parquet"]
.iter()
.map(|n| std::fs::metadata(dir.path().join(n)).unwrap().len())
.sum();
assert_eq!(entry.size, Some(own));
}
#[test]
fn a_hive_tree_of_another_format_is_still_laid_out() {
let dir = tempfile::tempdir().unwrap();
for year in ["year=2024", "year=2025"] {
let part = dir.path().join(year);
std::fs::create_dir_all(&part).unwrap();
std::fs::write(part.join("data.csv"), b"id\n1\n").unwrap();
}
std::fs::write(dir.path().join("summary.csv"), b"id\n1\n").unwrap();
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Hive);
assert!(entry.cost.partitions.is_some(), "the layout is named");
assert_eq!(entry.rows, None, "and nothing is invented about its rows");
}
#[test]
fn a_hive_dataset_is_described_despite_a_stray_file_at_its_root() {
let dir = tempfile::tempdir().unwrap();
for year in ["year=2024", "year=2025"] {
let part = dir.path().join(year);
std::fs::create_dir_all(&part).unwrap();
write(&part, "data.parquet", &["id"]);
}
std::fs::write(dir.path().join("schema.json"), b"{}").unwrap();
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Hive);
assert_eq!(entry.holds.one_format(), Some("json"), "its own only file");
assert_eq!(entry.rows, Some(2), "and the dataset is still counted");
assert_eq!(entry.cols, Some(2), "`id` and the partition column `year`");
}
#[test]
fn a_hive_dataset_is_described_despite_a_stray_file() {
let dir = tempfile::tempdir().unwrap();
for year in ["year=2024", "year=2025"] {
let part = dir.path().join(year);
std::fs::create_dir_all(&part).unwrap();
write(&part, "data.parquet", &["id"]);
}
std::fs::write(dir.path().join("year=2024/notes.csv"), b"x").unwrap();
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Hive);
assert_eq!(entry.rows, Some(2), "the dataset is still counted");
assert!(
entry.cost.partitions.is_some(),
"and its layout still named"
);
}
#[test]
fn a_directory_is_not_described_by_files_it_does_not_name() {
let dir = tempfile::tempdir().unwrap();
for name in ["a.json", "b.json", "c.json"] {
std::fs::write(dir.path().join(name), b"{}").unwrap();
}
let under = dir.path().join("derived");
std::fs::create_dir_all(&under).unwrap();
write(&under, "one.parquet", &["id", "ts", "amount"]);
write(&under, "two.parquet", &["id", "ts", "amount"]);
let entry = measured(dir.path());
assert_eq!(entry.label(), "3 json");
assert_eq!(
entry.cols, None,
"the Parquet below it is not this directory's shape"
);
assert_eq!(entry.rows, None);
assert!(entry.columns.is_empty());
}
#[test]
fn a_dataset_row_that_counted_nothing_is_still_described() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "a.parquet", &["id", "ts"]);
write(dir.path(), "b.parquet", &["id", "ts"]);
let mut entry = Entry {
kind: EntryKind::MultiFile,
..Entry::for_test(dir.path(), "data")
};
assert!(entry.holds.one_format().is_none(), "nothing counted");
enrich(&mut entry);
assert!(entry.size.is_some(), "the footers were read");
assert_eq!(entry.rows, Some(2));
assert_eq!(entry.cols, Some(2), "id and ts");
}
#[test]
fn a_parquet_file_named_like_a_writers_file_is_still_measured() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "_2024_sales.parquet", &["id", "amount"]);
let mut entry = Entry {
path: dir.path().join("_2024_sales.parquet"),
kind: EntryKind::File,
name: "_2024_sales.parquet".into(),
size: None,
modified: None,
rows: None,
cols: None,
cols_sampled: false,
columns: Vec::new(),
cost: Cost::default(),
holds: Default::default(),
opens_whole_directory: false,
format_spec: None,
table: None,
};
enrich(&mut entry);
assert_eq!(entry.rows, Some(1), "its footer was read");
assert_eq!(entry.cols, Some(2));
assert!(schema_preview(&entry).is_some(), "and the pane shows it");
assert!(is_bookkeeping("_2024_sales.parquet"));
assert_eq!(classify_directory(dir.path()), EntryKind::Directory);
}
#[test]
fn a_hive_preview_types_its_keys_as_the_scan_does() {
use polars::prelude::DataType;
let dir = tempfile::tempdir().unwrap();
let leaf = dir.path().join("day=2024-01-02/flag=true/n=3/x=1.5");
std::fs::create_dir_all(&leaf).unwrap();
write(&leaf, "part.parquet", &["id"]);
let mut entry = Entry::directory(dir.path());
entry.kind = EntryKind::Hive;
let preview = schema_preview(&entry).expect("a footer to read");
let types: Vec<(&str, &DataType)> = preview.iter().map(|(n, t)| (n.as_str(), t)).collect();
assert_eq!(
types[..4],
[
("day", &DataType::Date),
("flag", &DataType::Boolean),
("n", &DataType::Int64),
("x", &DataType::Float64),
]
);
assert_eq!(types[4].0, "id");
}
#[cfg(feature = "cloud")]
#[test]
fn extensionless_part_files_are_data_on_both_routes() {
let dir = tempfile::tempdir().unwrap();
let table = dir.path().join("occurrence.parquet");
std::fs::create_dir_all(&table).unwrap();
std::fs::write(table.join("000001"), b"PAR1").unwrap();
std::fs::write(table.join("000002"), b"PAR1").unwrap();
let objects: Vec<(String, u64)> = [
"gbif/occurrence.parquet/000001",
"gbif/occurrence.parquet/000002",
]
.iter()
.map(|k| ((*k).to_string(), 10u64))
.collect();
assert_eq!(
classify_directory(&table),
crate::cloud_browse::look_at_listing("gbif/occurrence.parquet/", &[], &objects).0,
"the two routes answer the same directory alike"
);
assert_eq!(classify_directory(&table), EntryKind::MultiFile);
}
#[test]
fn extensionless_part_files_are_measured_not_just_offered() {
let dir = tempfile::tempdir().unwrap();
let table = dir.path().join("occurrence.parquet");
std::fs::create_dir_all(&table).unwrap();
write(&table, "000001", &["id", "species"]);
write(&table, "000002", &["id", "species"]);
let entry = measured(&table);
assert_eq!(entry.kind, EntryKind::MultiFile);
assert_eq!(entry.rows, Some(2), "both footers were read");
assert_eq!(entry.cols, Some(2));
assert!(
schema_preview(&entry).is_some(),
"and the schema pane shows what those footers said, rather than asking \
for a full read of files already read"
);
let mut listed = scan_dir(&table);
assert_eq!(
listed.iter().map(|e| e.name.as_str()).collect::<Vec<_>>(),
vec!["000001", "000002"],
"the directory the label promises is not an empty listing"
);
let part = listed.first_mut().expect("a part file is listed");
enrich(part);
assert_eq!(part.rows, Some(1), "a part file counts its own rows");
assert_eq!(part.cols, Some(2));
}
#[cfg(feature = "cloud")]
#[test]
fn one_partition_beside_files_datui_cannot_read_answers_alike() {
let dir = tempfile::tempdir().unwrap();
std::fs::create_dir_all(dir.path().join("notes=old")).unwrap();
for note in ["README.md", "LICENSE", "logo.png"] {
std::fs::write(dir.path().join(note), b"x").unwrap();
}
let objects: Vec<(String, u64)> = ["out/README.md", "out/LICENSE", "out/logo.png"]
.iter()
.map(|k| ((*k).to_string(), 12u64))
.collect();
assert_eq!(
classify_directory(dir.path()),
crate::cloud_browse::look_at_listing("out/", &["out/notes=old/".to_string()], &objects)
.0,
"the two routes answer the same directory alike"
);
}
#[test]
fn a_partition_named_like_a_writers_file_is_still_a_partition() {
let dir = tempfile::tempdir().unwrap();
let mut directories = Vec::new();
for day in ["2024-01-01", "2024-01-02", "2024-01-03"] {
std::fs::create_dir_all(dir.path().join(format!("_date={day}"))).unwrap();
directories.push(format!("events/_date={day}/"));
}
#[cfg(feature = "cloud")]
assert_eq!(
classify_directory(dir.path()),
crate::cloud_browse::look_at_listing("events/", &directories, &[]).0,
"the two routes answer the same directory alike"
);
assert_eq!(classify_directory(dir.path()), EntryKind::Hive);
assert!(!is_bookkeeping("_date=2024-01-01"));
assert!(is_bookkeeping("_temporary"));
}
#[cfg(feature = "cloud")]
#[test]
fn a_writers_own_directory_is_skipped_on_both_routes() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "part-00000.parquet", &["id"]);
write(dir.path(), "part-00001.parquet", &["id"]);
std::fs::create_dir_all(dir.path().join("_temporary")).unwrap();
std::fs::create_dir_all(dir.path().join("notes")).unwrap();
std::fs::create_dir_all(dir.path().join("archive")).unwrap();
let local = classify_directory(dir.path());
let directories: Vec<String> = ["out/_temporary/", "out/notes/", "out/archive/"]
.iter()
.map(|f| (*f).to_string())
.collect();
let objects: Vec<(String, u64)> = [
("out/part-00000.parquet", 100u64),
("out/part-00001.parquet", 100),
]
.iter()
.map(|(k, s)| ((*k).to_string(), *s))
.collect();
let cloud = crate::cloud_browse::look_at_listing("out/", &directories, &objects).0;
assert_eq!(
local, cloud,
"the two routes answer the same directory alike"
);
assert_eq!(local, EntryKind::MultiFile);
}
#[cfg(unix)]
#[test]
fn every_entry_is_in_one_count() {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join("real.csv"), b"id\n1\n").unwrap();
std::os::unix::fs::symlink(dir.path().join("gone"), dir.path().join("broken.csv")).unwrap();
std::fs::write(dir.path().join("notes.md"), b"x").unwrap();
std::fs::write(dir.path().join("_SUCCESS"), b"").unwrap();
let holds = look_at_directory(dir.path()).1;
assert_eq!(holds.data_files(), 1);
assert_eq!(holds.not_read, 2, "the note and the broken link");
assert_eq!(holds.skipped, 1);
assert_eq!(holds.line(true).as_deref(), Some("1 csv"));
}
#[cfg(unix)]
#[test]
fn a_name_with_nothing_behind_it_is_not_a_data_file() {
let dir = tempfile::tempdir().unwrap();
std::os::unix::fs::symlink(dir.path().join("gone.csv"), dir.path().join("a.csv")).unwrap();
std::os::unix::fs::symlink(dir.path().join("gone.csv"), dir.path().join("b.csv")).unwrap();
assert_eq!(
classify_directory(dir.path()),
EntryKind::Directory,
"two broken symlinks are not a dataset"
);
}
#[test]
fn a_model_directory_is_its_weights() {
let dir = tempfile::tempdir().unwrap();
for name in [
"model-00001-of-00002.safetensors",
"model-00002-of-00002.safetensors",
"config.json",
"generation_config.json",
"tokenizer.json",
"tokenizer_config.json",
"model.safetensors.index.json",
] {
std::fs::write(dir.path().join(name), b"x").unwrap();
}
let DirectoryFormat::Mixed {
format,
files,
passed_over,
} = directory_format(dir.path())
else {
panic!("weights and JSON are two formats");
};
assert_eq!(format, crate::FileFormat::Safetensors);
assert_eq!(files.len(), 3, "the shards, and the index for its metadata");
assert_eq!(passed_over, [(crate::FileFormat::Json, 4)]);
let (kind, holds) = look_at_directory(dir.path());
assert_eq!(kind, EntryKind::MultiFile, "opened as one");
assert_eq!(holds.label(), "2 safetensors", "the shards, not the index");
assert_eq!(
data_format(Path::new("model.safetensors.index.json")),
Some(crate::FileFormat::Safetensors)
);
std::fs::write(dir.path().join("data.parquet"), b"x").unwrap();
let (kind, holds) = look_at_directory(dir.path());
assert_eq!(
(kind, holds.label().as_str()),
(EntryKind::Directory, "mixed")
);
}
#[test]
fn signed_files_are_sniffed_by_their_first_bytes() {
let dir = tempfile::tempdir().unwrap();
let gguf = dir.path().join("weights");
std::fs::write(&gguf, b"GGUF\x03\x00\x00\x00").unwrap();
let st = dir.path().join("checkpoint.bin");
let mut bytes = 2u64.to_le_bytes().to_vec();
bytes.extend_from_slice(b"{}");
std::fs::write(&st, &bytes).unwrap();
let text = dir.path().join("notes");
std::fs::write(&text, b"just some text").unwrap();
assert_eq!(sniff_format(&gguf), Some(crate::FileFormat::Gguf));
let opened = |path: &Path| crate::readers::sniff_open(path, None);
assert_eq!(opened(&st), Some(crate::FileFormat::Safetensors));
assert_eq!(opened(&text), None);
let midi = dir.path().join("song.bin");
std::fs::write(&midi, b"MThd\0\0\0\x06\0\0\0\x01\0\x60").unwrap();
assert_eq!(opened(&midi), Some(crate::FileFormat::Midi));
}
#[test]
fn a_format_that_cannot_be_read_as_many_is_not_offered_as_one() {
for ext in ["tsv", "psv", "xlsx", "xlsb"] {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join(format!("a.{ext}")), b"x").unwrap();
std::fs::write(dir.path().join(format!("b.{ext}")), b"x").unwrap();
assert_eq!(
classify_directory(dir.path()),
EntryKind::Directory,
"a directory of .{ext} has no reader that takes a list"
);
}
for ext in [
"parquet", "csv", "json", "jsonl", "ndjson", "arrow", "arrows", "ipc", "feather",
"avro", "orc",
] {
let dir = tempfile::tempdir().unwrap();
std::fs::write(dir.path().join(format!("a.{ext}")), b"x").unwrap();
std::fs::write(dir.path().join(format!("b.{ext}")), b"x").unwrap();
assert_eq!(
classify_directory(dir.path()),
EntryKind::MultiFile,
".{ext} reads as many files"
);
}
}
#[cfg(feature = "cloud")]
#[test]
fn a_writers_own_file_is_skipped_whatever_order_it_is_listed_in() {
let dir = tempfile::tempdir().unwrap();
for part in 0..8 {
write(dir.path(), &format!("{part}.parquet"), &["season"]);
}
std::fs::write(dir.path().join("_metadata.json"), b"{}").unwrap();
let local = classify_directory(dir.path());
let mut keys: Vec<(String, u64)> = vec![("jolpica/2000/_metadata.json".into(), 2)];
for part in 0..8 {
keys.push((format!("jolpica/2000/{part}.parquet"), 100));
}
keys.sort();
let cloud = crate::cloud_browse::look_at_listing("jolpica/2000/", &[], &keys).0;
assert_eq!(
local, cloud,
"the two routes answer the same directory alike"
);
assert_eq!(local, EntryKind::MultiFile);
}
#[test]
fn job_files_are_skipped_on_every_route() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "part-00000.parquet", &["id"]);
write(dir.path(), "part-00001.parquet", &["id"]);
for marker in [
"_SUCCESS",
"_committed_1727",
"_committed_1728",
"_started_1727",
".part.crc",
] {
std::fs::write(dir.path().join(marker), b"").unwrap();
}
assert_eq!(
classify_directory(dir.path()),
EntryKind::MultiFile,
"five markers beside two data files do not outvote them"
);
#[cfg(feature = "cloud")]
let keys: Vec<(String, u64)> = [
("out/_SUCCESS", 0u64),
("out/_committed_1727", 12),
("out/_committed_1728", 12),
("out/_started_1727", 12),
("out/.part.crc", 8),
("out/part-00000.parquet", 100),
("out/part-00001.parquet", 100),
]
.iter()
.map(|(k, s)| ((*k).to_string(), *s))
.collect();
#[cfg(feature = "cloud")]
assert_eq!(
crate::cloud_browse::look_at_listing("out/", &[], &keys).0,
EntryKind::MultiFile,
"and the same in a bucket"
);
}
fn write(dir: &Path, name: &str, columns: &[&str]) {
let mut frame = DataFrame::new(
1,
columns
.iter()
.map(|c| Column::new((*c).into(), &[1i32]))
.collect::<Vec<_>>(),
)
.unwrap();
let file = std::fs::File::create(dir.join(name)).unwrap();
ParquetWriter::new(file).finish(&mut frame).unwrap();
}
fn write_nested(dir: &Path, name: &str, struct_name: &str, fields: &[&str]) {
let inner = DataFrame::new(
1,
fields
.iter()
.map(|f| Column::new((*f).into(), &[1i32]))
.collect::<Vec<_>>(),
)
.unwrap();
let nested = inner
.into_struct(struct_name.into())
.into_series()
.into_column();
let mut frame = DataFrame::new(1, vec![Column::new("id".into(), &[1i32]), nested]).unwrap();
let file = std::fs::File::create(dir.join(name)).unwrap();
ParquetWriter::new(file).finish(&mut frame).unwrap();
}
fn measured(dir: &Path) -> Entry {
let (kind, holds) = look_at_directory(dir);
let mut entry = Entry {
path: dir.to_path_buf(),
kind,
name: dir.file_name().unwrap().to_string_lossy().into_owned(),
size: None,
modified: None,
rows: None,
cols: None,
cols_sampled: false,
columns: Vec::new(),
cost: Cost::default(),
holds,
opens_whole_directory: false,
format_spec: None,
table: None,
};
enrich(&mut entry);
entry
}
#[test]
fn holds_written_as_folders_still_reads() {
let old: Holds = serde_json::from_str(r#"{"folders":3,"partitions":2}"#).unwrap();
assert_eq!((old.directories, old.partitions), (3, 2));
let new = serde_json::to_string(&old).unwrap();
assert!(new.contains(r#""directories":3"#), "{new}");
}
#[test]
fn a_directory_of_separate_tables_is_not_a_dataset() {
let dir = tempfile::tempdir().unwrap();
write(
dir.path(),
"circuits.parquet",
&["circuit_id", "lat", "lng"],
);
write(
dir.path(),
"drivers.parquet",
&["driver_id", "code", "nationality"],
);
write(
dir.path(),
"laps.parquet",
&["lap", "position", "time_millis"],
);
assert_eq!(
classify_directory(dir.path()),
EntryKind::MultiFile,
"the filenames alone still say multi"
);
let entry = measured(dir.path());
assert_eq!(
entry.kind,
EntryKind::Directory,
"reading the footers says otherwise"
);
assert_eq!(
entry.rows, None,
"a sum across separate tables is not a row count"
);
assert_eq!(
entry.cols,
Some(9),
"the union of what the directory holds is still a true answer to what is in it"
);
assert_eq!(entry.label(), "3 parquet", "and the label counts the files");
}
#[test]
fn a_directory_whose_files_each_bring_a_column_is_a_place_to_look_inside() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "old.parquet", &["id", "ts", "amount"]);
write(dir.path(), "new.parquet", &["id", "ts", "amt"]);
assert_eq!(
classify_directory(dir.path()),
EntryKind::MultiFile,
"the names alone still say two Parquet files"
);
let entry = measured(dir.path());
assert_eq!(
entry.kind,
EntryKind::Directory,
"and the footers say neither file's columns are in the other's"
);
assert_eq!(entry.label(), "2 parquet", "which the label still reports");
assert_eq!(entry.rows, None, "a sum over two tables is not a number");
}
#[test]
fn a_directory_of_one_table_stays_a_dataset() {
let dir = tempfile::tempdir().unwrap();
for part in 0..3 {
write(
dir.path(),
&format!("part-0000{part}.parquet"),
&["id", "ts", "amount"],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile);
assert_eq!(entry.rows, Some(3));
assert_eq!(entry.cols, Some(3));
}
#[test]
fn a_dataset_that_gained_columns_stays_a_dataset() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "2009.parquet", &["id", "ts"]);
write(dir.path(), "2015.parquet", &["id", "ts", "fee"]);
write(
dir.path(),
"2025.parquet",
&["id", "ts", "fee", "witness", "address", "value"],
);
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile);
assert_eq!(entry.rows, Some(3));
assert_eq!(
entry.columns,
vec!["id", "ts", "fee", "witness", "address", "value"],
"every column any file has, in the order they first appear — not the \
2009 shape"
);
assert_eq!(entry.cols, Some(6), "and the count is of those");
}
#[test]
fn a_hive_dataset_that_gained_columns_reports_all_of_them() {
let dir = tempfile::tempdir().unwrap();
for (part, columns) in [
("year=2009", &["id", "ts"][..]),
("year=2025", &["id", "ts", "address"][..]),
] {
let sub = dir.path().join(part);
std::fs::create_dir_all(&sub).unwrap();
write(&sub, "part-0.parquet", columns);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Hive);
assert_eq!(entry.columns, vec!["id", "ts", "address"]);
assert_eq!(entry.cols, Some(4), "and the partition column `year`");
}
#[test]
fn a_directory_too_large_to_count_still_reports_the_columns_it_gained() {
let dir = tempfile::tempdir().unwrap();
for part in 0..MAX_FOOTERS_PER_DATASET + 1 {
let mut columns = vec!["id".to_string(), "ts".to_string()];
if part > MAX_FOOTERS_PER_DATASET / 2 {
columns.push("address".to_string());
}
let refs: Vec<&str> = columns.iter().map(String::as_str).collect();
write(dir.path(), &format!("part-{part:03}.parquet"), &refs);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile, "still one table");
assert_eq!(entry.rows, None, "too many files to count");
assert!(
entry.columns.contains(&"address".to_string()),
"the column the dataset gained is in the row: {:?}",
entry.columns
);
}
#[test]
fn a_lake_table_is_not_a_directory_of_parquet_files() {
for (marker, expected) in [
("_delta_log", EntryKind::Delta),
(".hoodie", EntryKind::Hudi),
] {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "part-0.parquet", &["id", "amount"]);
write(dir.path(), "part-1.parquet", &["id", "amount"]);
write(dir.path(), "part-2.parquet", &["id", "amount"]);
let log = dir.path().join(marker);
std::fs::create_dir_all(&log).unwrap();
std::fs::write(log.join("00000000000000000000.json"), b"{}").unwrap();
assert_eq!(
classify_directory(dir.path()),
expected,
"{marker} says what this directory is"
);
let entry = measured(dir.path());
assert_eq!(entry.kind, expected);
assert_eq!(
entry.rows, None,
"and no row count is claimed for it: summing the footers would count \
the rows the log says are gone"
);
assert!(!entry.kind.is_dataset(), "it does not open as one table");
}
}
#[test]
fn an_iceberg_root_is_metadata_beside_data() {
let iceberg = tempfile::tempdir().unwrap();
let data = iceberg.path().join("data");
let metadata = iceberg.path().join("metadata");
std::fs::create_dir_all(&data).unwrap();
std::fs::create_dir_all(&metadata).unwrap();
write(&data, "00000-0-abc.parquet", &["id", "amount"]);
write(&data, "00001-0-def.parquet", &["id", "amount"]);
std::fs::write(metadata.join("v2.metadata.json"), b"{}").unwrap();
std::fs::write(metadata.join("snap-1.avro"), b"x").unwrap();
assert_eq!(classify_directory(iceberg.path()), EntryKind::Iceberg);
let plain = tempfile::tempdir().unwrap();
std::fs::create_dir_all(plain.path().join("data")).unwrap();
std::fs::create_dir_all(plain.path().join("metadata")).unwrap();
std::fs::write(plain.path().join("metadata/notes.txt"), b"x").unwrap();
assert_eq!(
classify_directory(plain.path()),
EntryKind::Directory,
"no *.metadata.json, so no Iceberg table"
);
let no_data = tempfile::tempdir().unwrap();
let metadata = no_data.path().join("metadata");
std::fs::create_dir_all(&metadata).unwrap();
std::fs::write(metadata.join("v1.metadata.json"), b"{}").unwrap();
write(no_data.path(), "part-0.parquet", &["id"]);
write(no_data.path(), "part-1.parquet", &["id"]);
assert_eq!(
classify_directory(no_data.path()),
EntryKind::MultiFile,
"metadata with no data/ beside it is somebody's directory, not a table root"
);
}
#[test]
fn a_file_and_a_directory_of_it_count_the_same_columns() {
let dir = tempfile::tempdir().unwrap();
write_nested(dir.path(), "one.parquet", "inputs", &["address", "value"]);
let mut file = Entry::new(dir.path().join("one.parquet"), EntryKind::File);
enrich(&mut file);
assert_eq!(
file.cols,
Some(2),
"`id` and `inputs`, which is what opening it shows: {:?}",
file.columns
);
assert!(
file.columns.iter().any(|c| c == "inputs.address"),
"the leaves are still searchable: {:?}",
file.columns
);
write_nested(dir.path(), "two.parquet", "inputs", &["address", "value"]);
let directory = measured(dir.path());
assert_eq!(directory.kind, EntryKind::MultiFile);
assert_eq!(
directory.cols, file.cols,
"and a directory of them says the same number"
);
}
#[test]
fn a_dotted_column_name_is_its_own_column() {
let dir = tempfile::tempdir().unwrap();
write(dir.path(), "flat.parquet", &["id", "user.id", "user.name"]);
let mut file = Entry::new(dir.path().join("flat.parquet"), EntryKind::File);
enrich(&mut file);
assert_eq!(file.cols, Some(3), "three columns: {:?}", file.columns);
}
#[test]
fn a_writer_change_does_not_double_the_column_count() {
let dir = tempfile::tempdir().unwrap();
write_nested(dir.path(), "old.parquet", "inputs", &["address"]);
write_nested(dir.path(), "new.parquet", "inputs", &["address", "value"]);
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile, "still one table");
assert_eq!(
entry.cols,
Some(2),
"one `inputs`, not one per shape of it: {:?}",
entry.columns
);
assert!(
entry.columns.len() > 2,
"while every leaf stays searchable: {:?}",
entry.columns
);
}
#[test]
fn a_sampled_column_count_says_it_is_a_floor() {
let dir = tempfile::tempdir().unwrap();
for part in 0..MAX_FOOTERS_PER_DATASET * 2 {
write(
dir.path(),
&format!("part-{part:04}.parquet"),
&["id", "ts"],
);
}
let entry = measured(dir.path());
assert_eq!(entry.rows, None, "too many files to count");
assert!(entry.cols.is_some(), "but the width is still worth having");
assert!(
entry.cols_sampled,
"and it is marked as the floor it is, not presented as a total"
);
let small = tempfile::tempdir().unwrap();
write(small.path(), "a.parquet", &["id", "ts"]);
write(small.path(), "b.parquet", &["id", "ts"]);
assert!(!measured(small.path()).cols_sampled);
}
#[test]
fn the_files_a_directory_offers_come_back_in_order() {
let dir = tempfile::tempdir().unwrap();
for name in ["c.parquet", "a.parquet", "d.parquet", "b.parquet"] {
write(dir.path(), name, &["id"]);
}
let files = parquet_files_under(dir.path());
let names: Vec<String> = files
.iter()
.map(|p| p.file_name().unwrap().to_string_lossy().into_owned())
.collect();
assert_eq!(
names,
vec!["a.parquet", "b.parquet", "c.parquet", "d.parquet"],
"sorted, not in the order the directory was written"
);
}
#[test]
fn a_directory_past_the_budget_keeps_the_directories_first_files() {
let dir = tempfile::tempdir().unwrap();
for part in 0..MAX_FOOTERS_PER_DATASET * 3 {
write(dir.path(), &format!("part-{part:04}.parquet"), &["id"]);
}
let files = parquet_files_under(dir.path());
assert_eq!(
files.len(),
MAX_FOOTERS_PER_DATASET + 1,
"one past the budget, which is what says there are too many to count"
);
let names: Vec<String> = files
.iter()
.map(|p| p.file_name().unwrap().to_string_lossy().into_owned())
.collect();
let expected: Vec<String> = (0..=MAX_FOOTERS_PER_DATASET)
.map(|part| format!("part-{part:04}.parquet"))
.collect();
assert_eq!(
names, expected,
"the directory's first files, not the listing's"
);
}
#[test]
fn a_lake_table_is_recognized_among_its_data_files() {
let dir = tempfile::tempdir().unwrap();
for part in 0..32 {
write(dir.path(), &format!("part-{part:03}.parquet"), &["id"]);
}
std::fs::create_dir_all(dir.path().join("_delta_log")).unwrap();
assert_eq!(classify_directory(dir.path()), EntryKind::Delta);
}
#[test]
fn a_directory_too_large_to_count_is_still_checked() {
let dir = tempfile::tempdir().unwrap();
for table in 0..MAX_FOOTERS_PER_DATASET + 1 {
write(
dir.path(),
&format!("table_{table:03}.parquet"),
&[&format!("{table}_id"), &format!("{table}_value")],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Directory);
assert_eq!(entry.rows, None, "too many files to count either way");
}
#[test]
fn a_large_directory_of_one_table_stays_a_dataset() {
let dir = tempfile::tempdir().unwrap();
for part in 0..MAX_FOOTERS_PER_DATASET + 1 {
write(
dir.path(),
&format!("part-{part:05}.parquet"),
&["id", "ts"],
);
}
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::MultiFile);
}
#[test]
fn a_downgraded_directory_keeps_every_column_its_files_have() {
let dir = tempfile::tempdir().unwrap();
write(
dir.path(),
"circuits.parquet",
&["circuit_id", "lat", "lng"],
);
write(
dir.path(),
"drivers.parquet",
&["driver_id", "code", "nationality"],
);
let seasons = dir.path().join("seasons");
std::fs::create_dir_all(&seasons).unwrap();
write(&seasons, "2024.parquet", &["season_year", "round"]);
let entry = measured(dir.path());
assert_eq!(entry.kind, EntryKind::Directory);
for column in [
"circuit_id",
"lat",
"lng",
"driver_id",
"code",
"nationality",
"season_year",
"round",
] {
assert!(
entry.columns.iter().any(|c| c == column),
"{column} in {:?}",
entry.columns
);
}
}
#[test]
fn a_name_no_reader_takes_is_refused_before_opening() {
let refused = |name: &str| unreadable_by_name(std::path::Path::new(name));
assert!(refused("gs://b/ml/onnx/pipeline_rf.onnx"));
assert!(refused("model.onnx.gz"));
assert!(refused("README.md"));
for readable in [
"a.csv",
"a.CSV",
"a.csv.gz",
"a.parquet",
"a.xlsx",
"data.gz",
"part-0000",
] {
assert!(!refused(readable), "{readable}");
}
}
}