Skip to main content

datui_lib/loading/
open_scan.rs

1//! Opening a dataset: routing what was named, downloads, the scan and schema read
2//! for each kind of source, the load's phases as jobs, and installing the result.
3
4use crate::app::feedback::Confirm;
5use crate::app::jobs::{Answer, Job};
6use crate::app::modals::pivot_melt_modal::PivotMeltModal;
7use crate::app::modals::sort_filter_modal::SortFilterModal;
8use crate::cli::{CompressionFormat, FileFormat};
9#[cfg(feature = "cloud")]
10use crate::cloud::cloud_hive;
11use crate::loading::open_options::{OpenOptions, ReadReport, UnaskedDownload};
12use crate::loading::scan::Scan;
13use crate::table::{DataTableState, OpenFacts};
14#[cfg(feature = "cloud")]
15use crate::wait_on_runtime;
16use crate::{
17    App, AppEvent, UNSUPPORTED, cli, cloud::source, formats::dataset_files, home, home::catalog,
18    home::discover, loading,
19};
20use color_eyre::Result;
21#[cfg(feature = "cloud")]
22use polars::io::cloud::{AmazonS3ConfigKey, CloudOptions};
23use polars::prelude::{LazyFrame, Schema, col};
24#[cfg(feature = "cloud")]
25use polars::prelude::{PlRefPath, ScanArgsParquet};
26use std::path::{Path, PathBuf};
27use std::sync::{Arc, Mutex};
28
29/// Where the dataset on screen came from, and how it was opened.
30#[derive(Default)]
31pub struct OpenedSource {
32    pub(crate) original_file_format: Option<crate::export::export_modal::ExportFormat>,
33    pub(crate) original_file_delimiter: Option<u8>,
34    /// The paths and options the dataset on screen was opened with: what `H` reopens.
35    pub(crate) opened: Option<(Vec<PathBuf>, OpenOptions)>,
36    /// Whether the dataset was reached through home: `q` returns there, else quits.
37    pub(crate) opened_from_home: bool,
38    /// `--view NAME`, taken by the first install so later opens are not re-dressed.
39    pub(crate) startup_view: Option<String>,
40    /// The dataset whose downloaded shape is kept already. See
41    /// [`crate::App::remember_a_downloads_shape`].
42    pub(crate) shape_remembered: Option<u64>,
43}
44
45impl OpenedSource {
46    /// A new dataset is on screen, opened from `opened` (none for a frame handed over)
47    /// and read as `format`, with `delimiter`: what `H` reopens and an export starts
48    /// from. Reached through home, `q` goes back there; a reread (`H`) never unsets it,
49    /// since it is not a new place.
50    pub(crate) fn reset_for_dataset(
51        &mut self,
52        opened: Option<(Vec<PathBuf>, OpenOptions)>,
53        from_home: bool,
54        format: Option<crate::export::export_modal::ExportFormat>,
55        delimiter: Option<u8>,
56    ) {
57        self.opened = opened;
58        self.opened_from_home |= from_home;
59        self.original_file_format = format;
60        self.original_file_delimiter = delimiter;
61    }
62}
63
64/// What a cloud open was pointed at: the URL as given, the prefix to list, and the
65/// glob to keep.
66#[cfg(feature = "cloud")]
67struct CloudTarget<'a> {
68    /// The URL as typed, which is where the bucket and scheme come from.
69    full: &'a str,
70    /// The literal prefix to list: the whole key, or the part of a glob before its star.
71    key: String,
72    /// The glob the user named, if any. The listing keeps only matching keys, so
73    /// downstream sees plain files.
74    pattern: Option<&'a globset::GlobMatcher>,
75}
76
77/// Put hive partition columns first. `drifts` keeps the scan's hidden drift column,
78/// which the select would drop. Free so the table's state can reuse it when
79/// rebuilding a scan.
80pub(crate) fn hoist_partition_columns(
81    lf: LazyFrame,
82    schema: &Schema,
83    partition_columns: &[String],
84    drifts: bool,
85) -> LazyFrame {
86    if partition_columns.is_empty() {
87        return lf;
88    }
89    let mut exprs: Vec<_> = partition_columns
90        .iter()
91        .map(|s| col(s.as_str()))
92        .chain(
93            schema
94                .iter_names()
95                .map(|s| s.to_string())
96                .filter(|c| !partition_columns.contains(c))
97                .map(|s| col(s.as_str())),
98        )
99        .collect();
100    if drifts {
101        exprs.push(col(crate::formats::schema_union::DRIFT_COLUMN));
102    }
103    lf.select(exprs)
104}
105
106impl App {
107    /// Whether an open is on its way and not installed: `data_table_state` still holds
108    /// the outgoing dataset, so the main view shows the open's progress.
109    pub(crate) fn awaiting_dataset(&self) -> bool {
110        self.loading.awaiting_dataset()
111    }
112
113    /// The open's phase, percentage, path and size, for the loading screen and footer.
114    pub(crate) fn load_shown(&self) -> Option<(&str, u16, Option<&Path>, u64)> {
115        self.loading.current().map(|load| {
116            let (phase, percent) = load.phase().label();
117            (phase, percent, load.path(), load.size())
118        })
119    }
120
121    /// What the load is doing. While a footer pass runs its count stands in for the
122    /// phase, read once a frame from [`crate::loading::counting::Counting::footers_this_frame`] so callers agree.
123    pub(crate) fn loading_phase<'a>(&self, phase: &'a str) -> std::borrow::Cow<'a, str> {
124        match self.counting.footers_this_frame {
125            Some((read, total)) => std::borrow::Cow::Owned(format!(
126                "Reading footers: {} of {}",
127                crate::numfmt::group_chrome(read),
128                crate::numfmt::group_chrome(total)
129            )),
130            // A listing has no total to count towards, so it says how far it has got.
131            None => match self.counting.listed_this_frame {
132                Some(listed) => std::borrow::Cow::Owned(format!(
133                    "Listing files: {}",
134                    crate::numfmt::group_chrome(listed)
135                )),
136                None => std::borrow::Cow::Borrowed(phase),
137            },
138        }
139    }
140
141    /// An open is on its way: the loading screen shows `phase` and keys wait. Called
142    /// before the event that carries the open out (by `run`, or a key before its
143    /// `Open`), since a frame drawn between would show the outgoing dataset.
144    pub fn set_loading_phase(&mut self, phase: impl Into<String>, progress_percent: u16) {
145        self.announce_open(false, phase.into(), progress_percent);
146    }
147
148    /// As [`Self::set_loading_phase`]; `from_home` reports a failure on home.
149    pub(crate) fn announce_open(&mut self, from_home: bool, phase: String, percent: u16) {
150        self.put_down_load_in_flight();
151        self.loading.announce(from_home, phase, percent);
152    }
153
154    /// Put the path on the loading screen, so a wait says what it is waiting for.
155    pub(crate) fn name_what_is_loading(&mut self, path: PathBuf) {
156        self.loading.name(path);
157    }
158
159    /// Put down the load in flight unless it has started no work (the look or frame
160    /// leading to this open).
161    pub(crate) fn put_down_load_in_flight(&mut self) {
162        if let Some(retired) = self.loading.make_way() {
163            self.put_down_load(retired);
164        }
165    }
166
167    /// Start an open with its request: the load in flight goes, and so does the
168    /// screen's dataset's own reading.
169    pub(crate) fn begin_new_dataset(&mut self) {
170        self.put_down_load_in_flight();
171        // A preview's dataset this open did not take is unwanted.
172        self.home_app.previews.drop_prepared();
173        self.reset_chart_state();
174        self.jobs.advance();
175        // The open counts footers on its own counter, which the dataset takes over on
176        // install.
177        self.counting.stop_footer_pass();
178    }
179
180    /// Put down what the app keeps for a retired load: its jobs, their footer lines,
181    /// and the question about its download.
182    pub(crate) fn put_down_load(&mut self, retired: loading::Retired) {
183        let id = retired.id;
184        let lines = self.jobs.quiet(|job| job.load() == Some(id));
185        self.jobs.supersede(|job| job.load() == Some(id));
186        if self
187            .status_message
188            .as_ref()
189            .is_some_and(|status| lines.contains(status))
190        {
191            self.status_message = None;
192        }
193        if retired.asking {
194            self.confirmation_modal.hide();
195        }
196    }
197
198    /// Carry out what the open needs next.
199    pub(crate) fn run_load_step(&mut self, step: loading::Step) -> Option<AppEvent> {
200        use loading::Step;
201        let load = self.loading.id();
202        match step {
203            Step::Nothing => None,
204            Step::Crash(message) => Some(AppEvent::Crash(message)),
205            Step::Failed(failed) => {
206                self.load_failed(failed);
207                None
208            }
209            Step::Tables(tables) => {
210                self.land_on_tables(tables);
211                None
212            }
213            Step::Hex(hex) => {
214                self.land_on_hex(hex);
215                None
216            }
217            Step::Install(loaded) => {
218                // A view the open applies reads its own first rows.
219                if self.install_dataset(*loaded) {
220                    return None;
221                }
222                #[cfg(test)]
223                {
224                    self.counting.first_rows_asked += 1;
225                }
226                if !self.spawn_async_collect(Self::LOADING_BUFFER) {
227                    // Nothing to read: the buffer already serves the view.
228                    if self.status_message.as_deref() == Some(Self::LOADING_BUFFER) {
229                        self.status_message = None;
230                    }
231                    self.first_rows_settled();
232                }
233                None
234            }
235            #[cfg(any(feature = "http", feature = "cloud"))]
236            Step::Ask(pending) => {
237                // Nothing runs while the question is up (a spinner would read as progress); the
238                // loader holds the generation.
239                self.confirmation_modal.show(
240                    Self::download_confirmation_message(&pending, self.loading.download_note()),
241                    Confirm::Download,
242                );
243                None
244            }
245            Step::AskRead(read) => {
246                // As for a download: nothing runs, and the generation is held.
247                self.loading.hold_while_asking(self.jobs.hold());
248                self.confirmation_modal.show(
249                    Self::in_memory_confirmation_message(&read),
250                    Confirm::Download,
251                );
252                None
253            }
254            step => {
255                let load = load.expect("a step that runs work belongs to the open in flight");
256                self.spawn_load_phase(load, step);
257                None
258            }
259        }
260    }
261
262    /// Keep a downloaded dataset's shape under its URL once counted: nothing lists a
263    /// web file, so this is how its recent and catalog row show `344 × 9`. Once per
264    /// dataset.
265    pub(crate) fn remember_a_downloads_shape(&mut self) {
266        if self.source.shape_remembered == Some(self.dataset_generation) {
267            return;
268        }
269        let Some(url) = self.path.clone().filter(|p| source::is_remote_url(p)) else {
270            return;
271        };
272        let Some(state) = self.data_table_state.as_ref().filter(|s| s.fetched()) else {
273            return;
274        };
275        let Some(rows) = state.num_rows_if_valid().filter(|_| !state.changes_rows()) else {
276            return;
277        };
278        self.source.shape_remembered = Some(self.dataset_generation);
279        let columns: Vec<String> = state
280            .source_schema()
281            .iter_names()
282            .map(|name| name.to_string())
283            .collect();
284        let facts = crate::cache::DatasetFacts {
285            mtime: std::time::SystemTime::now()
286                .duration_since(std::time::UNIX_EPOCH)
287                .map(|d| d.as_secs())
288                .unwrap_or_default(),
289            size: 0,
290            rows: Some(rows),
291            cols: Some(columns.len()),
292            cols_sampled: false,
293            columns,
294            kind: Some(discover::EntryKind::File),
295            classified_by: discover::CLASSIFIER_VERSION,
296            cost: Default::default(),
297            holds: Default::default(),
298        };
299        // Off the UI thread: the index takes a lock other instances may hold.
300        let cache = self.cache.clone();
301        self.cache_writes
302            .spawn(move || cache.record_dataset_facts(&[(url, facts)]));
303    }
304
305    /// Install the dataset an open read and apply its view, if any. Returns whether
306    /// that view reads the first rows. Only the loader calls this, for the open in
307    /// flight.
308    pub(crate) fn install_dataset(&mut self, loaded: loading::Loaded) -> bool {
309        let loading::Loaded {
310            state,
311            path,
312            options,
313            debug_label,
314            paths,
315            recent,
316            from_home,
317            footers,
318        } = loaded;
319        let options = &options;
320        // One per dataset reaching the screen, not per open: a failed open leaves the last
321        // dataset up, and its footer pass must still finish into it.
322        self.dataset_generation = self.dataset_generation.wrapping_add(1);
323        // Each part of the app lets go of what it kept of the last dataset.
324        self.counting.reset_for_dataset(footers);
325        self.quality.reset_for_dataset();
326        self.analysis_modal.quality.reset_for_dataset();
327        self.prompt.reset_for_dataset();
328        self.sample.reset_for_dataset();
329        self.views.reset_for_dataset();
330        self.info
331            .reset_for_dataset(&self.app_config, path.as_deref());
332        self.sort_filter_modal = SortFilterModal::new();
333        self.pivot_melt_modal = PivotMeltModal::new();
334        // A view waiting on its pivot, a sample being drawn and a chart being prepared
335        // were over the replaced dataset.
336        self.jobs.supersede(|job| matches!(job, Job::ViewPivot(_)));
337        self.put_down_sample_draw();
338        self.reset_chart_state();
339        self.debug.schema_load = debug_label;
340        // A frame handed over has no path. The spec read is dropped: it holds the file's
341        // map, and a decompressed copy keeps its disk space while mapped.
342        let opened = paths.map(|paths| {
343            let options = OpenOptions {
344                format_read: None,
345                sqlite: None,
346                // Counted afresh by the next read.
347                tail: None,
348                prepared: None,
349                ..options.clone()
350            };
351            (paths, options)
352        });
353        let (format, delimiter) = match path.as_deref() {
354            Some(p) => (
355                Self::export_format_for(p, state.read_as().or(options.format)),
356                // CSV's delimiter: a comma unless the user named one. A `.tsv` exports as
357                // TSV (tab preset); a tab in a `.csv` would reopen as one column.
358                Some(options.separator_or(b',')),
359            ),
360            None => (None, None),
361        };
362        self.source
363            .reset_for_dataset(opened, from_home, format, delimiter);
364        // Recorded only once installed: a file that fails to load is not worth returning to.
365        if let Some(path) = recent {
366            // Off the opening path: the cache lock may be held by other instances. Only the
367            // next home listing waits on it.
368            let cache = self.cache.clone();
369            self.cache_writes.spawn(move || {
370                cache.push_recent(&path);
371            });
372        }
373        self.forget_the_rows_read();
374        self.data_table_state = Some(state);
375        // A followed file's watcher starts with its dataset and stops with it.
376        if options.follow
377            && let Some(state) = self.data_table_state.as_mut()
378        {
379            match options.tail.as_deref() {
380                Some(tail) => {
381                    let follow = crate::loading::follow::Follow::start(
382                        tail.clone(),
383                        self.app_config.read.follow_interval.duration(),
384                        self.events.clone(),
385                        options.spool.clone(),
386                    );
387                    state.start_following(if options.pipe {
388                        follow.as_pipe()
389                    } else {
390                        follow
391                    });
392                    // Already counted by the scan.
393                    state.follow_to(tail.rows(), false);
394                }
395                // A recording of a format that cannot be read as it grows.
396                None => self.flash_note(
397                    "Only text and Arrow streams are followed: this shows what had arrived, and recording goes on"
398                        .to_string(),
399                ),
400            }
401        }
402        // A count waiting on the last dataset's paint is no longer owed.
403        self.retire_a_count_the_rows_answered();
404        self.path = path.clone();
405        // Named for this file, so after its path is set.
406        self.open_info_documentation();
407        // A panel still up says what it says about the dataset on screen.
408        if self.overlay.shows(&crate::Overlay::Info) {
409            self.read_file_facts();
410            self.count_unfit();
411        }
412        // What the dataset still has to learn about itself is read behind it.
413        self.start_pending_footers();
414        self.index_lines();
415        // `#` for text and logs, unless the flag or the config said.
416        if options.row_numbers_auto
417            && let Some(state) = self.data_table_state.as_mut()
418            && state.numbered_by_default()
419        {
420            state.set_row_numbers(true);
421        }
422        self.status_message = Some(Self::LOADING_BUFFER.to_string());
423
424        // Where a view meets the dataset: `--view` names one for this first open only;
425        // `[views] auto_apply` dresses every open with a matching view. A fresh dataset
426        // starts with none applied, so the last file's view is not checked here.
427        let (view, reason) = match self.source.startup_view.take() {
428            Some(name) => match self.views.manager.get_view_by_name(&name).cloned() {
429                Some(view) => (Some(view), None),
430                None => {
431                    self.error_modal.show(format!("No view named \"{name}\""));
432                    (None, None)
433                }
434            },
435            None if self.app_config.views.auto_apply => self
436                .view_dataset()
437                .zip(self.data_table_state.as_ref())
438                .and_then(|(dataset, state)| {
439                    self.views
440                        .manager
441                        .get_most_relevant(dataset, state.source_schema())
442                })
443                .map_or((None, None), |(view, reason)| (Some(view), Some(reason))),
444            None => (None, None),
445        };
446        let Some(view) = view else {
447            return false;
448        };
449        let applied = match reason {
450            // Applied unasked, it says which view and why.
451            Some(why) => self.apply_matched_view(&view, why),
452            None => self.apply_view(&view),
453        };
454        match applied {
455            // The view reads its own first rows, so the dataset's are never read.
456            Ok(()) => true,
457            Err(e) => {
458                self.error_modal
459                    .show(format!("Error applying view \"{}\": {e}", view.name));
460                false
461            }
462        }
463    }
464
465    /// Enter the home screen, rebuilt, with the cursor on what is open, abandoning any
466    /// load. `loading::Loader::retire` supersedes the open's jobs (their answers are
467    /// dropped) and raises its stop flag, so downloads and cloud passes stop within a
468    /// wave. Work not the open's (an export, an analysis, the screen's footer pass) is
469    /// left running, with its progress and completion modal.
470    pub fn abandon_load(&mut self) {
471        let retired = self.loading.retire();
472        if let Some(retired) = retired {
473            self.put_down_load(retired);
474        }
475        // A chart being prepared would keep the throbber up at home, and could land later
476        // in another dataset with the same column names.
477        self.reset_chart_state();
478        // A look in flight may never return (a gone share): supersede it so nothing waits.
479        if self.jobs.supersede(|job| matches!(job, Job::Classify(_))) {
480            self.home.status = None;
481        }
482        // A collect owed behind this load goes too; it would read the abandoned dataset
483        // at home with every key held.
484        self.jobs.take_owed(Self::owed_rows);
485        // Only the open's wait is put down: an export holds keys too and keeps running.
486        // Rows the open's last step reads still land, unwaited.
487        if retired.is_some() {
488            self.busy = false;
489            let quieted = self.jobs.quiet(Self::reading_rows);
490            if self
491                .status_message
492                .as_ref()
493                .is_some_and(|status| quieted.contains(status))
494            {
495                self.status_message = None;
496            }
497        }
498        // Keys typed at the frozen screen were meant for the load; replayed at home they
499        // could open something unasked.
500        self.screen_generation = self.screen_generation.wrapping_add(1);
501    }
502
503    /// Whether finding out what a path is could hang on an unresponsive mount: only a
504    /// local path on what home calls a network mount. URLs are the scan's business,
505    /// and plain local paths answer at once.
506    pub(crate) fn looking_could_block(&self, path: &Path) -> bool {
507        // `cloud://<id>` is a place: `input_source` calls it local, `is_remote_path`
508        // remote, and a worker stat would report "No such path".
509        !home::is_cloud_place(path)
510            && matches!(source::input_source(path), source::InputSource::Local(_))
511            && (self.home.network_check)(path)
512    }
513
514    /// Browse into a path, report it as a lake table, or open it, as its kind calls for.
515    /// `jump` (a path typed at `~`) starts a new browse so Esc returns to the listing.
516    pub(crate) fn open_what_it_is(
517        &mut self,
518        path: PathBuf,
519        kind: discover::EntryKind,
520        jump: bool,
521    ) -> Option<AppEvent> {
522        let go_inside = |app: &mut Self, path: PathBuf| {
523            if jump {
524                app.home_jump_into(path);
525            } else {
526                app.home_browse_into(path);
527            }
528        };
529        if kind == discover::EntryKind::Directory {
530            go_inside(self, path);
531            return None;
532        }
533        // No reader: a local file opens in the hex view; a remote one is dimmed and its
534        // details pane says why.
535        if kind == discover::EntryKind::Other {
536            if matches!(source::input_source(&path), source::InputSource::Local(_)) {
537                self.open_hex(path, crate::app::hex_view::Origin::Home, true, None);
538            }
539            return None;
540        }
541        // A lake table's files are not its rows (tombstoned and rewritten files stay on
542        // disk), so datui goes inside rather than give a wrong answer.
543        if let Some(format) = kind.lake_name() {
544            self.home.lake_here = Some((path.clone(), format));
545            go_inside(self, path);
546            return None;
547        }
548        let directory = matches!(
549            kind,
550            discover::EntryKind::Hive | discover::EntryKind::MultiFile
551        );
552        // A directory typed at `~` is a place to go, as →: naming one never reads all of it.
553        if directory && jump {
554            go_inside(self, path);
555            return None;
556        }
557        // A cloud directory that is a dataset opens as a prefix, scanning every file under
558        // it.
559        if directory && home::is_object_store_url(&path) {
560            return Some(self.home_open_path(home::directory_dataset_url(&path), false));
561        }
562        // Said here, where the file was named, rather than after a download and load that
563        // could only fail; a typed path has no dimmed row to say it. A format spec may
564        // read it by glob or by magic.
565        let a_spec_may_read = !self.formats.by_glob(&path, false).is_empty()
566            || self.formats.specs.iter().any(|f| !f.spec.magic.is_empty());
567        // A table inside a file of tables (`flight.ulg/sensor_accel.1`) or a log found by
568        // magic (`00000042.BIN`) has a name that says nothing.
569        if kind == discover::EntryKind::File
570            && discover::unreadable_by_name(&path)
571            && !a_spec_may_read
572            && crate::formats::members::split(&path).is_none()
573            && crate::formats::members::holder(&path).is_none()
574            && crate::formats::members::split_variant(&path, &self.formats).is_none()
575            && crate::formats::hf_splits::split_place(&path).is_none()
576        {
577            self.home.status = Some(discover::NO_READER.to_string());
578            return None;
579        }
580        // The preview read this file's first page through the open's own steps: install
581        // that rather than read again.
582        let prepared = (!directory)
583            .then(|| self.home_app.previews.take_prepared(&path))
584            .flatten();
585        // A small built-in catalog file is fetched unasked (its row gave its size); a typed
586        // URL still asks.
587        let unasked = self
588            .home
589            .catalogs
590            .iter()
591            .filter(|c| c.origin == catalog::Origin::Bundled)
592            .flat_map(|c| c.datasets.iter())
593            .find(|d| d.location == path)
594            .filter(|_| {
595                !jump && matches!(source::input_source(&path), source::InputSource::Http(_))
596            })
597            .map(|dataset| UnaskedDownload {
598                limit: UnaskedDownload::LIMIT,
599                listed: dataset.size,
600            });
601        match self.home_open_path(path, directory) {
602            AppEvent::Open(paths, mut options) => {
603                options.prepared = prepared.map(|p| Arc::new(Mutex::new(Some(p))));
604                options.download_unasked = unasked;
605                Some(AppEvent::Open(paths, options))
606            }
607            event => Some(event),
608        }
609    }
610
611    /// What `datui <path>` does with a directory, by the same rule as `Enter` on its
612    /// row ([`home::look_into_as`]): a hive root or a one-table directory opens as one
613    /// table, any other opens home browsed into it. `--hive` still forces partition
614    /// columns. Stats the path, so `run` calls it on a worker
615    /// ([`AppEvent::OpenNamed`]); returns `LookThenOpenDirectory` or `Open`.
616    pub fn route_named_paths(paths: Vec<PathBuf>, options: OpenOptions) -> AppEvent {
617        Self::route_named_paths_with(paths, options, &crate::formats::Registry::default())
618    }
619
620    /// [`Self::route_named_paths`] with the format specs: a directory a spec reads as
621    /// column files opens rather than being looked at.
622    pub fn route_named_paths_with(
623        paths: Vec<PathBuf>,
624        options: OpenOptions,
625        formats: &crate::formats::Registry,
626    ) -> AppEvent {
627        if let Some(event) = Self::route_named_without_looking(&paths, &options) {
628            return event;
629        }
630        // Several paths are files read together, and `--hive` already answers; neither asks
631        // what one directory is.
632        let single = (paths.len() == 1 && !options.hive).then(|| paths[0].clone());
633        let Some(dir) = single.filter(|p| p.is_dir()) else {
634            return AppEvent::Open(paths, options);
635        };
636        // A format spec named for it, or one whose glob names it, reads it as columns.
637        if options.spec_file.is_some()
638            || options.spec_name.is_some()
639            || !formats.by_glob(&dir, true).is_empty()
640        {
641            return AppEvent::Open(paths, options);
642        }
643        // Looking reads footers or the front of a spread of files: seconds for large
644        // Parquet, so a worker does it.
645        AppEvent::LookThenOpenDirectory(dir, options)
646    }
647
648    /// The part of [`Self::route_named_paths`] that needs no filesystem. A cloud
649    /// directory is looked at by one listing page, whose contents pick the reader as
650    /// for `(all files)`. A glob, a file name or `--format` already says what to read.
651    pub(crate) fn route_named_without_looking(
652        paths: &[PathBuf],
653        options: &OpenOptions,
654    ) -> Option<AppEvent> {
655        #[cfg(feature = "cloud")]
656        if let [dir] = paths
657            && !options.hive
658            && home::is_object_store_url(dir)
659            && options.format.is_none()
660            && !dir.to_string_lossy().contains('*')
661            && !home::names_a_file(dir)
662        {
663            return Some(AppEvent::LookThenOpenDirectory(
664                dir.clone(),
665                options.clone(),
666            ));
667        }
668        let _ = (paths, options);
669        None
670    }
671
672    /// The first named local path that is not there. URLs and globs are left to the
673    /// open; stdin is no path.
674    pub fn missing_named_path(
675        paths: &[PathBuf],
676        formats: &crate::formats::Registry,
677    ) -> Option<PathBuf> {
678        paths
679            .iter()
680            .find(|path| {
681                !source::is_remote_url(path)
682                    && !crate::loading::stdin::is_stdin(path)
683                    && !source::expands_as_glob(path)
684                    && !path.exists()
685                    && crate::formats::members::split(path).is_none()
686                    && crate::formats::members::split_variant(path, formats).is_none()
687            })
688            .cloned()
689    }
690
691    /// Act on what the look at a directory named on the command line found; the rule is
692    /// at [`Self::route_named_paths`].
693    pub(crate) fn open_the_directory_looked_at(
694        &mut self,
695        dir: PathBuf,
696        kind: discover::EntryKind,
697        holds: Option<&discover::Holds>,
698        mut options: OpenOptions,
699    ) -> Option<AppEvent> {
700        #[cfg(feature = "cloud")]
701        if home::is_object_store_url(&dir) {
702            return self.open_the_cloud_directory_looked_at(dir, kind, holds, options);
703        }
704        let _ = holds;
705        // No override of the reader settings: the look read every file as this open will,
706        // so `--no-header` and skips are already accounted for. A lake table opens home
707        // with the same refusal its row gives.
708        if let Some(format) = kind.lake_name() {
709            self.enter_home();
710            self.home.lake_here = Some((dir.clone(), format));
711            self.home_jump_into(dir);
712            return None;
713        }
714        // One table: read it. `hive` puts the open on the directory route, where the
715        // contents pick the reader.
716        if matches!(
717            kind,
718            discover::EntryKind::Hive | discover::EntryKind::MultiFile
719        ) {
720            options.hive = true;
721            self.set_loading_phase("Scanning input", 10);
722            self.name_what_is_loading(dir.clone());
723            return Some(AppEvent::Open(vec![dir], options));
724        }
725        // A place to look inside: `datui .`, or a directory of separate tables (whose
726        // `(all files)` row unions them).
727        self.enter_home();
728        self.home_jump_into(dir);
729        None
730    }
731
732    /// As [`Self::open_the_directory_looked_at`] for a cloud directory: its `(all files)`
733    /// row's open, or a browse when nothing directly inside is readable.
734    #[cfg(feature = "cloud")]
735    fn open_the_cloud_directory_looked_at(
736        &mut self,
737        dir: PathBuf,
738        kind: discover::EntryKind,
739        holds: Option<&discover::Holds>,
740        options: OpenOptions,
741    ) -> Option<AppEvent> {
742        let open = |app: &mut Self, path: PathBuf, options: OpenOptions| {
743            app.set_loading_phase("Scanning input", 10);
744            app.name_what_is_loading(path.clone());
745            Some(AppEvent::Open(vec![path], options))
746        };
747        // The listing was refused: the open says why in the refuser's words.
748        let Some(holds) = holds else {
749            return open(self, dir, options);
750        };
751        if let Some(format) = kind.lake_name() {
752            self.enter_home();
753            self.home.lake_here = Some((dir.clone(), format));
754            self.home_jump_into(dir);
755            return None;
756        }
757        let directory = home::directory_dataset_url(&dir);
758        if matches!(
759            kind,
760            discover::EntryKind::Hive | discover::EntryKind::MultiFile
761        ) {
762            let options = OpenOptions {
763                hive: true,
764                ..options
765            };
766            return open(self, directory, options);
767        }
768        if let Some((format, left_out)) = Self::cloud_prefix_format(holds) {
769            let options = OpenOptions {
770                hive: true,
771                format: Some(format),
772                left_out,
773                ..options
774            };
775            return open(self, directory, options);
776        }
777        // Only directories, or nothing readable: browse inside, with the reason if any.
778        self.enter_home();
779        self.home_jump_into(dir);
780        self.home.status = Self::why_a_cloud_prefix_cannot_be_read(holds);
781        None
782    }
783
784    /// The options a home-started open reads with: the config's read and CSV settings,
785    /// as the command line would give them.
786    pub(crate) fn open_defaults(&self) -> OpenOptions {
787        match crate::cli::parse_args(["datui"]) {
788            Ok(args) => OpenOptions::from_args_and_config(&args, &self.app_config),
789            Err(_) => OpenOptions::default(),
790        }
791    }
792
793    /// `options` for the compressed delimited `file` in the dialect of the delimited
794    /// spec it matches, if any. Such files skip the scan that matches the others and
795    /// go straight to decompression.
796    fn with_delimited_spec(
797        file: &Path,
798        mut options: OpenOptions,
799        formats: &crate::formats::Registry,
800    ) -> Result<OpenOptions> {
801        if options.delimited.is_some() {
802            return Ok(options);
803        }
804        let asked = crate::formats::Asked {
805            spec_file: options.spec_file.clone(),
806            spec: options.spec_fetched.clone(),
807            spec_name: options.spec_name.clone(),
808            compression: options.compression,
809            ..Default::default()
810        };
811        let crate::formats::Route::Delimited(choice) =
812            crate::formats::route(file, &asked, formats).map_err(|e| color_eyre::eyre::eyre!(e))?
813        else {
814            return Ok(options);
815        };
816        let Some(delimited) = choice.spec.delimited.clone() else {
817            return Ok(options);
818        };
819        delimited.apply(&mut options);
820        let chosen = crate::formats::delimited_spec::DelimitedRead::chosen(
821            choice.spec,
822            choice.by,
823            choice.also,
824        );
825        let read =
826            crate::formats::delimited_spec::read_facts(&chosen, &[file.to_path_buf()], &options)?;
827        options.delimited = Some(Arc::new(read));
828        Ok(options)
829    }
830
831    /// Read a compressed CSV, TSV or PSV into a table state, split on its separator.
832    /// The one input that cannot be scanned lazily (minutes for a large export), so it
833    /// takes no `&self` and runs on a worker.
834    fn decompressed_delimited_state(
835        path: &Path,
836        options: &OpenOptions,
837        writer: &crate::loading::unfinished::Writer,
838    ) -> Result<DataTableState> {
839        let separator = options
840            .format
841            .and_then(FileFormat::separator)
842            .unwrap_or(b',');
843        DataTableState::from_read(
844            crate::formats::readers::csv::read_delimited(path, separator, options, writer)?,
845            options,
846        )
847    }
848
849    /// Polars' view of one source's S3 settings, for `scan_parquet`.
850    #[cfg(feature = "cloud")]
851    fn build_s3_cloud_options(settings: &crate::cloud::cloud_sources::S3Settings) -> CloudOptions {
852        let settings = settings.clone();
853        let virtual_hosted = (settings.endpoint.is_some() || settings.virtual_hosted.is_some())
854            .then(|| settings.virtual_hosted_style().to_string());
855        let configs: Vec<(AmazonS3ConfigKey, String)> = [
856            (AmazonS3ConfigKey::Endpoint, settings.endpoint),
857            (AmazonS3ConfigKey::AccessKeyId, settings.access_key_id),
858            (
859                AmazonS3ConfigKey::SecretAccessKey,
860                settings.secret_access_key,
861            ),
862            (AmazonS3ConfigKey::Token, settings.session_token),
863            (AmazonS3ConfigKey::Region, settings.region),
864            (AmazonS3ConfigKey::VirtualHostedStyleRequest, virtual_hosted),
865            (
866                AmazonS3ConfigKey::SkipSignature,
867                settings.skip_signature.then(|| "true".to_string()),
868            ),
869        ]
870        .into_iter()
871        .filter_map(|(key, value)| value.map(|v| (key, v)))
872        .chain([(
873            AmazonS3ConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
874            crate::cloud::user_agent::get(),
875        )])
876        .collect();
877        CloudOptions::default().with_aws(configs)
878    }
879
880    /// The bucket and key of an `s3://bucket/key` or `gs://bucket/key` URL. The key
881    /// is empty for a bucket root.
882    #[cfg(feature = "cloud")]
883    pub(crate) fn cloud_bucket_and_key(url: &str) -> Result<(String, String)> {
884        if let Some((_, container, key)) = source::azure_parts(url) {
885            return Ok((container, key.trim_matches('/').to_string()));
886        }
887        crate::cloud::cloud_browse::split_bucket_url(url)
888            .map(|(_, bucket, key)| (bucket, key))
889            .ok_or_else(|| {
890                color_eyre::eyre::eyre!("URL must be s3://bucket/key or gs://bucket/key")
891            })
892    }
893
894    /// The store Polars will scan `url` through, from its cache keyed on bucket and
895    /// `options`, so footer reads, size probes and downloads share its credentials,
896    /// TLS client and connection pool.
897    #[cfg(feature = "cloud")]
898    fn polars_object_store(
899        url: &str,
900        options: &CloudOptions,
901        runtime: &tokio::runtime::Handle,
902    ) -> Result<Arc<dyn object_store::ObjectStore>> {
903        let url = url.to_string();
904        let options = options.clone();
905        wait_on_runtime(runtime, async move {
906            let (_, store) = polars::io::cloud::build_object_store(
907                PlRefPath::new(url.as_str()),
908                Some(&options),
909                false,
910            )
911            .await?;
912            polars::prelude::PolarsResult::Ok(store.to_dyn_object_store().await.into_owned())
913        })
914        .ok_or_else(|| color_eyre::eyre::eyre!("cancelled"))?
915        .map_err(|e| color_eyre::eyre::eyre!("Object store config failed: {}", e))
916    }
917
918    /// An HTTP agent with a total time budget. ureq 3 bounds only through agent
919    /// config, and `timeout_global` covers the whole exchange, so a server dribbling a
920    /// byte every 29 seconds cannot hang the TUI.
921    #[cfg(feature = "http")]
922    fn http_agent(total: std::time::Duration) -> ureq::Agent {
923        crate::cloud::user_agent::ureq_config()
924            .timeout_global(Some(total))
925            .build()
926            .into()
927    }
928
929    /// What a HEAD says an HTTP(S) file weighs; `None` if unsaid. An error only when
930    /// the file cannot be had (404, no server); one refusing HEAD may still send it.
931    #[cfg(feature = "http")]
932    pub(crate) fn fetch_remote_size_http(
933        url: &str,
934    ) -> std::result::Result<Option<u64>, crate::error_display::HttpGone> {
935        let agent = Self::http_agent(std::time::Duration::from_secs(15));
936        // ureq asks for gzip and drops Content-Length from compressed answers (GitHub
937        // Pages compresses); identity gets the on-disk length.
938        match agent.head(url).header("Accept-Encoding", "identity").call() {
939            Ok(r) => Ok(r
940                .headers()
941                .get("Content-Length")
942                .and_then(|v| v.to_str().ok())
943                .and_then(|s| s.parse::<u64>().ok())),
944            Err(e) => crate::error_display::http_gone(url, &e).map_or(Ok(None), Err),
945        }
946    }
947
948    /// The size of one S3 or GCS object, from a HEAD through the shared store.
949    #[cfg(feature = "cloud")]
950    fn fetch_remote_size_cloud(
951        url: &str,
952        cloud: &crate::config::CloudConfig,
953        runtime: &tokio::runtime::Handle,
954    ) -> Result<Option<u64>> {
955        use object_store::ObjectStoreExt;
956
957        let (_bucket, key) = Self::cloud_bucket_and_key(url)?;
958        if key.is_empty() {
959            return Ok(None);
960        }
961        let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)?;
962        let path = crate::cloud::cloud_browse::object_path(&key);
963        let head = wait_on_runtime(runtime, async move { store.head(&path).await });
964        Ok(head.and_then(|r| r.ok()).map(|meta| meta.size))
965    }
966
967    /// Download `url` to a temp file; `stop` ends it early, even while the server is
968    /// silent, and any failure removes the file (see
969    /// [`crate::cloud::download::read_to_temp`]). Past `limit` bytes it fails with
970    /// [`crate::cloud::download::PastLimit`].
971    #[cfg(feature = "http")]
972    fn download_http_to_temp(
973        url: &str,
974        temp_dir: Option<&Path>,
975        extension: Option<&str>,
976        limit: Option<u64>,
977        writer: &crate::loading::unfinished::Writer,
978    ) -> Result<crate::cloud::download::TempDownload> {
979        use crate::cloud::download::StreamError;
980
981        let url = url.to_string();
982        let open = move || {
983            let agent = Self::http_agent(std::time::Duration::from_secs(300));
984            // ureq answers a 4xx or 5xx with an error, so every failure is said here.
985            let response = agent
986                .get(&url)
987                .call()
988                .map_err(|e| crate::error_display::http_message(&url, &e))?;
989            // No length: ureq decompresses, and the Content-Length was the wire's.
990            Ok((response.into_body().into_reader(), None))
991        };
992        crate::cloud::download::read_to_temp(temp_dir, extension, open, writer, limit).map_err(
993            |error| match error {
994                StreamError::Open(message) => color_eyre::eyre::eyre!(message),
995                StreamError::Read(e) => {
996                    color_eyre::eyre::eyre!("Download failed partway. Check your connection: {e}")
997                }
998                StreamError::Short { expected, got } => color_eyre::eyre::eyre!(
999                    "Download failed partway: it ended after {got} of {expected} bytes."
1000                ),
1001                StreamError::Write(report) => report,
1002                StreamError::Cut => color_eyre::eyre::eyre!("Download was cancelled."),
1003            },
1004        )
1005    }
1006
1007    /// Stream one S3, GCS or Azure object to a temp file, a few chunks in memory at a
1008    /// time ([`crate::cloud::download`]). Errors name its scheme; `writer`'s open stopping ends
1009    /// it, and any failure removes the file.
1010    #[cfg(feature = "cloud")]
1011    fn download_cloud_to_temp(
1012        url: &str,
1013        cloud: &crate::config::CloudConfig,
1014        options: &OpenOptions,
1015        runtime: &tokio::runtime::Handle,
1016        writer: &crate::loading::unfinished::Writer,
1017    ) -> Result<crate::cloud::download::TempDownload> {
1018        use crate::cloud::download::StreamError;
1019        use object_store::ObjectStoreExt;
1020
1021        let (label, example) = match source::input_source(Path::new(url)) {
1022            source::InputSource::Gcs(_) => ("GCS", "gs://bucket/path/file.csv"),
1023            source::InputSource::Azure(_) => (
1024                "Azure",
1025                "abfss://container@account.dfs.core.windows.net/path/file.csv",
1026            ),
1027            _ => ("S3", "s3://bucket/path/file.csv"),
1028        };
1029        let ext = source::download_suffix(url);
1030        let (_bucket, key) = Self::cloud_bucket_and_key(url)?;
1031        if key.is_empty() {
1032            return Err(crate::error_display::FileError::new(
1033                Path::new(url),
1034                format!("a {label} URL names an object here, such as {example}"),
1035            )
1036            .into());
1037        }
1038        let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)?;
1039
1040        let path = crate::cloud::cloud_browse::object_path(&key);
1041        let open = async move {
1042            let got = store
1043                .get(&path)
1044                .await
1045                .map_err(|e| crate::error_display::store_message(&e))?;
1046            let len = got.range.end - got.range.start;
1047            Ok((got.into_stream(), Some(len)))
1048        };
1049        let failed = |what: String| -> color_eyre::Report {
1050            crate::error_display::FileError::new(Path::new(url), what).into()
1051        };
1052        crate::cloud::download::stream_to_temp(
1053            runtime,
1054            options.temp_dir.as_deref(),
1055            ext.as_deref(),
1056            open,
1057            writer,
1058        )
1059        .map_err(|error| match error {
1060            StreamError::Open(e) => failed(e),
1061            StreamError::Read(e) => failed(format!("the download stopped: {e}")),
1062            StreamError::Short { expected, got } => {
1063                failed(format!("it ended after {got} of {expected} bytes"))
1064            }
1065            StreamError::Write(report) => report,
1066            StreamError::Cut => failed("the download was cancelled".to_string()),
1067        })
1068    }
1069
1070    /// The open's scan of `paths`, named `path`: the frame, or what must happen before
1071    /// there is one. Run by the open's `Scan` phase and by the home preview, which hands
1072    /// its result to the open.
1073    pub(crate) fn scan_for_open(
1074        cloud: &crate::config::CloudConfig,
1075        formats: &crate::formats::Registry,
1076        paths: &[PathBuf],
1077        options: OpenOptions,
1078        path: Option<PathBuf>,
1079    ) -> std::result::Result<loading::LoadAnswer, String> {
1080        use loading::LoadAnswer;
1081        let bytes_of = |files: &[PathBuf]| -> u64 {
1082            files
1083                .iter()
1084                .filter_map(|f| std::fs::metadata(f).ok())
1085                .map(|m| m.len())
1086                .sum()
1087        };
1088        // What the read passed over rides back with the options, so the dataset can say
1089        // what it left out. Seeded with the caller's knowledge and overwritten by the read:
1090        // on disk the read decides, but for an object-store prefix Polars lists and never
1091        // sees other formats, so home's listing is the only witness.
1092        let mut report = ReadReport {
1093            left_out: options.left_out.clone(),
1094            files_disagree: options.files_disagree,
1095            format: None,
1096            format_read: None,
1097            read_python: Vec::new(),
1098            sqlite: None,
1099            opened: None,
1100            splits: options.splits.clone(),
1101            delimited: None,
1102            table: None,
1103            guessed: false,
1104            read_notes: Vec::new(),
1105            typing: Default::default(),
1106        };
1107        // A followed file reads what it can and counts the rest: a misfit row never stops
1108        // the follow.
1109        let options = OpenOptions {
1110            ignore_errors: options.ignore_errors || options.follow,
1111            ..options
1112        };
1113        let named = |e: color_eyre::Report| {
1114            crate::error_display::user_message_from_report(&e, path.as_deref())
1115        };
1116        // A followed Arrow IPC stream gets its own scan, not a conversion; followed NDJSON
1117        // is scanned rather than read whole.
1118        let followed_stream = options.follow
1119            && crate::loading::follow::followed_stream(
1120                &paths[0],
1121                Some(crate::loading::follow::format_of(&paths[0], options.format)),
1122                &options,
1123            );
1124        let scan = if followed_stream {
1125            crate::loading::follow::stream::scan(&paths[0])
1126                .map(Scan::from)
1127                .map_err(|e| color_eyre::eyre::eyre!(e))
1128        } else {
1129            Self::build_lazyframe_from_paths_with(cloud, paths, &options, &mut report, formats)
1130        }
1131        // Named as the dataset is: a download by its URL, not its temp file.
1132        .map_err(named)?;
1133        let format = scan.format(report.format.or(options.format));
1134        // Bounded to complete records and counted for the watcher. A recording that cannot
1135        // be followed is read as it stands.
1136        let recording = options
1137            .spool
1138            .as_ref()
1139            .is_some_and(|handle| handle.spool().tee().is_some());
1140        let (scan, tail) = match scan {
1141            Scan::Frame(lf) if options.follow => {
1142                let format = crate::loading::follow::format_of(&paths[0], format);
1143                let refused = (!followed_stream)
1144                    .then(|| crate::loading::follow::refusal(Some(format), &options))
1145                    .flatten();
1146                match refused {
1147                    Some(_) if recording => (Scan::Frame(lf), None),
1148                    Some(refusal) => return Err(refusal),
1149                    None => {
1150                        let (lf, tail) = crate::loading::follow::bound_to_complete(
1151                            *lf, &paths[0], format, &options,
1152                        )
1153                        .map_err(named)?;
1154                        (Scan::Frame(Box::new(lf)), Some(Arc::new(tail)))
1155                    }
1156                }
1157            }
1158            _ if options.follow && !recording => {
1159                return Err(crate::loading::follow::refusal(format, &options)
1160                    .unwrap_or_else(|| "This file cannot be followed as it grows.".to_string()));
1161            }
1162            scan => (scan, None),
1163        };
1164        let read_mode = scan.read_mode(format, report.format_read.is_some(), &options);
1165        let mut options = OpenOptions {
1166            left_out: report.left_out,
1167            files_disagree: report.files_disagree,
1168            format,
1169            format_read: report.format_read,
1170            sqlite: report.sqlite,
1171            opened: report.opened,
1172            splits: report.splits,
1173            read_python: report.read_python,
1174            read_mode,
1175            tail,
1176            table: report.table.or_else(|| options.table.clone()),
1177            format_guessed: options.format_guessed || report.guessed,
1178            read_notes: report.read_notes,
1179            typing: report.typing,
1180            ..options
1181        };
1182        // The spec's dialect stays with the dataset, so a re-read (`H`, a decompressed
1183        // copy) reads the same.
1184        if let Some(read) = report.delimited {
1185            read.delimited().apply(&mut options);
1186            options.delimited = Some(read);
1187        }
1188        Ok(match scan {
1189            Scan::Frame(lf) => LoadAnswer::Scanned { lf, path, options },
1190            Scan::Decompress { file, .. } => LoadAnswer::Compressed {
1191                file,
1192                path,
1193                options,
1194            },
1195            Scan::Streams(files) => LoadAnswer::Convert {
1196                what: loading::Conversion::Streams,
1197                bytes: bytes_of(&files),
1198                files,
1199                path,
1200                options,
1201            },
1202            Scan::DecompressSpec { file, choice } => LoadAnswer::CompressedRecords {
1203                file,
1204                path,
1205                choice,
1206                options,
1207            },
1208            Scan::ReadInto { files, format } => LoadAnswer::Convert {
1209                what: loading::Conversion::Text(format),
1210                bytes: bytes_of(&files),
1211                files,
1212                path,
1213                options,
1214            },
1215            Scan::Tables { file, tables, .. } => LoadAnswer::Tables { file, tables, path },
1216            Scan::Unpack {
1217                file,
1218                member,
1219                format,
1220            } => LoadAnswer::Convert {
1221                what: loading::Conversion::Text(format),
1222                bytes: bytes_of(std::slice::from_ref(&file)),
1223                files: vec![file],
1224                path,
1225                options: OpenOptions {
1226                    table: Some(member),
1227                    ..options
1228                },
1229            },
1230            Scan::Hex { file, asked } => LoadAnswer::Hex {
1231                file,
1232                asked,
1233                record_size: options.record_size,
1234            },
1235        })
1236    }
1237
1238    /// The open's schema read of the scan's frame, building the dataset with everything
1239    /// the open `made`. Run by the `ReadSchema` phase and by the home preview.
1240    pub(crate) fn read_schema_for_open(
1241        lf: LazyFrame,
1242        path: Option<PathBuf>,
1243        options: OpenOptions,
1244        cloud: &crate::config::CloudConfig,
1245        runtime: &tokio::runtime::Handle,
1246        report: &crate::loading::measurements::OpenReport,
1247        made: loading::Made,
1248    ) -> std::result::Result<loading::LoadAnswer, String> {
1249        use loading::LoadAnswer;
1250        let (state, facts, debug_label) =
1251            Self::build_schema_state(lf, path.as_deref(), &options, cloud, runtime, report)
1252                .map_err(|e| crate::error_display::user_message_from_report(&e, path.as_deref()))?;
1253        // Everything the open found, given to the dataset as it is built.
1254        let loading::Made {
1255            download,
1256            converted,
1257            notes,
1258            other_tables,
1259            detail,
1260        } = made;
1261        let mut open_notes = facts.open_notes;
1262        open_notes.extend(notes);
1263        let mut other_tables_found = facts.other_tables;
1264        other_tables_found.extend(other_tables);
1265        let state = state.with_open(OpenFacts {
1266            fetched: Self::was_fetched(download.as_ref(), path.as_deref()),
1267            download,
1268            converted,
1269            other_tables: other_tables_found,
1270            open_notes,
1271            detail: detail.or(facts.detail),
1272            ..facts
1273        });
1274        Ok(LoadAnswer::SchemaRead {
1275            state: Box::new(state),
1276            path,
1277            options,
1278            debug_label: Some(debug_label),
1279        })
1280    }
1281
1282    /// Run the worker of an open's phase as `load`'s job; its answer goes to the
1283    /// loader. Every phase runs off the event thread: HEAD probes can stall for
1284    /// seconds, scans infer schemas, and directory schema reads touch every footer.
1285    fn spawn_load_phase(&mut self, load: loading::LoadId, step: loading::Step) {
1286        use loading::{LoadAnswer, Step};
1287        let job = Job::Load(load);
1288        // Every phase that reads, reads with the same cloud settings on the shared runtime.
1289        let (cloud, runtime) = (self.app_config.cloud.clone(), self.runtime.clone());
1290        match step {
1291            #[cfg(any(feature = "http", feature = "cloud"))]
1292            Step::ReadHeaders {
1293                url,
1294                format,
1295                options,
1296                writer,
1297            } => {
1298                self.spawn_job(job, Some("Reading headers..."), move |_| {
1299                    let read =
1300                        crate::cloud::remote_model::read(&url, format, &cloud, &runtime, &|| {
1301                            writer.stopped()
1302                        });
1303                    let crate::cloud::remote_model::Read { lf, summary, notes } = match read {
1304                        Ok(read) => read,
1305                        Err(crate::formats::model_files::RangeError::NoRanges) => {
1306                            return Ok(Answer::Load(Box::new(LoadAnswer::NoRanges { options })));
1307                        }
1308                        // The URL in the message may carry a password or a signature.
1309                        Err(crate::formats::model_files::RangeError::Failed(message)) => {
1310                            return Err(crate::logging::redact(&message, &[]));
1311                        }
1312                    };
1313                    let opened = Arc::new(crate::formats::model_files::opened(&summary));
1314                    let options = OpenOptions {
1315                        format: Some(format),
1316                        opened: Some(opened.clone()),
1317                        ..options
1318                    };
1319                    // The table is the headers, in memory: nothing is left to scan.
1320                    let state = Self::schema_state_from_full_scan(
1321                        lf,
1322                        None,
1323                        &OpenOptions {
1324                            hive: false,
1325                            ..options.clone()
1326                        },
1327                    )
1328                    .map_err(|e| crate::error_display::user_message_from_report(&e, Some(&url)))?
1329                    .with_open(OpenFacts {
1330                        detail: opened.detail.clone(),
1331                        open_notes: notes,
1332                        read_as: Some(format),
1333                        ..Default::default()
1334                    });
1335                    Ok(Answer::Load(Box::new(LoadAnswer::SchemaRead {
1336                        state: Box::new(state),
1337                        path: Some(url),
1338                        options,
1339                        debug_label: Some("model headers (ranged)".to_string()),
1340                    })))
1341                });
1342            }
1343            #[cfg(any(feature = "http", feature = "cloud"))]
1344            Step::Probe(pending) => {
1345                self.spawn_job(job, Some("Checking size..."), move |_| {
1346                    // Arrow in a store: its listing says which objects are streams to download.
1347                    #[cfg(feature = "cloud")]
1348                    if let loading::PendingDownload::Arrow { url, .. } = &pending {
1349                        let (_, _, options) = pending.parts();
1350                        let (objects, options) =
1351                            crate::cloud::cloud_arrow::list(url, options, &cloud, &runtime)
1352                                .map_err(|e| {
1353                                    crate::error_display::user_message_from_report(&e, None)
1354                                })?;
1355                        let size = crate::cloud::cloud_arrow::stream_bytes(&objects);
1356                        return Ok(Answer::Load(Box::new(LoadAnswer::Sized(
1357                            loading::PendingDownload::Arrow {
1358                                url: url.clone(),
1359                                objects,
1360                                size: Some(size),
1361                                options,
1362                            },
1363                        ))));
1364                    }
1365                    let size = match &pending {
1366                        #[cfg(feature = "http")]
1367                        loading::PendingDownload::Http { url, .. } => {
1368                            // A missing file or silent host ends the open here, before asking about a
1369                            // download.
1370                            Self::fetch_remote_size_http(url).map_err(|gone| gone.message)?
1371                        }
1372                        #[cfg(feature = "cloud")]
1373                        loading::PendingDownload::S3 { url, .. }
1374                        | loading::PendingDownload::Gcs { url, .. }
1375                        | loading::PendingDownload::Azure { url, .. } => {
1376                            Self::fetch_remote_size_cloud(url, &cloud, &runtime).unwrap_or(None)
1377                        }
1378                        #[cfg(feature = "cloud")]
1379                        loading::PendingDownload::Arrow { size, .. } => *size,
1380                    };
1381                    Ok(Answer::Load(Box::new(LoadAnswer::Sized(
1382                        pending.with_size(size),
1383                    ))))
1384                });
1385            }
1386            #[cfg(any(feature = "http", feature = "cloud"))]
1387            Step::Download { pending, writer } => {
1388                // The load's stop flag (abandon, replacement, app drop) stops the download at the
1389                // next chunk or during silence and removes its file; `ExitSweep` covers quitting.
1390                let status = match &pending {
1391                    #[cfg(feature = "http")]
1392                    loading::PendingDownload::Http { .. } => "Downloading...",
1393                    #[cfg(feature = "cloud")]
1394                    loading::PendingDownload::S3 { .. } => "Downloading from S3...",
1395                    #[cfg(feature = "cloud")]
1396                    loading::PendingDownload::Gcs { .. } => "Downloading from GCS...",
1397                    #[cfg(feature = "cloud")]
1398                    loading::PendingDownload::Azure { .. } => "Downloading from Azure...",
1399                    #[cfg(feature = "cloud")]
1400                    loading::PendingDownload::Arrow { url, .. } => {
1401                        match source::input_source(Path::new(url)) {
1402                            source::InputSource::Gcs(_) => "Downloading from GCS...",
1403                            source::InputSource::Azure(_) => "Downloading from Azure...",
1404                            _ => "Downloading from S3...",
1405                        }
1406                    }
1407                };
1408                // An unasked download says how much it fetches, when the server said.
1409                let sized = pending
1410                    .parts()
1411                    .1
1412                    .filter(|_| status == "Downloading...")
1413                    .map(|size| format!("Downloading {}...", crate::numfmt::bytes(size)));
1414                let status = sized.as_deref().unwrap_or(status);
1415                self.spawn_job(job, Some(status), move |_| {
1416                    let (url, _, options) = pending.parts();
1417                    let fetched = match &pending {
1418                        #[cfg(feature = "http")]
1419                        loading::PendingDownload::Http { .. } => {
1420                            let ext = source::download_suffix(url);
1421                            // An unasked download of unstated size stops at its limit; a stated size bounds the
1422                            // transfer itself, and the bytes counted here are decompressed.
1423                            let limit = options
1424                                .download_unasked
1425                                .filter(|_| pending.parts().1.is_none())
1426                                .map(|unasked| unasked.limit);
1427                            Self::download_http_to_temp(
1428                                url,
1429                                options.temp_dir.as_deref(),
1430                                ext.as_deref(),
1431                                limit,
1432                                &writer,
1433                            )
1434                            .map(|file| (file, options.clone()))
1435                        }
1436                        #[cfg(feature = "cloud")]
1437                        loading::PendingDownload::S3 { .. }
1438                        | loading::PendingDownload::Gcs { .. }
1439                        | loading::PendingDownload::Azure { .. } => {
1440                            Self::download_cloud_to_temp(url, &cloud, options, &runtime, &writer)
1441                                .map(|file| (file, options.clone()))
1442                        }
1443                        // Its streams, converted as they arrive; its IPC files stay put.
1444                        #[cfg(feature = "cloud")]
1445                        loading::PendingDownload::Arrow { objects, .. } => {
1446                            crate::cloud::cloud_arrow::download(
1447                                objects, options, &cloud, &runtime, &writer,
1448                            )
1449                            .map(|(file, parts)| {
1450                                let options = OpenOptions {
1451                                    format: Some(FileFormat::Arrow),
1452                                    hive: false,
1453                                    arrow_parts: Some(Arc::new(parts)),
1454                                    ..options.clone()
1455                                };
1456                                (file, options)
1457                            })
1458                        }
1459                    };
1460                    let (download, options) = match fetched {
1461                        Err(e)
1462                            if e.downcast_ref::<crate::cloud::download::PastLimit>()
1463                                .is_some() =>
1464                        {
1465                            return Ok(Answer::Load(Box::new(LoadAnswer::PastLimit(pending))));
1466                        }
1467                        fetched => fetched.map_err(|e| {
1468                            crate::error_display::user_message_from_report(&e, None)
1469                        })?,
1470                    };
1471                    Ok(Answer::Load(Box::new(LoadAnswer::Downloaded {
1472                        download,
1473                        options,
1474                    })))
1475                });
1476            }
1477            Step::Spool {
1478                options,
1479                writer,
1480                read,
1481            } => {
1482                // The read is its own thread, so a quiet producer does not hold up the stop:
1483                // Ctrl+O and quitting remove the partial file at once.
1484                let piped = self.pipes.stdin_reader.take();
1485                let stdout = self.pipes.stdout_pass.take();
1486                self.spawn_job(job, Some("Reading stdin..."), move |_| {
1487                    let open =
1488                        move || -> crate::cloud::download::Opened<Box<dyn std::io::Read + Send>> {
1489                            Ok((piped.unwrap_or_else(|| Box::new(std::io::stdin())), None))
1490                        };
1491                    // Followed, the copy goes on behind the first rows; recorded, it goes to the
1492                    // named file; and it is read as it arrives when its format allows.
1493                    let (download, options) = if options.follow
1494                        || options.tee.is_some()
1495                        || crate::loading::stdin::may_read_as_it_arrives(&options)
1496                    {
1497                        match crate::loading::follow::spool(open, options, &writer, &read, stdout)?
1498                        {
1499                            (crate::loading::follow::Spooled::Temp(download), options) => {
1500                                (download, options)
1501                            }
1502                            (crate::loading::follow::Spooled::Kept(file), options) => {
1503                                return Ok(Answer::Load(Box::new(LoadAnswer::Recorded {
1504                                    file,
1505                                    options,
1506                                })));
1507                            }
1508                        }
1509                    } else {
1510                        crate::loading::stdin::spool(open, options, &writer, &read)?
1511                    };
1512                    Ok(Answer::Load(Box::new(LoadAnswer::Spooled {
1513                        download,
1514                        options,
1515                    })))
1516                });
1517            }
1518            Step::FetchSpec {
1519                url,
1520                options,
1521                writer,
1522            } => {
1523                self.spawn_job(job, Some("Reading spec..."), move |_| {
1524                    #[cfg(any(feature = "http", feature = "cloud"))]
1525                    let fetched = crate::cloud::remote_model::fetch_small(
1526                        &url,
1527                        crate::formats::MAX_SPEC_BYTES,
1528                        &cloud,
1529                        &runtime,
1530                        &|| writer.stopped(),
1531                    );
1532                    #[cfg(not(any(feature = "http", feature = "cloud")))]
1533                    let fetched: std::result::Result<Option<Vec<u8>>, String> = {
1534                        let _ = &writer;
1535                        Err(crate::error_display::file_message(
1536                            &url,
1537                            "this build reads no URLs",
1538                        ))
1539                    };
1540                    // The URL in the message may carry a password or a signature.
1541                    let bytes = fetched
1542                        .map_err(|message| crate::logging::redact(&message, &[]))?
1543                        .ok_or_else(|| {
1544                            crate::logging::redact(
1545                                &crate::error_display::file_message(
1546                                    &url,
1547                                    &format!(
1548                                        "a format spec is at most {}",
1549                                        crate::formats::MAX_SPEC_SAID
1550                                    ),
1551                                ),
1552                                &[],
1553                            )
1554                        })?;
1555                    let spec = crate::formats::Spec::from_bytes(&bytes, &url)
1556                        .map_err(|e| crate::logging::redact(&e.to_string(), &[]))?;
1557                    Ok(Answer::Load(Box::new(LoadAnswer::SpecFetched {
1558                        spec: Arc::new(spec),
1559                        options,
1560                    })))
1561                });
1562            }
1563            Step::DecompressRecords {
1564                file,
1565                path,
1566                choice,
1567                options,
1568                writer,
1569            } => {
1570                self.spawn_job(job, Some("Decompressing..."), move |_| {
1571                    let failed = |e: color_eyre::Report| {
1572                        crate::error_display::user_message_from_report(&e, Some(path.as_path()))
1573                    };
1574                    let compression = options
1575                        .compression
1576                        .or_else(|| CompressionFormat::from_extension(&file))
1577                        .ok_or_else(|| format!("{} is not compressed", path.display()))?;
1578                    let temp_dir = options.temp_dir.clone().unwrap_or_else(std::env::temp_dir);
1579                    let copy = crate::formats::readers::csv::decompress_to_copy(
1580                        &file,
1581                        compression,
1582                        &temp_dir,
1583                        &writer,
1584                    )
1585                    .map_err(failed)?;
1586                    Ok(Answer::Load(Box::new(LoadAnswer::DecompressedRecords {
1587                        copy,
1588                        path,
1589                        choice,
1590                        options,
1591                    })))
1592                });
1593            }
1594            Step::ReadRecords {
1595                copy,
1596                path,
1597                choice,
1598                options,
1599            } => {
1600                self.spawn_job(job, Some("Reading records..."), move |_| {
1601                    let named = path
1602                        .file_name()
1603                        .map(|n| n.to_string_lossy().into_owned())
1604                        .unwrap_or_default();
1605                    let read = crate::formats::read(&copy, &named, choice)?;
1606                    let lf = Arc::clone(&read.records).into_lazy().map_err(|e| {
1607                        crate::error_display::user_message_from_report(
1608                            &color_eyre::eyre::eyre!(e),
1609                            Some(path.as_path()),
1610                        )
1611                    })?;
1612                    Ok(Answer::Load(Box::new(LoadAnswer::Scanned {
1613                        lf: Box::new(lf),
1614                        path: Some(path),
1615                        options: OpenOptions {
1616                            format_read: Some(Arc::new(read)),
1617                            ..options
1618                        },
1619                    })))
1620                });
1621            }
1622            Step::Decompress {
1623                file,
1624                path,
1625                options,
1626                writer,
1627                download,
1628            } => {
1629                // Only delimited text and lines come here, the format said by the loader or scan.
1630                let options = OpenOptions {
1631                    format: options.format.or(Some(FileFormat::TEXT)),
1632                    ..options
1633                };
1634                let formats = self.formats.clone();
1635                self.spawn_job(job, Some("Decompressing..."), move |_| {
1636                    let failed = |e: color_eyre::Report| {
1637                        crate::error_display::user_message_from_report(&e, Some(path.as_path()))
1638                    };
1639                    let options =
1640                        Self::with_delimited_spec(&file, options, &formats).map_err(failed)?;
1641                    let lines = options.delimited.is_none()
1642                        && options.format.is_some_and(FileFormat::is_lines);
1643                    let (state, opened) = if lines {
1644                        let (read, opened) = crate::formats::readers::csv::from_lines_decompressed(
1645                            &file, &options, &writer,
1646                        )
1647                        .map_err(failed)?;
1648                        let state = DataTableState::from_read(read, &options).map_err(failed)?;
1649                        (state, Some(opened))
1650                    } else {
1651                        let state = Self::decompressed_delimited_state(&file, &options, &writer)
1652                            .map_err(failed)?;
1653                        (state, None)
1654                    };
1655                    let mut open_notes = options
1656                        .delimited
1657                        .as_ref()
1658                        .map(|read| read.notes())
1659                        .unwrap_or_default();
1660                    open_notes.extend(opened.iter().flat_map(|o| o.notes.iter().cloned()));
1661                    let state = state.with_open(OpenFacts {
1662                        fetched: Self::was_fetched(download.as_ref(), Some(&path)),
1663                        download,
1664                        open_notes,
1665                        records: opened.and_then(|o| o.window),
1666                        delimited: options.delimited.clone(),
1667                        read_as: options.format,
1668                        // The loader sends a compressed file here without a scan.
1669                        read_mode: options.format.and_then(|f| {
1670                            f.read_mode(crate::Stored::Compressed {
1671                                in_memory: options.decompress_in_memory,
1672                            })
1673                        }),
1674                        ..Default::default()
1675                    });
1676                    Ok(Answer::Load(Box::new(LoadAnswer::SchemaRead {
1677                        state: Box::new(state),
1678                        path: Some(path),
1679                        options,
1680                        debug_label: Some("decompressed delimited".to_string()),
1681                    })))
1682                });
1683            }
1684            Step::Convert {
1685                what,
1686                files,
1687                path,
1688                options,
1689                writer,
1690                read,
1691            } => {
1692                // The load's stop flag ends it at the next batch or chunk, removing its files;
1693                // quitting removes them even if the process ends first.
1694                let formats = self.formats.clone();
1695                self.spawn_job(job, Some(what.status()), move |_| {
1696                    let named = |e: color_eyre::Report| {
1697                        crate::error_display::user_message_from_report(&e, path.as_deref())
1698                    };
1699                    let converted = match what {
1700                        loading::Conversion::Streams => {
1701                            let converted = crate::formats::ipc_stream::convert(
1702                                &files,
1703                                options.temp_dir.as_deref(),
1704                                &writer,
1705                                &read,
1706                            )
1707                            .map_err(named)?;
1708                            loading::Converted::Streams {
1709                                file: converted.file,
1710                                parts: converted.parts,
1711                            }
1712                        }
1713                        loading::Conversion::Text(format) => {
1714                            let display = path.clone().unwrap_or_else(|| files[0].clone());
1715                            let (converted, detail) = crate::formats::readers::convert(
1716                                &crate::formats::readers::ConvertIn {
1717                                    files: &files,
1718                                    display: &display,
1719                                    format,
1720                                    options: &options,
1721                                    formats: &formats,
1722                                    writer: &writer,
1723                                    read: &read,
1724                                },
1725                            )
1726                            .map_err(named)?;
1727                            loading::Converted::Frame {
1728                                files: converted.files,
1729                                lf: Box::new(converted.lf),
1730                                notes: converted.notes,
1731                                other_tables: converted.other_tables,
1732                                detail,
1733                            }
1734                        }
1735                    };
1736                    Ok(Answer::Load(Box::new(LoadAnswer::Converted {
1737                        converted,
1738                        path,
1739                        options,
1740                    })))
1741                });
1742            }
1743            Step::Scan {
1744                paths,
1745                options,
1746                display,
1747                status,
1748            } => {
1749                let formats = self.formats.clone();
1750                // A download is scanned from a temp path; the URL the user typed names the dataset.
1751                let path = display.or_else(|| paths.first().cloned());
1752                self.home_app.reads.scans += 1;
1753                self.spawn_job(job, Some(status), move |_| {
1754                    Self::scan_for_open(&cloud, &formats, &paths, options, path)
1755                        .map(|answer| Answer::Load(Box::new(answer)))
1756                });
1757            }
1758            Step::ReadSchema {
1759                lf,
1760                path,
1761                options,
1762                progress,
1763                made,
1764            } => {
1765                self.debug.schema_load = None;
1766                let report = crate::loading::measurements::OpenReport {
1767                    progress,
1768                    meter: Arc::new(crate::loading::measurements::Meter::default()),
1769                    remembered: Some(self.cache.clone()),
1770                    writes: self.cache_writes.clone(),
1771                };
1772                self.spawn_job(job, Some("Reading schema..."), move |_| {
1773                    Self::read_schema_for_open(*lf, path, options, &cloud, &runtime, &report, made)
1774                        .map(|answer| Answer::Load(Box::new(answer)))
1775                });
1776            }
1777            Step::Nothing
1778            | Step::Crash(_)
1779            | Step::Install(_)
1780            | Step::Failed(_)
1781            | Step::Tables(_)
1782            | Step::Hex(_) => {
1783                unreachable!("not a phase with a worker")
1784            }
1785            #[cfg(any(feature = "http", feature = "cloud"))]
1786            Step::Ask(_) => unreachable!("not a phase with a worker"),
1787            Step::AskRead(_) => unreachable!("not a phase with a worker"),
1788        }
1789    }
1790
1791    /// The question before files past `[read] memory_warning` are read whole:
1792    /// `big.json: JSON reads 2.0 GiB into memory`.
1793    fn in_memory_confirmation_message(read: &loading::InMemory) -> String {
1794        let what = match read.files {
1795            1 => format!(
1796                "{}: {} reads",
1797                read.name
1798                    .file_name()
1799                    .map(|n| n.to_string_lossy().into_owned())
1800                    .unwrap_or_else(|| read.name.display().to_string()),
1801                read.format.title()
1802            ),
1803            n => format!("{n} {} files read", read.format.title()),
1804        };
1805        format!(
1806            "{what} {} into memory before the table appears.\n\nRead it?",
1807            crate::numfmt::bytes(read.bytes)
1808        )
1809    }
1810
1811    /// The question asked before a remote file is downloaded.
1812    #[cfg(any(feature = "http", feature = "cloud"))]
1813    fn download_confirmation_message(
1814        pending: &loading::PendingDownload,
1815        note: Option<&str>,
1816    ) -> String {
1817        let (url, size, options) = pending.parts();
1818        let size_str = size
1819            .map(crate::numfmt::bytes)
1820            .unwrap_or_else(|| "unknown".to_string());
1821        let dest_dir = options
1822            .temp_dir
1823            .as_deref()
1824            .map(|p| p.display().to_string())
1825            .unwrap_or_else(|| std::env::temp_dir().display().to_string());
1826        let note = note.map(|note| format!("{note}\n\n")).unwrap_or_default();
1827        // Only a store's Arrow streams are downloaded, and converted as they arrive.
1828        let files = match pending.arrow_files() {
1829            Some((1, 0)) => "Arrow stream: converted as it downloads\n".to_string(),
1830            Some((streams, 0)) => {
1831                format!("Files: {streams} Arrow streams, converted as they download\n")
1832            }
1833            Some((streams, in_place)) => {
1834                let streams = match streams {
1835                    1 => "1 Arrow stream, converted as it downloads".to_string(),
1836                    n => format!("{n} Arrow streams, converted as they download"),
1837                };
1838                let in_place = match in_place {
1839                    1 => "1 IPC file read in place".to_string(),
1840                    n => format!("{n} IPC files read in place"),
1841                };
1842                format!("Files: {streams}; {in_place}\n")
1843            }
1844            None => String::new(),
1845        };
1846        format!(
1847            "{note}URL: {url}\n{files}File size: {size_str}\nDestination: {dest_dir} (temporary file)\n\nContinue with download?"
1848        )
1849    }
1850}
1851
1852/// Scans and schema reads for each kind of source, used by the load's phases.
1853impl App {
1854    fn hoist_partition_columns(
1855        lf: LazyFrame,
1856        schema: &Schema,
1857        partition_columns: &[String],
1858        drifts: bool,
1859    ) -> LazyFrame {
1860        hoist_partition_columns(lf, schema, partition_columns, drifts)
1861    }
1862
1863    /// Schema for a local directory of Parquet: the union of every footer's columns,
1864    /// not `collect_schema()` or one file's. Like a cloud prefix, past one wave of
1865    /// footers the dataset opens and the rest join behind it; an unchanged listing
1866    /// reopens from its cached footers. `None` when not that shape or nothing could
1867    /// be read: the caller falls back to the general scan, which reports any error.
1868    pub(crate) fn schema_state_from_local_hive(
1869        path: Option<&Path>,
1870        options: &OpenOptions,
1871        report: &crate::loading::measurements::OpenReport,
1872    ) -> Option<(DataTableState, OpenFacts)> {
1873        if !options.single_spine_schema {
1874            return None;
1875        }
1876        let p = path.filter(|p| p.is_dir() && options.hive)?;
1877        dataset_files::open(Arc::new(dataset_files::LocalFiles::new(p)), options, report)
1878    }
1879
1880    /// The same for an object store, off the UI thread.
1881    #[cfg(feature = "cloud")]
1882    fn schema_state_from_cloud_hive(
1883        path: Option<&Path>,
1884        options: &OpenOptions,
1885        cloud: &crate::config::CloudConfig,
1886        runtime: &tokio::runtime::Handle,
1887        report: &crate::loading::measurements::OpenReport,
1888    ) -> Option<(DataTableState, OpenFacts)> {
1889        // A prefix of Arrow files is read from its download (`cloud_arrow`).
1890        if !options.single_spine_schema || options.format == Some(FileFormat::Arrow) {
1891            return None;
1892        }
1893        // Unlike the local path this needs no --hive: a directory or glob URL is a hive
1894        // scan by shape.
1895        let p = path.filter(|p| {
1896            let s = p.as_os_str().to_string_lossy();
1897            home::is_object_store_url(p) && (options.hive || source::is_prefix_or_glob(&s))
1898        })?;
1899
1900        let (full, cloud_opts, store) = Self::cloud_store_for(p, cloud, runtime).ok()?;
1901        let (_bucket, key) = Self::cloud_bucket_and_key(&full).ok()?;
1902        Self::schema_state_from_cloud_hive_with(
1903            full, key, store, cloud_opts, options, runtime, report,
1904        )
1905    }
1906
1907    /// The same against a built store; a test hands it an in-memory one to cover the
1908    /// route choice and that each route gets the caller's counter.
1909    #[cfg(feature = "cloud")]
1910    pub(crate) fn schema_state_from_cloud_hive_with(
1911        full: String,
1912        key: String,
1913        store: Arc<dyn object_store::ObjectStore>,
1914        cloud_opts: CloudOptions,
1915        options: &OpenOptions,
1916        runtime: &tokio::runtime::Handle,
1917        report: &crate::loading::measurements::OpenReport,
1918    ) -> Option<(DataTableState, OpenFacts)> {
1919        // List every file once; scan, schema and count work from that list. A listing
1920        // prefix is literal, so datui expands a glob itself: it lists the literal part
1921        // and matches the rest, and the glob gets the schema union, count, notes and
1922        // measurements a prefix gets.
1923        let pattern = full.contains('*').then(|| {
1924            globset::GlobBuilder::new(&key)
1925                .literal_separator(true)
1926                .build()
1927                .map(|g| g.compile_matcher())
1928        });
1929        let pattern = match pattern {
1930            // A pattern datui cannot read is not one it should guess at.
1931            Some(Err(_)) => return None,
1932            Some(Ok(matcher)) => Some(matcher),
1933            None => None,
1934        };
1935        let listed = cloud_hive::prefix_of_glob(&key).to_string();
1936        Self::schema_state_from_cloud_files(
1937            CloudTarget {
1938                full: &full,
1939                key: listed,
1940                pattern: pattern.as_ref(),
1941            },
1942            store,
1943            cloud_opts,
1944            options,
1945            runtime,
1946            report,
1947        )
1948    }
1949
1950    /// A cloud prefix of Parquet files as one dataset from a single listing. Files are
1951    /// scanned by name (no second listing) and leniently (`cloud_hive::lenient_scan`),
1952    /// since files written years apart differ. The state keeps the list so counts read
1953    /// footers and buffers read only the files holding their rows (`RemoteFiles`).
1954    #[cfg(feature = "cloud")]
1955    fn schema_state_from_cloud_files(
1956        target: CloudTarget<'_>,
1957        store: Arc<dyn object_store::ObjectStore>,
1958        cloud_opts: CloudOptions,
1959        options: &OpenOptions,
1960        runtime: &tokio::runtime::Handle,
1961        report: &crate::loading::measurements::OpenReport,
1962    ) -> Option<(DataTableState, OpenFacts)> {
1963        let CloudTarget { full, key, pattern } = target;
1964        dataset_files::open(
1965            Arc::new(dataset_files::StoreFiles::new(
1966                full,
1967                key,
1968                pattern.cloned(),
1969                store,
1970                cloud_opts,
1971                runtime,
1972            )),
1973            options,
1974            report,
1975        )
1976    }
1977
1978    /// Record one object opened from a bucket for home: rows and columns from its
1979    /// footer under the resolved URL. No size (the footer is read from the tail).
1980    /// `mtime` is the open time: remote records are not fingerprinted by it, and the
1981    /// index evicts oldest `mtime` first, so zero would evict these first.
1982    #[cfg(feature = "cloud")]
1983    pub(crate) fn record_cloud_object_facts(
1984        cache: Option<&crate::cache::CacheManager>,
1985        full: &str,
1986        footer: &cloud_hive::FileFooter,
1987    ) {
1988        let Some(cache) = cache else {
1989            return;
1990        };
1991        let columns: Vec<String> = footer
1992            .schema
1993            .iter_names()
1994            .map(|name| name.to_string())
1995            .collect();
1996        cache.record_dataset_facts(&[(
1997            PathBuf::from(full),
1998            crate::cache::DatasetFacts {
1999                mtime: std::time::SystemTime::now()
2000                    .duration_since(std::time::UNIX_EPOCH)
2001                    .map(|d| d.as_secs())
2002                    .unwrap_or_default(),
2003                size: 0,
2004                rows: Some(footer.rows()),
2005                cols: Some(columns.len()),
2006                cols_sampled: false,
2007                columns,
2008                kind: Some(discover::EntryKind::File),
2009                classified_by: discover::CLASSIFIER_VERSION,
2010                cost: discover::Cost {
2011                    row_groups: Some(footer.row_group_rows.len()),
2012                    ..Default::default()
2013                },
2014                holds: Default::default(),
2015            },
2016        )]);
2017    }
2018
2019    /// General schema route: ask the frame. Slow for a wide hive dataset, hence off the
2020    /// UI thread.
2021    fn schema_state_from_full_scan(
2022        mut lf: LazyFrame,
2023        path: Option<&Path>,
2024        options: &OpenOptions,
2025    ) -> Result<DataTableState> {
2026        let schema = lf
2027            .collect_schema()
2028            .map_err(color_eyre::eyre::Report::from)?;
2029        let partition_columns =
2030            match path.filter(|p| options.hive && (p.is_dir() || source::expands_as_glob(p))) {
2031                Some(p) => crate::formats::readers::hive::discover_hive_partition_columns(p)
2032                    .into_iter()
2033                    .filter(|c| schema.contains(c.as_str()))
2034                    .collect::<Vec<_>>(),
2035                None => Vec::new(),
2036            };
2037        let lf = Self::hoist_partition_columns(lf, &schema, &partition_columns, false);
2038        let part_cols = (!partition_columns.is_empty()).then_some(partition_columns);
2039        DataTableState::from_schema_and_lazyframe(schema, lf, options, part_cols)
2040    }
2041
2042    /// Build the table state by the cheapest route that applies, with a label naming
2043    /// the route for the debug overlay. Takes config by value so it runs off-thread.
2044    fn build_schema_state(
2045        lf: LazyFrame,
2046        path: Option<&Path>,
2047        options: &OpenOptions,
2048        cloud: &crate::config::CloudConfig,
2049        runtime: &tokio::runtime::Handle,
2050        report: &crate::loading::measurements::OpenReport,
2051    ) -> Result<(DataTableState, OpenFacts, String)> {
2052        // The facts carry the meter of the route that built the dataset, installed with it.
2053        // A failed open never gets here, so the dataset on screen keeps its own figures.
2054        let (state, mut facts, label) =
2055            Self::schema_state_by_route(lf, path, options, cloud, runtime, report)?;
2056        // What the open did, beside what it found: the one place the scan's report, the
2057        // lake-table flag and the state are all known (see
2058        // `DataTableState::open_notes`). A delimited header whose names are all numbers
2059        // is likely a data row, which `H` reads as data.
2060        let names_look_like_data = options.format.and_then(FileFormat::separator).is_some()
2061            && options.has_header != Some(false)
2062            && !crate::formats::schema_union::names_are_names(
2063                &state
2064                    .schema()
2065                    .iter_names()
2066                    .map(|n| n.to_string())
2067                    .collect::<Vec<_>>(),
2068            );
2069        facts.open_notes = crate::notes::from_the_open(
2070            &options.left_out,
2071            options.read_as_plain_files_of,
2072            options.files_disagree,
2073            names_look_like_data,
2074        );
2075        // The row count on screen is a true count of the files and a wrong one of the
2076        // table.
2077        facts.not_the_table = options.read_as_plain_files_of;
2078        if let Some(splits) = &options.splits {
2079            facts.other_tables = splits.others.clone();
2080            facts
2081                .open_notes
2082                .extend(crate::notes::map_caches(splits.caches));
2083        }
2084        if let Some(read) = &options.format_read {
2085            facts.open_notes.extend(read.notes());
2086            facts.format_read = Some(read.clone());
2087        }
2088        if let Some(opened) = &options.opened {
2089            facts.records = opened.window.clone();
2090            facts.detail = opened.detail.clone();
2091            facts.other_tables = opened.other_tables.clone();
2092            facts.open_notes.extend(opened.notes.iter().cloned());
2093            facts.units = opened.units.clone();
2094            facts.indexing = opened.indexing.clone();
2095            facts.numbering = opened.numbering.clone();
2096        }
2097        if let Some(sqlite) = &options.sqlite {
2098            facts.pushdown = Some(sqlite.pushdown.clone());
2099            facts.hold = sqlite.hold.lock().ok().and_then(|mut hold| hold.take());
2100            facts.other_tables = sqlite.other_tables.clone();
2101        }
2102        if let Some(read) = &options.delimited {
2103            facts.open_notes.extend(read.notes());
2104            facts.delimited = Some(read.clone());
2105        }
2106        facts.open_notes.extend(options.read_notes.iter().cloned());
2107        facts.typing = options.typing.clone();
2108        facts.read_mode = options.read_mode;
2109        facts.read_as = options.format;
2110        // A downloaded object's display path is its URL too; only a scan reading the store
2111        // in place buffers like one. Arrow in a store reads IPC files in place and streams
2112        // from their download (`cloud_arrow`).
2113        facts.remote_source = match &options.arrow_parts {
2114            Some(parts) => parts.iter().any(|part| {
2115                matches!(part, crate::formats::ipc_stream::Part::InPlace(p) if source::is_remote_url(p))
2116            }),
2117            None => path.is_some_and(source::scans_in_place),
2118        };
2119        // The footer-sum row count for a local Parquet hive directory, asked here because a
2120        // stat on a dead mount hangs its thread. A directory of another format counts by
2121        // a scan; one read by file has a counter for the footers the open did not read.
2122        if options.hive
2123            && facts.remote_files.is_none()
2124            && options.format.is_none_or(|f| f == FileFormat::Parquet)
2125            && let Some(dir) = path.filter(|p| !source::is_remote_url(p) && p.is_dir())
2126        {
2127            facts.parquet_count_dir = Some(dir.to_path_buf());
2128        }
2129        Ok((state, facts, label))
2130    }
2131
2132    /// Scan an object-store prefix with the reader its format calls for; Polars'
2133    /// other scans take the same `CloudOptions`. Parquet keeps its own branch at each
2134    /// call site (the only format with hive partitioning here). `None` for other
2135    /// formats, sending the caller to its Parquet scan.
2136    #[cfg(feature = "cloud")]
2137    pub(crate) fn scan_cloud_prefix(
2138        url: &str,
2139        cloud_opts: CloudOptions,
2140        format: FileFormat,
2141        glob: bool,
2142        options: &OpenOptions,
2143    ) -> Option<Result<LazyFrame>> {
2144        // The formats the docs say a prefix reads in place. Parquet takes the caller's
2145        // scan; model files are read by their headers before this.
2146        if !format.reads_bucket_prefix() {
2147            return None;
2148        }
2149        // Narrow a plain prefix to keys with an extension: a console's folder marker lists
2150        // as `data` for `data/`, and Polars refuses the whole prefix over it.
2151        let pl_path = if url.ends_with('/') && !url.contains('*') {
2152            PlRefPath::new(format!("{url}**/*.*").as_str())
2153        } else {
2154            PlRefPath::new(url)
2155        };
2156        let scan = crate::formats::readers::of(format).bucket_scan?;
2157        Some(scan(crate::formats::readers::BucketIn {
2158            url,
2159            path: pl_path,
2160            cloud: cloud_opts,
2161            glob,
2162            options,
2163            format,
2164        }))
2165    }
2166
2167    /// The non-Parquet format a store prefix or glob reads as: the listing's, else a
2168    /// glob's extension (`*.arrow`).
2169    #[cfg(feature = "cloud")]
2170    fn cloud_glob_format(url: &str, options: &OpenOptions) -> Option<FileFormat> {
2171        options
2172            .format
2173            .or_else(|| {
2174                url.contains('*')
2175                    .then(|| FileFormat::from_path(Path::new(url)))
2176                    .flatten()
2177            })
2178            .filter(|f| *f != FileFormat::Parquet)
2179    }
2180
2181    /// The plain URL and Polars options for one object-store path, through the source
2182    /// it names or belongs to (`cloud_sources::resolve`).
2183    #[cfg(feature = "cloud")]
2184    fn resolve_cloud_url(
2185        path: &Path,
2186        cloud: &crate::config::CloudConfig,
2187    ) -> Result<(String, CloudOptions)> {
2188        let text = path.to_string_lossy();
2189        let resolved = crate::cloud::cloud_sources::resolve_for_open(&text, cloud)
2190            .map_err(|e| color_eyre::eyre::eyre!(e))?;
2191        use object_store::azure::AzureConfigKey;
2192        use polars::io::cloud::GoogleConfigKey;
2193        let gcs_agent = (
2194            GoogleConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
2195            crate::cloud::user_agent::get(),
2196        );
2197        let options = match resolved.kind {
2198            crate::cloud::source::ProviderKind::S3 => Self::build_s3_cloud_options(&resolved.s3),
2199            crate::cloud::source::ProviderKind::Gcs
2200                if resolved.signing == crate::cloud::cloud_sources::Signing::Unsigned =>
2201            {
2202                CloudOptions::default()
2203                    .with_gcp([(GoogleConfigKey::SkipSignature, "true".into()), gcs_agent])
2204            }
2205            crate::cloud::source::ProviderKind::Gcs => match &resolved.gcloud {
2206                // The token comes from `gcloud` whenever Polars asks, so long scans outlive it.
2207                Some((configuration, _)) => CloudOptions::default()
2208                    .with_gcp([gcs_agent])
2209                    .with_credential_provider(Some(crate::cloud::gcloud::polars_provider(
2210                        configuration,
2211                    ))),
2212                None => match &resolved.google_credentials {
2213                    Some(file) => CloudOptions::default().with_gcp([
2214                        (
2215                            GoogleConfigKey::ApplicationCredentials,
2216                            file.to_string_lossy().into_owned(),
2217                        ),
2218                        gcs_agent,
2219                    ]),
2220                    None => CloudOptions::default().with_gcp([gcs_agent]),
2221                },
2222            },
2223            crate::cloud::source::ProviderKind::Azure => {
2224                let (account, _, _) = source::azure_parts(&resolved.url)
2225                    .ok_or_else(|| color_eyre::eyre::eyre!("not an Azure URL"))?;
2226                let mut azure = crate::cloud::azure::polars_options(&account, &resolved.azure);
2227                azure.push((
2228                    AzureConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
2229                    crate::cloud::user_agent::get(),
2230                ));
2231                CloudOptions::default().with_azure(azure)
2232            }
2233        };
2234        Ok((resolved.url, options))
2235    }
2236
2237    /// The URL, Polars options and store for one object-store path. `cloud` is the
2238    /// effective config (see `OpenOptions::effective_cloud`).
2239    #[cfg(feature = "cloud")]
2240    pub(crate) fn cloud_store_for(
2241        path: &Path,
2242        cloud: &crate::config::CloudConfig,
2243        runtime: &tokio::runtime::Handle,
2244    ) -> Result<(String, CloudOptions, Arc<dyn object_store::ObjectStore>)> {
2245        let (full, cloud_opts) = Self::resolve_cloud_url(path, cloud)?;
2246        let store = Self::polars_object_store(&full, &cloud_opts, runtime)?;
2247        Ok((full, cloud_opts, store))
2248    }
2249
2250    /// One Parquet object read in place: schema and row count from its footer in one
2251    /// 256 KiB tail read. Polars answers a cloud `len()` by reading row group 0, so the
2252    /// schema is handed to the scan and the count known before the first frame. A
2253    /// failure is returned: the caller falls back to asking the frame and labels why.
2254    #[cfg(feature = "cloud")]
2255    fn schema_state_from_cloud_object(
2256        path: &Path,
2257        options: &OpenOptions,
2258        cloud: &crate::config::CloudConfig,
2259        runtime: &tokio::runtime::Handle,
2260        report: &crate::loading::measurements::OpenReport,
2261    ) -> Result<(DataTableState, OpenFacts)> {
2262        let (full, cloud_opts, store) = Self::cloud_store_for(path, cloud, runtime)?;
2263        let (_bucket, key) = Self::cloud_bucket_and_key(&full)?;
2264        if key.is_empty() {
2265            return Err(color_eyre::eyre::eyre!("a bucket, not an object"));
2266        }
2267        let meter = report.meter.clone();
2268        let (footer, etag) = wait_on_runtime(runtime, async move {
2269            cloud_hive::footer_of_cloud_parquet(store, &key, &meter).await
2270        })
2271        .ok_or_else(|| color_eyre::eyre::eyre!("cancelled"))??;
2272        let args = ScanArgsParquet {
2273            schema: Some(footer.schema.clone()),
2274            cloud_options: Some(cloud_opts),
2275            hive_options: polars::io::HiveOptions::default(),
2276            glob: false,
2277            ..Default::default()
2278        };
2279        let lf = LazyFrame::scan_parquet(PlRefPath::new(full.as_str()), args)?;
2280        let state =
2281            DataTableState::from_schema_and_lazyframe(footer.schema.clone(), lf, options, None)?;
2282        // Record it in the dataset index, as the prefix route does.
2283        Self::record_cloud_object_facts(report.remembered.as_ref(), &full, &footer);
2284        let column_bytes =
2285            crate::formats::schema_union::column_bytes_per_row(&[Some(footer.clone())]);
2286        let facts = OpenFacts {
2287            remote_objects: vec![crate::cloud::local_copy::RemoteObject {
2288                url: full,
2289                size: footer.file_bytes as u64,
2290                etag,
2291            }],
2292            row_groups: vec![footer.row_group_rows],
2293            column_bytes,
2294            ..Default::default()
2295        };
2296        Ok((state, facts))
2297    }
2298
2299    /// The schema routes, cheapest first: one local footer, one cloud footer (a hive
2300    /// prefix or a single object), then asking the frame.
2301    fn schema_state_by_route(
2302        lf: LazyFrame,
2303        path: Option<&Path>,
2304        options: &OpenOptions,
2305        cloud: &crate::config::CloudConfig,
2306        runtime: &tokio::runtime::Handle,
2307        report: &crate::loading::measurements::OpenReport,
2308    ) -> Result<(DataTableState, OpenFacts, String)> {
2309        #[cfg(not(feature = "cloud"))]
2310        let _ = (cloud, runtime);
2311
2312        // A meter per attempt; the winner's becomes the open's. Earlier routes measure
2313        // before finding they cannot finish, and a shared meter would leave their figures
2314        // on a dataset a later route built. Both full-scan arms also return an empty one;
2315        // `test_a_route_that_gave_up_leaves_no_figures_on_the_dataset_that_opened` holds
2316        // the two guards together.
2317        let attempt = |report: &crate::loading::measurements::OpenReport| {
2318            crate::loading::measurements::OpenReport {
2319                progress: report.progress.clone(),
2320                meter: Arc::new(crate::loading::measurements::Meter::default()),
2321                remembered: report.remembered.clone(),
2322                writes: report.writes.clone(),
2323            }
2324        };
2325
2326        let local = attempt(report);
2327        if let Some((state, facts)) = Self::schema_state_from_local_hive(path, options, &local) {
2328            let facts = OpenFacts {
2329                measurements: local.meter,
2330                ..facts
2331            };
2332            return Ok((state, facts, "one-file (local)".to_string()));
2333        }
2334        // An open abandoned mid-listing is not one for the routes below to scan whole.
2335        if report.progress.is_cancelled() {
2336            return Err(color_eyre::eyre::eyre!("cancelled"));
2337        }
2338        #[cfg(feature = "cloud")]
2339        let cloud_hive_attempt = attempt(report);
2340        #[cfg(feature = "cloud")]
2341        if let Some((state, facts)) =
2342            Self::schema_state_from_cloud_hive(path, options, cloud, runtime, &cloud_hive_attempt)
2343        {
2344            let facts = OpenFacts {
2345                measurements: cloud_hive_attempt.meter,
2346                ..facts
2347            };
2348            return Ok((state, facts, "one-file (cloud)".to_string()));
2349        }
2350        #[cfg(feature = "cloud")]
2351        if let Some(p) = path.filter(|p| {
2352            source::scans_in_place(p)
2353                && !options.hive
2354                && !source::is_prefix_or_glob(&p.to_string_lossy())
2355        }) {
2356            let object = attempt(report);
2357            match Self::schema_state_from_cloud_object(p, options, cloud, runtime, &object) {
2358                Ok((state, facts)) => {
2359                    let facts = OpenFacts {
2360                        measurements: object.meter,
2361                        ..facts
2362                    };
2363                    return Ok((state, facts, "footer (cloud)".to_string()));
2364                }
2365                // Shown in the debug overlay: the fallback costs a row group for the count.
2366                Err(e) => {
2367                    // A fresh meter: the full scan measures nothing, and the failed attempt's figures
2368                    // would describe a route not taken.
2369                    return Self::schema_state_from_full_scan(lf, path, options).map(|state| {
2370                        (
2371                            state,
2372                            OpenFacts::default(),
2373                            format!("full scan (cloud footer: {e})"),
2374                        )
2375                    });
2376                }
2377            }
2378        }
2379        Self::schema_state_from_full_scan(lf, path, options)
2380            .map(|state| (state, OpenFacts::default(), "full scan".to_string()))
2381    }
2382
2383    /// The files of one split when `dir` is a Hugging Face `datasets` cache (its
2384    /// `dataset_info.json` or `state.json` beside Arrow files); choices go in `report`.
2385    /// Any other directory reads every file.
2386    fn hugging_face_split(
2387        dir: &Path,
2388        format: FileFormat,
2389        files: Vec<PathBuf>,
2390        options: &OpenOptions,
2391        report: &mut ReadReport,
2392    ) -> Result<Vec<PathBuf>> {
2393        let metadata = || {
2394            ["dataset_info.json", "state.json"]
2395                .iter()
2396                .any(|name| dir.join(name).is_file())
2397        };
2398        if format != FileFormat::Arrow || !metadata() {
2399            return Ok(files);
2400        }
2401        let names: Vec<&str> = files
2402            .iter()
2403            .map(|f| f.file_name().and_then(|n| n.to_str()).unwrap_or_default())
2404            .collect();
2405        let (chosen, splits) = crate::formats::hf_splits::choose(&names, options.table.as_deref())
2406            .map_err(|e| color_eyre::eyre::eyre!("{}: {e}", dir.display()))?;
2407        report.splits = Some(Arc::new(splits));
2408        Ok(chosen.into_iter().map(|i| files[i].clone()).collect())
2409    }
2410
2411    /// Why `--table` was refused for `path`, a file of `format`, which holds one table.
2412    fn one_table(path: Option<&Path>, format: Option<FileFormat>) -> color_eyre::Report {
2413        match path {
2414            Some(path) => crate::error_display::FileError::new(path, cli::one_table(format)).into(),
2415            None => color_eyre::eyre::eyre!(cli::one_table(format)),
2416        }
2417    }
2418
2419    /// A scan of `url` in an object store that Polars refused, with what to check.
2420    #[cfg(feature = "cloud")]
2421    fn cloud_scan_failed(url: &str, e: &polars::prelude::PolarsError) -> color_eyre::Report {
2422        let said = crate::error_display::user_message_from_polars(e);
2423        let (first, rest) = said.split_once('\n').unwrap_or((&said, ""));
2424        let first = first.trim_end().trim_end_matches('.');
2425        let what = format!("could not read it: {first}. Check the credentials and the URL.");
2426        let what = match rest {
2427            "" => what,
2428            rest => format!("{what}\n{rest}"),
2429        };
2430        crate::error_display::FileError::new(Path::new(url), what).into()
2431    }
2432
2433    /// The inputs of an Arrow read as one table, in order: IPC files scanned in place,
2434    /// and each run of streams as its rows of the `converted` IPC file. Stacked as a
2435    /// directory's files are ([`crate::formats::readers::polars::union_of_files`]).
2436    fn scan_arrow_parts(
2437        cloud: &crate::config::CloudConfig,
2438        converted: Option<&PathBuf>,
2439        parts: &[crate::formats::ipc_stream::Part],
2440    ) -> Result<LazyFrame> {
2441        use crate::formats::ipc_stream::Part;
2442        #[cfg(not(feature = "cloud"))]
2443        let _ = cloud;
2444        let scan = |path: &Path| -> Result<LazyFrame> {
2445            #[cfg(feature = "cloud")]
2446            if source::is_remote_url(path) {
2447                let (url, cloud_options) = Self::resolve_cloud_url(path, cloud)?;
2448                let args = polars::prelude::UnifiedScanArgs {
2449                    cloud_options: Some(cloud_options),
2450                    ..Default::default()
2451                };
2452                return Ok(LazyFrame::scan_ipc(
2453                    PlRefPath::new(url.as_str()),
2454                    Default::default(),
2455                    args,
2456                )?);
2457            }
2458            // Converted streams sit in a user-named temp directory (`[` and all), and a file
2459            // read in place may be called `d[1].arrow`.
2460            let args = polars::prelude::UnifiedScanArgs {
2461                glob: source::expands_as_glob(path),
2462                ..Default::default()
2463            };
2464            Ok(LazyFrame::scan_ipc(
2465                polars::prelude::PlRefPath::try_from_path(path)?,
2466                Default::default(),
2467                args,
2468            )?)
2469        };
2470        let streams = |offset: u64, rows: u64| -> Result<LazyFrame> {
2471            let file = converted
2472                .ok_or_else(|| color_eyre::eyre::eyre!("No converted Arrow file to read."))?;
2473            let lf = scan(file)?;
2474            // The whole file needs no slice, which would hide its row count.
2475            let whole = offset == 0
2476                && parts
2477                    .iter()
2478                    .all(|part| matches!(part, Part::Converted { .. }));
2479            Ok(if whole {
2480                lf
2481            } else {
2482                lf.slice(offset as i64, rows as polars::prelude::IdxSize)
2483            })
2484        };
2485        let mut frames = Vec::new();
2486        let mut run: Option<(u64, u64)> = None;
2487        for part in parts {
2488            match part {
2489                Part::Converted { offset, rows, .. } => {
2490                    run = Some(match run {
2491                        Some((start, n)) if start + n == *offset => (start, n + rows),
2492                        Some((start, n)) => {
2493                            frames.push(streams(start, n)?);
2494                            (*offset, *rows)
2495                        }
2496                        None => (*offset, *rows),
2497                    });
2498                }
2499                Part::InPlace(path) => {
2500                    if let Some((start, n)) = run.take() {
2501                        frames.push(streams(start, n)?);
2502                    }
2503                    frames.push(scan(path)?);
2504                }
2505            }
2506        }
2507        if let Some((start, n)) = run {
2508            frames.push(streams(start, n)?);
2509        }
2510        match frames.len() {
2511            0 => Err(color_eyre::eyre::eyre!("No Arrow files to read.")),
2512            1 => Ok(frames.remove(0)),
2513            _ => Ok(polars::prelude::concat(
2514                frames.as_slice(),
2515                crate::formats::readers::polars::union_of_files(),
2516            )?),
2517        }
2518    }
2519
2520    /// One split of a `save_to_disk` DatasetDict whose `dataset_dict.json` names
2521    /// `splits`: the one `--table` names, else the first, read as any directory. The
2522    /// others are listed.
2523    fn dataset_dict_split(
2524        dir: &Path,
2525        splits: &[String],
2526        options: &OpenOptions,
2527        report: &mut ReadReport,
2528        formats: &crate::formats::Registry,
2529    ) -> Result<Scan> {
2530        let listed: Vec<&str> = splits.iter().map(String::as_str).collect();
2531        let mut picked = crate::formats::hf_splits::pick(&listed, options.table.as_deref())
2532            .map_err(|e| color_eyre::eyre::eyre!("{}: {e}", dir.display()))?;
2533        let split = dir.join(picked.split.as_deref().unwrap_or_default());
2534        let inner = OpenOptions {
2535            table: None,
2536            splits: None,
2537            ..options.clone()
2538        };
2539        let scan = Self::build_local_lazyframe(&[split], &inner, report, formats)?;
2540        // The split's own directory names no splits; its `map()` files are still counted.
2541        picked.caches = report.splits.as_ref().map_or(0, |inner| inner.caches);
2542        report.splits = Some(Arc::new(picked));
2543        Ok(scan)
2544    }
2545
2546    /// Whether the files about to be read as one table differ in columns, for a note.
2547    /// Only footerless formats (Parquet's footers give the exact version), and only
2548    /// the [`crate::formats::schema_union::sample_files`] spread, since the user is waiting.
2549    fn files_disagree(
2550        files: &[PathBuf],
2551        options: &OpenOptions,
2552        found: FileFormat,
2553    ) -> crate::formats::schema_union::Disagreement {
2554        // The format the read will use: `--format` outranks the names, and judging with
2555        // another reader would describe a read that never happened.
2556        let format = options.format.unwrap_or(found);
2557        if format == FileFormat::Parquet {
2558            return Default::default();
2559        }
2560        // `--null` takes `COL=VAL` forms resolved per file, which a sample cannot mirror;
2561        // guessing would report a widening the table never did.
2562        if options.null_values.is_some() {
2563            return Default::default();
2564        }
2565        crate::formats::schema_union::sample_files(files, format, &Self::read_as(options))
2566            .disagreement()
2567    }
2568
2569    /// The reader settings a sample copies to describe what the open will do, from the
2570    /// open's actual options (`from_args_and_config` always fills
2571    /// `infer_schema_length` and `parse_strings`, so "did the user set anything" is
2572    /// always true).
2573    pub(crate) fn read_as(options: &OpenOptions) -> crate::formats::schema_union::ReadAs {
2574        crate::formats::schema_union::ReadAs {
2575            delimiter: options.delimiter,
2576            has_header: options.has_header,
2577            skip_rows: options.skip_rows,
2578            skip_lines: options.skip_lines,
2579            infer_schema_length: options.infer_schema_length,
2580            ignore_errors: options.ignore_errors,
2581            try_parse_dates: options.csv_try_parse_dates(),
2582            comment_char: options.comment_char.clone(),
2583            header_rows: options.header_rows.clone(),
2584            header_join: options.header_join.clone(),
2585        }
2586    }
2587
2588    /// A format spec reads a local file or one downloaded remote object
2589    /// (`loading::remote_download`); a remote prefix or glob here would be scanned in
2590    /// place without the spec.
2591    pub(crate) fn refuse_spec_in_place(path: &Path, options: &OpenOptions) -> Result<()> {
2592        if source::is_remote_url(path)
2593            && (options.spec_file.is_some() || options.spec_name.is_some())
2594        {
2595            return Err(crate::error_display::FileError::new(
2596                path,
2597                "a format spec reads one remote object at a time, not a prefix or a glob; name the object",
2598            )
2599            .into());
2600        }
2601        Ok(())
2602    }
2603
2604    /// Build the LazyFrame for `paths`. Takes the cloud config by reference so it runs
2605    /// on a worker. `found` collects what the read says of itself for the notes (data
2606    /// files passed over, whether columns differ): this pass decides, so a second
2607    /// directory read cannot disagree.
2608    pub(crate) fn build_lazyframe_from_paths_with(
2609        cloud: &crate::config::CloudConfig,
2610        paths: &[PathBuf],
2611        options: &OpenOptions,
2612        report: &mut ReadReport,
2613        formats: &crate::formats::Registry,
2614    ) -> Result<Scan> {
2615        // Arrow streams converted or a bucket's Arrow listed: the load says where each
2616        // input's rows are.
2617        if let Some(parts) = &options.arrow_parts {
2618            if options.table.is_some() && options.splits.is_none() {
2619                // Named by an input the user knows, never the converted copy.
2620                let named = parts.first().map(|part| match part {
2621                    crate::formats::ipc_stream::Part::InPlace(path) => path.as_path(),
2622                    crate::formats::ipc_stream::Part::Converted { source, .. } => source.as_path(),
2623                });
2624                return Err(Self::one_table(named, Some(FileFormat::Arrow)));
2625            }
2626            return Self::scan_arrow_parts(cloud, paths.first(), parts).map(Scan::from);
2627        }
2628        // Only the cloud readers below take the settings.
2629        #[cfg(not(feature = "cloud"))]
2630        let _ = cloud;
2631        let path = &paths[0];
2632        Self::refuse_spec_in_place(path, options)?;
2633        match source::input_source(path) {
2634            source::InputSource::Http(_url) => {
2635                #[cfg(feature = "http")]
2636                {
2637                    return Err(color_eyre::eyre::eyre!(
2638                        "HTTP/HTTPS load is handled in the event loop; this path should not be reached."
2639                    ));
2640                }
2641                #[cfg(not(feature = "http"))]
2642                {
2643                    return Err(color_eyre::eyre::eyre!(
2644                        "HTTP/HTTPS URLs are not supported in this build. Rebuild with default features."
2645                    ));
2646                }
2647            }
2648            source::InputSource::S3(url) => {
2649                #[cfg(feature = "cloud")]
2650                {
2651                    let (full, cloud_opts) =
2652                        Self::resolve_cloud_url(Path::new(&format!("s3://{url}")), cloud)?;
2653                    let is_glob = source::is_prefix_or_glob(&full);
2654                    // The reader the prefix's own format calls for, if the listing said; only Parquet
2655                    // falls through.
2656                    if let Some(format) = Self::cloud_glob_format(&full, options)
2657                        && let Some(lf) = Self::scan_cloud_prefix(
2658                            &full,
2659                            cloud_opts.clone(),
2660                            format,
2661                            is_glob,
2662                            options,
2663                        )
2664                    {
2665                        return lf.map(Scan::from);
2666                    }
2667                    let pl_path = PlRefPath::new(full.as_str());
2668                    let hive_options = if is_glob {
2669                        polars::io::HiveOptions::new_enabled()
2670                    } else {
2671                        polars::io::HiveOptions::default()
2672                    };
2673                    let args = ScanArgsParquet {
2674                        cloud_options: Some(cloud_opts),
2675                        hive_options,
2676                        glob: is_glob,
2677                        ..Default::default()
2678                    };
2679                    let lf = LazyFrame::scan_parquet(pl_path, args)
2680                        .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2681                    // The frame alone: a state would make Polars list every file under the prefix,
2682                    // and the schema phase lists them again.
2683                    return Ok(lf.into());
2684                }
2685                #[cfg(not(feature = "cloud"))]
2686                {
2687                    let _ = url;
2688                    return Err(color_eyre::eyre::eyre!(
2689                        "S3 is not supported in this build. Rebuild with default features and set AWS credentials (e.g. AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, AWS_REGION)."
2690                    ));
2691                }
2692            }
2693            source::InputSource::Gcs(url) => {
2694                #[cfg(feature = "cloud")]
2695                {
2696                    let (full, cloud_opts) =
2697                        Self::resolve_cloud_url(Path::new(&format!("gs://{url}")), cloud)?;
2698                    let is_glob = source::is_prefix_or_glob(&full);
2699                    // The reader the prefix's own format calls for, if the listing said; only Parquet
2700                    // falls through.
2701                    if let Some(format) = Self::cloud_glob_format(&full, options)
2702                        && let Some(lf) = Self::scan_cloud_prefix(
2703                            &full,
2704                            cloud_opts.clone(),
2705                            format,
2706                            is_glob,
2707                            options,
2708                        )
2709                    {
2710                        return lf.map(Scan::from);
2711                    }
2712                    let pl_path = PlRefPath::new(full.as_str());
2713                    let hive_options = if is_glob {
2714                        polars::io::HiveOptions::new_enabled()
2715                    } else {
2716                        polars::io::HiveOptions::default()
2717                    };
2718                    let args = ScanArgsParquet {
2719                        cloud_options: Some(cloud_opts),
2720                        hive_options,
2721                        glob: is_glob,
2722                        ..Default::default()
2723                    };
2724                    let lf = LazyFrame::scan_parquet(pl_path, args)
2725                        .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2726                    return Ok(lf.into());
2727                }
2728                #[cfg(not(feature = "cloud"))]
2729                {
2730                    let _ = url;
2731                    return Err(color_eyre::eyre::eyre!(
2732                        "GCS (gs://) is not supported in this build. Rebuild with default features."
2733                    ));
2734                }
2735            }
2736            source::InputSource::Azure(url) => {
2737                #[cfg(feature = "cloud")]
2738                {
2739                    let (full, cloud_opts) = Self::resolve_cloud_url(Path::new(&url), cloud)?;
2740                    let is_glob = source::is_prefix_or_glob(&full);
2741                    // The reader the prefix's own format calls for, if the listing said; only Parquet
2742                    // falls through.
2743                    if let Some(format) = Self::cloud_glob_format(&full, options)
2744                        && let Some(lf) = Self::scan_cloud_prefix(
2745                            &full,
2746                            cloud_opts.clone(),
2747                            format,
2748                            is_glob,
2749                            options,
2750                        )
2751                    {
2752                        return lf.map(Scan::from);
2753                    }
2754                    let args = ScanArgsParquet {
2755                        cloud_options: Some(cloud_opts),
2756                        hive_options: if is_glob {
2757                            polars::io::HiveOptions::new_enabled()
2758                        } else {
2759                            polars::io::HiveOptions::default()
2760                        },
2761                        glob: is_glob,
2762                        ..Default::default()
2763                    };
2764                    let lf = LazyFrame::scan_parquet(PlRefPath::new(full.as_str()), args)
2765                        .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2766                    return Ok(lf.into());
2767                }
2768                #[cfg(not(feature = "cloud"))]
2769                {
2770                    let _ = url;
2771                    return Err(color_eyre::eyre::eyre!(
2772                        "Azure is not supported in this build. Rebuild with default features."
2773                    ));
2774                }
2775            }
2776            source::InputSource::Local(_) => {}
2777        }
2778        Self::build_local_lazyframe(paths, options, report, formats)
2779    }
2780
2781    /// A directory's files read as `found`, through the delimited spec the first
2782    /// matches, if any.
2783    fn read_directory_files(
2784        files: &[PathBuf],
2785        options: &OpenOptions,
2786        found: FileFormat,
2787        report: &mut ReadReport,
2788        formats: &crate::formats::Registry,
2789    ) -> Result<Scan> {
2790        if options.delimited.is_none()
2791            && options.format.is_none()
2792            && found.separator().is_some()
2793            && let Some(first) = files
2794                .iter()
2795                .find(|f| !crate::formats::nul_tail::holds_nothing(f))
2796            && let Some(choice) = Self::delimited_spec_of(first, options, formats)?
2797        {
2798            let nested = OpenOptions {
2799                hive: false,
2800                format: Some(found),
2801                splits: report.splits.clone(),
2802                ..options.clone()
2803            };
2804            return Self::read_with_delimited_spec(files, &nested, report, formats, choice);
2805        }
2806        // A spec's read reports differing files itself, from their header lines.
2807        report.files_disagree = Self::files_disagree(files, options, found);
2808        let nested = OpenOptions {
2809            hive: false,
2810            format: Some(options.format.unwrap_or(found)),
2811            splits: report.splits.clone(),
2812            ..options.clone()
2813        };
2814        Self::build_local_lazyframe(files, &nested, report, formats)
2815    }
2816
2817    /// The delimited spec whose glob or magic `file` matches, if one does.
2818    fn delimited_spec_of(
2819        file: &Path,
2820        options: &OpenOptions,
2821        formats: &crate::formats::Registry,
2822    ) -> Result<Option<crate::formats::Choice>> {
2823        let asked = crate::formats::Asked {
2824            compression: options.compression,
2825            text_only: true,
2826            ..Default::default()
2827        };
2828        match crate::formats::route(file, &asked, formats)
2829            .map_err(|e| color_eyre::eyre::eyre!(e))?
2830        {
2831            crate::formats::Route::Delimited(choice) => Ok(Some(choice)),
2832            _ => Ok(None),
2833        }
2834    }
2835
2836    /// `paths` read with the CSV reader in the dialect of `choice`'s delimited spec.
2837    fn read_with_delimited_spec(
2838        paths: &[PathBuf],
2839        options: &OpenOptions,
2840        report: &mut ReadReport,
2841        formats: &crate::formats::Registry,
2842        choice: crate::formats::Choice,
2843    ) -> Result<Scan> {
2844        let mut nested = options.clone();
2845        let Some(delimited) = choice.spec.delimited.clone() else {
2846            return Err(color_eyre::eyre::eyre!(
2847                "{} is not a delimited spec",
2848                choice.spec.name
2849            ));
2850        };
2851        delimited.apply(&mut nested);
2852        nested.delimited = Some(Arc::new(
2853            crate::formats::delimited_spec::DelimitedRead::chosen(
2854                choice.spec,
2855                choice.by,
2856                choice.also,
2857            ),
2858        ));
2859        Self::build_local_lazyframe(paths, &nested, report, formats)
2860    }
2861
2862    /// The local half of `build_lazyframe_from_paths_with`: directories resolve to local
2863    /// files, so no cloud settings.
2864    pub(crate) fn build_local_lazyframe(
2865        paths: &[PathBuf],
2866        options: &OpenOptions,
2867        report: &mut ReadReport,
2868        formats: &crate::formats::Registry,
2869    ) -> Result<Scan> {
2870        let path = &paths[0];
2871
2872        // `--hex`: the file's bytes, whatever it holds.
2873        if options.hex
2874            && let [one] = paths
2875            && one.is_file()
2876        {
2877            return Ok(Scan::Hex {
2878                file: one.clone(),
2879                asked: true,
2880            });
2881        }
2882
2883        // A local glob whose first file a delimited spec reads goes through the spec file
2884        // by file; Polars' scan would read the spec's header lines as data.
2885        if let [pattern] = paths
2886            && !options.hive
2887            && options.delimited.is_none()
2888            && options.format.is_none()
2889            && source::expands_as_glob(pattern)
2890        {
2891            let files = crate::loading::local_glob::expand(pattern);
2892            if let Some(first) = files
2893                .iter()
2894                .find(|f| !crate::formats::nul_tail::holds_nothing(f))
2895                && let Some(choice) = Self::delimited_spec_of(first, options, formats)?
2896            {
2897                let format = FileFormat::from_path(first).filter(|f| f.separator().is_some());
2898                let nested = OpenOptions {
2899                    format: format.or(Some(FileFormat::Csv)),
2900                    ..options.clone()
2901                };
2902                return Self::read_with_delimited_spec(&files, &nested, report, formats, choice);
2903            }
2904        }
2905
2906        // A format spec, asked for or matched by glob or magic. A path whose name or bytes
2907        // already say its format opens as such, except for delimited specs' text. Several
2908        // files match by the first, and only to a delimited spec.
2909        if !options.hive && options.delimited.is_none() {
2910            let asked = crate::formats::Asked {
2911                spec_file: options.spec_file.clone(),
2912                spec_name: options.spec_name.clone(),
2913                variant: options.table.clone(),
2914                spec: options.spec_fetched.clone(),
2915                builtin: options.format.is_some(),
2916                compression: options.compression,
2917                text_only: paths.len() > 1,
2918            };
2919            match crate::formats::route(path, &asked, formats)
2920                .map_err(|e| color_eyre::eyre::eyre!(e))?
2921            {
2922                crate::formats::Route::Elsewhere => {}
2923                crate::formats::Route::Delimited(choice) => {
2924                    return Self::read_with_delimited_spec(paths, options, report, formats, choice);
2925                }
2926                crate::formats::Route::Read(read) => {
2927                    let lf = Arc::clone(&read.records).into_lazy()?;
2928                    report.format_read = Some(Arc::new(*read));
2929                    return Ok(lf.into());
2930                }
2931                crate::formats::Route::Decompress(choice) => {
2932                    return Ok(Scan::DecompressSpec {
2933                        file: path.clone(),
2934                        choice,
2935                    });
2936                }
2937            }
2938        } else if options.hive && (options.spec_file.is_some() || options.spec_name.is_some()) {
2939            return Err(color_eyre::eyre::eyre!(
2940                "a format spec reads one file, or one directory of column files"
2941            ));
2942        }
2943
2944        // The header lines of a delimited spec's first file: its units and metadata.
2945        if let Some(read) = &options.delimited
2946            && report.delimited.is_none()
2947            && path.is_file()
2948        {
2949            report.delimited = Some(Arc::new(crate::formats::delimited_spec::read_facts(
2950                read, paths, options,
2951            )?));
2952        }
2953
2954        // One directory path, with or without `--hive`: naming a directory asks to read
2955        // it, and the dispatch below picks the reader for its contents.
2956        if paths.len() == 1 && (options.hive || path.is_dir()) {
2957            // A file is a file whatever its name holds: `a*b.parquet` is not a glob.
2958            let is_single_file = path.is_file();
2959            if !is_single_file {
2960                // The directory's contents pick the reader.
2961                if path.is_dir()
2962                    && let Some(splits) = crate::formats::hf_splits::dataset_dict(path)
2963                {
2964                    return Self::dataset_dict_split(path, &splits, options, report, formats);
2965                }
2966                if path.is_dir() {
2967                    match crate::home::discover::directory_format(path) {
2968                        // Flat and Parquet: the scan below is already right for it.
2969                        crate::home::discover::DirectoryFormat::One(FileFormat::Parquet, _) => {}
2970                        // Partitions, or empty. Only the hive scan walks `key=value` trees, and hive
2971                        // partitioning is Parquet-only here (`HiveOptions::new_disabled()` for CSV and
2972                        // NDJSON), so other formats in partitions are refused with their file names.
2973                        crate::home::discover::DirectoryFormat::Deeper => {
2974                            if let crate::home::discover::DirectoryFormat::One(found, files) =
2975                                crate::home::discover::hive_leaf_format(path)
2976                                && found != FileFormat::Parquet
2977                            {
2978                                // The extension, as the user sees it on the files, not the format's name.
2979                                let named = files
2980                                    .first()
2981                                    .and_then(|f| crate::home::discover::data_extension(f))
2982                                    .unwrap_or_else(|| format!("{found:?}").to_lowercase());
2983                                return Err(color_eyre::eyre::eyre!(
2984                                    "{} is partitioned into key=value directories of .{} \
2985                                     files. datui reads hive partitioning for Parquet \
2986                                     only — open one partition instead.",
2987                                    path.display(),
2988                                    named
2989                                ));
2990                            }
2991                        }
2992                        crate::home::discover::DirectoryFormat::One(found, files) => {
2993                            // Read as a list of files typed on the command line would be; `--format`
2994                            // outranks the names.
2995                            let format = options.format.unwrap_or(found);
2996                            let files =
2997                                Self::hugging_face_split(path, format, files, options, report)?;
2998                            return Self::read_directory_files(
2999                                &files, options, found, report, formats,
3000                            );
3001                        }
3002                        crate::home::discover::DirectoryFormat::Mixed {
3003                            format: found,
3004                            files,
3005                            passed_over,
3006                        } => {
3007                            // The commonest format is the table: a thousand CSVs and one stray JSON are a
3008                            // directory of CSVs.
3009                            let format = options.format.unwrap_or(found);
3010                            let files =
3011                                Self::hugging_face_split(path, format, files, options, report)?;
3012                            let lf = Self::read_directory_files(
3013                                &files, options, found, report, formats,
3014                            )?;
3015                            // A model's config and tokenizer JSON are not data passed over, so a model
3016                            // directory reports nothing left out.
3017                            if !matches!(found, FileFormat::Safetensors | FileFormat::Gguf) {
3018                                report.left_out = passed_over;
3019                            }
3020                            return Ok(lf);
3021                        }
3022                    }
3023                }
3024                let use_parquet_hive =
3025                    path.is_dir() || path.as_os_str().to_string_lossy().contains(".parquet");
3026                if use_parquet_hive {
3027                    // Only the LazyFrame here; schema and partition discovery are the schema phase's.
3028                    return crate::formats::readers::hive::scan_parquet_hive(path).map(Scan::from);
3029                }
3030                return Err(color_eyre::eyre::eyre!(
3031                    "With --hive use a directory or a glob pattern for Parquet (e.g. path/to/dir or path/**/*.parquet)"
3032                ));
3033            }
3034        }
3035
3036        // A file with no extension may still be Parquet (a part file in a `.parquet`
3037        // directory). A text name (`.log`, `.txt`) reads as lines unless its bytes say a
3038        // format: candump writes `.log`.
3039        let compressed = options
3040            .compression
3041            .or_else(|| CompressionFormat::from_extension(path))
3042            .is_some();
3043        // Under a compression suffix, the name before it says delimited text or lines
3044        // (`x.tsv.gz`, `app.log.gz`).
3045        let named = FileFormat::from_path(path).or_else(|| {
3046            compressed
3047                .then(|| FileFormat::from_path(Path::new(path.file_stem()?)))
3048                .flatten()
3049                .filter(|f| f.decompressed_once())
3050        });
3051        let mut effective_format = options
3052            .format
3053            // A format another refines asks the bytes: journal JSON in a `.json` file.
3054            .or_else(|| {
3055                named.filter(|f| !f.is_lines()).map(|f| {
3056                    (!compressed)
3057                        .then(|| crate::formats::readers::refined(path, f))
3058                        .flatten()
3059                        .unwrap_or(f)
3060                })
3061            })
3062            .or_else(|| {
3063                (path.extension().is_none()
3064                    && crate::home::discover::is_parquet_key(&path.to_string_lossy()))
3065                .then_some(FileFormat::Parquet)
3066            })
3067            // Any other unnamed format, by its first bytes (see `crate::formats::readers`).
3068            .or_else(|| crate::formats::readers::sniff_open(path, options.compression))
3069            .or(named);
3070        // Text no signature claims: JSON, CSV or TSV on evidence, lines otherwise; non-text
3071        // bytes are shown raw.
3072        if effective_format.is_none()
3073            && let [file] = paths
3074            && file.is_file()
3075        {
3076            effective_format = crate::formats::lines::guess_file(file, options.compression)
3077                .map(|f| crate::formats::lines::as_asked(f, options));
3078            report.guessed = effective_format.is_some();
3079        }
3080        report.format = effective_format;
3081
3082        // Refused, not ignored: a one-table file with `--table` would pass for the table
3083        // asked for.
3084        if options.table.is_some()
3085            && !effective_format.is_some_and(FileFormat::takes_table)
3086            && options.splits.is_none()
3087        {
3088            return Err(Self::one_table(Some(path), effective_format));
3089        }
3090
3091        // One compressed CSV, TSV, PSV or text file: the load decompresses it
3092        // (`Step::Decompress`) into a copy the dataset holds; a copy made here would be
3093        // dropped with the state below.
3094        if let [file] = paths
3095            && compressed
3096            && let Some(format) = effective_format.filter(|f| f.decompressed_once())
3097        {
3098            return Ok(Scan::Decompress {
3099                file: file.clone(),
3100                format,
3101            });
3102        }
3103
3104        let Some(format) = effective_format else {
3105            if !path.exists() {
3106                return Err(std::io::Error::new(
3107                    std::io::ErrorKind::NotFound,
3108                    format!("File not found: {}", path.display()),
3109                )
3110                .into());
3111            }
3112            // A local file nothing reads is shown as its bytes (#588).
3113            if paths.len() == 1 && path.is_file() {
3114                return Ok(Scan::Hex {
3115                    file: path.clone(),
3116                    asked: false,
3117                });
3118            }
3119            return Err(color_eyre::eyre::eyre!(match paths.len() {
3120                1 => UNSUPPORTED.to_string(),
3121                _ => crate::formats::readers::many_files_refused(),
3122            }));
3123        };
3124        // The same question home asks (`reads_many_files`) before offering a directory, so
3125        // it never offers one this refuses.
3126        if paths.len() > 1 && !format.reads_many_files() {
3127            if !path.exists() {
3128                return Err(std::io::Error::new(
3129                    std::io::ErrorKind::NotFound,
3130                    format!("File not found: {}", path.display()),
3131                )
3132                .into());
3133            }
3134            return Err(color_eyre::eyre::eyre!(
3135                crate::formats::readers::many_files_refused()
3136            ));
3137        }
3138        let guessed;
3139        let options = if report.guessed {
3140            guessed = OpenOptions {
3141                format_guessed: true,
3142                ..options.clone()
3143            };
3144            &guessed
3145        } else {
3146            options
3147        };
3148        crate::formats::readers::scan(crate::formats::readers::ScanIn {
3149            format,
3150            paths,
3151            options,
3152            report,
3153            formats,
3154        })
3155    }
3156}