pub(crate) mod counting;
pub(crate) mod first_rows_trace;
pub mod follow;
pub(crate) mod local_glob;
pub mod measurements;
pub(crate) mod open_options;
pub(crate) mod open_scan;
pub(crate) mod scan;
pub mod stdin;
pub mod tee;
pub(crate) mod unfinished;
use std::path::{Path, PathBuf};
use std::sync::Arc;
use std::sync::atomic::{AtomicU64, Ordering};
use polars::prelude::LazyFrame;
use crate::cloud::download::TempDownload;
use crate::formats::schema_union::FooterProgress;
use crate::loading::unfinished::{Unfinished, Writer};
use crate::table::DataTableState;
use crate::{CompressionFormat, FileFormat, OpenOptions, cloud::source};
use crate::app::jobs::Hold;
#[cfg(any(feature = "http", feature = "cloud"))]
use crate::app::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::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::formats::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::formats::members::holder(&first).is_some()
{
table = Some(crate::formats::members::place(&first, name));
} else if local
&& paths.len() == 1
&& options.table.is_none()
&& let Some((dir, split)) = crate::formats::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::formats::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::formats::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::formats::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::formats::ipc_stream::Part>>,
splits: Option<Arc<crate::formats::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::formats::ipc_stream::Part>,
},
Frame {
files: Vec<TempDownload>,
lf: Box<LazyFrame>,
notes: Vec<crate::notes::Note>,
other_tables: Vec<String>,
detail: Option<Arc<crate::formats::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::formats::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::home_preview::Prepared) -> Step {
let crate::home::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::loading::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::cloud::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::cloud_arrow::in_place(&objects).is_some() => {
load.phase = Phase::Scanning { downloaded: false };
let parts = crate::cloud::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::formats::readers::refined(path, named).unwrap_or(named)),
None => crate::formats::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 MiB, 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;