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