use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::AtomicU64;
use color_eyre::Result;
use color_eyre::eyre::eyre;
use crate::export_modal::ExportFormat;
use crate::members::Table;
use crate::scan::Scan;
use crate::segments::Converted;
use crate::text_formats::Detail;
use crate::unfinished::Writer;
use crate::{FileFormat, OpenOptions, ReadReport};
pub(crate) mod facts;
pub(crate) mod polars;
pub(crate) type ListTables = fn(&Path) -> Result<Vec<Table>>;
pub(crate) type FactsFn = fn(&Path) -> Result<FormatFacts>;
#[derive(Clone, Copy)]
pub(crate) struct Facts {
pub read: FactsFn,
pub footer: bool,
}
#[derive(Debug, Clone, Default)]
pub struct FormatFacts {
pub detail: Option<Arc<Detail>>,
pub footer: Option<crate::parquet_footer::Footer>,
}
pub(crate) type TableSchema = fn(&Path, Option<&str>) -> Option<crate::discover::SchemaPreview>;
pub(crate) struct ScanIn<'a> {
pub format: FileFormat,
pub paths: &'a [PathBuf],
pub options: &'a OpenOptions,
pub report: &'a mut ReadReport,
pub formats: &'a crate::formats::Registry,
}
impl ScanIn<'_> {
pub fn path(&self) -> &Path {
&self.paths[0]
}
}
pub(crate) type ScanFn = fn(ScanIn<'_>) -> Result<Scan>;
pub(crate) struct ConvertIn<'a> {
pub files: &'a [PathBuf],
pub display: &'a Path,
pub format: FileFormat,
pub options: &'a OpenOptions,
pub formats: &'a crate::formats::Registry,
pub writer: &'a Writer,
pub read: &'a AtomicU64,
}
pub(crate) type ConvertOut = Result<(Converted, Option<Arc<Detail>>)>;
pub(crate) type ConvertFn = fn(&ConvertIn<'_>) -> ConvertOut;
#[cfg(feature = "cloud")]
pub(crate) struct BucketIn<'a> {
pub url: &'a str,
pub path: ::polars::prelude::PlRefPath,
pub cloud: ::polars::io::cloud::CloudOptions,
pub glob: bool,
pub options: &'a OpenOptions,
pub format: FileFormat,
}
#[cfg(feature = "cloud")]
pub(crate) type BucketScan = fn(BucketIn<'_>) -> Result<::polars::prelude::LazyFrame>;
pub(crate) struct Reader {
pub scan: ScanFn,
#[cfg(feature = "cloud")]
pub bucket_scan: Option<BucketScan>,
pub convert: Option<ConvertFn>,
pub signatures: &'static [Signature],
pub refines: &'static [FileFormat],
pub tables: Option<ListTables>,
pub facts: Option<Facts>,
pub table_schema: Option<TableSchema>,
pub bytes_decide: bool,
pub preview: Option<Preview>,
pub python: Option<crate::python_script::Python>,
pub export: Option<ExportFormat>,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Preview {
Scan,
RowGroup,
}
pub(crate) const BASE: Reader = Reader {
scan: |input| {
Err(eyre!(
"datui has no reader for {} files.",
input.format.name()
))
},
#[cfg(feature = "cloud")]
bucket_scan: None,
convert: None,
signatures: &[],
refines: &[],
tables: None,
facts: None,
table_schema: None,
bytes_decide: false,
preview: None,
python: None,
export: None,
};
pub(crate) fn scan(input: ScanIn<'_>) -> Result<Scan> {
let one = (input.paths.len() == 1).then(|| input.paths[0].clone());
(of(input.format).scan)(input).map_err(|e| match one {
Some(path) => crate::error_display::in_file(&path, e),
None => e,
})
}
pub(crate) fn of(format: FileFormat) -> &'static Reader {
match format {
FileFormat::Parquet => &polars::PARQUET,
FileFormat::Csv => &polars::CSV,
FileFormat::Tsv => &polars::TSV,
FileFormat::Psv => &polars::PSV,
FileFormat::Json => &polars::JSON,
FileFormat::Jsonl => &polars::JSONL,
FileFormat::Arrow => &polars::ARROW,
FileFormat::Avro => &polars::AVRO,
FileFormat::Orc => &polars::ORC,
FileFormat::Excel => &polars::EXCEL,
FileFormat::Safetensors => &crate::model_files::SAFETENSORS,
FileFormat::Gguf => &crate::model_files::GGUF,
FileFormat::Nmea => &crate::gps::NMEA,
FileFormat::Gpx => &crate::gps::GPX,
FileFormat::Audio => &crate::audio::READER,
FileFormat::Midi => &crate::midi::READER,
FileFormat::Sqlite => &crate::sqlite::READER,
FileFormat::Vcd => &crate::vcd::READER,
FileFormat::Fix => &crate::fix::READER,
FileFormat::Sdf => &crate::sdf::READER,
FileFormat::Numpy => &crate::numpy::READER,
FileFormat::Elf => &crate::elf::READER,
FileFormat::Ulog => &crate::ulog::READER,
FileFormat::Dataflash => &crate::dataflash::READER,
FileFormat::Candump => &crate::candump::READER,
FileFormat::Text => &crate::lines::READER,
FileFormat::Journal => &crate::journal::READER,
}
}
pub const HEAD: usize = 4096;
pub(crate) struct Signature {
pub says: fn(&[u8], Option<&Path>) -> bool,
pub kind: Kind,
pub trusted: Trusted,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Kind {
Magic,
Structure,
Text,
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct Trusted {
pub pipe: bool,
pub open: Unnamed,
pub listing: bool,
pub tables: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Unnamed {
Never,
Any,
NoExtension,
}
pub(crate) const EVERYWHERE: Trusted = Trusted {
pipe: true,
open: Unnamed::Any,
listing: true,
tables: false,
};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Asked {
Pipe,
Open { extension: bool },
Listing,
Tables,
}
impl Asked {
fn believes(self, format: FileFormat, trusted: Trusted) -> bool {
match self {
Asked::Pipe => trusted.pipe,
Asked::Open { extension } => match trusted.open {
Unnamed::Never => false,
Unnamed::Any => true,
Unnamed::NoExtension => !extension,
},
Asked::Listing => trusted.listing,
Asked::Tables => trusted.tables && format.holds_tables(),
}
}
}
pub(crate) fn sniff(
head: &[u8],
file: Option<&Path>,
asked: Asked,
among: impl Fn(FileFormat) -> bool,
) -> Option<FileFormat> {
[Kind::Magic, Kind::Structure, Kind::Text]
.into_iter()
.find_map(|kind| {
FileFormat::ALL.into_iter().find(|&format| {
among(format)
&& of(format).signatures.iter().any(|sig| {
sig.kind == kind
&& asked.believes(format, sig.trusted)
&& (sig.says)(head, file)
})
})
})
}
pub(crate) fn refined(path: &Path, named: FileFormat) -> Option<FileFormat> {
if !FileFormat::ALL
.into_iter()
.any(|f| of(f).refines.contains(&named))
{
return None;
}
let head = head_of(path)?;
sniff(&head, Some(path), Asked::Open { extension: true }, |f| {
of(f).refines.contains(&named)
})
}
pub(crate) fn sniff_file(path: &Path, asked: Asked) -> Option<FileFormat> {
let head = head_of(path)?;
sniff(&head, Some(path), asked, |_| true)
}
pub(crate) fn sniff_open(
path: &Path,
compression: Option<crate::CompressionFormat>,
) -> Option<FileFormat> {
let read_through = |f: FileFormat| f.reads_into();
let asked = Asked::Open {
extension: path.extension().is_some(),
};
let is_file = path.is_file();
if is_file
&& let Some(found) =
head_of(path).and_then(|head| sniff(&head, Some(path), asked, |_| true))
{
return Some(found);
}
let compression = compression.or_else(|| crate::CompressionFormat::from_extension(path))?;
if let Some(named) = path
.file_stem()
.and_then(|stem| FileFormat::from_path(Path::new(stem)))
.filter(|f| read_through(*f))
{
return Some(named);
}
if !is_file {
return None;
}
let head = crate::formats::head_of(path, Some(compression), HEAD as u64)?;
sniff(&head, Some(path), asked, read_through)
}
pub(crate) fn head_of(path: &Path) -> Option<Vec<u8>> {
use std::io::Read;
let mut head = Vec::with_capacity(HEAD);
std::fs::File::open(path)
.ok()?
.take(HEAD as u64)
.read_to_end(&mut head)
.ok()?;
Some(head)
}
pub(crate) fn read_into(input: ScanIn<'_>) -> Result<Scan> {
Ok(Scan::ReadInto {
files: input.paths.to_vec(),
format: input.format,
})
}
pub(crate) fn convert(input: &ConvertIn<'_>) -> ConvertOut {
match of(input.format).convert {
Some(convert) => convert(input),
None => Err(eyre!(
"{} files are not read into files of their own.",
input.format.name()
)),
}
.map_err(|e| crate::error_display::in_file(input.display, e))
}
pub(crate) fn many_files_refused() -> String {
let (many, one): (Vec<FileFormat>, Vec<FileFormat>) = FileFormat::ALL
.into_iter()
.partition(|f| f.reads_many_files());
let names = |formats: &[FileFormat]| {
formats
.iter()
.map(|f| f.title())
.collect::<Vec<_>>()
.join(", ")
};
format!(
"Unsupported file type for multiple files: {} files are read as one table; open {} files one at a time.",
names(&many),
names(&one)
)
}
pub(crate) fn export_default(format: FileFormat) -> Option<ExportFormat> {
of(format).export
}
#[cfg(test)]
pub(crate) mod bad_input {
use std::path::Path;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, AtomicU64};
use super::{ConvertIn, ScanIn};
use crate::scan::Scan;
use crate::{FileFormat, OpenOptions, ReadReport};
pub(crate) fn opening(
dir: &Path,
name: &str,
bytes: &[u8],
format: FileFormat,
options: &OpenOptions,
) -> Option<String> {
let path = dir.join(name);
std::fs::write(&path, bytes).unwrap();
let said =
|e: color_eyre::Report| crate::error_display::user_message_from_report(&e, Some(&path));
let formats = crate::formats::Registry::of(Vec::new());
let mut report = ReadReport::default();
let paths = [path.clone()];
let scan = super::scan(ScanIn {
format,
paths: &paths,
options,
report: &mut report,
formats: &formats,
});
let mut written = Vec::new();
let lf = match scan {
Err(e) => return Some(said(e)),
Ok(Scan::Frame(lf)) => *lf,
Ok(Scan::ReadInto { files, format }) => {
let unfinished = crate::unfinished::Unfinished::default();
let writer = unfinished.writer(Arc::new(AtomicBool::new(false)));
let read = AtomicU64::new(0);
match super::convert(&ConvertIn {
files: &files,
display: &path,
format,
options,
formats: &formats,
writer: &writer,
read: &read,
}) {
Err(e) => return Some(said(e)),
Ok((converted, _)) => {
written = converted.files;
converted.lf
}
}
}
Ok(_) => return None,
};
let rows = lf.limit(100).collect();
drop(written);
rows.err().map(|e| said(color_eyre::Report::new(e)))
}
pub(crate) fn assert_shape(message: &str, path: &Path) {
let quoted = format!("\"{}\":", path.display());
assert!(message.starts_with("ed), "names the file: {message}");
let after = &message[quoted.len()..];
let what = match after.strip_prefix(' ') {
Some(what) => what,
None => {
let mut parts = after.splitn(3, ':');
let (line, column) = (parts.next().unwrap(), parts.next().unwrap_or_default());
for n in [line, column] {
assert!(
n.parse::<usize>().is_ok_and(|n| n > 0),
"a place in the file: {message}"
);
}
parts
.next()
.and_then(|what| what.strip_prefix(' '))
.unwrap_or_else(|| panic!("a place, then a space: {message}"))
}
};
let first = message.lines().next().unwrap_or_default();
assert!(first.ends_with('.'), "ends with a full stop: {message}");
assert!(
crate::error_display::starts_with_a_key(what)
|| what.chars().next().is_some_and(|c| !c.is_lowercase()),
"sentence case: {message}"
);
assert_eq!(
what.matches(path.to_string_lossy().as_ref()).count(),
0,
"named once: {message}"
);
for rust in [
"Some(",
"None",
"Error {",
"Kind(",
"PolarsError",
"ComputeError",
"os error",
] {
assert!(!message.contains(rust), "no Rust ({rust}): {message}");
}
}
pub(crate) fn each_names_its_file(format: FileFormat, bad: &[(&str, &[u8], &str)]) {
let dir = tempfile::tempdir().unwrap();
for (name, bytes, says) in bad {
let message = opening(dir.path(), name, bytes, format, &OpenOptions::default())
.unwrap_or_else(|| panic!("{name} opens"));
eprintln!("{message}");
assert_shape(&message, &dir.path().join(name));
assert!(message.contains(says), "{name}: {message}");
}
}
}
#[cfg(test)]
mod tests {
use super::*;
fn piped(head: &[u8]) -> Option<FileFormat> {
sniff(head, None, Asked::Pipe, |_| true)
}
#[test]
fn text_formats_by_their_first_bytes() {
let said = [
(
&b"8=FIX.4.4\x019=5\x0135=0\x0110=000\x01\n"[..],
FileFormat::Fix,
),
(
b"$timescale 1ns $end\n$scope module top $end\n",
FileFormat::Vcd,
),
(
b"aspirin\n RDKit\n\n 0 0 0 0 0 0 0 0 0 0999 V2000\nM END\n$$$$\n",
FileFormat::Sdf,
),
(b"$GPGGA,1,2", FileFormat::Nmea),
(b"<?xml version=\"1.0\"?>\n<gpx>", FileFormat::Gpx),
];
for (head, format) in said {
assert_eq!(piped(head), Some(format), "{format:?}");
}
assert_eq!(piped(b"a,b\n1,2\n"), None);
}
#[test]
fn a_compressed_file_by_its_name_or_what_it_holds() {
use std::io::Write;
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("session.log.gz");
let mut gz = flate2::write::GzEncoder::new(
std::fs::File::create(&path).unwrap(),
Default::default(),
);
gz.write_all(b"20260101-00:00:00 : 8=FIX.4.2|9=5|35=0|10=000|\n")
.unwrap();
gz.finish().unwrap();
assert_eq!(sniff_open(&path, None), Some(FileFormat::Fix));
for (name, format) in [
("lib.sdf.gz", Some(FileFormat::Sdf)),
("a.nmea.gz", Some(FileFormat::Nmea)),
("a.csv.gz", None),
] {
assert_eq!(sniff_open(&dir.path().join(name), None), format, "{name}");
}
}
#[test]
fn a_signature_is_believed_where_it_says() {
let elf = b"\x7fELF\x02\x01\x01\0\0\0\0\0\0\0\0\0";
assert_eq!(piped(elf), Some(FileFormat::Elf));
assert_eq!(sniff(elf, None, Asked::Listing, |_| true), None);
assert_eq!(
sniff(elf, None, Asked::Tables, |_| true),
Some(FileFormat::Elf)
);
assert_eq!(piped(b"ORC\x00"), None);
assert_eq!(
sniff(b"ORC\x00", None, Asked::Listing, |_| true),
Some(FileFormat::Orc)
);
let npy = b"\x93NUMPY\x01\x00";
assert_eq!(piped(npy), Some(FileFormat::Numpy));
assert_eq!(sniff(npy, None, Asked::Tables, |_| true), None);
}
#[test]
fn the_copy_docs_name_each_format_s_reader() {
let page = std::path::Path::new(env!("CARGO_MANIFEST_DIR"))
.join("../../docs/user-guide/copying.md");
let text = std::fs::read_to_string(&page).expect("the copying page");
let start = text.find("| Format | Reader |").expect("the reader table");
let rows: Vec<(String, String)> = text[start..]
.lines()
.skip(2)
.take_while(|l| l.starts_with('|'))
.map(|l| {
let (format, reader) = l.trim_matches('|').split_once(" | ").expect("two cells");
(format.trim().to_string(), reader.trim().to_string())
})
.collect();
let said: Vec<&str> = rows.iter().map(|(f, _)| f.as_str()).collect();
let titles: Vec<&str> = FileFormat::ALL.iter().map(|f| f.title()).collect();
assert_eq!(said, titles, "one row per format, in --format's order");
for ((title, reader), format) in rows.iter().zip(FileFormat::ALL) {
let expected = of(format)
.python
.as_ref()
.map_or("`df = ...`".to_string(), |p| format!("`{}`", p.call));
assert_eq!(*reader, expected, "{title}");
}
}
#[test]
fn readers_agree_with_their_descriptors() {
for format in FileFormat::ALL {
let reader = of(format);
assert_eq!(
reader.tables.is_some(),
format.holds_tables(),
"{}",
format.name()
);
if format.reads_into() {
assert!(reader.convert.is_some(), "{}", format.name());
}
if reader.facts.is_some() {
assert!(format.summary_tab().is_some(), "{}", format.name());
}
#[cfg(feature = "cloud")]
if reader.bucket_scan.is_some() {
assert!(format.reads_bucket_prefix(), "{}", format.name());
}
}
}
}