1use 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#[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 pub(crate) opened: Option<(Vec<PathBuf>, OpenOptions)>,
36 pub(crate) opened_from_home: bool,
38 pub(crate) startup_view: Option<String>,
40 pub(crate) shape_remembered: Option<u64>,
43}
44
45impl OpenedSource {
46 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#[cfg(feature = "cloud")]
67struct CloudTarget<'a> {
68 full: &'a str,
70 key: String,
72 pattern: Option<&'a globset::GlobMatcher>,
75}
76
77pub(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 pub(crate) fn awaiting_dataset(&self) -> bool {
110 self.loading.awaiting_dataset()
111 }
112
113 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 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 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 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 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 pub(crate) fn name_what_is_loading(&mut self, path: PathBuf) {
156 self.loading.name(path);
157 }
158
159 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 pub(crate) fn begin_new_dataset(&mut self) {
170 self.put_down_load_in_flight();
171 self.home_app.previews.drop_prepared();
173 self.reset_chart_state();
174 self.jobs.advance();
175 self.counting.stop_footer_pass();
178 }
179
180 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 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 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 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 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 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 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 let cache = self.cache.clone();
301 self.cache_writes
302 .spawn(move || cache.record_dataset_facts(&[(url, facts)]));
303 }
304
305 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 self.dataset_generation = self.dataset_generation.wrapping_add(1);
323 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 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 let opened = paths.map(|paths| {
343 let options = OpenOptions {
344 format_read: None,
345 sqlite: None,
346 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 Some(options.separator_or(b',')),
360 ),
361 None => (None, None),
362 };
363 self.source
364 .reset_for_dataset(opened, from_home, format, delimiter);
365 if let Some(path) = recent {
367 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 if let Some(closed) = self.data_table_state.replace(state) {
377 crate::app::background::release(closed);
378 }
379 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 state.follow_to(tail.rows(), false);
398 }
399 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 self.retire_a_count_the_rows_answered();
408 self.path = path.clone();
409 self.open_info_documentation();
411 if self.overlay.shows(&crate::Overlay::Info) {
413 self.read_file_facts();
414 self.count_unfit();
415 }
416 self.start_pending_footers();
418 self.index_lines();
419 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 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 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 Some(why) => self.apply_matched_view(&view, why),
472 None => self.apply_view(&view),
473 };
474 match applied {
475 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 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 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 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 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 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 self.reset_chart_state();
544 if self.jobs.supersede(|job| matches!(job, Job::Classify(_))) {
546 self.home.status = None;
547 }
548 self.jobs.take_owed(Self::owed_rows);
551 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 self.screen_generation = self.screen_generation.wrapping_add(1);
567 }
568
569 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 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 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 if directory && jump {
609 go_inside(self, path);
610 return None;
611 }
612 if directory && home::is_object_store_url(&path) {
615 return Some(self.home_open_path(home::directory_dataset_url(&path), false));
616 }
617 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 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 let prepared = (!directory)
638 .then(|| self.home_app.previews.take_prepared(&path))
639 .flatten();
640 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 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 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 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 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 AppEvent::LookThenOpenDirectory(dir, options)
701 }
702
703 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 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 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 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 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 self.rest_at_start();
784 self.enter_home();
785 self.home_jump_into(dir);
786 None
787 }
788
789 #[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 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 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 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 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 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 #[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 #[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 #[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 #[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 #[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 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 #[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 #[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 let response = agent
1045 .get(&url)
1046 .call()
1047 .map_err(|e| crate::error_display::http_message(&url, &e))?;
1048 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 #[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 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 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 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 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 .map_err(named)?;
1192 let format = scan.format(report.format.or(options.format));
1193 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 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 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 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 fn spawn_load_phase(&mut self, load: loading::LoadId, step: loading::Step) {
1345 use loading::{LoadAnswer, Step};
1346 let job = Job::Load(load);
1347 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 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 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 #[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 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 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 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 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 #[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 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 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 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(©, &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 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 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 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 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 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 #[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 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
1911impl 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 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 #[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 if !options.single_spine_schema || options.format == Some(FileFormat::Arrow) {
1950 return None;
1951 }
1952 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 #[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 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 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 #[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 #[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 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 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 let (state, mut facts, label) =
2114 Self::schema_state_by_route(lf, path, options, cloud, runtime, report)?;
2115 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 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 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 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 #[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 if !format.reads_bucket_prefix() {
2206 return None;
2207 }
2208 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 #[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 #[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 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 #[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 #[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 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 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 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 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 Err(e) => {
2426 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 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 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 #[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 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 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 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 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 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 fn files_disagree(
2609 files: &[PathBuf],
2610 options: &OpenOptions,
2611 found: FileFormat,
2612 ) -> crate::formats::schema_union::Disagreement {
2613 let format = options.format.unwrap_or(found);
2616 if format == FileFormat::Parquet {
2617 return Default::default();
2618 }
2619 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 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 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 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 if let Some(parts) = &options.arrow_parts {
2677 if options.table.is_some() && options.splits.is_none() {
2678 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 #[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 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 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 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 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 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 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 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 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 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 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 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 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 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 if paths.len() == 1 && (options.hive || path.is_dir()) {
3016 let is_single_file = path.is_file();
3018 if !is_single_file {
3019 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 crate::home::discover::DirectoryFormat::One(FileFormat::Parquet, _) => {}
3029 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 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 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 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 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 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 let compressed = options
3099 .compression
3100 .or_else(|| CompressionFormat::from_extension(path))
3101 .is_some();
3102 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 .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 .or_else(|| crate::formats::readers::sniff_open(path, options.compression))
3128 .or(named);
3129 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 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 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 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 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}