use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use polars::prelude::LazyFrame;
use crate::download::TempDownload;
use crate::schema_union::FooterProgress;
use crate::unfinished::{Unfinished, Writer};
use crate::widgets::datatable::DataTableState;
use crate::{CompressionFormat, FileFormat, OpenOptions, source, stdin};
use crate::jobs::Hold;
#[cfg(any(feature = "http", feature = "cloud"))]
use crate::jobs::Jobs;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub(crate) struct LoadId(u64);
impl LoadId {
#[cfg(test)]
pub(crate) fn for_tests(n: u64) -> Self {
Self(n)
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
#[derive(Clone)]
pub(crate) enum PendingDownload {
#[cfg(feature = "http")]
Http {
url: String,
size: Option<u64>,
options: OpenOptions,
},
#[cfg(feature = "cloud")]
S3 {
url: String,
size: Option<u64>,
options: OpenOptions,
},
#[cfg(feature = "cloud")]
Gcs {
url: String,
size: Option<u64>,
options: OpenOptions,
},
#[cfg(feature = "cloud")]
Azure {
url: String,
size: Option<u64>,
options: OpenOptions,
},
#[cfg(feature = "cloud")]
Arrow {
url: String,
objects: Vec<crate::cloud_arrow::Object>,
size: Option<u64>,
options: OpenOptions,
},
}
#[cfg(any(feature = "http", feature = "cloud"))]
impl PendingDownload {
pub(crate) fn parts(&self) -> (&str, Option<u64>, &OpenOptions) {
match self {
#[cfg(feature = "http")]
PendingDownload::Http { url, size, options } => (url, *size, options),
#[cfg(feature = "cloud")]
PendingDownload::S3 { url, size, options } => (url, *size, options),
#[cfg(feature = "cloud")]
PendingDownload::Gcs { url, size, options } => (url, *size, options),
#[cfg(feature = "cloud")]
PendingDownload::Azure { url, size, options } => (url, *size, options),
#[cfg(feature = "cloud")]
PendingDownload::Arrow {
url, size, options, ..
} => (url, *size, options),
}
}
pub(crate) fn arrow_files(&self) -> Option<(usize, usize)> {
match self {
#[cfg(feature = "cloud")]
PendingDownload::Arrow { objects, .. } => {
let streams = objects.iter().filter(|o| o.stream).count();
Some((streams, objects.len() - streams))
}
_ => None,
}
}
pub(crate) fn asked(mut self) -> Self {
match &mut self {
#[cfg(feature = "http")]
PendingDownload::Http { options, .. } => options.download_unasked = None,
#[cfg(feature = "cloud")]
PendingDownload::S3 { options, .. }
| PendingDownload::Gcs { options, .. }
| PendingDownload::Azure { options, .. }
| PendingDownload::Arrow { options, .. } => options.download_unasked = None,
}
self
}
pub(crate) fn with_size(mut self, found: Option<u64>) -> Self {
match &mut self {
#[cfg(feature = "http")]
PendingDownload::Http { size, .. } => *size = found,
#[cfg(feature = "cloud")]
PendingDownload::S3 { size, .. } => *size = found,
#[cfg(feature = "cloud")]
PendingDownload::Gcs { size, .. } => *size = found,
#[cfg(feature = "cloud")]
PendingDownload::Azure { size, .. } => *size = found,
#[cfg(feature = "cloud")]
PendingDownload::Arrow { size, .. } => *size = found,
}
self
}
}
#[derive(Clone)]
pub(crate) struct OpenRequest {
pub(crate) paths: Vec<PathBuf>,
pub(crate) options: OpenOptions,
pub(crate) size: u64,
pub(crate) recent: Option<PathBuf>,
pub(crate) shown: Option<PathBuf>,
pub(crate) warn_in_memory_above: Option<u64>,
}
impl OpenRequest {
pub(crate) fn named(
mut paths: Vec<PathBuf>,
mut options: OpenOptions,
formats: &crate::formats::Registry,
) -> Self {
let first = paths[0].clone();
let piped = stdin::is_stdin(&first);
let local = !piped && matches!(source::input_source(&first), source::InputSource::Local(_));
let mut table = None;
if local
&& paths.len() == 1
&& let Some((db, name)) = crate::members::split(&first)
{
table = Some(first.clone());
options.table = Some(name);
paths = vec![db];
} else if local
&& first.is_file()
&& let Some(name) = options.table.as_deref()
&& crate::members::holder(&first).is_some()
{
table = Some(crate::members::place(&first, name));
} else if local
&& paths.len() == 1
&& options.table.is_none()
&& let Some((dir, split)) = crate::hf_splits::split_place(&first)
{
table = Some(first.clone());
options.table = Some(split);
paths = vec![dir];
} else if local
&& first.is_dir()
&& let Some(split) = options.table.as_deref()
&& crate::hf_splits::cache_splits(&first)
.iter()
.any(|s| s == split)
{
table = Some(first.join(split));
} else if local
&& paths.len() == 1
&& options.table.is_none()
&& let Some((file, variant)) = crate::members::split_variant(&first, formats)
{
table = Some(first.clone());
options.table = Some(variant);
paths = vec![file];
} else if local
&& let Some(variant) = options.table.as_deref()
&& first.is_file()
&& formats.variants_of(&first).is_some()
{
table = Some(crate::members::place(&first, variant));
}
let first = &paths[0];
let size = if local {
std::fs::metadata(first).map(|m| m.len()).unwrap_or(0)
} else {
0
};
let recent = table
.clone()
.or_else(|| (!piped && (!local || first.exists())).then(|| first.clone()));
Self {
paths,
options,
size,
recent,
shown: table,
warn_in_memory_above: None,
}
}
}
pub(crate) enum Phase {
Starting {
label: String,
percent: u16,
},
LookingAtPaths,
ReadingSpec,
LookingAtDirectory,
#[cfg(any(feature = "http", feature = "cloud"))]
ReadingHeaders,
#[cfg(any(feature = "http", feature = "cloud"))]
CheckingSize {
note: Option<&'static str>,
},
ConfirmingRead {
scan: Box<Scan>,
_hold: Option<Hold>,
},
#[cfg(any(feature = "http", feature = "cloud"))]
Confirming {
pending: Box<PendingDownload>,
note: Option<&'static str>,
_hold: Hold,
},
#[cfg(any(feature = "http", feature = "cloud"))]
Downloading,
Spooling {
read: Arc<AtomicU64>,
},
Decompressing,
DecompressingRecords,
ReadingRecords,
Converting {
what: Conversion,
read: Arc<AtomicU64>,
total: u64,
},
ScanningStrings,
CountingFooter,
Scanning {
downloaded: bool,
},
ReadingLines,
ReadingSchema,
FirstRows,
}
impl Phase {
pub(crate) fn label(&self) -> (&str, u16) {
match self {
Phase::Starting { label, percent } => (label, *percent),
Phase::LookingAtPaths => ("Scanning input", 10),
Phase::ReadingSpec => ("Reading spec", 5),
Phase::LookingAtDirectory => (crate::App::LOOKING_AT_A_DIRECTORY, 5),
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::ReadingHeaders => ("Reading headers", 20),
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::CheckingSize { .. } | Phase::Confirming { .. } => ("Checking size", 0),
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::Downloading => ("Downloading", 20),
Phase::Spooling { .. } => ("Reading stdin", 5),
Phase::ConfirmingRead { .. } => ("Scanning input", 0),
Phase::Decompressing | Phase::DecompressingRecords => ("Decompressing", 30),
Phase::ReadingRecords => ("Reading records", 35),
Phase::Converting { what, read, total } => {
let done = read.load(Ordering::Relaxed).min(*total);
let share = (done * 20).checked_div(*total).unwrap_or(0);
(what.label(), 10 + share as u16)
}
Phase::ScanningStrings => ("Scanning string columns", 55),
Phase::CountingFooter => (COUNTING_FOOTER, 10),
Phase::Scanning { downloaded: false } => ("Scanning input", 10),
Phase::Scanning { downloaded: true } => ("Scanning", 30),
Phase::ReadingLines => (READING_LINES, 10),
Phase::ReadingSchema => ("Reading schema", 40),
Phase::FirstRows => ("Loading buffer", 70),
}
}
fn asks(&self) -> bool {
match self {
Phase::ConfirmingRead { .. } => true,
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::Confirming { .. } => true,
_ => false,
}
}
fn builds_the_dataset(&self) -> bool {
match self {
Phase::ReadingSchema | Phase::Decompressing => true,
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::ReadingHeaders => true,
_ => false,
}
}
fn starting(&self) -> bool {
matches!(
self,
Phase::Starting { .. } | Phase::LookingAtPaths | Phase::LookingAtDirectory
)
}
}
#[derive(Clone)]
struct Fetched {
url: PathBuf,
file: TempDownload,
arrow: Option<KeptArrow>,
}
#[derive(Clone)]
struct KeptArrow {
parts: Arc<Vec<crate::ipc_stream::Part>>,
splits: Option<Arc<crate::hf_splits::Splits>>,
table: Option<String>,
}
impl Fetched {
fn serves(&self, options: &OpenOptions) -> bool {
self.arrow
.as_ref()
.is_none_or(|arrow| arrow.table == options.table)
}
}
pub(crate) struct Load {
id: LoadId,
from_home: bool,
phase: Phase,
path: Option<PathBuf>,
size: u64,
paths: Option<Vec<PathBuf>>,
recent: Option<PathBuf>,
progress: Arc<FooterProgress>,
writer: Writer,
download: Option<Fetched>,
converted: Vec<TempDownload>,
warn_in_memory_above: Option<u64>,
}
pub(crate) struct Scan {
paths: Vec<PathBuf>,
options: OpenOptions,
display: Option<PathBuf>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub(crate) struct InMemory {
pub(crate) bytes: u64,
pub(crate) format: FileFormat,
pub(crate) files: usize,
pub(crate) name: PathBuf,
}
impl Load {
pub(crate) fn phase(&self) -> &Phase {
&self.phase
}
fn made(&self) -> Made {
Made {
download: self.download.as_ref().map(|fetched| fetched.file.clone()),
converted: self.converted.clone(),
..Made::default()
}
}
pub(crate) fn path(&self) -> Option<&Path> {
self.path.as_deref()
}
pub(crate) fn size(&self) -> u64 {
match &self.phase {
Phase::Spooling { read } => read.load(Ordering::Relaxed),
_ => self.size,
}
}
}
pub(crate) enum Step {
Nothing,
Crash(String),
#[cfg(any(feature = "http", feature = "cloud"))]
ReadHeaders {
url: PathBuf,
format: FileFormat,
options: OpenOptions,
writer: Writer,
},
#[cfg(any(feature = "http", feature = "cloud"))]
Probe(PendingDownload),
#[cfg(any(feature = "http", feature = "cloud"))]
Ask(PendingDownload),
AskRead(InMemory),
#[cfg(any(feature = "http", feature = "cloud"))]
Download {
pending: PendingDownload,
writer: Writer,
},
FetchSpec {
url: PathBuf,
options: OpenOptions,
writer: Writer,
},
Spool {
options: OpenOptions,
writer: Writer,
read: Arc<AtomicU64>,
},
Decompress {
file: PathBuf,
path: PathBuf,
options: OpenOptions,
writer: Writer,
download: Option<TempDownload>,
},
DecompressRecords {
file: PathBuf,
path: PathBuf,
choice: crate::formats::Choice,
options: OpenOptions,
writer: Writer,
},
ReadRecords {
copy: PathBuf,
path: PathBuf,
choice: crate::formats::Choice,
options: OpenOptions,
},
Convert {
what: Conversion,
files: Vec<PathBuf>,
path: Option<PathBuf>,
options: OpenOptions,
writer: Writer,
read: Arc<AtomicU64>,
},
Scan {
paths: Vec<PathBuf>,
options: OpenOptions,
display: Option<PathBuf>,
status: &'static str,
},
ReadSchema {
lf: Box<LazyFrame>,
path: Option<PathBuf>,
options: OpenOptions,
progress: Arc<FooterProgress>,
made: Made,
},
Install(Box<Loaded>),
Failed(Failed),
Tables(Tables),
Hex(Hex),
}
#[derive(Debug)]
pub(crate) struct Hex {
pub(crate) file: PathBuf,
pub(crate) from_home: bool,
pub(crate) asked: bool,
pub(crate) record_size: Option<usize>,
}
#[derive(Debug)]
pub(crate) struct Tables {
pub(crate) database: PathBuf,
pub(crate) from_home: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub(crate) enum Conversion {
Streams,
Text(FileFormat),
}
impl Conversion {
pub(crate) fn label(self) -> &'static str {
match self {
Conversion::Streams => "Converting Arrow stream",
Conversion::Text(format) => format.conversion().label,
}
}
pub(crate) fn status(self) -> &'static str {
match self {
Conversion::Streams => "Converting Arrow stream...",
Conversion::Text(format) => format.conversion().status,
}
}
}
pub(crate) enum Converted {
Streams {
file: TempDownload,
parts: Vec<crate::ipc_stream::Part>,
},
Frame {
files: Vec<TempDownload>,
lf: Box<LazyFrame>,
notes: Vec<crate::notes::Note>,
other_tables: Vec<String>,
detail: Option<Arc<crate::text_formats::Detail>>,
},
}
#[derive(Default)]
pub(crate) struct Made {
pub(crate) download: Option<TempDownload>,
pub(crate) converted: Vec<TempDownload>,
pub(crate) notes: Vec<crate::notes::Note>,
pub(crate) other_tables: Vec<String>,
pub(crate) detail: Option<Arc<crate::text_formats::Detail>>,
}
#[derive(Debug)]
pub(crate) struct Failed {
pub(crate) message: String,
pub(crate) from_home: bool,
}
pub(crate) struct Loaded {
pub(crate) state: DataTableState,
pub(crate) path: Option<PathBuf>,
pub(crate) options: OpenOptions,
pub(crate) debug_label: Option<String>,
pub(crate) paths: Option<Vec<PathBuf>>,
pub(crate) recent: Option<PathBuf>,
pub(crate) from_home: bool,
pub(crate) footers: Arc<FooterProgress>,
}
pub(crate) enum LoadAnswer {
Scanned {
lf: Box<LazyFrame>,
path: Option<PathBuf>,
options: OpenOptions,
},
Compressed {
file: PathBuf,
path: Option<PathBuf>,
options: OpenOptions,
},
SpecFetched {
spec: Arc<crate::formats::Spec>,
options: OpenOptions,
},
CompressedRecords {
file: PathBuf,
path: Option<PathBuf>,
choice: crate::formats::Choice,
options: OpenOptions,
},
DecompressedRecords {
copy: TempDownload,
path: PathBuf,
choice: crate::formats::Choice,
options: OpenOptions,
},
Convert {
what: Conversion,
files: Vec<PathBuf>,
bytes: u64,
path: Option<PathBuf>,
options: OpenOptions,
},
Converted {
converted: Converted,
path: Option<PathBuf>,
options: OpenOptions,
},
Tables {
file: PathBuf,
tables: Vec<String>,
path: Option<PathBuf>,
},
Hex {
file: PathBuf,
asked: bool,
record_size: Option<usize>,
},
SchemaRead {
state: Box<DataTableState>,
path: Option<PathBuf>,
options: OpenOptions,
debug_label: Option<String>,
},
#[cfg(any(feature = "http", feature = "cloud"))]
NoRanges { options: OpenOptions },
#[cfg(any(feature = "http", feature = "cloud"))]
Sized(PendingDownload),
#[cfg(any(feature = "http", feature = "cloud"))]
PastLimit(PendingDownload),
#[cfg(any(feature = "http", feature = "cloud"))]
Downloaded {
download: TempDownload,
options: OpenOptions,
},
Spooled {
download: TempDownload,
options: OpenOptions,
},
Recorded { file: PathBuf, options: OpenOptions },
}
#[derive(Debug, Clone, Copy)]
pub(crate) struct Retired {
pub(crate) id: LoadId,
pub(crate) asking: bool,
}
#[derive(Default)]
pub(crate) struct Loader {
next_id: u64,
load: Option<Load>,
kept: Option<Fetched>,
unfinished: Unfinished,
}
impl Loader {
pub(crate) fn unfinished(&self) -> &Unfinished {
&self.unfinished
}
pub(crate) fn current(&self) -> Option<&Load> {
self.load.as_ref()
}
pub(crate) fn id(&self) -> Option<LoadId> {
self.load.as_ref().map(|load| load.id)
}
fn current_in(&self, id: LoadId, phase: impl Fn(&Phase) -> bool) -> bool {
self.load
.as_ref()
.is_some_and(|load| load.id == id && phase(&load.phase))
}
pub(crate) fn looking_at_paths(&self, id: LoadId) -> bool {
self.current_in(id, |phase| matches!(phase, Phase::LookingAtPaths))
}
pub(crate) fn looking_at_directory(&self, id: LoadId) -> bool {
self.current_in(id, |phase| matches!(phase, Phase::LookingAtDirectory))
}
pub(crate) fn awaiting_dataset(&self) -> bool {
self.load
.as_ref()
.is_some_and(|load| !matches!(load.phase, Phase::FirstRows))
}
pub(crate) fn waits(&self) -> bool {
self.awaiting_dataset() && !self.asking()
}
pub(crate) fn asking(&self) -> bool {
self.load.as_ref().is_some_and(|load| load.phase.asks())
}
pub(crate) fn hold_while_asking(&mut self, hold: Hold) {
if let Some(Load {
phase: Phase::ConfirmingRead { _hold, .. },
..
}) = self.load.as_mut()
{
*_hold = Some(hold);
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
pub(crate) fn download_note(&self) -> Option<&'static str> {
match self.load.as_ref().map(|load| &load.phase) {
Some(Phase::Confirming { note, .. }) => *note,
_ => None,
}
}
pub(crate) fn progress(&self) -> Option<&Arc<FooterProgress>> {
self.load
.as_ref()
.filter(|load| !matches!(load.phase, Phase::FirstRows))
.map(|load| &load.progress)
}
pub(crate) fn retire(&mut self) -> Option<Retired> {
let load = self.load.take()?;
if !matches!(load.phase, Phase::FirstRows) {
load.progress.cancel();
}
Some(Retired {
id: load.id,
asking: load.phase.asks(),
})
}
pub(crate) fn make_way(&mut self) -> Option<Retired> {
if self.load.as_ref().is_some_and(|load| load.phase.starting()) {
return None;
}
self.retire()
}
fn start(&mut self, from_home: bool) -> &mut Load {
if self
.load
.as_ref()
.is_some_and(|load| !load.phase.starting())
{
debug_assert!(false, "an open started without making way");
self.retire();
}
if self.load.is_none() {
self.next_id = self.next_id.wrapping_add(1);
let progress = Arc::<FooterProgress>::default();
self.load = Some(Load {
id: LoadId(self.next_id),
from_home,
phase: Phase::Starting {
label: "Loading".to_string(),
percent: 0,
},
path: None,
size: 0,
paths: None,
recent: None,
writer: self.unfinished.writer(progress.cancel_flag()),
progress,
download: None,
converted: Vec::new(),
warn_in_memory_above: None,
});
}
let load = self.load.as_mut().expect("started just above");
load.from_home |= from_home;
load
}
pub(crate) fn announce(&mut self, from_home: bool, label: String, percent: u16) -> LoadId {
let load = self.start(from_home);
load.phase = Phase::Starting { label, percent };
load.id
}
#[cfg(test)]
pub(crate) fn first_rows_for_tests(&mut self) {
self.retire();
self.start(false).phase = Phase::FirstRows;
}
#[cfg(test)]
pub(crate) fn size_for_tests(&mut self, size: u64) {
if let Some(load) = self.load.as_mut() {
load.size = size;
}
}
pub(crate) fn name(&mut self, path: PathBuf) {
if let Some(load) = self
.load
.as_mut()
.filter(|load| !matches!(load.phase, Phase::FirstRows))
{
load.path = Some(stdin::named(&path));
}
}
pub(crate) fn look_at_paths(&mut self) -> LoadId {
let load = self.start(false);
load.phase = Phase::LookingAtPaths;
load.id
}
pub(crate) fn look_at_directory(&mut self, dir: PathBuf) -> LoadId {
let load = self.start(false);
load.phase = Phase::LookingAtDirectory;
load.path = Some(dir);
load.id
}
pub(crate) fn open(&mut self, request: OpenRequest) -> Step {
if !(request.paths.len() == 1
&& self
.kept
.as_ref()
.is_some_and(|kept| kept.url == request.paths[0]))
{
self.kept = None;
}
let OpenRequest {
paths,
mut options,
size,
recent,
shown,
warn_in_memory_above,
} = request;
options.arrow_parts = None;
let prepared = options
.prepared
.take()
.and_then(|handoff| handoff.lock().ok()?.take());
let load = self.start(false);
load.path = Some(shown.unwrap_or_else(|| stdin::named(&paths[0])));
load.size = size;
load.recent = recent;
load.paths = Some(paths.clone());
load.warn_in_memory_above = warn_in_memory_above;
match prepared {
Some(prepared) => self.install_prepared(*prepared),
None => self.first_step(paths, options),
}
}
fn install_prepared(&mut self, prepared: crate::home_preview::Prepared) -> Step {
let crate::home_preview::Prepared {
state,
options,
debug_label,
progress,
} = prepared;
let writer = self.unfinished.writer(progress.cancel_flag());
let load = self.load.as_mut().expect("started by the open");
load.phase = Phase::FirstRows;
load.progress = progress;
load.writer = writer;
Step::Install(Box::new(Loaded {
state: *state,
path: load.path.clone(),
options,
debug_label,
paths: load.paths.clone(),
recent: load.recent.clone(),
from_home: load.from_home,
footers: load.progress.clone(),
}))
}
pub(crate) fn open_frame(&mut self, lf: LazyFrame, options: OpenOptions) -> Step {
let load = self.start(false);
load.path = None;
load.size = 0;
load.paths = None;
load.recent = None;
load.phase = Phase::ReadingSchema;
Step::ReadSchema {
lf: Box::new(lf),
path: None,
options,
progress: load.progress.clone(),
made: Made::default(),
}
}
fn first_step(&mut self, paths: Vec<PathBuf>, options: OpenOptions) -> Step {
let first = paths[0].clone();
let src = source::input_source(&first);
if paths.len() > 1 {
let only_one = match &src {
source::InputSource::S3(_) => {
Some("Only one S3 URL at a time. Open a single s3:// path.")
}
source::InputSource::Gcs(_) => {
Some("Only one GCS URL at a time. Open a single gs:// path.")
}
source::InputSource::Azure(_) => {
Some("Only one Azure URL at a time. Open a single abfss:// path.")
}
source::InputSource::Http(_) => {
Some("Only one HTTP/HTTPS URL at a time. Open a single URL.")
}
source::InputSource::Local(_) => None,
};
if let Some(message) = only_one {
self.load = None;
return Step::Crash(message.to_string());
}
}
if let Some(message) = stdin::refuse(&paths, true) {
self.load = None;
return Step::Crash(message.to_string());
}
if options.follow
&& let Some(message) = crate::follow::refuse_paths(&paths, &options)
{
self.load = None;
return Step::Crash(message);
}
if let Some(url) = options
.spec_file
.clone()
.filter(|file| options.spec_fetched.is_none() && source::is_remote_url(file))
{
let load = self.load.as_mut().expect("an open has a load");
load.phase = Phase::ReadingSpec;
return Step::FetchSpec {
url,
options,
writer: load.writer.clone(),
};
}
if stdin::is_stdin(&first) {
if let Some(kept) = self.kept.clone().filter(|kept| {
kept.url == first && kept.file.path().exists() && kept.serves(&options)
}) {
return self.read_download(kept, options);
}
let load = self.load.as_mut().expect("an open has a load");
let read = Arc::<AtomicU64>::default();
load.phase = Phase::Spooling { read: read.clone() };
return Step::Spool {
options,
writer: load.writer.clone(),
read,
};
}
let load = self.load.as_mut().expect("an open has a load");
let compression = options
.compression
.or_else(|| CompressionFormat::from_extension(&first));
let delimited = delimited_format(&first, &options);
if matches!(src, source::InputSource::Local(_))
&& paths.len() == 1
&& compression.is_some()
&& let Some(format) = delimited
{
load.phase = Phase::Decompressing;
return Step::Decompress {
file: first.clone(),
path: first,
options: OpenOptions {
format: Some(format),
..options
},
writer: load.writer.clone(),
download: None,
};
}
#[cfg(any(feature = "http", feature = "cloud"))]
if let Some(kept) = self
.kept
.clone()
.filter(|kept| paths.len() == 1 && kept.file.path().exists() && kept.serves(&options))
{
return self.read_download(kept, options);
}
#[cfg(any(feature = "http", feature = "cloud"))]
if paths.len() == 1
&& let Some(format) = crate::remote_model::model_format(&first, options.format)
{
let load = self.load.as_mut().expect("an open has a load");
load.phase = Phase::ReadingHeaders;
return Step::ReadHeaders {
url: first,
format,
options,
writer: load.writer.clone(),
};
}
#[cfg(any(feature = "http", feature = "cloud"))]
if let Some(pending) = remote_download(&src, &options) {
let load = self.load.as_mut().expect("an open has a load");
load.phase = Phase::CheckingSize { note: None };
return Step::Probe(pending);
}
let load = self.load.as_mut().expect("an open has a load");
if counts_footer(&paths, &options) {
load.phase = Phase::CountingFooter;
let display = load.path.clone().filter(|shown| *shown != paths[0]);
return Step::Scan {
paths,
options,
display,
status: COUNTING_FOOTER_STATUS,
};
}
if paths.len() == 1 && delimited.is_some() && options.parse_strings.is_some() {
load.phase = Phase::ScanningStrings;
return Step::Scan {
paths,
options,
display: None,
status: "Scanning string columns...",
};
}
let display = load.path.clone().filter(|shown| *shown != paths[0]);
if let Some(limit) = load.warn_in_memory_above
&& let Some(read) = in_memory(&paths, &options)
&& read.bytes > limit
{
load.phase = Phase::ConfirmingRead {
scan: Box::new(Scan {
paths,
options,
display,
}),
_hold: None,
};
return Step::AskRead(read);
}
if reads_lines(&paths, &options) {
load.phase = Phase::ReadingLines;
return Step::Scan {
paths,
options,
display,
status: READING_LINES_STATUS,
};
}
load.phase = Phase::Scanning { downloaded: false };
Step::Scan {
paths,
options,
display,
status: "Scanning input...",
}
}
fn read_download(&mut self, fetched: Fetched, mut options: OpenOptions) -> Step {
if let Some(arrow) = &fetched.arrow {
options.format = Some(FileFormat::Arrow);
options.hive = false;
options.arrow_parts = Some(arrow.parts.clone());
options.splits = arrow.splits.clone();
}
let load = self.load.as_mut().expect("a download read has a load");
let file = fetched.file.path().to_path_buf();
let url = stdin::named(&fetched.url);
let download = fetched.file.clone();
load.download = Some(fetched);
let compressed = options
.compression
.or_else(|| CompressionFormat::from_extension(&file))
.is_some();
if compressed && let Some(format) = delimited_format(&file, &options) {
load.phase = Phase::Decompressing;
return Step::Decompress {
file,
path: url,
options: OpenOptions {
format: Some(format),
..options
},
writer: load.writer.clone(),
download: Some(download),
};
}
let paths = vec![file];
let status = if counts_footer(&paths, &options) {
load.phase = Phase::CountingFooter;
COUNTING_FOOTER_STATUS
} else if reads_lines(&paths, &options) {
load.phase = Phase::ReadingLines;
READING_LINES_STATUS
} else {
load.phase = Phase::Scanning { downloaded: true };
"Scanning..."
};
Step::Scan {
paths,
options,
display: Some(url),
status,
}
}
pub(crate) fn answered(
&mut self,
id: LoadId,
answer: LoadAnswer,
#[cfg(any(feature = "http", feature = "cloud"))] jobs: &Jobs,
) -> Step {
let Some(load) = self.load.as_mut().filter(|load| load.id == id) else {
return Step::Nothing;
};
match (answer, &load.phase) {
(
LoadAnswer::Scanned { lf, path, options },
Phase::Scanning { .. }
| Phase::ScanningStrings
| Phase::CountingFooter
| Phase::ReadingRecords
| Phase::ReadingLines,
) => {
load.phase = Phase::ReadingSchema;
Step::ReadSchema {
lf,
path,
options,
progress: load.progress.clone(),
made: load.made(),
}
}
(
LoadAnswer::Convert {
what,
files,
bytes,
path,
options,
},
Phase::Scanning { .. }
| Phase::ScanningStrings
| Phase::CountingFooter
| Phase::ReadingLines,
) if load.converted.is_empty() => {
let read = Arc::<AtomicU64>::default();
load.phase = Phase::Converting {
what,
read: read.clone(),
total: bytes,
};
Step::Convert {
what,
files,
path,
options,
writer: load.writer.clone(),
read,
}
}
(
LoadAnswer::Converted {
converted: Converted::Streams { file, parts },
path,
options,
},
Phase::Converting { .. },
) => {
let paths = vec![file.path().to_path_buf()];
if let Some(fetched) = load.download.as_mut() {
fetched.file = file.clone();
}
load.converted = vec![file];
load.phase = Phase::Scanning { downloaded: true };
Step::Scan {
paths,
options: OpenOptions {
format: Some(FileFormat::Arrow),
hive: false,
arrow_parts: Some(Arc::new(parts)),
..options
},
display: path,
status: "Scanning...",
}
}
(
LoadAnswer::Converted {
converted:
Converted::Frame {
files,
lf,
notes,
other_tables,
detail,
},
path,
options,
},
Phase::Converting { .. },
) => {
load.converted = files;
load.phase = Phase::ReadingSchema;
Step::ReadSchema {
lf,
path,
options,
progress: load.progress.clone(),
made: Made {
notes,
other_tables,
detail,
..load.made()
},
}
}
(
LoadAnswer::Compressed {
file,
path,
options,
},
Phase::Scanning { .. }
| Phase::ScanningStrings
| Phase::CountingFooter
| Phase::ReadingLines,
) => {
load.phase = Phase::Decompressing;
Step::Decompress {
path: path.unwrap_or_else(|| file.clone()),
file,
options,
writer: load.writer.clone(),
download: load.download.as_ref().map(|fetched| fetched.file.clone()),
}
}
(LoadAnswer::SpecFetched { spec, options }, Phase::ReadingSpec) => {
let paths = load.paths.clone().unwrap_or_default();
if paths.is_empty() {
return Step::Nothing;
}
self.first_step(
paths,
OpenOptions {
spec_fetched: Some(spec),
..options
},
)
}
(
LoadAnswer::CompressedRecords {
file,
path,
choice,
options,
},
Phase::Scanning { .. } | Phase::ScanningStrings | Phase::ReadingLines,
) => {
load.phase = Phase::DecompressingRecords;
Step::DecompressRecords {
path: path.unwrap_or_else(|| file.clone()),
file,
choice,
options,
writer: load.writer.clone(),
}
}
(
LoadAnswer::DecompressedRecords {
copy,
path,
choice,
options,
},
Phase::DecompressingRecords,
) => {
let file = copy.path().to_path_buf();
load.converted = vec![copy];
load.phase = Phase::ReadingRecords;
Step::ReadRecords {
copy: file,
path,
choice,
options,
}
}
(
LoadAnswer::Hex {
file,
asked,
record_size,
},
Phase::Scanning { .. }
| Phase::ScanningStrings
| Phase::CountingFooter
| Phase::ReadingLines,
) => {
let from_home = load.from_home;
let fetched = load.download.is_some();
let named = load.path.clone();
self.retire();
if fetched {
let message = match named {
Some(path) => crate::error_display::file_message(&path, crate::UNSUPPORTED),
None => crate::UNSUPPORTED.to_string(),
};
return Step::Failed(Failed { message, from_home });
}
Step::Hex(Hex {
file,
from_home,
asked,
record_size,
})
}
(
LoadAnswer::Tables { file, tables, path },
Phase::Scanning { .. }
| Phase::ScanningStrings
| Phase::CountingFooter
| Phase::ReadingLines,
) => {
let from_home = load.from_home;
let database = path.unwrap_or(file);
let fetched = load.download.is_some();
self.retire();
if fetched {
let shown = tables.iter().take(20).cloned().collect::<Vec<_>>();
let more = tables.len().saturating_sub(shown.len());
let more = match more {
0 => String::new(),
n => format!(" and {n} more"),
};
return Step::Failed(Failed {
message: format!(
"{} holds {} tables: {}{more}. Open one with --table NAME.",
database.display(),
tables.len(),
shown.join(", ")
),
from_home,
});
}
Step::Tables(Tables {
database,
from_home,
})
}
(
LoadAnswer::SchemaRead {
state,
path,
options,
debug_label,
},
phase,
) if phase.builds_the_dataset() => {
load.phase = Phase::FirstRows;
if let Some(fetched) = load.download.take() {
self.kept = Some(fetched);
}
let state = *state;
let load = self.load.as_ref().expect("installing its load");
Step::Install(Box::new(Loaded {
state,
path,
options,
debug_label,
paths: load.paths.clone(),
recent: load.recent.clone(),
from_home: load.from_home,
footers: load.progress.clone(),
}))
}
#[cfg(feature = "cloud")]
(
LoadAnswer::Sized(PendingDownload::Arrow {
url,
objects,
options,
..
}),
Phase::CheckingSize { .. },
) if crate::cloud_arrow::in_place(&objects).is_some() => {
load.phase = Phase::Scanning { downloaded: false };
let parts = crate::cloud_arrow::in_place(&objects).unwrap_or_default();
Step::Scan {
paths: vec![PathBuf::from(url)],
options: OpenOptions {
format: Some(FileFormat::Arrow),
hive: false,
arrow_parts: Some(Arc::new(parts)),
..options
},
display: None,
status: "Scanning...",
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
(LoadAnswer::NoRanges { options }, Phase::ReadingHeaders) => {
let first = load.paths.as_ref().and_then(|paths| paths.first().cloned());
match first
.and_then(|first| remote_download(&source::input_source(&first), &options))
{
Some(pending) => {
load.phase = Phase::CheckingSize {
note: Some(NO_RANGES),
};
Step::Probe(pending)
}
None => self.failed(id, NO_RANGES),
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
(LoadAnswer::Sized(pending), Phase::CheckingSize { note: None })
if pending
.parts()
.2
.download_unasked
.is_some_and(|unasked| unasked.covers(pending.parts().1)) =>
{
load.phase = Phase::Downloading;
Step::Download {
pending,
writer: load.writer.clone(),
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
(LoadAnswer::Sized(pending), Phase::CheckingSize { note }) => {
load.phase = Phase::Confirming {
pending: Box::new(pending.clone()),
note: *note,
_hold: jobs.hold(),
};
Step::Ask(pending)
}
#[cfg(any(feature = "http", feature = "cloud"))]
(LoadAnswer::PastLimit(pending), Phase::Downloading) => {
let pending = pending.with_size(None);
load.phase = Phase::Confirming {
pending: Box::new(pending.clone()),
note: Some(PAST_LIMIT),
_hold: jobs.hold(),
};
Step::Ask(pending)
}
#[cfg(any(feature = "http", feature = "cloud"))]
(LoadAnswer::Downloaded { download, options }, Phase::Downloading) => {
let fetched = Fetched {
url: load.path.clone().unwrap_or_default(),
file: download,
arrow: options.arrow_parts.clone().map(|parts| KeptArrow {
parts,
splits: options.splits.clone(),
table: options.table.clone(),
}),
};
self.read_download(fetched, options)
}
(LoadAnswer::Recorded { file, options }, Phase::Spooling { read }) => {
load.size = read.load(Ordering::Relaxed);
self.kept = None;
load.path = Some(file.clone());
load.paths = Some(vec![file.clone()]);
load.recent = Some(file.clone());
let paths = vec![file];
let status = if counts_footer(&paths, &options) {
load.phase = Phase::CountingFooter;
COUNTING_FOOTER_STATUS
} else {
load.phase = Phase::Scanning { downloaded: false };
"Scanning input..."
};
Step::Scan {
paths,
options,
display: None,
status,
}
}
(LoadAnswer::Spooled { download, options }, Phase::Spooling { read }) => {
load.size = read.load(Ordering::Relaxed);
let fetched = Fetched {
url: PathBuf::from(stdin::PATH),
file: download,
arrow: None,
};
self.read_download(fetched, options)
}
_ => Step::Nothing,
}
}
pub(crate) fn failed(&mut self, id: LoadId, message: &str) -> Step {
let Some(load) = self.load.as_ref().filter(|load| load.id == id) else {
return Step::Nothing;
};
if matches!(load.phase, Phase::FirstRows) {
return Step::Nothing;
}
let from_home = load.from_home;
let mut message = match &load.download {
Some(fetched) => crate::error_display::named_by_source(
message,
fetched.file.path(),
&stdin::named(&fetched.url),
),
None => message.to_string(),
};
if let Some(path) = &load.path {
for converted in &load.converted {
message = crate::error_display::named_by_source(&message, converted.path(), path);
}
}
self.retire();
Step::Failed(Failed { message, from_home })
}
pub(crate) fn confirmed(&mut self) -> Step {
let Some(load) = self.load.as_mut().filter(|load| load.phase.asks()) else {
return Step::Nothing;
};
match std::mem::replace(&mut load.phase, Phase::Scanning { downloaded: false }) {
Phase::ConfirmingRead { scan, .. } => {
let Scan {
paths,
options,
display,
} = *scan;
Step::Scan {
paths,
options,
display,
status: "Scanning input...",
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
Phase::Confirming { pending, .. } => {
load.phase = Phase::Downloading;
Step::Download {
pending: pending.asked(),
writer: load.writer.clone(),
}
}
_ => unreachable!("a phase that asks"),
}
}
pub(crate) fn first_rows_settled(&mut self) {
if self
.load
.as_ref()
.is_some_and(|load| matches!(load.phase, Phase::FirstRows))
{
self.load = None;
}
}
}
impl Drop for Loader {
fn drop(&mut self) {
self.retire();
}
}
const READING_LINES: &str = "Reading as lines";
const READING_LINES_STATUS: &str = "Reading as lines...";
fn reads_lines(paths: &[PathBuf], options: &OpenOptions) -> bool {
options
.format
.or_else(|| match paths {
[one] => FileFormat::from_path(one),
_ => None,
})
.is_some_and(FileFormat::is_lines)
}
pub(crate) fn in_memory(paths: &[PathBuf], options: &OpenOptions) -> Option<InMemory> {
let mut found: Option<InMemory> = None;
for path in paths {
if !matches!(source::input_source(path), source::InputSource::Local(_)) {
continue;
}
let compression = options
.compression
.or_else(|| CompressionFormat::from_extension(path));
let format = options.format.or_else(|| match compression {
Some(_) => path
.file_stem()
.and_then(|stem| FileFormat::from_path(Path::new(stem))),
None => match FileFormat::from_path(path) {
Some(named) => Some(crate::readers::refined(path, named).unwrap_or(named)),
None => crate::readers::sniff_open(path, None),
},
});
let Some(format) = format else {
continue;
};
let stored = match compression {
Some(_) => crate::Stored::Compressed {
in_memory: options.decompress_in_memory,
},
None => crate::Stored::Plain,
};
if format.read_mode(stored) != Some(crate::ReadMode::InMemory)
|| format.http_file() == crate::RemoteRead::InPlace
{
continue;
}
let Some(bytes) = std::fs::metadata(path)
.ok()
.filter(|m| m.is_file())
.map(|m| m.len())
else {
continue;
};
match &mut found {
Some(read) => {
read.bytes += bytes;
read.files += 1;
}
None => {
found = Some(InMemory {
bytes,
format,
files: 1,
name: path.clone(),
})
}
}
}
found
}
const COUNTING_FOOTER: &str = "Counting rows to skip the footer";
const COUNTING_FOOTER_STATUS: &str = "Counting rows to skip the footer...";
fn counts_footer(paths: &[PathBuf], options: &OpenOptions) -> bool {
options.skip_tail_rows.is_some_and(|n| n > 0)
&& paths.iter().all(|p| delimited_format(p, options).is_some())
}
pub(crate) fn delimited_format(path: &Path, options: &OpenOptions) -> Option<FileFormat> {
let format = options.format.or_else(|| {
FileFormat::from_path(path).or_else(|| {
CompressionFormat::from_extension(path)?;
FileFormat::from_path(Path::new(path.file_stem()?))
})
})?;
format.decompressed_once().then_some(format)
}
#[cfg(any(feature = "http", feature = "cloud"))]
pub(crate) const PAST_LIMIT: &str =
"This catalog file passed 50 MB, more than it may download without asking.";
#[cfg(any(feature = "http", feature = "cloud"))]
pub(crate) const NO_RANGES: &str = "The server does not send byte ranges, so the model's header cannot be read without downloading the whole file.";
#[cfg(any(feature = "http", feature = "cloud"))]
fn remote_download(src: &source::InputSource, options: &OpenOptions) -> Option<PendingDownload> {
#[cfg(feature = "cloud")]
let spec = options.spec_file.is_some() || options.spec_name.is_some();
#[cfg(feature = "cloud")]
let should_download =
|url: &str| should_download(url) || (spec && !source::is_prefix_or_glob(url));
#[cfg(feature = "cloud")]
if !spec && let Some(arrow) = cloud_arrow(src, options) {
return Some(arrow);
}
let options = options.clone();
match src {
#[cfg(feature = "http")]
source::InputSource::Http(url) => Some(PendingDownload::Http {
url: url.clone(),
size: None,
options,
}),
#[cfg(feature = "cloud")]
source::InputSource::S3(url) => {
let full = format!("s3://{url}");
should_download(&full).then_some(PendingDownload::S3 {
url: full,
size: None,
options,
})
}
#[cfg(feature = "cloud")]
source::InputSource::Gcs(url) => {
let full = format!("gs://{url}");
should_download(&full).then_some(PendingDownload::Gcs {
url: full,
size: None,
options,
})
}
#[cfg(feature = "cloud")]
source::InputSource::Azure(url) => should_download(url).then(|| PendingDownload::Azure {
url: url.clone(),
size: None,
options,
}),
_ => None,
}
}
#[cfg(feature = "cloud")]
fn cloud_arrow(src: &source::InputSource, options: &OpenOptions) -> Option<PendingDownload> {
let url = match src {
source::InputSource::S3(url) => format!("s3://{url}"),
source::InputSource::Gcs(url) => format!("gs://{url}"),
source::InputSource::Azure(url) => url.clone(),
_ => return None,
};
if url.contains('*') {
return None;
}
let arrow = if url.ends_with('/') {
options.format == Some(FileFormat::Arrow)
} else {
options.compression.is_none()
&& options
.format
.or_else(|| FileFormat::from_path(Path::new(&url)))
== Some(FileFormat::Arrow)
};
arrow.then(|| PendingDownload::Arrow {
url,
objects: Vec::new(),
size: None,
options: options.clone(),
})
}
#[cfg(feature = "cloud")]
fn should_download(url: &str) -> bool {
let (_, ext) = source::url_path_extension(url);
source::cloud_path_should_download(ext.as_deref(), source::is_prefix_or_glob(url))
}
#[cfg(test)]
mod tests {
use super::*;
use polars::prelude::IntoLazy;
#[cfg(all(feature = "http", feature = "cloud"))]
#[test]
fn remote_files_are_routed_as_their_format_says() {
use crate::{FileFormat, RemoteRead, Stored};
let in_place = |url: &str, stream: bool| {
let path = Path::new(url);
crate::remote_model::model_format(path, None).is_some()
|| match remote_download(&source::input_source(path), &OpenOptions::default()) {
None => true,
Some(PendingDownload::Arrow { .. }) => {
let object = crate::cloud_arrow::Object {
url: url.to_string(),
size: 10,
stream,
};
crate::cloud_arrow::in_place(&[object]).is_some()
}
Some(_) => false,
}
};
let mut seen = Vec::new();
for ext in [
"parquet",
"csv",
"tsv",
"psv",
"json",
"jsonl",
"arrow",
"avro",
"orc",
"xlsx",
"safetensors",
"gguf",
"nmea",
"gpx",
"wav",
"mid",
"db",
"vcd",
"sdf",
"log",
"npy",
"elf",
"ulg",
] {
let format = FileFormat::from_extension(ext).expect(ext);
seen.push(format);
let https = format!("https://example.com/d/x.{ext}");
assert_eq!(
in_place(&https, false),
format.http_file() == RemoteRead::InPlace,
"{https}"
);
for url in [format!("s3://b/d/x.{ext}"), format!("gs://b/d/x.{ext}")] {
let said = |stored| format.bucket_object(stored) == RemoteRead::InPlace;
assert_eq!(in_place(&url, false), said(Stored::Plain), "{url}");
if format == FileFormat::Arrow {
assert_eq!(in_place(&url, true), said(Stored::Stream), "{url}: stream");
}
let gz = format!("{url}.gz");
let compressed = Stored::Compressed { in_memory: false };
assert_eq!(in_place(&gz, false), said(compressed), "{gz}");
}
}
assert_eq!(FileFormat::Fix.http_file(), RemoteRead::Downloaded);
assert_eq!(
FileFormat::Fix.bucket_object(Stored::Plain),
RemoteRead::Downloaded
);
seen.push(FileFormat::Fix);
assert_eq!(FileFormat::Dataflash.http_file(), RemoteRead::Downloaded);
assert_eq!(
FileFormat::Dataflash.bucket_object(Stored::Plain),
RemoteRead::Downloaded
);
seen.push(FileFormat::Dataflash);
assert_eq!(FileFormat::Candump.http_file(), RemoteRead::Downloaded);
assert_eq!(
FileFormat::Candump.bucket_object(Stored::Plain),
RemoteRead::Downloaded
);
seen.push(FileFormat::Candump);
assert_eq!(FileFormat::Journal.http_file(), RemoteRead::Downloaded);
assert_eq!(
FileFormat::Journal.bucket_object(Stored::Plain),
RemoteRead::Downloaded
);
seen.push(FileFormat::Journal);
for f in FileFormat::ALL {
assert!(seen.contains(&f), "{} is checked", f.name());
}
}
fn frame() -> LazyFrame {
polars::df!("a" => [1i32, 2, 3]).unwrap().lazy()
}
fn state() -> Box<DataTableState> {
Box::new(DataTableState::from_lazyframe(frame(), &OpenOptions::default()).unwrap())
}
fn request(path: &str) -> OpenRequest {
OpenRequest {
paths: vec![PathBuf::from(path)],
options: OpenOptions::default(),
size: 7,
recent: Some(PathBuf::from(path)),
shown: None,
warn_in_memory_above: None,
}
}
#[cfg(any(feature = "http", feature = "cloud"))]
fn jobs() -> Jobs {
Jobs::new(std::sync::mpsc::channel().0)
}
fn answer(loader: &mut Loader, id: LoadId, answer: LoadAnswer) -> Step {
loader.answered(
id,
answer,
#[cfg(any(feature = "http", feature = "cloud"))]
&jobs(),
)
}
fn scanned(path: &str) -> LoadAnswer {
LoadAnswer::Scanned {
lf: Box::new(frame()),
path: Some(PathBuf::from(path)),
options: OpenOptions::default(),
}
}
fn schema_read(path: &str) -> LoadAnswer {
LoadAnswer::SchemaRead {
state: state(),
path: Some(PathBuf::from(path)),
options: OpenOptions::default(),
debug_label: None,
}
}
#[test]
fn a_local_open_runs_its_phases_in_order() {
let mut loader = Loader::default();
let step = loader.open(request("data.parquet"));
let id = loader.id().expect("a load");
assert!(matches!(
step,
Step::Scan { ref paths, display: None, status: "Scanning input...", .. }
if paths == &[PathBuf::from("data.parquet")]
));
let load = loader.current().unwrap();
assert_eq!(load.phase().label(), ("Scanning input", 10));
assert_eq!(load.path(), Some(Path::new("data.parquet")));
assert_eq!(load.size(), 7);
assert!(loader.awaiting_dataset() && loader.waits());
let Step::ReadSchema { progress, .. } = answer(&mut loader, id, scanned("data.parquet"))
else {
panic!("the scan goes on to the schema");
};
assert!(Arc::ptr_eq(&progress, loader.progress().unwrap()));
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading schema", 40)
);
let Step::Install(loaded) = answer(&mut loader, id, schema_read("data.parquet")) else {
panic!("the schema goes on to the install");
};
assert_eq!(loaded.path.as_deref(), Some(Path::new("data.parquet")));
assert_eq!(
loaded.paths.as_deref(),
Some(&[PathBuf::from("data.parquet")][..])
);
assert_eq!(loaded.recent.as_deref(), Some(Path::new("data.parquet")));
assert!(!loaded.from_home);
assert!(
Arc::ptr_eq(&loaded.footers, &progress),
"the dataset takes the load's footer counter"
);
assert!(!loader.awaiting_dataset(), "the dataset is up");
assert!(!loader.waits(), "the read of its rows holds the keys");
assert!(loader.progress().is_none(), "the counter is the dataset's");
assert_eq!(
loader.current().unwrap().phase().label(),
("Loading buffer", 70)
);
loader.first_rows_settled();
assert!(loader.current().is_none(), "the open is done");
}
#[test]
fn streams_are_converted_then_scanned() {
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let _ = loader.open(request("cache.arrow"));
let id = loader.id().unwrap();
let Step::Convert {
files, path, read, ..
} = answer(
&mut loader,
id,
LoadAnswer::Convert {
what: Conversion::Streams,
files: vec![PathBuf::from("cache.arrow")],
bytes: 200,
path: Some(PathBuf::from("cache.arrow")),
options: OpenOptions::default(),
},
)
else {
panic!("the streams are converted");
};
assert_eq!(files, [PathBuf::from("cache.arrow")]);
assert_eq!(path.as_deref(), Some(Path::new("cache.arrow")));
assert_eq!(
loader.current().unwrap().phase().label(),
("Converting Arrow stream", 10)
);
read.store(100, Ordering::Relaxed);
assert_eq!(
loader.current().unwrap().phase().label(),
("Converting Arrow stream", 20),
"the bar moves with the bytes read"
);
assert!(loader.waits());
let converted =
TempDownload::keep(TempDownload::create(Some(dir.path()), Some("arrow")).unwrap());
let temp = converted.path().to_path_buf();
let Step::Scan {
paths,
options,
display,
..
} = answer(
&mut loader,
id,
LoadAnswer::Converted {
converted: Converted::Streams {
file: converted,
parts: Vec::new(),
},
path: Some(PathBuf::from("cache.arrow")),
options: OpenOptions::default(),
},
)
else {
panic!("the IPC file is scanned");
};
assert_eq!(paths, std::slice::from_ref(&temp));
assert_eq!(options.format, Some(FileFormat::Arrow));
assert_eq!(display.as_deref(), Some(Path::new("cache.arrow")));
assert!(
matches!(
answer(
&mut loader,
id,
LoadAnswer::Convert {
what: Conversion::Streams,
files: vec![temp.clone()],
bytes: 1,
path: None,
options: OpenOptions::default(),
},
),
Step::Nothing
),
"converted once"
);
let Step::ReadSchema { made, .. } = answer(&mut loader, id, scanned("cache.arrow")) else {
panic!("the schema is read");
};
assert!(made.download.is_none(), "nothing was downloaded");
assert_eq!(
made.converted
.iter()
.map(|d| d.path().to_path_buf())
.collect::<Vec<_>>(),
std::slice::from_ref(&temp),
"the dataset holds the converted file"
);
drop(made);
let Step::Failed(failed) = loader.failed(id, &format!("could not read {}", temp.display()))
else {
panic!("the open fails");
};
assert_eq!(failed.message, "could not read cache.arrow");
assert!(!temp.exists(), "the retired load let the file go");
}
#[test]
fn an_old_load_s_answers_change_nothing() {
let mut loader = Loader::default();
let _ = loader.open(request("first.csv"));
let first = loader.id().unwrap();
assert!(loader.make_way().is_some(), "the first is doing work");
let _ = loader.open(request("second.csv"));
let second = loader.id().unwrap();
assert_ne!(first, second);
assert!(matches!(
answer(&mut loader, first, scanned("first.csv")),
Step::Nothing
));
assert!(matches!(
answer(&mut loader, first, schema_read("first.csv")),
Step::Nothing
));
assert!(matches!(loader.failed(first, "gone"), Step::Nothing));
assert_eq!(loader.id(), Some(second));
let load = loader.current().unwrap();
assert_eq!(load.path(), Some(Path::new("second.csv")));
assert_eq!(load.phase().label(), ("Scanning input", 10));
assert!(matches!(
answer(&mut loader, second, schema_read("second.csv")),
Step::Nothing
));
assert!(loader.awaiting_dataset());
}
#[test]
fn a_converted_file_nobody_wants_is_removed() {
let dir = tempfile::tempdir().unwrap();
let converted = || {
let file =
TempDownload::keep(TempDownload::create(Some(dir.path()), Some("arrow")).unwrap());
let path = file.path().to_path_buf();
(
LoadAnswer::Converted {
converted: Converted::Streams {
file,
parts: Vec::new(),
},
path: Some(PathBuf::from("cache.arrow")),
options: OpenOptions::default(),
},
path,
)
};
let mut loader = Loader::default();
let _ = loader.open(request("cache.arrow"));
let first = loader.id().unwrap();
let (early, early_path) = converted();
assert!(matches!(answer(&mut loader, first, early), Step::Nothing));
assert!(!early_path.exists(), "not converting yet");
let Step::Convert { writer, .. } = answer(
&mut loader,
first,
LoadAnswer::Convert {
what: Conversion::Streams,
files: vec![PathBuf::from("cache.arrow")],
bytes: 1,
path: Some(PathBuf::from("cache.arrow")),
options: OpenOptions::default(),
},
) else {
panic!("the streams are converted");
};
assert!(!writer.stopped());
assert!(loader.make_way().is_some(), "the conversion is work");
assert!(writer.stopped(), "Ctrl+O or another open stops it");
let _ = loader.open(request("other.csv"));
let (late, late_path) = converted();
assert!(matches!(answer(&mut loader, first, late), Step::Nothing));
assert!(!late_path.exists(), "the replaced load's file");
assert_eq!(
loader.current().unwrap().path(),
Some(Path::new("other.csv"))
);
}
#[test]
fn a_superseding_open_stops_the_one_it_replaces() {
let mut loader = Loader::default();
let _ = loader.open(request("big"));
let id = loader.id().unwrap();
let Step::ReadSchema { progress, .. } = answer(&mut loader, id, scanned("big")) else {
panic!("reading the schema");
};
let retired = loader.make_way().expect("the load is replaced");
assert!(!retired.asking);
assert!(progress.is_cancelled(), "its pass stops issuing reads");
assert!(loader.current().is_none());
}
#[test]
fn an_open_carries_on_the_load_that_announced_it() {
let mut loader = Loader::default();
let announced = loader.announce(true, "Scanning input".to_string(), 10);
loader.name(PathBuf::from("chosen.csv"));
assert_eq!(
loader.current().unwrap().path(),
Some(Path::new("chosen.csv"))
);
assert!(loader.waits(), "the keys wait from the key that asked");
assert!(loader.make_way().is_none(), "nothing to replace yet");
let _ = loader.open(request("chosen.csv"));
assert_eq!(loader.id(), Some(announced));
let Step::Install(loaded) = ({
let _ = answer(&mut loader, announced, scanned("chosen.csv"));
answer(&mut loader, announced, schema_read("chosen.csv"))
}) else {
panic!("installs");
};
assert!(loaded.from_home, "a failure would be reported at home");
let mut loader = Loader::default();
let look = loader.look_at_paths();
assert!(loader.looking_at_paths(look));
assert!(loader.make_way().is_none());
let directory = loader.look_at_directory(PathBuf::from("dir"));
assert_eq!(directory, look);
assert!(!loader.looking_at_paths(look), "that look is answered");
assert!(loader.looking_at_directory(look));
assert_eq!(
loader.current().unwrap().phase().label(),
(crate::App::LOOKING_AT_A_DIRECTORY, 5)
);
let _ = loader.open(request("dir"));
assert_eq!(loader.id(), Some(look));
assert!(
!loader.looking_at_directory(look),
"nor is a late look taken"
);
}
#[test]
fn a_failed_phase_retires_its_load() {
let mut loader = Loader::default();
let id = loader.announce(true, "Scanning input".to_string(), 10);
let _ = loader.open(request("broken.parquet"));
let Step::ReadSchema { progress, .. } = answer(&mut loader, id, scanned("broken.parquet"))
else {
panic!("reading the schema");
};
let Step::Failed(failed) = loader.failed(id, "not parquet") else {
panic!("the open fails");
};
assert_eq!(failed.message, "not parquet");
assert!(failed.from_home);
assert!(loader.current().is_none());
assert!(progress.is_cancelled());
assert!(matches!(loader.failed(id, "again"), Step::Nothing));
let id = {
let _ = loader.open(request("good.parquet"));
loader.id().unwrap()
};
let _ = answer(&mut loader, id, scanned("good.parquet"));
let Step::Install(loaded) = answer(&mut loader, id, schema_read("good.parquet")) else {
panic!("installs");
};
assert!(
matches!(loader.failed(id, "the rows"), Step::Nothing),
"the rows' failure is the table's"
);
assert!(!loaded.footers.is_cancelled());
assert!(
loader.retire().is_some(),
"going home puts down the read of the first rows"
);
assert!(
!loaded.footers.is_cancelled(),
"but not the dataset's footer pass"
);
}
#[test]
fn a_large_read_into_memory_is_asked_about_first() {
let dir = tempfile::tempdir().unwrap();
let json = dir.path().join("big.json");
std::fs::write(&json, "[{\"a\": 1}]").unwrap();
let csv = dir.path().join("big.csv");
std::fs::write(&csv, "a\n1\n").unwrap();
let open = |loader: &mut Loader, path: &Path, limit: u64| {
loader.open(OpenRequest {
warn_in_memory_above: Some(limit),
..request(&path.to_string_lossy())
})
};
let mut loader = Loader::default();
let Step::AskRead(read) = open(&mut loader, &json, 4) else {
panic!("a JSON file past the size is asked about");
};
assert_eq!(read.format, FileFormat::Json);
assert_eq!((read.bytes, read.files), (10, 1));
assert!(loader.asking() && loader.awaiting_dataset() && !loader.waits());
assert!(matches!(loader.confirmed(), Step::Scan { ref paths, .. } if paths[0] == json));
assert!(!loader.asking());
let mut loader = Loader::default();
assert!(matches!(open(&mut loader, &json, 10), Step::Scan { .. }));
let mut loader = Loader::default();
assert!(matches!(open(&mut loader, &csv, 0), Step::Scan { .. }));
let mut loader = Loader::default();
let _ = open(&mut loader, &json, 0);
assert!(loader.retire().is_some_and(|retired| retired.asking));
}
#[test]
fn a_large_journal_is_asked_about_first() {
let dir = tempfile::tempdir().unwrap();
let entry = "{\"__CURSOR\":\"s=1\",\"__REALTIME_TIMESTAMP\":\"1\",\"MESSAGE\":\"m\"}\n";
for name in ["journal", "journal.json"] {
let path = dir.path().join(name);
std::fs::write(&path, entry).unwrap();
let mut loader = Loader::default();
let step = loader.open(OpenRequest {
warn_in_memory_above: Some(4),
..request(&path.to_string_lossy())
});
let Step::AskRead(read) = step else {
panic!("{name}: a journal past the size is asked about");
};
assert_eq!(read.format, FileFormat::Journal, "{name}");
}
}
#[test]
fn a_footer_to_drop_says_the_file_is_counted() {
let footer = |n| OpenOptions {
skip_tail_rows: Some(n),
parse_strings: Some(crate::ParseStringsTarget::All),
..OpenOptions::default()
};
let mut loader = Loader::default();
let step = loader.open(OpenRequest {
options: footer(2),
..request("vendor.csv")
});
assert!(
matches!(step, Step::Scan { status, .. } if status == COUNTING_FOOTER_STATUS),
"{:?}",
loader.current().map(|l| l.phase().label())
);
assert_eq!(
loader.current().unwrap().phase().label(),
(COUNTING_FOOTER, 10)
);
for (path, n) in [("vendor.csv", 0), ("vendor.parquet", 2)] {
let mut loader = Loader::default();
let step = loader.open(OpenRequest {
options: footer(n),
..request(path)
});
assert!(
matches!(step, Step::Scan { status, .. } if status != COUNTING_FOOTER_STATUS),
"{path}"
);
}
}
#[test]
fn several_urls_end_the_session() {
let mut loader = Loader::default();
let step = loader.open(OpenRequest {
paths: vec![PathBuf::from("s3://a/x.csv"), PathBuf::from("s3://a/y.csv")],
options: OpenOptions::default(),
size: 0,
recent: None,
shown: None,
warn_in_memory_above: None,
});
assert!(matches!(step, Step::Crash(message) if message.contains("S3")));
assert!(loader.current().is_none());
}
#[test]
fn a_compressed_delimited_file_is_decompressed_as_its_format() {
let opts = |format: Option<FileFormat>| OpenOptions {
format,
..OpenOptions::default()
};
for (name, format, expected) in [
("x.tsv.gz", None, Some(FileFormat::Tsv)),
("x.PSV.xz", None, Some(FileFormat::Psv)),
("x.csv.zst", None, Some(FileFormat::Csv)),
("x.tsv", None, Some(FileFormat::Tsv)),
("x.gz", Some(FileFormat::Tsv), Some(FileFormat::Tsv)),
("x.csv.gz", Some(FileFormat::Psv), Some(FileFormat::Psv)),
("x.json.gz", None, None),
("x.gz", None, None),
("x.tsv.gz", Some(FileFormat::Parquet), None),
] {
assert_eq!(
delimited_format(Path::new(name), &opts(format)),
expected,
"{name} {format:?}"
);
}
let mut loader = Loader::default();
assert!(matches!(
loader.open(request("logs.tsv.gz")),
Step::Decompress { ref options, .. } if options.format == Some(FileFormat::Tsv)
));
}
#[test]
fn a_compressed_file_the_scan_found_is_decompressed() {
let compressed = || LoadAnswer::Compressed {
file: PathBuf::from("dir/data.tsv.gz"),
path: Some(PathBuf::from("dir")),
options: OpenOptions {
format: Some(FileFormat::Tsv),
..OpenOptions::default()
},
};
let mut loader = Loader::default();
assert!(matches!(loader.open(request("dir")), Step::Scan { .. }));
let id = loader.id().unwrap();
assert!(matches!(
answer(&mut loader, id, compressed()),
Step::Decompress { ref file, ref path, ref options, download: None, .. }
if file == Path::new("dir/data.tsv.gz")
&& path == Path::new("dir")
&& options.format == Some(FileFormat::Tsv)
));
assert_eq!(
loader.current().unwrap().phase().label(),
("Decompressing", 30)
);
assert!(matches!(
answer(&mut loader, id, compressed()),
Step::Nothing
));
assert!(matches!(
answer(&mut loader, id, schema_read("dir")),
Step::Install(_)
));
}
#[cfg(feature = "cloud")]
#[test]
fn a_remote_object_a_spec_reads_is_downloaded_first() {
let asked = [
OpenOptions {
spec_name: Some("vendor.feed".into()),
..OpenOptions::default()
},
OpenOptions {
spec_file: Some(PathBuf::from("feed.toml")),
..OpenOptions::default()
},
];
let open = |url: &str, options: &OpenOptions| {
Loader::default().open(OpenRequest {
options: options.clone(),
..request(url)
})
};
for options in &asked {
for url in [
"s3://b/day",
"s3://b/day.parquet",
"gs://b/day.arrow",
"abfss://c@acct.dfs.core.windows.net/day.l2",
] {
assert!(
matches!(open(url, options), Step::Probe(_)),
"{url} is downloaded"
);
}
for url in ["s3://b/days/", "gs://b/days/*.l2"] {
assert!(
matches!(open(url, options), Step::Scan { .. }),
"{url} goes on to the scan"
);
}
}
assert!(matches!(
open("s3://b/day.parquet", &OpenOptions::default()),
Step::Scan { .. }
));
}
#[test]
fn a_remote_spec_is_fetched_before_the_open_goes_on() {
let spec = Arc::new(
crate::formats::Spec::parse(
"name = \"t.rec\"\n[records]\nfields = [{ name = \"v\", type = \"u1\" }]\n",
None,
)
.unwrap(),
);
let options = OpenOptions {
spec_file: Some(PathBuf::from("https://example.com/t.toml")),
..OpenOptions::default()
};
let mut loader = Loader::default();
let Step::FetchSpec { url, options, .. } = loader.open(OpenRequest {
options,
..request("day.rec")
}) else {
panic!("the spec is fetched first");
};
assert_eq!(url, Path::new("https://example.com/t.toml"));
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading spec", 5)
);
let id = loader.id().unwrap();
let fetched = || LoadAnswer::SpecFetched {
spec: spec.clone(),
options: options.clone(),
};
let Step::Scan { paths, options, .. } = answer(&mut loader, id, fetched()) else {
panic!("the open goes on to its scan");
};
assert_eq!(paths, [PathBuf::from("day.rec")]);
assert!(options.spec_fetched.is_some_and(|s| s.name == "t.rec"));
assert!(matches!(answer(&mut loader, id, fetched()), Step::Nothing));
}
#[test]
fn a_compressed_file_a_spec_reads_is_decompressed_then_its_records_read() {
let spec = crate::formats::Spec::parse(
"name = \"t.rec\"\n[records]\nfields = [{ name = \"v\", type = \"u1\" }]\n",
None,
)
.unwrap();
let choice = crate::formats::Choice {
spec: Arc::new(spec),
by: crate::formats::Chosen::Glob,
also: Vec::new(),
};
let compressed = || LoadAnswer::CompressedRecords {
file: PathBuf::from("day.rec.zst"),
path: None,
choice: choice.clone(),
options: OpenOptions::default(),
};
let mut loader = Loader::default();
let _ = loader.open(request("day.rec.zst"));
let id = loader.id().unwrap();
let Step::DecompressRecords {
file, path, writer, ..
} = answer(&mut loader, id, compressed())
else {
panic!("the file is decompressed");
};
assert_eq!(
(file.as_path(), path.as_path()),
(Path::new("day.rec.zst"), Path::new("day.rec.zst"))
);
assert_eq!(
loader.current().unwrap().phase().label(),
("Decompressing", 30)
);
assert!(loader.waits());
assert!(matches!(
answer(&mut loader, id, compressed()),
Step::Nothing
));
loader.retire();
assert!(writer.stopped(), "Esc or another open stops the copy");
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let _ = loader.open(request("day.rec.zst"));
let id = loader.id().unwrap();
let _ = answer(&mut loader, id, compressed());
let copy = TempDownload::keep(TempDownload::create(Some(dir.path()), None).unwrap());
let temp = copy.path().to_path_buf();
let decompressed = || LoadAnswer::DecompressedRecords {
copy: copy.clone(),
path: PathBuf::from("day.rec.zst"),
choice: choice.clone(),
options: OpenOptions::default(),
};
let Step::ReadRecords { copy: read, .. } = answer(&mut loader, id, decompressed()) else {
panic!("the copy's records are read");
};
assert_eq!(read, temp);
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading records", 35)
);
assert!(matches!(
answer(&mut loader, id, decompressed()),
Step::Nothing
));
let Step::ReadSchema { made, .. } = answer(&mut loader, id, scanned("day.rec.zst")) else {
panic!("the frame's schema is read");
};
assert_eq!(made.converted.len(), 1, "the dataset holds the copy");
drop((made, copy));
loader.retire();
assert!(!temp.exists(), "the retired load let the copy go");
}
#[test]
fn gps_logs_the_scan_found_are_converted() {
let convert = || LoadAnswer::Convert {
what: Conversion::Text(FileFormat::Nmea),
files: vec![PathBuf::from("logs/a.nmea"), PathBuf::from("logs/b.nmea")],
bytes: 200,
path: Some(PathBuf::from("logs")),
options: OpenOptions {
format: Some(FileFormat::Nmea),
..OpenOptions::default()
},
};
let mut loader = Loader::default();
let _ = loader.open(request("logs"));
let id = loader.id().unwrap();
let Step::Convert {
what,
files,
writer,
read,
..
} = answer(&mut loader, id, convert())
else {
panic!("the logs are converted");
};
assert_eq!(what, Conversion::Text(FileFormat::Nmea));
assert_eq!(what.status(), "Reading NMEA...");
assert_eq!(files.len(), 2);
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading NMEA log", 10)
);
read.store(100, Ordering::Relaxed);
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading NMEA log", 20),
"the bar moves with the bytes read"
);
assert!(loader.waits());
assert!(matches!(answer(&mut loader, id, convert()), Step::Nothing));
assert!(!writer.stopped());
loader.retire();
assert!(
writer.stopped(),
"Ctrl+O or another open stops the conversion"
);
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let _ = loader.open(request("logs"));
let id = loader.id().unwrap();
assert!(matches!(
answer(&mut loader, id, convert()),
Step::Convert { .. }
));
let file =
TempDownload::keep(TempDownload::create(Some(dir.path()), Some("arrow")).unwrap());
let temp = file.path().to_path_buf();
let note = crate::notes::Note {
summary: "1 line left out: not NMEA".to_string(),
scope: "of 9 lines in 2 logs".to_string(),
read_as_text: None,
passed_over: None,
};
let Step::ReadSchema { made, path, .. } = answer(
&mut loader,
id,
LoadAnswer::Converted {
converted: Converted::Frame {
files: vec![file],
lf: Box::new(frame()),
notes: vec![note],
other_tables: vec!["GSV 3".to_string()],
detail: None,
},
path: Some(PathBuf::from("logs")),
options: OpenOptions::default(),
},
) else {
panic!("the frame's schema is read");
};
assert_eq!(path.as_deref(), Some(Path::new("logs")));
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading schema", 40)
);
assert_eq!(made.converted.len(), 1);
assert_eq!(made.notes.len(), 1);
assert_eq!(made.other_tables, ["GSV 3"]);
drop(made);
let Step::Failed(failed) = loader.failed(id, &format!("could not read {}", temp.display()))
else {
panic!("the open fails");
};
assert_eq!(failed.message, "could not read logs", "named as asked for");
assert!(!temp.exists(), "the retired load let the files go");
}
#[test]
fn a_database_of_several_tables_lands_on_them() {
let tables = |file: &str| LoadAnswer::Tables {
file: PathBuf::from(file),
tables: vec!["users".to_string(), "orders".to_string()],
path: Some(PathBuf::from(file)),
};
let mut loader = Loader::default();
let _ = loader.open(request("shop.db"));
let id = loader.id().unwrap();
let Step::Tables(landed) = answer(&mut loader, id, tables("shop.db")) else {
panic!("the tables are listed");
};
assert_eq!(landed.database, Path::new("shop.db"));
assert!(!landed.from_home);
assert!(loader.current().is_none(), "the load is put down");
assert!(matches!(
answer(&mut loader, id, tables("shop.db")),
Step::Nothing
));
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let Step::Spool { .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
OpenOptions::default(),
&Default::default(),
)) else {
panic!("standard input is read first");
};
let id = loader.id().unwrap();
let spool = TempDownload::keep(TempDownload::create(Some(dir.path()), None).unwrap());
let options = OpenOptions {
format: Some(FileFormat::Sqlite),
..Default::default()
};
let Step::Scan { paths, .. } = answer(
&mut loader,
id,
LoadAnswer::Spooled {
download: spool,
options,
},
) else {
panic!("the spool is scanned");
};
let Step::Failed(failed) = answer(
&mut loader,
id,
LoadAnswer::Tables {
file: paths[0].clone(),
tables: vec!["users".to_string(), "orders".to_string()],
path: Some(PathBuf::from("stdin")),
},
) else {
panic!("piped in, the table has to be named");
};
assert_eq!(
failed.message,
"stdin holds 2 tables: users, orders. Open one with --table NAME."
);
}
#[test]
fn a_path_inside_a_database_opens_the_database() {
let dir = tempfile::tempdir().unwrap();
let db = dir.path().join("shop.db");
let mut header = crate::sqlite::MAGIC.to_vec();
header.resize(512, 0);
std::fs::write(&db, header).unwrap();
let request = OpenRequest::named(
vec![db.join("orders")],
OpenOptions::default(),
&Default::default(),
);
assert_eq!(request.paths, std::slice::from_ref(&db));
assert_eq!(request.options.table.as_deref(), Some("orders"));
assert_eq!(request.recent, Some(db.join("orders")));
assert_eq!(request.shown, Some(db.join("orders")));
let named = OpenRequest::named(
vec![db.clone()],
OpenOptions {
table: Some("users".to_string()),
..Default::default()
},
&Default::default(),
);
assert_eq!(named.paths, std::slice::from_ref(&db));
assert_eq!(
named.recent,
Some(db.join("users")),
"--table is recorded too"
);
let plain = OpenRequest::named(
vec![db.clone()],
OpenOptions::default(),
&Default::default(),
);
assert_eq!(plain.recent, Some(db));
assert_eq!(plain.shown, None);
}
#[test]
fn the_first_step_follows_what_was_asked_for() {
let mut loader = Loader::default();
assert!(matches!(
loader.open(request("logs.csv.gz")),
Step::Decompress { ref file, ref path, .. }
if file == Path::new("logs.csv.gz") && path == Path::new("logs.csv.gz")
));
assert_eq!(
loader.current().unwrap().phase().label(),
("Decompressing", 30)
);
let id = loader.id().unwrap();
assert!(matches!(
answer(&mut loader, id, schema_read("logs.csv.gz")),
Step::Install(_)
));
let mut loader = Loader::default();
let mut parsed = request("table.csv");
parsed.options.parse_strings = Some(crate::ParseStringsTarget::All);
assert!(matches!(
loader.open(parsed),
Step::Scan {
status: "Scanning string columns...",
..
}
));
let mut loader = Loader::default();
let Step::ReadSchema { path: None, .. } =
loader.open_frame(frame(), OpenOptions::default())
else {
panic!("a frame's schema is read");
};
let id = loader.id().unwrap();
let Step::Install(loaded) = answer(
&mut loader,
id,
LoadAnswer::SchemaRead {
state: state(),
path: None,
options: OpenOptions::default(),
debug_label: None,
},
) else {
panic!("installs");
};
assert!(loaded.paths.is_none() && loaded.recent.is_none());
}
fn downloaded(dir: &Path, body: &str, extension: &str) -> TempDownload {
use std::io::Write;
let mut file = TempDownload::create(Some(dir), Some(extension)).unwrap();
file.write_all(body.as_bytes()).unwrap();
TempDownload::keep(file)
}
#[cfg(feature = "http")]
#[test]
fn a_remote_model_reads_its_headers_or_falls_back_to_the_download() {
let url = "https://example.com/m/model.safetensors";
let jobs = jobs();
let mut loader = Loader::default();
let Step::ReadHeaders {
url: read, format, ..
} = loader.open(request(url))
else {
panic!("the headers are read, not downloaded");
};
assert_eq!(
(read.as_path(), format),
(Path::new(url), FileFormat::Safetensors)
);
let id = loader.id().unwrap();
assert_eq!(
loader.current().unwrap().phase().label(),
("Reading headers", 20)
);
assert!(loader.waits());
let Step::Install(loaded) = loader.answered(id, schema_read(url), &jobs) else {
panic!("the headers are the dataset");
};
assert_eq!(loaded.recent.as_deref(), Some(Path::new(url)));
let mut loader = Loader::default();
let _ = loader.open(request(url));
let id = loader.id().unwrap();
let Step::Probe(pending) = loader.answered(
id,
LoadAnswer::NoRanges {
options: OpenOptions::default(),
},
&jobs,
) else {
panic!("no ranges: the file is sized for its download");
};
assert_eq!(pending.parts().0, url);
assert_eq!(loader.download_note(), None, "not asked yet");
assert!(matches!(
loader.answered(id, LoadAnswer::Sized(pending), &jobs),
Step::Ask(_)
));
assert_eq!(loader.download_note(), Some(NO_RANGES));
assert!(matches!(
loader.answered(id, schema_read(url), &jobs),
Step::Nothing
));
}
#[cfg(feature = "cloud")]
#[test]
fn arrow_in_a_bucket_is_listed_first() {
let open = |url: &str, format: Option<FileFormat>| {
let mut loader = Loader::default();
let mut request = request(url);
request.options.format = format;
request.options.hive = url.ends_with('/');
let step = loader.open(request);
(loader, step)
};
for url in ["s3://lake/hf/", "gs://lake/hf/", "s3://lake/one.arrow"] {
let (_, step) = open(url, url.ends_with('/').then_some(FileFormat::Arrow));
let Step::Probe(PendingDownload::Arrow { url: listed, .. }) = step else {
panic!("{url}: listed first");
};
assert_eq!(listed, url);
}
assert!(matches!(
open("s3://lake/hf/", Some(FileFormat::Csv)).1,
Step::Scan { .. }
));
assert!(matches!(
open("s3://lake/hf/*.arrow", Some(FileFormat::Arrow)).1,
Step::Scan { .. }
));
let object = |name: &str, stream: bool| crate::cloud_arrow::Object {
url: format!("s3://lake/hf/{name}"),
size: 10,
stream,
};
let jobs = jobs();
let (mut loader, _) = open("s3://lake/hf/", Some(FileFormat::Arrow));
let id = loader.id().unwrap();
let listed = |objects| {
LoadAnswer::Sized(PendingDownload::Arrow {
url: "s3://lake/hf/".to_string(),
objects,
size: Some(0),
options: OpenOptions::default(),
})
};
let Step::Scan { paths, options, .. } = loader.answered(
id,
listed(vec![object("a.arrow", false), object("b.arrow", false)]),
&jobs,
) else {
panic!("IPC files only: scanned in place");
};
assert_eq!(paths, [PathBuf::from("s3://lake/hf/")]);
assert_eq!(
options.arrow_parts.as_deref(),
Some(&vec![
crate::ipc_stream::Part::InPlace(PathBuf::from("s3://lake/hf/a.arrow")),
crate::ipc_stream::Part::InPlace(PathBuf::from("s3://lake/hf/b.arrow")),
])
);
let (mut loader, _) = open("s3://lake/hf/", Some(FileFormat::Arrow));
let id = loader.id().unwrap();
assert!(matches!(
loader.answered(
id,
listed(vec![object("a.arrow", false), object("s.arrow", true)]),
&jobs
),
Step::Ask(PendingDownload::Arrow { .. })
));
}
#[cfg(feature = "http")]
#[test]
fn a_small_catalog_file_downloads_without_asking() {
let url = "https://example.com/penguins.csv";
let jobs = jobs();
let unasked = |listed| crate::UnaskedDownload {
limit: 1_000,
listed,
};
let ask = |options: Option<crate::UnaskedDownload>, size: Option<u64>| {
let mut loader = Loader::default();
let mut request = request(url);
request.options.download_unasked = options;
let Step::Probe(pending) = loader.open(request) else {
panic!("the size is asked first");
};
let id = loader.id().unwrap();
let step = loader.answered(id, LoadAnswer::Sized(pending.with_size(size)), &jobs);
(matches!(step, Step::Download { .. }), loader)
};
let (downloads, loader) = ask(Some(unasked(None)), Some(800));
assert!(downloads, "under the limit: no question");
assert!(!loader.asking());
assert_eq!(
loader.current().unwrap().phase().label(),
("Downloading", 20)
);
assert!(
ask(Some(unasked(Some(800))), None).0,
"the catalog's size stands in"
);
assert!(
!ask(Some(unasked(Some(800))), Some(5_000)).0,
"the server's size wins"
);
assert!(!ask(Some(unasked(None)), None).0, "nothing says how big");
assert!(!ask(None, Some(10)).0, "a URL from anywhere else asks");
}
#[cfg(feature = "http")]
#[test]
fn an_unasked_download_past_its_limit_asks_once() {
let url = "https://example.com/penguins.csv";
let jobs = jobs();
let unasked = Some(crate::UnaskedDownload {
limit: 1_000,
listed: Some(800),
});
let mut loader = Loader::default();
let mut asked_for = request(url);
asked_for.options.download_unasked = unasked;
let Step::Probe(pending) = loader.open(asked_for) else {
panic!("the size is asked first");
};
let id = loader.id().unwrap();
let Step::Download { pending, .. } =
loader.answered(id, LoadAnswer::Sized(pending.with_size(None)), &jobs)
else {
panic!("small by the catalog: no question");
};
assert!(
pending.parts().2.download_unasked.is_some(),
"it has a limit"
);
let Step::Ask(asked) = loader.answered(id, LoadAnswer::PastLimit(pending), &jobs) else {
panic!("past the limit, it asks");
};
assert_eq!(asked.parts().1, None, "its size is not what anyone said");
assert_eq!(loader.download_note(), Some(PAST_LIMIT));
assert!(jobs.would_strand(), "the generation is held while it asks");
let Step::Download { pending, .. } = loader.confirmed() else {
panic!("agreed to, it downloads");
};
assert!(pending.parts().2.download_unasked.is_none(), "asked once");
let mut loader = Loader::default();
let mut asked_for = request(url);
asked_for.options.download_unasked = unasked;
let Step::Probe(pending) = loader.open(asked_for) else {
panic!("the size is asked first");
};
let id = loader.id().unwrap();
let step = loader.answered(id, LoadAnswer::Sized(pending.with_size(Some(5_000))), &jobs);
assert!(matches!(step, Step::Ask(_)));
let Step::Download { pending, .. } = loader.confirmed() else {
panic!("agreed to, it downloads");
};
assert!(pending.parts().2.download_unasked.is_none());
}
#[cfg(any(feature = "http", feature = "cloud"))]
#[test]
fn the_past_limit_note_names_the_limit() {
assert_eq!(crate::UnaskedDownload::LIMIT, 50 * 1024 * 1024);
assert!(PAST_LIMIT.contains("50 MB"));
}
#[cfg(feature = "http")]
#[test]
fn a_remote_file_is_asked_about_then_downloaded_and_kept() {
let url = "https://example.com/data.csv";
let dir = tempfile::tempdir().unwrap();
let jobs = jobs();
let mut loader = Loader::default();
let Step::Probe(pending) = loader.open(request(url)) else {
panic!("the size is asked first");
};
let id = loader.id().unwrap();
assert_eq!(
loader.current().unwrap().phase().label(),
("Checking size", 0)
);
let sized = pending.with_size(Some(42));
assert!(matches!(
loader.answered(id, LoadAnswer::Sized(sized), &jobs),
Step::Ask(ref p) if p.parts().1 == Some(42)
));
assert!(loader.asking() && loader.awaiting_dataset());
assert!(!loader.waits(), "the question has the keys");
assert!(jobs.would_strand(), "the generation is held while it asks");
let Step::Download { .. } = loader.confirmed() else {
panic!("agreed to, it downloads");
};
assert!(!jobs.would_strand(), "the hold goes with the question");
assert!(matches!(loader.confirmed(), Step::Nothing), "asked once");
assert_eq!(
loader.current().unwrap().phase().label(),
("Downloading", 20)
);
let file = downloaded(dir.path(), "a\n1\n", "csv");
let at = file.path().to_path_buf();
let Step::Scan {
paths,
display,
status,
..
} = loader.answered(
id,
LoadAnswer::Downloaded {
download: file,
options: OpenOptions::default(),
},
&jobs,
)
else {
panic!("the download is scanned");
};
assert_eq!(paths, vec![at.clone()]);
assert_eq!(display.as_deref(), Some(Path::new(url)), "named by its URL");
assert_eq!(status, "Scanning...");
let Step::ReadSchema {
made: Made {
download: Some(download),
..
},
..
} = loader.answered(id, scanned(url), &jobs)
else {
panic!("the schema read is handed the download");
};
assert_eq!(download.path(), at, "the file the scan reads");
let state = DataTableState::from_lazyframe(frame(), &OpenOptions::default())
.unwrap()
.with_open(crate::widgets::datatable::OpenFacts {
download: Some(download),
..Default::default()
});
let read = LoadAnswer::SchemaRead {
state: Box::new(state),
path: Some(PathBuf::from(url)),
options: OpenOptions::default(),
debug_label: None,
};
let Step::Install(loaded) = loader.answered(id, read, &jobs) else {
panic!("installs");
};
assert!(
loaded.state.scans_a_download(),
"the dataset holds its file"
);
loader.first_rows_settled();
drop(loaded);
assert!(at.exists(), "and the loader keeps it to read again");
let Step::Scan { paths, display, .. } = loader.open(request(url)) else {
panic!("the copy on hand is scanned, with no question");
};
assert_eq!(paths, vec![at.clone()]);
assert_eq!(display.as_deref(), Some(Path::new(url)));
assert!(loader.make_way().is_some());
let _ = loader.open(request("local.csv"));
assert!(!at.exists(), "nothing else held it");
}
#[cfg(feature = "http")]
#[test]
fn quitting_mid_schema_read_sweeps_the_download_the_worker_holds() {
use std::time::{Duration, Instant};
let url = "https://example.com/data.csv";
let dir = tempfile::tempdir().unwrap();
let jobs = jobs();
let mut loader = Loader::default();
let Step::Probe(pending) = loader.open(request(url)) else {
panic!("the size is asked first");
};
let id = loader.id().unwrap();
let _ = loader.answered(id, LoadAnswer::Sized(pending.with_size(Some(4))), &jobs);
let Step::Download { writer, .. } = loader.confirmed() else {
panic!("agreed to, it downloads");
};
let file = crate::download::read_to_temp(
Some(dir.path()),
Some("csv"),
|| Ok((std::io::Cursor::new(b"a\n1\n".to_vec()), Some(4))),
&writer,
None,
)
.unwrap();
let at = file.path().to_path_buf();
let downloaded = LoadAnswer::Downloaded {
download: file,
options: OpenOptions::default(),
};
let _ = loader.answered(id, downloaded, &jobs);
let Step::ReadSchema {
made: Made {
download: Some(held),
..
},
..
} = loader.answered(id, scanned(url), &jobs)
else {
panic!("the schema read is handed the download");
};
let unfinished = loader.unfinished().clone();
drop(loader);
assert!(at.exists(), "the worker's copy keeps it");
unfinished.sweep(Instant::now() + Duration::from_millis(50));
assert!(!at.exists(), "the sweep removes it");
drop(held);
}
#[cfg(feature = "http")]
#[test]
fn a_remote_open_put_down_leaves_nothing_behind() {
let url = "https://example.com/data.csv";
let dir = tempfile::tempdir().unwrap();
let jobs = jobs();
let mut loader = Loader::default();
let Step::Probe(pending) = loader.open(request(url)) else {
panic!("probe");
};
let id = loader.id().unwrap();
let _ = loader.answered(id, LoadAnswer::Sized(pending), &jobs);
let retired = loader.retire().expect("declined");
assert!(retired.asking, "the question goes with it");
assert!(!jobs.would_strand(), "and so does its hold");
let Step::Probe(pending) = loader.open(request(url)) else {
panic!("probe");
};
let old = loader.id().unwrap();
let _ = loader.answered(old, LoadAnswer::Sized(pending), &jobs);
let Step::Download { writer, .. } = loader.confirmed() else {
panic!("download");
};
assert!(loader.make_way().is_some());
assert!(writer.stopped(), "the download in flight is told to stop");
let _ = loader.open(request("local.csv"));
let file = downloaded(dir.path(), "a\n1\n", "csv");
let at = file.path().to_path_buf();
assert!(matches!(
loader.answered(
old,
LoadAnswer::Downloaded {
download: file,
options: OpenOptions::default()
},
&jobs
),
Step::Nothing
));
assert!(!at.exists(), "the stale download's file went with it");
assert!(loader.make_way().is_some());
let Step::Probe(pending) = loader.open(request(url)) else {
panic!("probe");
};
let id = loader.id().unwrap();
let _ = loader.answered(id, LoadAnswer::Sized(pending), &jobs);
let _ = loader.confirmed();
let file = downloaded(dir.path(), "a\n1\n", "csv");
let at = file.path().to_path_buf();
let _ = loader.answered(
id,
LoadAnswer::Downloaded {
download: file,
options: OpenOptions::default(),
},
&jobs,
);
let reason = format!(
"\"{}\": Bad csv.\nIt stopped at {}.",
at.display(),
at.display()
);
let Step::Failed(failed) = loader.failed(id, &reason) else {
panic!("the open fails");
};
assert_eq!(
failed.message,
format!("\"{url}\": Bad csv.\nIt stopped at {url}."),
"named by the URL opened, never the temp file (#511)"
);
assert!(!at.exists(), "a failed open keeps no download");
}
#[cfg(feature = "http")]
#[test]
fn a_downloaded_compressed_csv_is_decompressed() {
for (ext, format) in [
("csv.gz", FileFormat::Csv),
("tsv.zst", FileFormat::Tsv),
("psv.xz", FileFormat::Psv),
] {
let url = format!("https://example.com/data.{ext}");
let dir = tempfile::tempdir().unwrap();
let jobs = jobs();
let mut loader = Loader::default();
let Step::Probe(pending) = loader.open(request(&url)) else {
panic!("probe");
};
let id = loader.id().unwrap();
let _ = loader.answered(id, LoadAnswer::Sized(pending), &jobs);
let _ = loader.confirmed();
let file = downloaded(dir.path(), "", ext);
let step = loader.answered(
id,
LoadAnswer::Downloaded {
download: file,
options: OpenOptions::default(),
},
&jobs,
);
assert!(
matches!(
step,
Step::Decompress { ref path, ref options, .. }
if path == Path::new(&url) && options.format == Some(format)
),
"{ext}"
);
}
}
#[test]
fn stdin_is_spooled_then_read_as_stdin() {
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let request = OpenRequest::named(
vec![PathBuf::from("-")],
OpenOptions::default(),
&Default::default(),
);
assert_eq!(request.recent, None, "not recorded in recents");
let Step::Spool { writer, read, .. } = loader.open(request) else {
panic!("standard input is read first");
};
let id = loader.id().unwrap();
let load = loader.current().unwrap();
assert_eq!(load.phase().label(), ("Reading stdin", 5));
assert_eq!(load.path(), Some(Path::new("stdin")));
read.store(1234, Ordering::Relaxed);
assert_eq!(loader.current().unwrap().size(), 1234, "what has come in");
assert!(loader.waits());
let file = {
use std::io::Write;
let mut file = TempDownload::create(Some(dir.path()), None).unwrap();
file.write_all(b"a\n1\n").unwrap();
TempDownload::keep(file)
};
let at = file.path().to_path_buf();
let options = OpenOptions {
format: Some(FileFormat::Csv),
..Default::default()
};
let Step::Scan { paths, display, .. } = answer(
&mut loader,
id,
LoadAnswer::Spooled {
download: file,
options: options.clone(),
},
) else {
panic!("the file is scanned");
};
assert_eq!(paths, vec![at.clone()]);
assert_eq!(display.as_deref(), Some(Path::new("stdin")));
let Step::ReadSchema {
made: Made {
download: Some(_), ..
},
..
} = answer(&mut loader, id, scanned("stdin"))
else {
panic!("the dataset holds the file");
};
let Step::Install(loaded) = answer(&mut loader, id, schema_read("stdin")) else {
panic!("installs");
};
assert_eq!(loaded.recent, None);
assert_eq!(loaded.paths.as_deref(), Some(&[PathBuf::from("-")][..]));
loader.first_rows_settled();
drop(writer);
let Step::Scan { paths, display, .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
options,
&Default::default(),
)) else {
panic!("the copy on hand is read again");
};
assert_eq!(paths, vec![at]);
assert_eq!(display.as_deref(), Some(Path::new("stdin")));
assert!(loader.make_way().is_some());
let step = loader.open(OpenRequest::named(
vec![PathBuf::from("-"), PathBuf::from("a.csv")],
OpenOptions::default(),
&Default::default(),
));
assert!(matches!(step, Step::Crash(_)));
}
#[test]
fn a_converted_pipe_keeps_the_copy_not_the_stream() {
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let options = OpenOptions {
format: Some(FileFormat::Arrow),
..Default::default()
};
let Step::Spool { .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
OpenOptions::default(),
&Default::default(),
)) else {
panic!("standard input is read first");
};
let id = loader.id().unwrap();
let spool = downloaded(dir.path(), "stream", "tmp");
let spooled = spool.path().to_path_buf();
let Step::Scan { paths, .. } = answer(
&mut loader,
id,
LoadAnswer::Spooled {
download: spool,
options: options.clone(),
},
) else {
panic!("the spool is scanned");
};
let Step::Convert { .. } = answer(
&mut loader,
id,
LoadAnswer::Convert {
what: Conversion::Streams,
files: paths,
bytes: 6,
path: Some(PathBuf::from("stdin")),
options: options.clone(),
},
) else {
panic!("the stream is converted");
};
let copy = downloaded(dir.path(), "ipc file", "arrow");
let converted = copy.path().to_path_buf();
let Step::Scan { paths, display, .. } = answer(
&mut loader,
id,
LoadAnswer::Converted {
converted: Converted::Streams {
file: copy,
parts: Vec::new(),
},
path: Some(PathBuf::from("stdin")),
options: options.clone(),
},
) else {
panic!("the copy is scanned");
};
assert_eq!(paths, std::slice::from_ref(&converted));
assert_eq!(display.as_deref(), Some(Path::new("stdin")));
assert!(!spooled.exists(), "the spool goes once converted");
let Step::ReadSchema {
made: Made {
download: Some(_), ..
},
..
} = answer(&mut loader, id, scanned("stdin"))
else {
panic!("the dataset holds the copy");
};
let Step::Install(loaded) = answer(&mut loader, id, schema_read("stdin")) else {
panic!("installs");
};
loader.first_rows_settled();
drop(loaded);
assert!(converted.exists(), "kept to read again");
let Step::Scan { paths, .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
options,
&Default::default(),
)) else {
panic!("the copy on hand is read again");
};
assert_eq!(paths, [converted]);
}
#[test]
fn piped_compressed_csv_is_decompressed_and_a_stop_reaches_the_spool() {
let dir = tempfile::tempdir().unwrap();
let mut loader = Loader::default();
let Step::Spool { writer, .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
OpenOptions::default(),
&Default::default(),
)) else {
panic!("spool");
};
assert!(loader.make_way().is_some(), "Ctrl+O puts it down");
assert!(writer.stopped(), "and the spool stops at its next chunk");
let Step::Spool { .. } = loader.open(OpenRequest::named(
vec![PathBuf::from("-")],
OpenOptions::default(),
&Default::default(),
)) else {
panic!("spool");
};
let id = loader.id().unwrap();
let file = TempDownload::keep(TempDownload::create(Some(dir.path()), None).unwrap());
let step = answer(
&mut loader,
id,
LoadAnswer::Spooled {
download: file,
options: OpenOptions {
format: Some(FileFormat::Csv),
compression: Some(CompressionFormat::Gzip),
..Default::default()
},
},
);
assert!(matches!(
step,
Step::Decompress { ref path, download: Some(_), .. } if path == Path::new("stdin")
));
}
}