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 ..options.clone()
350 };
351 (paths, options)
352 });
353 let (format, delimiter) = match path.as_deref() {
354 Some(p) => (
355 Self::export_format_for(p, state.read_as().or(options.format)),
356 Some(options.separator_or(b',')),
359 ),
360 None => (None, None),
361 };
362 self.source
363 .reset_for_dataset(opened, from_home, format, delimiter);
364 if let Some(path) = recent {
366 let cache = self.cache.clone();
369 self.cache_writes.spawn(move || {
370 cache.push_recent(&path);
371 });
372 }
373 self.forget_the_rows_read();
374 self.data_table_state = Some(state);
375 if options.follow
377 && let Some(state) = self.data_table_state.as_mut()
378 {
379 match options.tail.as_deref() {
380 Some(tail) => {
381 let follow = crate::loading::follow::Follow::start(
382 tail.clone(),
383 self.app_config.read.follow_interval.duration(),
384 self.events.clone(),
385 options.spool.clone(),
386 );
387 state.start_following(if options.pipe {
388 follow.as_pipe()
389 } else {
390 follow
391 });
392 state.follow_to(tail.rows(), false);
394 }
395 None => self.flash_note(
397 "Only text and Arrow streams are followed: this shows what had arrived, and recording goes on"
398 .to_string(),
399 ),
400 }
401 }
402 self.retire_a_count_the_rows_answered();
404 self.path = path.clone();
405 self.open_info_documentation();
407 if self.overlay.shows(&crate::Overlay::Info) {
409 self.read_file_facts();
410 self.count_unfit();
411 }
412 self.start_pending_footers();
414 self.index_lines();
415 if options.row_numbers_auto
417 && let Some(state) = self.data_table_state.as_mut()
418 && state.numbered_by_default()
419 {
420 state.set_row_numbers(true);
421 }
422 self.status_message = Some(Self::LOADING_BUFFER.to_string());
423
424 let (view, reason) = match self.source.startup_view.take() {
428 Some(name) => match self.views.manager.get_view_by_name(&name).cloned() {
429 Some(view) => (Some(view), None),
430 None => {
431 self.error_modal.show(format!("No view named \"{name}\""));
432 (None, None)
433 }
434 },
435 None if self.app_config.views.auto_apply => self
436 .view_dataset()
437 .zip(self.data_table_state.as_ref())
438 .and_then(|(dataset, state)| {
439 self.views
440 .manager
441 .get_most_relevant(dataset, state.source_schema())
442 })
443 .map_or((None, None), |(view, reason)| (Some(view), Some(reason))),
444 None => (None, None),
445 };
446 let Some(view) = view else {
447 return false;
448 };
449 let applied = match reason {
450 Some(why) => self.apply_matched_view(&view, why),
452 None => self.apply_view(&view),
453 };
454 match applied {
455 Ok(()) => true,
457 Err(e) => {
458 self.error_modal
459 .show(format!("Error applying view \"{}\": {e}", view.name));
460 false
461 }
462 }
463 }
464
465 pub fn abandon_load(&mut self) {
471 let retired = self.loading.retire();
472 if let Some(retired) = retired {
473 self.put_down_load(retired);
474 }
475 self.reset_chart_state();
478 if self.jobs.supersede(|job| matches!(job, Job::Classify(_))) {
480 self.home.status = None;
481 }
482 self.jobs.take_owed(Self::owed_rows);
485 if retired.is_some() {
488 self.busy = false;
489 let quieted = self.jobs.quiet(Self::reading_rows);
490 if self
491 .status_message
492 .as_ref()
493 .is_some_and(|status| quieted.contains(status))
494 {
495 self.status_message = None;
496 }
497 }
498 self.screen_generation = self.screen_generation.wrapping_add(1);
501 }
502
503 pub(crate) fn looking_could_block(&self, path: &Path) -> bool {
507 !home::is_cloud_place(path)
510 && matches!(source::input_source(path), source::InputSource::Local(_))
511 && (self.home.network_check)(path)
512 }
513
514 pub(crate) fn open_what_it_is(
517 &mut self,
518 path: PathBuf,
519 kind: discover::EntryKind,
520 jump: bool,
521 ) -> Option<AppEvent> {
522 let go_inside = |app: &mut Self, path: PathBuf| {
523 if jump {
524 app.home_jump_into(path);
525 } else {
526 app.home_browse_into(path);
527 }
528 };
529 if kind == discover::EntryKind::Directory {
530 go_inside(self, path);
531 return None;
532 }
533 if kind == discover::EntryKind::Other {
536 if matches!(source::input_source(&path), source::InputSource::Local(_)) {
537 self.open_hex(path, crate::app::hex_view::Origin::Home, true, None);
538 }
539 return None;
540 }
541 if let Some(format) = kind.lake_name() {
544 self.home.lake_here = Some((path.clone(), format));
545 go_inside(self, path);
546 return None;
547 }
548 let directory = matches!(
549 kind,
550 discover::EntryKind::Hive | discover::EntryKind::MultiFile
551 );
552 if directory && jump {
554 go_inside(self, path);
555 return None;
556 }
557 if directory && home::is_object_store_url(&path) {
560 return Some(self.home_open_path(home::directory_dataset_url(&path), false));
561 }
562 let a_spec_may_read = !self.formats.by_glob(&path, false).is_empty()
566 || self.formats.specs.iter().any(|f| !f.spec.magic.is_empty());
567 if kind == discover::EntryKind::File
570 && discover::unreadable_by_name(&path)
571 && !a_spec_may_read
572 && crate::formats::members::split(&path).is_none()
573 && crate::formats::members::holder(&path).is_none()
574 && crate::formats::members::split_variant(&path, &self.formats).is_none()
575 && crate::formats::hf_splits::split_place(&path).is_none()
576 {
577 self.home.status = Some(discover::NO_READER.to_string());
578 return None;
579 }
580 let prepared = (!directory)
583 .then(|| self.home_app.previews.take_prepared(&path))
584 .flatten();
585 let unasked = self
588 .home
589 .catalogs
590 .iter()
591 .filter(|c| c.origin == catalog::Origin::Bundled)
592 .flat_map(|c| c.datasets.iter())
593 .find(|d| d.location == path)
594 .filter(|_| {
595 !jump && matches!(source::input_source(&path), source::InputSource::Http(_))
596 })
597 .map(|dataset| UnaskedDownload {
598 limit: UnaskedDownload::LIMIT,
599 listed: dataset.size,
600 });
601 match self.home_open_path(path, directory) {
602 AppEvent::Open(paths, mut options) => {
603 options.prepared = prepared.map(|p| Arc::new(Mutex::new(Some(p))));
604 options.download_unasked = unasked;
605 Some(AppEvent::Open(paths, options))
606 }
607 event => Some(event),
608 }
609 }
610
611 pub fn route_named_paths(paths: Vec<PathBuf>, options: OpenOptions) -> AppEvent {
617 Self::route_named_paths_with(paths, options, &crate::formats::Registry::default())
618 }
619
620 pub fn route_named_paths_with(
623 paths: Vec<PathBuf>,
624 options: OpenOptions,
625 formats: &crate::formats::Registry,
626 ) -> AppEvent {
627 if let Some(event) = Self::route_named_without_looking(&paths, &options) {
628 return event;
629 }
630 let single = (paths.len() == 1 && !options.hive).then(|| paths[0].clone());
633 let Some(dir) = single.filter(|p| p.is_dir()) else {
634 return AppEvent::Open(paths, options);
635 };
636 if options.spec_file.is_some()
638 || options.spec_name.is_some()
639 || !formats.by_glob(&dir, true).is_empty()
640 {
641 return AppEvent::Open(paths, options);
642 }
643 AppEvent::LookThenOpenDirectory(dir, options)
646 }
647
648 pub(crate) fn route_named_without_looking(
652 paths: &[PathBuf],
653 options: &OpenOptions,
654 ) -> Option<AppEvent> {
655 #[cfg(feature = "cloud")]
656 if let [dir] = paths
657 && !options.hive
658 && home::is_object_store_url(dir)
659 && options.format.is_none()
660 && !dir.to_string_lossy().contains('*')
661 && !home::names_a_file(dir)
662 {
663 return Some(AppEvent::LookThenOpenDirectory(
664 dir.clone(),
665 options.clone(),
666 ));
667 }
668 let _ = (paths, options);
669 None
670 }
671
672 pub fn missing_named_path(
675 paths: &[PathBuf],
676 formats: &crate::formats::Registry,
677 ) -> Option<PathBuf> {
678 paths
679 .iter()
680 .find(|path| {
681 !source::is_remote_url(path)
682 && !crate::loading::stdin::is_stdin(path)
683 && !source::expands_as_glob(path)
684 && !path.exists()
685 && crate::formats::members::split(path).is_none()
686 && crate::formats::members::split_variant(path, formats).is_none()
687 })
688 .cloned()
689 }
690
691 pub(crate) fn open_the_directory_looked_at(
694 &mut self,
695 dir: PathBuf,
696 kind: discover::EntryKind,
697 holds: Option<&discover::Holds>,
698 mut options: OpenOptions,
699 ) -> Option<AppEvent> {
700 #[cfg(feature = "cloud")]
701 if home::is_object_store_url(&dir) {
702 return self.open_the_cloud_directory_looked_at(dir, kind, holds, options);
703 }
704 let _ = holds;
705 if let Some(format) = kind.lake_name() {
709 self.enter_home();
710 self.home.lake_here = Some((dir.clone(), format));
711 self.home_jump_into(dir);
712 return None;
713 }
714 if matches!(
717 kind,
718 discover::EntryKind::Hive | discover::EntryKind::MultiFile
719 ) {
720 options.hive = true;
721 self.set_loading_phase("Scanning input", 10);
722 self.name_what_is_loading(dir.clone());
723 return Some(AppEvent::Open(vec![dir], options));
724 }
725 self.enter_home();
728 self.home_jump_into(dir);
729 None
730 }
731
732 #[cfg(feature = "cloud")]
735 fn open_the_cloud_directory_looked_at(
736 &mut self,
737 dir: PathBuf,
738 kind: discover::EntryKind,
739 holds: Option<&discover::Holds>,
740 options: OpenOptions,
741 ) -> Option<AppEvent> {
742 let open = |app: &mut Self, path: PathBuf, options: OpenOptions| {
743 app.set_loading_phase("Scanning input", 10);
744 app.name_what_is_loading(path.clone());
745 Some(AppEvent::Open(vec![path], options))
746 };
747 let Some(holds) = holds else {
749 return open(self, dir, options);
750 };
751 if let Some(format) = kind.lake_name() {
752 self.enter_home();
753 self.home.lake_here = Some((dir.clone(), format));
754 self.home_jump_into(dir);
755 return None;
756 }
757 let directory = home::directory_dataset_url(&dir);
758 if matches!(
759 kind,
760 discover::EntryKind::Hive | discover::EntryKind::MultiFile
761 ) {
762 let options = OpenOptions {
763 hive: true,
764 ..options
765 };
766 return open(self, directory, options);
767 }
768 if let Some((format, left_out)) = Self::cloud_prefix_format(holds) {
769 let options = OpenOptions {
770 hive: true,
771 format: Some(format),
772 left_out,
773 ..options
774 };
775 return open(self, directory, options);
776 }
777 self.enter_home();
779 self.home_jump_into(dir);
780 self.home.status = Self::why_a_cloud_prefix_cannot_be_read(holds);
781 None
782 }
783
784 pub(crate) fn open_defaults(&self) -> OpenOptions {
787 match crate::cli::parse_args(["datui"]) {
788 Ok(args) => OpenOptions::from_args_and_config(&args, &self.app_config),
789 Err(_) => OpenOptions::default(),
790 }
791 }
792
793 fn with_delimited_spec(
797 file: &Path,
798 mut options: OpenOptions,
799 formats: &crate::formats::Registry,
800 ) -> Result<OpenOptions> {
801 if options.delimited.is_some() {
802 return Ok(options);
803 }
804 let asked = crate::formats::Asked {
805 spec_file: options.spec_file.clone(),
806 spec: options.spec_fetched.clone(),
807 spec_name: options.spec_name.clone(),
808 compression: options.compression,
809 ..Default::default()
810 };
811 let crate::formats::Route::Delimited(choice) =
812 crate::formats::route(file, &asked, formats).map_err(|e| color_eyre::eyre::eyre!(e))?
813 else {
814 return Ok(options);
815 };
816 let Some(delimited) = choice.spec.delimited.clone() else {
817 return Ok(options);
818 };
819 delimited.apply(&mut options);
820 let chosen = crate::formats::delimited_spec::DelimitedRead::chosen(
821 choice.spec,
822 choice.by,
823 choice.also,
824 );
825 let read =
826 crate::formats::delimited_spec::read_facts(&chosen, &[file.to_path_buf()], &options)?;
827 options.delimited = Some(Arc::new(read));
828 Ok(options)
829 }
830
831 fn decompressed_delimited_state(
835 path: &Path,
836 options: &OpenOptions,
837 writer: &crate::loading::unfinished::Writer,
838 ) -> Result<DataTableState> {
839 let separator = options
840 .format
841 .and_then(FileFormat::separator)
842 .unwrap_or(b',');
843 DataTableState::from_read(
844 crate::formats::readers::csv::read_delimited(path, separator, options, writer)?,
845 options,
846 )
847 }
848
849 #[cfg(feature = "cloud")]
851 fn build_s3_cloud_options(settings: &crate::cloud::cloud_sources::S3Settings) -> CloudOptions {
852 let settings = settings.clone();
853 let virtual_hosted = (settings.endpoint.is_some() || settings.virtual_hosted.is_some())
854 .then(|| settings.virtual_hosted_style().to_string());
855 let configs: Vec<(AmazonS3ConfigKey, String)> = [
856 (AmazonS3ConfigKey::Endpoint, settings.endpoint),
857 (AmazonS3ConfigKey::AccessKeyId, settings.access_key_id),
858 (
859 AmazonS3ConfigKey::SecretAccessKey,
860 settings.secret_access_key,
861 ),
862 (AmazonS3ConfigKey::Token, settings.session_token),
863 (AmazonS3ConfigKey::Region, settings.region),
864 (AmazonS3ConfigKey::VirtualHostedStyleRequest, virtual_hosted),
865 (
866 AmazonS3ConfigKey::SkipSignature,
867 settings.skip_signature.then(|| "true".to_string()),
868 ),
869 ]
870 .into_iter()
871 .filter_map(|(key, value)| value.map(|v| (key, v)))
872 .chain([(
873 AmazonS3ConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
874 crate::cloud::user_agent::get(),
875 )])
876 .collect();
877 CloudOptions::default().with_aws(configs)
878 }
879
880 #[cfg(feature = "cloud")]
883 pub(crate) fn cloud_bucket_and_key(url: &str) -> Result<(String, String)> {
884 if let Some((_, container, key)) = source::azure_parts(url) {
885 return Ok((container, key.trim_matches('/').to_string()));
886 }
887 crate::cloud::cloud_browse::split_bucket_url(url)
888 .map(|(_, bucket, key)| (bucket, key))
889 .ok_or_else(|| {
890 color_eyre::eyre::eyre!("URL must be s3://bucket/key or gs://bucket/key")
891 })
892 }
893
894 #[cfg(feature = "cloud")]
898 fn polars_object_store(
899 url: &str,
900 options: &CloudOptions,
901 runtime: &tokio::runtime::Handle,
902 ) -> Result<Arc<dyn object_store::ObjectStore>> {
903 let url = url.to_string();
904 let options = options.clone();
905 wait_on_runtime(runtime, async move {
906 let (_, store) = polars::io::cloud::build_object_store(
907 PlRefPath::new(url.as_str()),
908 Some(&options),
909 false,
910 )
911 .await?;
912 polars::prelude::PolarsResult::Ok(store.to_dyn_object_store().await.into_owned())
913 })
914 .ok_or_else(|| color_eyre::eyre::eyre!("cancelled"))?
915 .map_err(|e| color_eyre::eyre::eyre!("Object store config failed: {}", e))
916 }
917
918 #[cfg(feature = "http")]
922 fn http_agent(total: std::time::Duration) -> ureq::Agent {
923 crate::cloud::user_agent::ureq_config()
924 .timeout_global(Some(total))
925 .build()
926 .into()
927 }
928
929 #[cfg(feature = "http")]
932 pub(crate) fn fetch_remote_size_http(
933 url: &str,
934 ) -> std::result::Result<Option<u64>, crate::error_display::HttpGone> {
935 let agent = Self::http_agent(std::time::Duration::from_secs(15));
936 match agent.head(url).header("Accept-Encoding", "identity").call() {
939 Ok(r) => Ok(r
940 .headers()
941 .get("Content-Length")
942 .and_then(|v| v.to_str().ok())
943 .and_then(|s| s.parse::<u64>().ok())),
944 Err(e) => crate::error_display::http_gone(url, &e).map_or(Ok(None), Err),
945 }
946 }
947
948 #[cfg(feature = "cloud")]
950 fn fetch_remote_size_cloud(
951 url: &str,
952 cloud: &crate::config::CloudConfig,
953 runtime: &tokio::runtime::Handle,
954 ) -> Result<Option<u64>> {
955 use object_store::ObjectStoreExt;
956
957 let (_bucket, key) = Self::cloud_bucket_and_key(url)?;
958 if key.is_empty() {
959 return Ok(None);
960 }
961 let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)?;
962 let path = crate::cloud::cloud_browse::object_path(&key);
963 let head = wait_on_runtime(runtime, async move { store.head(&path).await });
964 Ok(head.and_then(|r| r.ok()).map(|meta| meta.size))
965 }
966
967 #[cfg(feature = "http")]
972 fn download_http_to_temp(
973 url: &str,
974 temp_dir: Option<&Path>,
975 extension: Option<&str>,
976 limit: Option<u64>,
977 writer: &crate::loading::unfinished::Writer,
978 ) -> Result<crate::cloud::download::TempDownload> {
979 use crate::cloud::download::StreamError;
980
981 let url = url.to_string();
982 let open = move || {
983 let agent = Self::http_agent(std::time::Duration::from_secs(300));
984 let response = agent
986 .get(&url)
987 .call()
988 .map_err(|e| crate::error_display::http_message(&url, &e))?;
989 Ok((response.into_body().into_reader(), None))
991 };
992 crate::cloud::download::read_to_temp(temp_dir, extension, open, writer, limit).map_err(
993 |error| match error {
994 StreamError::Open(message) => color_eyre::eyre::eyre!(message),
995 StreamError::Read(e) => {
996 color_eyre::eyre::eyre!("Download failed partway. Check your connection: {e}")
997 }
998 StreamError::Short { expected, got } => color_eyre::eyre::eyre!(
999 "Download failed partway: it ended after {got} of {expected} bytes."
1000 ),
1001 StreamError::Write(report) => report,
1002 StreamError::Cut => color_eyre::eyre::eyre!("Download was cancelled."),
1003 },
1004 )
1005 }
1006
1007 #[cfg(feature = "cloud")]
1011 fn download_cloud_to_temp(
1012 url: &str,
1013 cloud: &crate::config::CloudConfig,
1014 options: &OpenOptions,
1015 runtime: &tokio::runtime::Handle,
1016 writer: &crate::loading::unfinished::Writer,
1017 ) -> Result<crate::cloud::download::TempDownload> {
1018 use crate::cloud::download::StreamError;
1019 use object_store::ObjectStoreExt;
1020
1021 let (label, example) = match source::input_source(Path::new(url)) {
1022 source::InputSource::Gcs(_) => ("GCS", "gs://bucket/path/file.csv"),
1023 source::InputSource::Azure(_) => (
1024 "Azure",
1025 "abfss://container@account.dfs.core.windows.net/path/file.csv",
1026 ),
1027 _ => ("S3", "s3://bucket/path/file.csv"),
1028 };
1029 let ext = source::download_suffix(url);
1030 let (_bucket, key) = Self::cloud_bucket_and_key(url)?;
1031 if key.is_empty() {
1032 return Err(crate::error_display::FileError::new(
1033 Path::new(url),
1034 format!("a {label} URL names an object here, such as {example}"),
1035 )
1036 .into());
1037 }
1038 let (_, _, store) = Self::cloud_store_for(Path::new(url), cloud, runtime)?;
1039
1040 let path = crate::cloud::cloud_browse::object_path(&key);
1041 let open = async move {
1042 let got = store
1043 .get(&path)
1044 .await
1045 .map_err(|e| crate::error_display::store_message(&e))?;
1046 let len = got.range.end - got.range.start;
1047 Ok((got.into_stream(), Some(len)))
1048 };
1049 let failed = |what: String| -> color_eyre::Report {
1050 crate::error_display::FileError::new(Path::new(url), what).into()
1051 };
1052 crate::cloud::download::stream_to_temp(
1053 runtime,
1054 options.temp_dir.as_deref(),
1055 ext.as_deref(),
1056 open,
1057 writer,
1058 )
1059 .map_err(|error| match error {
1060 StreamError::Open(e) => failed(e),
1061 StreamError::Read(e) => failed(format!("the download stopped: {e}")),
1062 StreamError::Short { expected, got } => {
1063 failed(format!("it ended after {got} of {expected} bytes"))
1064 }
1065 StreamError::Write(report) => report,
1066 StreamError::Cut => failed("the download was cancelled".to_string()),
1067 })
1068 }
1069
1070 pub(crate) fn scan_for_open(
1074 cloud: &crate::config::CloudConfig,
1075 formats: &crate::formats::Registry,
1076 paths: &[PathBuf],
1077 options: OpenOptions,
1078 path: Option<PathBuf>,
1079 ) -> std::result::Result<loading::LoadAnswer, String> {
1080 use loading::LoadAnswer;
1081 let bytes_of = |files: &[PathBuf]| -> u64 {
1082 files
1083 .iter()
1084 .filter_map(|f| std::fs::metadata(f).ok())
1085 .map(|m| m.len())
1086 .sum()
1087 };
1088 let mut report = ReadReport {
1093 left_out: options.left_out.clone(),
1094 files_disagree: options.files_disagree,
1095 format: None,
1096 format_read: None,
1097 read_python: Vec::new(),
1098 sqlite: None,
1099 opened: None,
1100 splits: options.splits.clone(),
1101 delimited: None,
1102 table: None,
1103 guessed: false,
1104 read_notes: Vec::new(),
1105 typing: Default::default(),
1106 };
1107 let options = OpenOptions {
1110 ignore_errors: options.ignore_errors || options.follow,
1111 ..options
1112 };
1113 let named = |e: color_eyre::Report| {
1114 crate::error_display::user_message_from_report(&e, path.as_deref())
1115 };
1116 let followed_stream = options.follow
1119 && crate::loading::follow::followed_stream(
1120 &paths[0],
1121 Some(crate::loading::follow::format_of(&paths[0], options.format)),
1122 &options,
1123 );
1124 let scan = if followed_stream {
1125 crate::loading::follow::stream::scan(&paths[0])
1126 .map(Scan::from)
1127 .map_err(|e| color_eyre::eyre::eyre!(e))
1128 } else {
1129 Self::build_lazyframe_from_paths_with(cloud, paths, &options, &mut report, formats)
1130 }
1131 .map_err(named)?;
1133 let format = scan.format(report.format.or(options.format));
1134 let recording = options
1137 .spool
1138 .as_ref()
1139 .is_some_and(|handle| handle.spool().tee().is_some());
1140 let (scan, tail) = match scan {
1141 Scan::Frame(lf) if options.follow => {
1142 let format = crate::loading::follow::format_of(&paths[0], format);
1143 let refused = (!followed_stream)
1144 .then(|| crate::loading::follow::refusal(Some(format), &options))
1145 .flatten();
1146 match refused {
1147 Some(_) if recording => (Scan::Frame(lf), None),
1148 Some(refusal) => return Err(refusal),
1149 None => {
1150 let (lf, tail) = crate::loading::follow::bound_to_complete(
1151 *lf, &paths[0], format, &options,
1152 )
1153 .map_err(named)?;
1154 (Scan::Frame(Box::new(lf)), Some(Arc::new(tail)))
1155 }
1156 }
1157 }
1158 _ if options.follow && !recording => {
1159 return Err(crate::loading::follow::refusal(format, &options)
1160 .unwrap_or_else(|| "This file cannot be followed as it grows.".to_string()));
1161 }
1162 scan => (scan, None),
1163 };
1164 let read_mode = scan.read_mode(format, report.format_read.is_some(), &options);
1165 let mut options = OpenOptions {
1166 left_out: report.left_out,
1167 files_disagree: report.files_disagree,
1168 format,
1169 format_read: report.format_read,
1170 sqlite: report.sqlite,
1171 opened: report.opened,
1172 splits: report.splits,
1173 read_python: report.read_python,
1174 read_mode,
1175 tail,
1176 table: report.table.or_else(|| options.table.clone()),
1177 format_guessed: options.format_guessed || report.guessed,
1178 read_notes: report.read_notes,
1179 typing: report.typing,
1180 ..options
1181 };
1182 if let Some(read) = report.delimited {
1185 read.delimited().apply(&mut options);
1186 options.delimited = Some(read);
1187 }
1188 Ok(match scan {
1189 Scan::Frame(lf) => LoadAnswer::Scanned { lf, path, options },
1190 Scan::Decompress { file, .. } => LoadAnswer::Compressed {
1191 file,
1192 path,
1193 options,
1194 },
1195 Scan::Streams(files) => LoadAnswer::Convert {
1196 what: loading::Conversion::Streams,
1197 bytes: bytes_of(&files),
1198 files,
1199 path,
1200 options,
1201 },
1202 Scan::DecompressSpec { file, choice } => LoadAnswer::CompressedRecords {
1203 file,
1204 path,
1205 choice,
1206 options,
1207 },
1208 Scan::ReadInto { files, format } => LoadAnswer::Convert {
1209 what: loading::Conversion::Text(format),
1210 bytes: bytes_of(&files),
1211 files,
1212 path,
1213 options,
1214 },
1215 Scan::Tables { file, tables, .. } => LoadAnswer::Tables { file, tables, path },
1216 Scan::Unpack {
1217 file,
1218 member,
1219 format,
1220 } => LoadAnswer::Convert {
1221 what: loading::Conversion::Text(format),
1222 bytes: bytes_of(std::slice::from_ref(&file)),
1223 files: vec![file],
1224 path,
1225 options: OpenOptions {
1226 table: Some(member),
1227 ..options
1228 },
1229 },
1230 Scan::Hex { file, asked } => LoadAnswer::Hex {
1231 file,
1232 asked,
1233 record_size: options.record_size,
1234 },
1235 })
1236 }
1237
1238 pub(crate) fn read_schema_for_open(
1241 lf: LazyFrame,
1242 path: Option<PathBuf>,
1243 options: OpenOptions,
1244 cloud: &crate::config::CloudConfig,
1245 runtime: &tokio::runtime::Handle,
1246 report: &crate::loading::measurements::OpenReport,
1247 made: loading::Made,
1248 ) -> std::result::Result<loading::LoadAnswer, String> {
1249 use loading::LoadAnswer;
1250 let (state, facts, debug_label) =
1251 Self::build_schema_state(lf, path.as_deref(), &options, cloud, runtime, report)
1252 .map_err(|e| crate::error_display::user_message_from_report(&e, path.as_deref()))?;
1253 let loading::Made {
1255 download,
1256 converted,
1257 notes,
1258 other_tables,
1259 detail,
1260 } = made;
1261 let mut open_notes = facts.open_notes;
1262 open_notes.extend(notes);
1263 let mut other_tables_found = facts.other_tables;
1264 other_tables_found.extend(other_tables);
1265 let state = state.with_open(OpenFacts {
1266 fetched: Self::was_fetched(download.as_ref(), path.as_deref()),
1267 download,
1268 converted,
1269 other_tables: other_tables_found,
1270 open_notes,
1271 detail: detail.or(facts.detail),
1272 ..facts
1273 });
1274 Ok(LoadAnswer::SchemaRead {
1275 state: Box::new(state),
1276 path,
1277 options,
1278 debug_label: Some(debug_label),
1279 })
1280 }
1281
1282 fn spawn_load_phase(&mut self, load: loading::LoadId, step: loading::Step) {
1286 use loading::{LoadAnswer, Step};
1287 let job = Job::Load(load);
1288 let (cloud, runtime) = (self.app_config.cloud.clone(), self.runtime.clone());
1290 match step {
1291 #[cfg(any(feature = "http", feature = "cloud"))]
1292 Step::ReadHeaders {
1293 url,
1294 format,
1295 options,
1296 writer,
1297 } => {
1298 self.spawn_job(job, Some("Reading headers..."), move |_| {
1299 let read =
1300 crate::cloud::remote_model::read(&url, format, &cloud, &runtime, &|| {
1301 writer.stopped()
1302 });
1303 let crate::cloud::remote_model::Read { lf, summary, notes } = match read {
1304 Ok(read) => read,
1305 Err(crate::formats::model_files::RangeError::NoRanges) => {
1306 return Ok(Answer::Load(Box::new(LoadAnswer::NoRanges { options })));
1307 }
1308 Err(crate::formats::model_files::RangeError::Failed(message)) => {
1310 return Err(crate::logging::redact(&message, &[]));
1311 }
1312 };
1313 let opened = Arc::new(crate::formats::model_files::opened(&summary));
1314 let options = OpenOptions {
1315 format: Some(format),
1316 opened: Some(opened.clone()),
1317 ..options
1318 };
1319 let state = Self::schema_state_from_full_scan(
1321 lf,
1322 None,
1323 &OpenOptions {
1324 hive: false,
1325 ..options.clone()
1326 },
1327 )
1328 .map_err(|e| crate::error_display::user_message_from_report(&e, Some(&url)))?
1329 .with_open(OpenFacts {
1330 detail: opened.detail.clone(),
1331 open_notes: notes,
1332 read_as: Some(format),
1333 ..Default::default()
1334 });
1335 Ok(Answer::Load(Box::new(LoadAnswer::SchemaRead {
1336 state: Box::new(state),
1337 path: Some(url),
1338 options,
1339 debug_label: Some("model headers (ranged)".to_string()),
1340 })))
1341 });
1342 }
1343 #[cfg(any(feature = "http", feature = "cloud"))]
1344 Step::Probe(pending) => {
1345 self.spawn_job(job, Some("Checking size..."), move |_| {
1346 #[cfg(feature = "cloud")]
1348 if let loading::PendingDownload::Arrow { url, .. } = &pending {
1349 let (_, _, options) = pending.parts();
1350 let (objects, options) =
1351 crate::cloud::cloud_arrow::list(url, options, &cloud, &runtime)
1352 .map_err(|e| {
1353 crate::error_display::user_message_from_report(&e, None)
1354 })?;
1355 let size = crate::cloud::cloud_arrow::stream_bytes(&objects);
1356 return Ok(Answer::Load(Box::new(LoadAnswer::Sized(
1357 loading::PendingDownload::Arrow {
1358 url: url.clone(),
1359 objects,
1360 size: Some(size),
1361 options,
1362 },
1363 ))));
1364 }
1365 let size = match &pending {
1366 #[cfg(feature = "http")]
1367 loading::PendingDownload::Http { url, .. } => {
1368 Self::fetch_remote_size_http(url).map_err(|gone| gone.message)?
1371 }
1372 #[cfg(feature = "cloud")]
1373 loading::PendingDownload::S3 { url, .. }
1374 | loading::PendingDownload::Gcs { url, .. }
1375 | loading::PendingDownload::Azure { url, .. } => {
1376 Self::fetch_remote_size_cloud(url, &cloud, &runtime).unwrap_or(None)
1377 }
1378 #[cfg(feature = "cloud")]
1379 loading::PendingDownload::Arrow { size, .. } => *size,
1380 };
1381 Ok(Answer::Load(Box::new(LoadAnswer::Sized(
1382 pending.with_size(size),
1383 ))))
1384 });
1385 }
1386 #[cfg(any(feature = "http", feature = "cloud"))]
1387 Step::Download { pending, writer } => {
1388 let status = match &pending {
1391 #[cfg(feature = "http")]
1392 loading::PendingDownload::Http { .. } => "Downloading...",
1393 #[cfg(feature = "cloud")]
1394 loading::PendingDownload::S3 { .. } => "Downloading from S3...",
1395 #[cfg(feature = "cloud")]
1396 loading::PendingDownload::Gcs { .. } => "Downloading from GCS...",
1397 #[cfg(feature = "cloud")]
1398 loading::PendingDownload::Azure { .. } => "Downloading from Azure...",
1399 #[cfg(feature = "cloud")]
1400 loading::PendingDownload::Arrow { url, .. } => {
1401 match source::input_source(Path::new(url)) {
1402 source::InputSource::Gcs(_) => "Downloading from GCS...",
1403 source::InputSource::Azure(_) => "Downloading from Azure...",
1404 _ => "Downloading from S3...",
1405 }
1406 }
1407 };
1408 let sized = pending
1410 .parts()
1411 .1
1412 .filter(|_| status == "Downloading...")
1413 .map(|size| format!("Downloading {}...", crate::numfmt::bytes(size)));
1414 let status = sized.as_deref().unwrap_or(status);
1415 self.spawn_job(job, Some(status), move |_| {
1416 let (url, _, options) = pending.parts();
1417 let fetched = match &pending {
1418 #[cfg(feature = "http")]
1419 loading::PendingDownload::Http { .. } => {
1420 let ext = source::download_suffix(url);
1421 let limit = options
1424 .download_unasked
1425 .filter(|_| pending.parts().1.is_none())
1426 .map(|unasked| unasked.limit);
1427 Self::download_http_to_temp(
1428 url,
1429 options.temp_dir.as_deref(),
1430 ext.as_deref(),
1431 limit,
1432 &writer,
1433 )
1434 .map(|file| (file, options.clone()))
1435 }
1436 #[cfg(feature = "cloud")]
1437 loading::PendingDownload::S3 { .. }
1438 | loading::PendingDownload::Gcs { .. }
1439 | loading::PendingDownload::Azure { .. } => {
1440 Self::download_cloud_to_temp(url, &cloud, options, &runtime, &writer)
1441 .map(|file| (file, options.clone()))
1442 }
1443 #[cfg(feature = "cloud")]
1445 loading::PendingDownload::Arrow { objects, .. } => {
1446 crate::cloud::cloud_arrow::download(
1447 objects, options, &cloud, &runtime, &writer,
1448 )
1449 .map(|(file, parts)| {
1450 let options = OpenOptions {
1451 format: Some(FileFormat::Arrow),
1452 hive: false,
1453 arrow_parts: Some(Arc::new(parts)),
1454 ..options.clone()
1455 };
1456 (file, options)
1457 })
1458 }
1459 };
1460 let (download, options) = match fetched {
1461 Err(e)
1462 if e.downcast_ref::<crate::cloud::download::PastLimit>()
1463 .is_some() =>
1464 {
1465 return Ok(Answer::Load(Box::new(LoadAnswer::PastLimit(pending))));
1466 }
1467 fetched => fetched.map_err(|e| {
1468 crate::error_display::user_message_from_report(&e, None)
1469 })?,
1470 };
1471 Ok(Answer::Load(Box::new(LoadAnswer::Downloaded {
1472 download,
1473 options,
1474 })))
1475 });
1476 }
1477 Step::Spool {
1478 options,
1479 writer,
1480 read,
1481 } => {
1482 let piped = self.pipes.stdin_reader.take();
1485 let stdout = self.pipes.stdout_pass.take();
1486 self.spawn_job(job, Some("Reading stdin..."), move |_| {
1487 let open =
1488 move || -> crate::cloud::download::Opened<Box<dyn std::io::Read + Send>> {
1489 Ok((piped.unwrap_or_else(|| Box::new(std::io::stdin())), None))
1490 };
1491 let (download, options) = if options.follow
1494 || options.tee.is_some()
1495 || crate::loading::stdin::may_read_as_it_arrives(&options)
1496 {
1497 match crate::loading::follow::spool(open, options, &writer, &read, stdout)?
1498 {
1499 (crate::loading::follow::Spooled::Temp(download), options) => {
1500 (download, options)
1501 }
1502 (crate::loading::follow::Spooled::Kept(file), options) => {
1503 return Ok(Answer::Load(Box::new(LoadAnswer::Recorded {
1504 file,
1505 options,
1506 })));
1507 }
1508 }
1509 } else {
1510 crate::loading::stdin::spool(open, options, &writer, &read)?
1511 };
1512 Ok(Answer::Load(Box::new(LoadAnswer::Spooled {
1513 download,
1514 options,
1515 })))
1516 });
1517 }
1518 Step::FetchSpec {
1519 url,
1520 options,
1521 writer,
1522 } => {
1523 self.spawn_job(job, Some("Reading spec..."), move |_| {
1524 #[cfg(any(feature = "http", feature = "cloud"))]
1525 let fetched = crate::cloud::remote_model::fetch_small(
1526 &url,
1527 crate::formats::MAX_SPEC_BYTES,
1528 &cloud,
1529 &runtime,
1530 &|| writer.stopped(),
1531 );
1532 #[cfg(not(any(feature = "http", feature = "cloud")))]
1533 let fetched: std::result::Result<Option<Vec<u8>>, String> = {
1534 let _ = &writer;
1535 Err(crate::error_display::file_message(
1536 &url,
1537 "this build reads no URLs",
1538 ))
1539 };
1540 let bytes = fetched
1542 .map_err(|message| crate::logging::redact(&message, &[]))?
1543 .ok_or_else(|| {
1544 crate::logging::redact(
1545 &crate::error_display::file_message(
1546 &url,
1547 &format!(
1548 "a format spec is at most {}",
1549 crate::formats::MAX_SPEC_SAID
1550 ),
1551 ),
1552 &[],
1553 )
1554 })?;
1555 let spec = crate::formats::Spec::from_bytes(&bytes, &url)
1556 .map_err(|e| crate::logging::redact(&e.to_string(), &[]))?;
1557 Ok(Answer::Load(Box::new(LoadAnswer::SpecFetched {
1558 spec: Arc::new(spec),
1559 options,
1560 })))
1561 });
1562 }
1563 Step::DecompressRecords {
1564 file,
1565 path,
1566 choice,
1567 options,
1568 writer,
1569 } => {
1570 self.spawn_job(job, Some("Decompressing..."), move |_| {
1571 let failed = |e: color_eyre::Report| {
1572 crate::error_display::user_message_from_report(&e, Some(path.as_path()))
1573 };
1574 let compression = options
1575 .compression
1576 .or_else(|| CompressionFormat::from_extension(&file))
1577 .ok_or_else(|| format!("{} is not compressed", path.display()))?;
1578 let temp_dir = options.temp_dir.clone().unwrap_or_else(std::env::temp_dir);
1579 let copy = crate::formats::readers::csv::decompress_to_copy(
1580 &file,
1581 compression,
1582 &temp_dir,
1583 &writer,
1584 )
1585 .map_err(failed)?;
1586 Ok(Answer::Load(Box::new(LoadAnswer::DecompressedRecords {
1587 copy,
1588 path,
1589 choice,
1590 options,
1591 })))
1592 });
1593 }
1594 Step::ReadRecords {
1595 copy,
1596 path,
1597 choice,
1598 options,
1599 } => {
1600 self.spawn_job(job, Some("Reading records..."), move |_| {
1601 let named = path
1602 .file_name()
1603 .map(|n| n.to_string_lossy().into_owned())
1604 .unwrap_or_default();
1605 let read = crate::formats::read(©, &named, choice)?;
1606 let lf = Arc::clone(&read.records).into_lazy().map_err(|e| {
1607 crate::error_display::user_message_from_report(
1608 &color_eyre::eyre::eyre!(e),
1609 Some(path.as_path()),
1610 )
1611 })?;
1612 Ok(Answer::Load(Box::new(LoadAnswer::Scanned {
1613 lf: Box::new(lf),
1614 path: Some(path),
1615 options: OpenOptions {
1616 format_read: Some(Arc::new(read)),
1617 ..options
1618 },
1619 })))
1620 });
1621 }
1622 Step::Decompress {
1623 file,
1624 path,
1625 options,
1626 writer,
1627 download,
1628 } => {
1629 let options = OpenOptions {
1631 format: options.format.or(Some(FileFormat::TEXT)),
1632 ..options
1633 };
1634 let formats = self.formats.clone();
1635 self.spawn_job(job, Some("Decompressing..."), move |_| {
1636 let failed = |e: color_eyre::Report| {
1637 crate::error_display::user_message_from_report(&e, Some(path.as_path()))
1638 };
1639 let options =
1640 Self::with_delimited_spec(&file, options, &formats).map_err(failed)?;
1641 let lines = options.delimited.is_none()
1642 && options.format.is_some_and(FileFormat::is_lines);
1643 let (state, opened) = if lines {
1644 let (read, opened) = crate::formats::readers::csv::from_lines_decompressed(
1645 &file, &options, &writer,
1646 )
1647 .map_err(failed)?;
1648 let state = DataTableState::from_read(read, &options).map_err(failed)?;
1649 (state, Some(opened))
1650 } else {
1651 let state = Self::decompressed_delimited_state(&file, &options, &writer)
1652 .map_err(failed)?;
1653 (state, None)
1654 };
1655 let mut open_notes = options
1656 .delimited
1657 .as_ref()
1658 .map(|read| read.notes())
1659 .unwrap_or_default();
1660 open_notes.extend(opened.iter().flat_map(|o| o.notes.iter().cloned()));
1661 let state = state.with_open(OpenFacts {
1662 fetched: Self::was_fetched(download.as_ref(), Some(&path)),
1663 download,
1664 open_notes,
1665 records: opened.and_then(|o| o.window),
1666 delimited: options.delimited.clone(),
1667 read_as: options.format,
1668 read_mode: options.format.and_then(|f| {
1670 f.read_mode(crate::Stored::Compressed {
1671 in_memory: options.decompress_in_memory,
1672 })
1673 }),
1674 ..Default::default()
1675 });
1676 Ok(Answer::Load(Box::new(LoadAnswer::SchemaRead {
1677 state: Box::new(state),
1678 path: Some(path),
1679 options,
1680 debug_label: Some("decompressed delimited".to_string()),
1681 })))
1682 });
1683 }
1684 Step::Convert {
1685 what,
1686 files,
1687 path,
1688 options,
1689 writer,
1690 read,
1691 } => {
1692 let formats = self.formats.clone();
1695 self.spawn_job(job, Some(what.status()), move |_| {
1696 let named = |e: color_eyre::Report| {
1697 crate::error_display::user_message_from_report(&e, path.as_deref())
1698 };
1699 let converted = match what {
1700 loading::Conversion::Streams => {
1701 let converted = crate::formats::ipc_stream::convert(
1702 &files,
1703 options.temp_dir.as_deref(),
1704 &writer,
1705 &read,
1706 )
1707 .map_err(named)?;
1708 loading::Converted::Streams {
1709 file: converted.file,
1710 parts: converted.parts,
1711 }
1712 }
1713 loading::Conversion::Text(format) => {
1714 let display = path.clone().unwrap_or_else(|| files[0].clone());
1715 let (converted, detail) = crate::formats::readers::convert(
1716 &crate::formats::readers::ConvertIn {
1717 files: &files,
1718 display: &display,
1719 format,
1720 options: &options,
1721 formats: &formats,
1722 writer: &writer,
1723 read: &read,
1724 },
1725 )
1726 .map_err(named)?;
1727 loading::Converted::Frame {
1728 files: converted.files,
1729 lf: Box::new(converted.lf),
1730 notes: converted.notes,
1731 other_tables: converted.other_tables,
1732 detail,
1733 }
1734 }
1735 };
1736 Ok(Answer::Load(Box::new(LoadAnswer::Converted {
1737 converted,
1738 path,
1739 options,
1740 })))
1741 });
1742 }
1743 Step::Scan {
1744 paths,
1745 options,
1746 display,
1747 status,
1748 } => {
1749 let formats = self.formats.clone();
1750 let path = display.or_else(|| paths.first().cloned());
1752 self.home_app.reads.scans += 1;
1753 self.spawn_job(job, Some(status), move |_| {
1754 Self::scan_for_open(&cloud, &formats, &paths, options, path)
1755 .map(|answer| Answer::Load(Box::new(answer)))
1756 });
1757 }
1758 Step::ReadSchema {
1759 lf,
1760 path,
1761 options,
1762 progress,
1763 made,
1764 } => {
1765 self.debug.schema_load = None;
1766 let report = crate::loading::measurements::OpenReport {
1767 progress,
1768 meter: Arc::new(crate::loading::measurements::Meter::default()),
1769 remembered: Some(self.cache.clone()),
1770 writes: self.cache_writes.clone(),
1771 };
1772 self.spawn_job(job, Some("Reading schema..."), move |_| {
1773 Self::read_schema_for_open(*lf, path, options, &cloud, &runtime, &report, made)
1774 .map(|answer| Answer::Load(Box::new(answer)))
1775 });
1776 }
1777 Step::Nothing
1778 | Step::Crash(_)
1779 | Step::Install(_)
1780 | Step::Failed(_)
1781 | Step::Tables(_)
1782 | Step::Hex(_) => {
1783 unreachable!("not a phase with a worker")
1784 }
1785 #[cfg(any(feature = "http", feature = "cloud"))]
1786 Step::Ask(_) => unreachable!("not a phase with a worker"),
1787 Step::AskRead(_) => unreachable!("not a phase with a worker"),
1788 }
1789 }
1790
1791 fn in_memory_confirmation_message(read: &loading::InMemory) -> String {
1794 let what = match read.files {
1795 1 => format!(
1796 "{}: {} reads",
1797 read.name
1798 .file_name()
1799 .map(|n| n.to_string_lossy().into_owned())
1800 .unwrap_or_else(|| read.name.display().to_string()),
1801 read.format.title()
1802 ),
1803 n => format!("{n} {} files read", read.format.title()),
1804 };
1805 format!(
1806 "{what} {} into memory before the table appears.\n\nRead it?",
1807 crate::numfmt::bytes(read.bytes)
1808 )
1809 }
1810
1811 #[cfg(any(feature = "http", feature = "cloud"))]
1813 fn download_confirmation_message(
1814 pending: &loading::PendingDownload,
1815 note: Option<&str>,
1816 ) -> String {
1817 let (url, size, options) = pending.parts();
1818 let size_str = size
1819 .map(crate::numfmt::bytes)
1820 .unwrap_or_else(|| "unknown".to_string());
1821 let dest_dir = options
1822 .temp_dir
1823 .as_deref()
1824 .map(|p| p.display().to_string())
1825 .unwrap_or_else(|| std::env::temp_dir().display().to_string());
1826 let note = note.map(|note| format!("{note}\n\n")).unwrap_or_default();
1827 let files = match pending.arrow_files() {
1829 Some((1, 0)) => "Arrow stream: converted as it downloads\n".to_string(),
1830 Some((streams, 0)) => {
1831 format!("Files: {streams} Arrow streams, converted as they download\n")
1832 }
1833 Some((streams, in_place)) => {
1834 let streams = match streams {
1835 1 => "1 Arrow stream, converted as it downloads".to_string(),
1836 n => format!("{n} Arrow streams, converted as they download"),
1837 };
1838 let in_place = match in_place {
1839 1 => "1 IPC file read in place".to_string(),
1840 n => format!("{n} IPC files read in place"),
1841 };
1842 format!("Files: {streams}; {in_place}\n")
1843 }
1844 None => String::new(),
1845 };
1846 format!(
1847 "{note}URL: {url}\n{files}File size: {size_str}\nDestination: {dest_dir} (temporary file)\n\nContinue with download?"
1848 )
1849 }
1850}
1851
1852impl App {
1854 fn hoist_partition_columns(
1855 lf: LazyFrame,
1856 schema: &Schema,
1857 partition_columns: &[String],
1858 drifts: bool,
1859 ) -> LazyFrame {
1860 hoist_partition_columns(lf, schema, partition_columns, drifts)
1861 }
1862
1863 pub(crate) fn schema_state_from_local_hive(
1869 path: Option<&Path>,
1870 options: &OpenOptions,
1871 report: &crate::loading::measurements::OpenReport,
1872 ) -> Option<(DataTableState, OpenFacts)> {
1873 if !options.single_spine_schema {
1874 return None;
1875 }
1876 let p = path.filter(|p| p.is_dir() && options.hive)?;
1877 dataset_files::open(Arc::new(dataset_files::LocalFiles::new(p)), options, report)
1878 }
1879
1880 #[cfg(feature = "cloud")]
1882 fn schema_state_from_cloud_hive(
1883 path: Option<&Path>,
1884 options: &OpenOptions,
1885 cloud: &crate::config::CloudConfig,
1886 runtime: &tokio::runtime::Handle,
1887 report: &crate::loading::measurements::OpenReport,
1888 ) -> Option<(DataTableState, OpenFacts)> {
1889 if !options.single_spine_schema || options.format == Some(FileFormat::Arrow) {
1891 return None;
1892 }
1893 let p = path.filter(|p| {
1896 let s = p.as_os_str().to_string_lossy();
1897 home::is_object_store_url(p) && (options.hive || source::is_prefix_or_glob(&s))
1898 })?;
1899
1900 let (full, cloud_opts, store) = Self::cloud_store_for(p, cloud, runtime).ok()?;
1901 let (_bucket, key) = Self::cloud_bucket_and_key(&full).ok()?;
1902 Self::schema_state_from_cloud_hive_with(
1903 full, key, store, cloud_opts, options, runtime, report,
1904 )
1905 }
1906
1907 #[cfg(feature = "cloud")]
1910 pub(crate) fn schema_state_from_cloud_hive_with(
1911 full: String,
1912 key: String,
1913 store: Arc<dyn object_store::ObjectStore>,
1914 cloud_opts: CloudOptions,
1915 options: &OpenOptions,
1916 runtime: &tokio::runtime::Handle,
1917 report: &crate::loading::measurements::OpenReport,
1918 ) -> Option<(DataTableState, OpenFacts)> {
1919 let pattern = full.contains('*').then(|| {
1924 globset::GlobBuilder::new(&key)
1925 .literal_separator(true)
1926 .build()
1927 .map(|g| g.compile_matcher())
1928 });
1929 let pattern = match pattern {
1930 Some(Err(_)) => return None,
1932 Some(Ok(matcher)) => Some(matcher),
1933 None => None,
1934 };
1935 let listed = cloud_hive::prefix_of_glob(&key).to_string();
1936 Self::schema_state_from_cloud_files(
1937 CloudTarget {
1938 full: &full,
1939 key: listed,
1940 pattern: pattern.as_ref(),
1941 },
1942 store,
1943 cloud_opts,
1944 options,
1945 runtime,
1946 report,
1947 )
1948 }
1949
1950 #[cfg(feature = "cloud")]
1955 fn schema_state_from_cloud_files(
1956 target: CloudTarget<'_>,
1957 store: Arc<dyn object_store::ObjectStore>,
1958 cloud_opts: CloudOptions,
1959 options: &OpenOptions,
1960 runtime: &tokio::runtime::Handle,
1961 report: &crate::loading::measurements::OpenReport,
1962 ) -> Option<(DataTableState, OpenFacts)> {
1963 let CloudTarget { full, key, pattern } = target;
1964 dataset_files::open(
1965 Arc::new(dataset_files::StoreFiles::new(
1966 full,
1967 key,
1968 pattern.cloned(),
1969 store,
1970 cloud_opts,
1971 runtime,
1972 )),
1973 options,
1974 report,
1975 )
1976 }
1977
1978 #[cfg(feature = "cloud")]
1983 pub(crate) fn record_cloud_object_facts(
1984 cache: Option<&crate::cache::CacheManager>,
1985 full: &str,
1986 footer: &cloud_hive::FileFooter,
1987 ) {
1988 let Some(cache) = cache else {
1989 return;
1990 };
1991 let columns: Vec<String> = footer
1992 .schema
1993 .iter_names()
1994 .map(|name| name.to_string())
1995 .collect();
1996 cache.record_dataset_facts(&[(
1997 PathBuf::from(full),
1998 crate::cache::DatasetFacts {
1999 mtime: std::time::SystemTime::now()
2000 .duration_since(std::time::UNIX_EPOCH)
2001 .map(|d| d.as_secs())
2002 .unwrap_or_default(),
2003 size: 0,
2004 rows: Some(footer.rows()),
2005 cols: Some(columns.len()),
2006 cols_sampled: false,
2007 columns,
2008 kind: Some(discover::EntryKind::File),
2009 classified_by: discover::CLASSIFIER_VERSION,
2010 cost: discover::Cost {
2011 row_groups: Some(footer.row_group_rows.len()),
2012 ..Default::default()
2013 },
2014 holds: Default::default(),
2015 },
2016 )]);
2017 }
2018
2019 fn schema_state_from_full_scan(
2022 mut lf: LazyFrame,
2023 path: Option<&Path>,
2024 options: &OpenOptions,
2025 ) -> Result<DataTableState> {
2026 let schema = lf
2027 .collect_schema()
2028 .map_err(color_eyre::eyre::Report::from)?;
2029 let partition_columns =
2030 match path.filter(|p| options.hive && (p.is_dir() || source::expands_as_glob(p))) {
2031 Some(p) => crate::formats::readers::hive::discover_hive_partition_columns(p)
2032 .into_iter()
2033 .filter(|c| schema.contains(c.as_str()))
2034 .collect::<Vec<_>>(),
2035 None => Vec::new(),
2036 };
2037 let lf = Self::hoist_partition_columns(lf, &schema, &partition_columns, false);
2038 let part_cols = (!partition_columns.is_empty()).then_some(partition_columns);
2039 DataTableState::from_schema_and_lazyframe(schema, lf, options, part_cols)
2040 }
2041
2042 fn build_schema_state(
2045 lf: LazyFrame,
2046 path: Option<&Path>,
2047 options: &OpenOptions,
2048 cloud: &crate::config::CloudConfig,
2049 runtime: &tokio::runtime::Handle,
2050 report: &crate::loading::measurements::OpenReport,
2051 ) -> Result<(DataTableState, OpenFacts, String)> {
2052 let (state, mut facts, label) =
2055 Self::schema_state_by_route(lf, path, options, cloud, runtime, report)?;
2056 let names_look_like_data = options.format.and_then(FileFormat::separator).is_some()
2061 && options.has_header != Some(false)
2062 && !crate::formats::schema_union::names_are_names(
2063 &state
2064 .schema()
2065 .iter_names()
2066 .map(|n| n.to_string())
2067 .collect::<Vec<_>>(),
2068 );
2069 facts.open_notes = crate::notes::from_the_open(
2070 &options.left_out,
2071 options.read_as_plain_files_of,
2072 options.files_disagree,
2073 names_look_like_data,
2074 );
2075 facts.not_the_table = options.read_as_plain_files_of;
2078 if let Some(splits) = &options.splits {
2079 facts.other_tables = splits.others.clone();
2080 facts
2081 .open_notes
2082 .extend(crate::notes::map_caches(splits.caches));
2083 }
2084 if let Some(read) = &options.format_read {
2085 facts.open_notes.extend(read.notes());
2086 facts.format_read = Some(read.clone());
2087 }
2088 if let Some(opened) = &options.opened {
2089 facts.records = opened.window.clone();
2090 facts.detail = opened.detail.clone();
2091 facts.other_tables = opened.other_tables.clone();
2092 facts.open_notes.extend(opened.notes.iter().cloned());
2093 facts.units = opened.units.clone();
2094 facts.indexing = opened.indexing.clone();
2095 facts.numbering = opened.numbering.clone();
2096 }
2097 if let Some(sqlite) = &options.sqlite {
2098 facts.pushdown = Some(sqlite.pushdown.clone());
2099 facts.hold = sqlite.hold.lock().ok().and_then(|mut hold| hold.take());
2100 facts.other_tables = sqlite.other_tables.clone();
2101 }
2102 if let Some(read) = &options.delimited {
2103 facts.open_notes.extend(read.notes());
2104 facts.delimited = Some(read.clone());
2105 }
2106 facts.open_notes.extend(options.read_notes.iter().cloned());
2107 facts.typing = options.typing.clone();
2108 facts.read_mode = options.read_mode;
2109 facts.read_as = options.format;
2110 facts.remote_source = match &options.arrow_parts {
2114 Some(parts) => parts.iter().any(|part| {
2115 matches!(part, crate::formats::ipc_stream::Part::InPlace(p) if source::is_remote_url(p))
2116 }),
2117 None => path.is_some_and(source::scans_in_place),
2118 };
2119 if options.hive
2123 && facts.remote_files.is_none()
2124 && options.format.is_none_or(|f| f == FileFormat::Parquet)
2125 && let Some(dir) = path.filter(|p| !source::is_remote_url(p) && p.is_dir())
2126 {
2127 facts.parquet_count_dir = Some(dir.to_path_buf());
2128 }
2129 Ok((state, facts, label))
2130 }
2131
2132 #[cfg(feature = "cloud")]
2137 pub(crate) fn scan_cloud_prefix(
2138 url: &str,
2139 cloud_opts: CloudOptions,
2140 format: FileFormat,
2141 glob: bool,
2142 options: &OpenOptions,
2143 ) -> Option<Result<LazyFrame>> {
2144 if !format.reads_bucket_prefix() {
2147 return None;
2148 }
2149 let pl_path = if url.ends_with('/') && !url.contains('*') {
2152 PlRefPath::new(format!("{url}**/*.*").as_str())
2153 } else {
2154 PlRefPath::new(url)
2155 };
2156 let scan = crate::formats::readers::of(format).bucket_scan?;
2157 Some(scan(crate::formats::readers::BucketIn {
2158 url,
2159 path: pl_path,
2160 cloud: cloud_opts,
2161 glob,
2162 options,
2163 format,
2164 }))
2165 }
2166
2167 #[cfg(feature = "cloud")]
2170 fn cloud_glob_format(url: &str, options: &OpenOptions) -> Option<FileFormat> {
2171 options
2172 .format
2173 .or_else(|| {
2174 url.contains('*')
2175 .then(|| FileFormat::from_path(Path::new(url)))
2176 .flatten()
2177 })
2178 .filter(|f| *f != FileFormat::Parquet)
2179 }
2180
2181 #[cfg(feature = "cloud")]
2184 fn resolve_cloud_url(
2185 path: &Path,
2186 cloud: &crate::config::CloudConfig,
2187 ) -> Result<(String, CloudOptions)> {
2188 let text = path.to_string_lossy();
2189 let resolved = crate::cloud::cloud_sources::resolve_for_open(&text, cloud)
2190 .map_err(|e| color_eyre::eyre::eyre!(e))?;
2191 use object_store::azure::AzureConfigKey;
2192 use polars::io::cloud::GoogleConfigKey;
2193 let gcs_agent = (
2194 GoogleConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
2195 crate::cloud::user_agent::get(),
2196 );
2197 let options = match resolved.kind {
2198 crate::cloud::source::ProviderKind::S3 => Self::build_s3_cloud_options(&resolved.s3),
2199 crate::cloud::source::ProviderKind::Gcs
2200 if resolved.signing == crate::cloud::cloud_sources::Signing::Unsigned =>
2201 {
2202 CloudOptions::default()
2203 .with_gcp([(GoogleConfigKey::SkipSignature, "true".into()), gcs_agent])
2204 }
2205 crate::cloud::source::ProviderKind::Gcs => match &resolved.gcloud {
2206 Some((configuration, _)) => CloudOptions::default()
2208 .with_gcp([gcs_agent])
2209 .with_credential_provider(Some(crate::cloud::gcloud::polars_provider(
2210 configuration,
2211 ))),
2212 None => match &resolved.google_credentials {
2213 Some(file) => CloudOptions::default().with_gcp([
2214 (
2215 GoogleConfigKey::ApplicationCredentials,
2216 file.to_string_lossy().into_owned(),
2217 ),
2218 gcs_agent,
2219 ]),
2220 None => CloudOptions::default().with_gcp([gcs_agent]),
2221 },
2222 },
2223 crate::cloud::source::ProviderKind::Azure => {
2224 let (account, _, _) = source::azure_parts(&resolved.url)
2225 .ok_or_else(|| color_eyre::eyre::eyre!("not an Azure URL"))?;
2226 let mut azure = crate::cloud::azure::polars_options(&account, &resolved.azure);
2227 azure.push((
2228 AzureConfigKey::Client(crate::cloud::user_agent::CLIENT_KEY),
2229 crate::cloud::user_agent::get(),
2230 ));
2231 CloudOptions::default().with_azure(azure)
2232 }
2233 };
2234 Ok((resolved.url, options))
2235 }
2236
2237 #[cfg(feature = "cloud")]
2240 pub(crate) fn cloud_store_for(
2241 path: &Path,
2242 cloud: &crate::config::CloudConfig,
2243 runtime: &tokio::runtime::Handle,
2244 ) -> Result<(String, CloudOptions, Arc<dyn object_store::ObjectStore>)> {
2245 let (full, cloud_opts) = Self::resolve_cloud_url(path, cloud)?;
2246 let store = Self::polars_object_store(&full, &cloud_opts, runtime)?;
2247 Ok((full, cloud_opts, store))
2248 }
2249
2250 #[cfg(feature = "cloud")]
2255 fn schema_state_from_cloud_object(
2256 path: &Path,
2257 options: &OpenOptions,
2258 cloud: &crate::config::CloudConfig,
2259 runtime: &tokio::runtime::Handle,
2260 report: &crate::loading::measurements::OpenReport,
2261 ) -> Result<(DataTableState, OpenFacts)> {
2262 let (full, cloud_opts, store) = Self::cloud_store_for(path, cloud, runtime)?;
2263 let (_bucket, key) = Self::cloud_bucket_and_key(&full)?;
2264 if key.is_empty() {
2265 return Err(color_eyre::eyre::eyre!("a bucket, not an object"));
2266 }
2267 let meter = report.meter.clone();
2268 let (footer, etag) = wait_on_runtime(runtime, async move {
2269 cloud_hive::footer_of_cloud_parquet(store, &key, &meter).await
2270 })
2271 .ok_or_else(|| color_eyre::eyre::eyre!("cancelled"))??;
2272 let args = ScanArgsParquet {
2273 schema: Some(footer.schema.clone()),
2274 cloud_options: Some(cloud_opts),
2275 hive_options: polars::io::HiveOptions::default(),
2276 glob: false,
2277 ..Default::default()
2278 };
2279 let lf = LazyFrame::scan_parquet(PlRefPath::new(full.as_str()), args)?;
2280 let state =
2281 DataTableState::from_schema_and_lazyframe(footer.schema.clone(), lf, options, None)?;
2282 Self::record_cloud_object_facts(report.remembered.as_ref(), &full, &footer);
2284 let column_bytes =
2285 crate::formats::schema_union::column_bytes_per_row(&[Some(footer.clone())]);
2286 let facts = OpenFacts {
2287 remote_objects: vec![crate::cloud::local_copy::RemoteObject {
2288 url: full,
2289 size: footer.file_bytes as u64,
2290 etag,
2291 }],
2292 row_groups: vec![footer.row_group_rows],
2293 column_bytes,
2294 ..Default::default()
2295 };
2296 Ok((state, facts))
2297 }
2298
2299 fn schema_state_by_route(
2302 lf: LazyFrame,
2303 path: Option<&Path>,
2304 options: &OpenOptions,
2305 cloud: &crate::config::CloudConfig,
2306 runtime: &tokio::runtime::Handle,
2307 report: &crate::loading::measurements::OpenReport,
2308 ) -> Result<(DataTableState, OpenFacts, String)> {
2309 #[cfg(not(feature = "cloud"))]
2310 let _ = (cloud, runtime);
2311
2312 let attempt = |report: &crate::loading::measurements::OpenReport| {
2318 crate::loading::measurements::OpenReport {
2319 progress: report.progress.clone(),
2320 meter: Arc::new(crate::loading::measurements::Meter::default()),
2321 remembered: report.remembered.clone(),
2322 writes: report.writes.clone(),
2323 }
2324 };
2325
2326 let local = attempt(report);
2327 if let Some((state, facts)) = Self::schema_state_from_local_hive(path, options, &local) {
2328 let facts = OpenFacts {
2329 measurements: local.meter,
2330 ..facts
2331 };
2332 return Ok((state, facts, "one-file (local)".to_string()));
2333 }
2334 if report.progress.is_cancelled() {
2336 return Err(color_eyre::eyre::eyre!("cancelled"));
2337 }
2338 #[cfg(feature = "cloud")]
2339 let cloud_hive_attempt = attempt(report);
2340 #[cfg(feature = "cloud")]
2341 if let Some((state, facts)) =
2342 Self::schema_state_from_cloud_hive(path, options, cloud, runtime, &cloud_hive_attempt)
2343 {
2344 let facts = OpenFacts {
2345 measurements: cloud_hive_attempt.meter,
2346 ..facts
2347 };
2348 return Ok((state, facts, "one-file (cloud)".to_string()));
2349 }
2350 #[cfg(feature = "cloud")]
2351 if let Some(p) = path.filter(|p| {
2352 source::scans_in_place(p)
2353 && !options.hive
2354 && !source::is_prefix_or_glob(&p.to_string_lossy())
2355 }) {
2356 let object = attempt(report);
2357 match Self::schema_state_from_cloud_object(p, options, cloud, runtime, &object) {
2358 Ok((state, facts)) => {
2359 let facts = OpenFacts {
2360 measurements: object.meter,
2361 ..facts
2362 };
2363 return Ok((state, facts, "footer (cloud)".to_string()));
2364 }
2365 Err(e) => {
2367 return Self::schema_state_from_full_scan(lf, path, options).map(|state| {
2370 (
2371 state,
2372 OpenFacts::default(),
2373 format!("full scan (cloud footer: {e})"),
2374 )
2375 });
2376 }
2377 }
2378 }
2379 Self::schema_state_from_full_scan(lf, path, options)
2380 .map(|state| (state, OpenFacts::default(), "full scan".to_string()))
2381 }
2382
2383 fn hugging_face_split(
2387 dir: &Path,
2388 format: FileFormat,
2389 files: Vec<PathBuf>,
2390 options: &OpenOptions,
2391 report: &mut ReadReport,
2392 ) -> Result<Vec<PathBuf>> {
2393 let metadata = || {
2394 ["dataset_info.json", "state.json"]
2395 .iter()
2396 .any(|name| dir.join(name).is_file())
2397 };
2398 if format != FileFormat::Arrow || !metadata() {
2399 return Ok(files);
2400 }
2401 let names: Vec<&str> = files
2402 .iter()
2403 .map(|f| f.file_name().and_then(|n| n.to_str()).unwrap_or_default())
2404 .collect();
2405 let (chosen, splits) = crate::formats::hf_splits::choose(&names, options.table.as_deref())
2406 .map_err(|e| color_eyre::eyre::eyre!("{}: {e}", dir.display()))?;
2407 report.splits = Some(Arc::new(splits));
2408 Ok(chosen.into_iter().map(|i| files[i].clone()).collect())
2409 }
2410
2411 fn one_table(path: Option<&Path>, format: Option<FileFormat>) -> color_eyre::Report {
2413 match path {
2414 Some(path) => crate::error_display::FileError::new(path, cli::one_table(format)).into(),
2415 None => color_eyre::eyre::eyre!(cli::one_table(format)),
2416 }
2417 }
2418
2419 #[cfg(feature = "cloud")]
2421 fn cloud_scan_failed(url: &str, e: &polars::prelude::PolarsError) -> color_eyre::Report {
2422 let said = crate::error_display::user_message_from_polars(e);
2423 let (first, rest) = said.split_once('\n').unwrap_or((&said, ""));
2424 let first = first.trim_end().trim_end_matches('.');
2425 let what = format!("could not read it: {first}. Check the credentials and the URL.");
2426 let what = match rest {
2427 "" => what,
2428 rest => format!("{what}\n{rest}"),
2429 };
2430 crate::error_display::FileError::new(Path::new(url), what).into()
2431 }
2432
2433 fn scan_arrow_parts(
2437 cloud: &crate::config::CloudConfig,
2438 converted: Option<&PathBuf>,
2439 parts: &[crate::formats::ipc_stream::Part],
2440 ) -> Result<LazyFrame> {
2441 use crate::formats::ipc_stream::Part;
2442 #[cfg(not(feature = "cloud"))]
2443 let _ = cloud;
2444 let scan = |path: &Path| -> Result<LazyFrame> {
2445 #[cfg(feature = "cloud")]
2446 if source::is_remote_url(path) {
2447 let (url, cloud_options) = Self::resolve_cloud_url(path, cloud)?;
2448 let args = polars::prelude::UnifiedScanArgs {
2449 cloud_options: Some(cloud_options),
2450 ..Default::default()
2451 };
2452 return Ok(LazyFrame::scan_ipc(
2453 PlRefPath::new(url.as_str()),
2454 Default::default(),
2455 args,
2456 )?);
2457 }
2458 let args = polars::prelude::UnifiedScanArgs {
2461 glob: source::expands_as_glob(path),
2462 ..Default::default()
2463 };
2464 Ok(LazyFrame::scan_ipc(
2465 polars::prelude::PlRefPath::try_from_path(path)?,
2466 Default::default(),
2467 args,
2468 )?)
2469 };
2470 let streams = |offset: u64, rows: u64| -> Result<LazyFrame> {
2471 let file = converted
2472 .ok_or_else(|| color_eyre::eyre::eyre!("No converted Arrow file to read."))?;
2473 let lf = scan(file)?;
2474 let whole = offset == 0
2476 && parts
2477 .iter()
2478 .all(|part| matches!(part, Part::Converted { .. }));
2479 Ok(if whole {
2480 lf
2481 } else {
2482 lf.slice(offset as i64, rows as polars::prelude::IdxSize)
2483 })
2484 };
2485 let mut frames = Vec::new();
2486 let mut run: Option<(u64, u64)> = None;
2487 for part in parts {
2488 match part {
2489 Part::Converted { offset, rows, .. } => {
2490 run = Some(match run {
2491 Some((start, n)) if start + n == *offset => (start, n + rows),
2492 Some((start, n)) => {
2493 frames.push(streams(start, n)?);
2494 (*offset, *rows)
2495 }
2496 None => (*offset, *rows),
2497 });
2498 }
2499 Part::InPlace(path) => {
2500 if let Some((start, n)) = run.take() {
2501 frames.push(streams(start, n)?);
2502 }
2503 frames.push(scan(path)?);
2504 }
2505 }
2506 }
2507 if let Some((start, n)) = run {
2508 frames.push(streams(start, n)?);
2509 }
2510 match frames.len() {
2511 0 => Err(color_eyre::eyre::eyre!("No Arrow files to read.")),
2512 1 => Ok(frames.remove(0)),
2513 _ => Ok(polars::prelude::concat(
2514 frames.as_slice(),
2515 crate::formats::readers::polars::union_of_files(),
2516 )?),
2517 }
2518 }
2519
2520 fn dataset_dict_split(
2524 dir: &Path,
2525 splits: &[String],
2526 options: &OpenOptions,
2527 report: &mut ReadReport,
2528 formats: &crate::formats::Registry,
2529 ) -> Result<Scan> {
2530 let listed: Vec<&str> = splits.iter().map(String::as_str).collect();
2531 let mut picked = crate::formats::hf_splits::pick(&listed, options.table.as_deref())
2532 .map_err(|e| color_eyre::eyre::eyre!("{}: {e}", dir.display()))?;
2533 let split = dir.join(picked.split.as_deref().unwrap_or_default());
2534 let inner = OpenOptions {
2535 table: None,
2536 splits: None,
2537 ..options.clone()
2538 };
2539 let scan = Self::build_local_lazyframe(&[split], &inner, report, formats)?;
2540 picked.caches = report.splits.as_ref().map_or(0, |inner| inner.caches);
2542 report.splits = Some(Arc::new(picked));
2543 Ok(scan)
2544 }
2545
2546 fn files_disagree(
2550 files: &[PathBuf],
2551 options: &OpenOptions,
2552 found: FileFormat,
2553 ) -> crate::formats::schema_union::Disagreement {
2554 let format = options.format.unwrap_or(found);
2557 if format == FileFormat::Parquet {
2558 return Default::default();
2559 }
2560 if options.null_values.is_some() {
2563 return Default::default();
2564 }
2565 crate::formats::schema_union::sample_files(files, format, &Self::read_as(options))
2566 .disagreement()
2567 }
2568
2569 pub(crate) fn read_as(options: &OpenOptions) -> crate::formats::schema_union::ReadAs {
2574 crate::formats::schema_union::ReadAs {
2575 delimiter: options.delimiter,
2576 has_header: options.has_header,
2577 skip_rows: options.skip_rows,
2578 skip_lines: options.skip_lines,
2579 infer_schema_length: options.infer_schema_length,
2580 ignore_errors: options.ignore_errors,
2581 try_parse_dates: options.csv_try_parse_dates(),
2582 comment_char: options.comment_char.clone(),
2583 header_rows: options.header_rows.clone(),
2584 header_join: options.header_join.clone(),
2585 }
2586 }
2587
2588 pub(crate) fn refuse_spec_in_place(path: &Path, options: &OpenOptions) -> Result<()> {
2592 if source::is_remote_url(path)
2593 && (options.spec_file.is_some() || options.spec_name.is_some())
2594 {
2595 return Err(crate::error_display::FileError::new(
2596 path,
2597 "a format spec reads one remote object at a time, not a prefix or a glob; name the object",
2598 )
2599 .into());
2600 }
2601 Ok(())
2602 }
2603
2604 pub(crate) fn build_lazyframe_from_paths_with(
2609 cloud: &crate::config::CloudConfig,
2610 paths: &[PathBuf],
2611 options: &OpenOptions,
2612 report: &mut ReadReport,
2613 formats: &crate::formats::Registry,
2614 ) -> Result<Scan> {
2615 if let Some(parts) = &options.arrow_parts {
2618 if options.table.is_some() && options.splits.is_none() {
2619 let named = parts.first().map(|part| match part {
2621 crate::formats::ipc_stream::Part::InPlace(path) => path.as_path(),
2622 crate::formats::ipc_stream::Part::Converted { source, .. } => source.as_path(),
2623 });
2624 return Err(Self::one_table(named, Some(FileFormat::Arrow)));
2625 }
2626 return Self::scan_arrow_parts(cloud, paths.first(), parts).map(Scan::from);
2627 }
2628 #[cfg(not(feature = "cloud"))]
2630 let _ = cloud;
2631 let path = &paths[0];
2632 Self::refuse_spec_in_place(path, options)?;
2633 match source::input_source(path) {
2634 source::InputSource::Http(_url) => {
2635 #[cfg(feature = "http")]
2636 {
2637 return Err(color_eyre::eyre::eyre!(
2638 "HTTP/HTTPS load is handled in the event loop; this path should not be reached."
2639 ));
2640 }
2641 #[cfg(not(feature = "http"))]
2642 {
2643 return Err(color_eyre::eyre::eyre!(
2644 "HTTP/HTTPS URLs are not supported in this build. Rebuild with default features."
2645 ));
2646 }
2647 }
2648 source::InputSource::S3(url) => {
2649 #[cfg(feature = "cloud")]
2650 {
2651 let (full, cloud_opts) =
2652 Self::resolve_cloud_url(Path::new(&format!("s3://{url}")), cloud)?;
2653 let is_glob = source::is_prefix_or_glob(&full);
2654 if let Some(format) = Self::cloud_glob_format(&full, options)
2657 && let Some(lf) = Self::scan_cloud_prefix(
2658 &full,
2659 cloud_opts.clone(),
2660 format,
2661 is_glob,
2662 options,
2663 )
2664 {
2665 return lf.map(Scan::from);
2666 }
2667 let pl_path = PlRefPath::new(full.as_str());
2668 let hive_options = if is_glob {
2669 polars::io::HiveOptions::new_enabled()
2670 } else {
2671 polars::io::HiveOptions::default()
2672 };
2673 let args = ScanArgsParquet {
2674 cloud_options: Some(cloud_opts),
2675 hive_options,
2676 glob: is_glob,
2677 ..Default::default()
2678 };
2679 let lf = LazyFrame::scan_parquet(pl_path, args)
2680 .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2681 return Ok(lf.into());
2684 }
2685 #[cfg(not(feature = "cloud"))]
2686 {
2687 let _ = url;
2688 return Err(color_eyre::eyre::eyre!(
2689 "S3 is not supported in this build. Rebuild with default features and set AWS credentials (e.g. AWS_ACCESS_KEY_ID, AWS_SECRET_ACCESS_KEY, AWS_REGION)."
2690 ));
2691 }
2692 }
2693 source::InputSource::Gcs(url) => {
2694 #[cfg(feature = "cloud")]
2695 {
2696 let (full, cloud_opts) =
2697 Self::resolve_cloud_url(Path::new(&format!("gs://{url}")), cloud)?;
2698 let is_glob = source::is_prefix_or_glob(&full);
2699 if let Some(format) = Self::cloud_glob_format(&full, options)
2702 && let Some(lf) = Self::scan_cloud_prefix(
2703 &full,
2704 cloud_opts.clone(),
2705 format,
2706 is_glob,
2707 options,
2708 )
2709 {
2710 return lf.map(Scan::from);
2711 }
2712 let pl_path = PlRefPath::new(full.as_str());
2713 let hive_options = if is_glob {
2714 polars::io::HiveOptions::new_enabled()
2715 } else {
2716 polars::io::HiveOptions::default()
2717 };
2718 let args = ScanArgsParquet {
2719 cloud_options: Some(cloud_opts),
2720 hive_options,
2721 glob: is_glob,
2722 ..Default::default()
2723 };
2724 let lf = LazyFrame::scan_parquet(pl_path, args)
2725 .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2726 return Ok(lf.into());
2727 }
2728 #[cfg(not(feature = "cloud"))]
2729 {
2730 let _ = url;
2731 return Err(color_eyre::eyre::eyre!(
2732 "GCS (gs://) is not supported in this build. Rebuild with default features."
2733 ));
2734 }
2735 }
2736 source::InputSource::Azure(url) => {
2737 #[cfg(feature = "cloud")]
2738 {
2739 let (full, cloud_opts) = Self::resolve_cloud_url(Path::new(&url), cloud)?;
2740 let is_glob = source::is_prefix_or_glob(&full);
2741 if let Some(format) = Self::cloud_glob_format(&full, options)
2744 && let Some(lf) = Self::scan_cloud_prefix(
2745 &full,
2746 cloud_opts.clone(),
2747 format,
2748 is_glob,
2749 options,
2750 )
2751 {
2752 return lf.map(Scan::from);
2753 }
2754 let args = ScanArgsParquet {
2755 cloud_options: Some(cloud_opts),
2756 hive_options: if is_glob {
2757 polars::io::HiveOptions::new_enabled()
2758 } else {
2759 polars::io::HiveOptions::default()
2760 },
2761 glob: is_glob,
2762 ..Default::default()
2763 };
2764 let lf = LazyFrame::scan_parquet(PlRefPath::new(full.as_str()), args)
2765 .map_err(|e| Self::cloud_scan_failed(&full, &e))?;
2766 return Ok(lf.into());
2767 }
2768 #[cfg(not(feature = "cloud"))]
2769 {
2770 let _ = url;
2771 return Err(color_eyre::eyre::eyre!(
2772 "Azure is not supported in this build. Rebuild with default features."
2773 ));
2774 }
2775 }
2776 source::InputSource::Local(_) => {}
2777 }
2778 Self::build_local_lazyframe(paths, options, report, formats)
2779 }
2780
2781 fn read_directory_files(
2784 files: &[PathBuf],
2785 options: &OpenOptions,
2786 found: FileFormat,
2787 report: &mut ReadReport,
2788 formats: &crate::formats::Registry,
2789 ) -> Result<Scan> {
2790 if options.delimited.is_none()
2791 && options.format.is_none()
2792 && found.separator().is_some()
2793 && let Some(first) = files
2794 .iter()
2795 .find(|f| !crate::formats::nul_tail::holds_nothing(f))
2796 && let Some(choice) = Self::delimited_spec_of(first, options, formats)?
2797 {
2798 let nested = OpenOptions {
2799 hive: false,
2800 format: Some(found),
2801 splits: report.splits.clone(),
2802 ..options.clone()
2803 };
2804 return Self::read_with_delimited_spec(files, &nested, report, formats, choice);
2805 }
2806 report.files_disagree = Self::files_disagree(files, options, found);
2808 let nested = OpenOptions {
2809 hive: false,
2810 format: Some(options.format.unwrap_or(found)),
2811 splits: report.splits.clone(),
2812 ..options.clone()
2813 };
2814 Self::build_local_lazyframe(files, &nested, report, formats)
2815 }
2816
2817 fn delimited_spec_of(
2819 file: &Path,
2820 options: &OpenOptions,
2821 formats: &crate::formats::Registry,
2822 ) -> Result<Option<crate::formats::Choice>> {
2823 let asked = crate::formats::Asked {
2824 compression: options.compression,
2825 text_only: true,
2826 ..Default::default()
2827 };
2828 match crate::formats::route(file, &asked, formats)
2829 .map_err(|e| color_eyre::eyre::eyre!(e))?
2830 {
2831 crate::formats::Route::Delimited(choice) => Ok(Some(choice)),
2832 _ => Ok(None),
2833 }
2834 }
2835
2836 fn read_with_delimited_spec(
2838 paths: &[PathBuf],
2839 options: &OpenOptions,
2840 report: &mut ReadReport,
2841 formats: &crate::formats::Registry,
2842 choice: crate::formats::Choice,
2843 ) -> Result<Scan> {
2844 let mut nested = options.clone();
2845 let Some(delimited) = choice.spec.delimited.clone() else {
2846 return Err(color_eyre::eyre::eyre!(
2847 "{} is not a delimited spec",
2848 choice.spec.name
2849 ));
2850 };
2851 delimited.apply(&mut nested);
2852 nested.delimited = Some(Arc::new(
2853 crate::formats::delimited_spec::DelimitedRead::chosen(
2854 choice.spec,
2855 choice.by,
2856 choice.also,
2857 ),
2858 ));
2859 Self::build_local_lazyframe(paths, &nested, report, formats)
2860 }
2861
2862 pub(crate) fn build_local_lazyframe(
2865 paths: &[PathBuf],
2866 options: &OpenOptions,
2867 report: &mut ReadReport,
2868 formats: &crate::formats::Registry,
2869 ) -> Result<Scan> {
2870 let path = &paths[0];
2871
2872 if options.hex
2874 && let [one] = paths
2875 && one.is_file()
2876 {
2877 return Ok(Scan::Hex {
2878 file: one.clone(),
2879 asked: true,
2880 });
2881 }
2882
2883 if let [pattern] = paths
2886 && !options.hive
2887 && options.delimited.is_none()
2888 && options.format.is_none()
2889 && source::expands_as_glob(pattern)
2890 {
2891 let files = crate::loading::local_glob::expand(pattern);
2892 if let Some(first) = files
2893 .iter()
2894 .find(|f| !crate::formats::nul_tail::holds_nothing(f))
2895 && let Some(choice) = Self::delimited_spec_of(first, options, formats)?
2896 {
2897 let format = FileFormat::from_path(first).filter(|f| f.separator().is_some());
2898 let nested = OpenOptions {
2899 format: format.or(Some(FileFormat::Csv)),
2900 ..options.clone()
2901 };
2902 return Self::read_with_delimited_spec(&files, &nested, report, formats, choice);
2903 }
2904 }
2905
2906 if !options.hive && options.delimited.is_none() {
2910 let asked = crate::formats::Asked {
2911 spec_file: options.spec_file.clone(),
2912 spec_name: options.spec_name.clone(),
2913 variant: options.table.clone(),
2914 spec: options.spec_fetched.clone(),
2915 builtin: options.format.is_some(),
2916 compression: options.compression,
2917 text_only: paths.len() > 1,
2918 };
2919 match crate::formats::route(path, &asked, formats)
2920 .map_err(|e| color_eyre::eyre::eyre!(e))?
2921 {
2922 crate::formats::Route::Elsewhere => {}
2923 crate::formats::Route::Delimited(choice) => {
2924 return Self::read_with_delimited_spec(paths, options, report, formats, choice);
2925 }
2926 crate::formats::Route::Read(read) => {
2927 let lf = Arc::clone(&read.records).into_lazy()?;
2928 report.format_read = Some(Arc::new(*read));
2929 return Ok(lf.into());
2930 }
2931 crate::formats::Route::Decompress(choice) => {
2932 return Ok(Scan::DecompressSpec {
2933 file: path.clone(),
2934 choice,
2935 });
2936 }
2937 }
2938 } else if options.hive && (options.spec_file.is_some() || options.spec_name.is_some()) {
2939 return Err(color_eyre::eyre::eyre!(
2940 "a format spec reads one file, or one directory of column files"
2941 ));
2942 }
2943
2944 if let Some(read) = &options.delimited
2946 && report.delimited.is_none()
2947 && path.is_file()
2948 {
2949 report.delimited = Some(Arc::new(crate::formats::delimited_spec::read_facts(
2950 read, paths, options,
2951 )?));
2952 }
2953
2954 if paths.len() == 1 && (options.hive || path.is_dir()) {
2957 let is_single_file = path.is_file();
2959 if !is_single_file {
2960 if path.is_dir()
2962 && let Some(splits) = crate::formats::hf_splits::dataset_dict(path)
2963 {
2964 return Self::dataset_dict_split(path, &splits, options, report, formats);
2965 }
2966 if path.is_dir() {
2967 match crate::home::discover::directory_format(path) {
2968 crate::home::discover::DirectoryFormat::One(FileFormat::Parquet, _) => {}
2970 crate::home::discover::DirectoryFormat::Deeper => {
2974 if let crate::home::discover::DirectoryFormat::One(found, files) =
2975 crate::home::discover::hive_leaf_format(path)
2976 && found != FileFormat::Parquet
2977 {
2978 let named = files
2980 .first()
2981 .and_then(|f| crate::home::discover::data_extension(f))
2982 .unwrap_or_else(|| format!("{found:?}").to_lowercase());
2983 return Err(color_eyre::eyre::eyre!(
2984 "{} is partitioned into key=value directories of .{} \
2985 files. datui reads hive partitioning for Parquet \
2986 only — open one partition instead.",
2987 path.display(),
2988 named
2989 ));
2990 }
2991 }
2992 crate::home::discover::DirectoryFormat::One(found, files) => {
2993 let format = options.format.unwrap_or(found);
2996 let files =
2997 Self::hugging_face_split(path, format, files, options, report)?;
2998 return Self::read_directory_files(
2999 &files, options, found, report, formats,
3000 );
3001 }
3002 crate::home::discover::DirectoryFormat::Mixed {
3003 format: found,
3004 files,
3005 passed_over,
3006 } => {
3007 let format = options.format.unwrap_or(found);
3010 let files =
3011 Self::hugging_face_split(path, format, files, options, report)?;
3012 let lf = Self::read_directory_files(
3013 &files, options, found, report, formats,
3014 )?;
3015 if !matches!(found, FileFormat::Safetensors | FileFormat::Gguf) {
3018 report.left_out = passed_over;
3019 }
3020 return Ok(lf);
3021 }
3022 }
3023 }
3024 let use_parquet_hive =
3025 path.is_dir() || path.as_os_str().to_string_lossy().contains(".parquet");
3026 if use_parquet_hive {
3027 return crate::formats::readers::hive::scan_parquet_hive(path).map(Scan::from);
3029 }
3030 return Err(color_eyre::eyre::eyre!(
3031 "With --hive use a directory or a glob pattern for Parquet (e.g. path/to/dir or path/**/*.parquet)"
3032 ));
3033 }
3034 }
3035
3036 let compressed = options
3040 .compression
3041 .or_else(|| CompressionFormat::from_extension(path))
3042 .is_some();
3043 let named = FileFormat::from_path(path).or_else(|| {
3046 compressed
3047 .then(|| FileFormat::from_path(Path::new(path.file_stem()?)))
3048 .flatten()
3049 .filter(|f| f.decompressed_once())
3050 });
3051 let mut effective_format = options
3052 .format
3053 .or_else(|| {
3055 named.filter(|f| !f.is_lines()).map(|f| {
3056 (!compressed)
3057 .then(|| crate::formats::readers::refined(path, f))
3058 .flatten()
3059 .unwrap_or(f)
3060 })
3061 })
3062 .or_else(|| {
3063 (path.extension().is_none()
3064 && crate::home::discover::is_parquet_key(&path.to_string_lossy()))
3065 .then_some(FileFormat::Parquet)
3066 })
3067 .or_else(|| crate::formats::readers::sniff_open(path, options.compression))
3069 .or(named);
3070 if effective_format.is_none()
3073 && let [file] = paths
3074 && file.is_file()
3075 {
3076 effective_format = crate::formats::lines::guess_file(file, options.compression)
3077 .map(|f| crate::formats::lines::as_asked(f, options));
3078 report.guessed = effective_format.is_some();
3079 }
3080 report.format = effective_format;
3081
3082 if options.table.is_some()
3085 && !effective_format.is_some_and(FileFormat::takes_table)
3086 && options.splits.is_none()
3087 {
3088 return Err(Self::one_table(Some(path), effective_format));
3089 }
3090
3091 if let [file] = paths
3095 && compressed
3096 && let Some(format) = effective_format.filter(|f| f.decompressed_once())
3097 {
3098 return Ok(Scan::Decompress {
3099 file: file.clone(),
3100 format,
3101 });
3102 }
3103
3104 let Some(format) = effective_format else {
3105 if !path.exists() {
3106 return Err(std::io::Error::new(
3107 std::io::ErrorKind::NotFound,
3108 format!("File not found: {}", path.display()),
3109 )
3110 .into());
3111 }
3112 if paths.len() == 1 && path.is_file() {
3114 return Ok(Scan::Hex {
3115 file: path.clone(),
3116 asked: false,
3117 });
3118 }
3119 return Err(color_eyre::eyre::eyre!(match paths.len() {
3120 1 => UNSUPPORTED.to_string(),
3121 _ => crate::formats::readers::many_files_refused(),
3122 }));
3123 };
3124 if paths.len() > 1 && !format.reads_many_files() {
3127 if !path.exists() {
3128 return Err(std::io::Error::new(
3129 std::io::ErrorKind::NotFound,
3130 format!("File not found: {}", path.display()),
3131 )
3132 .into());
3133 }
3134 return Err(color_eyre::eyre::eyre!(
3135 crate::formats::readers::many_files_refused()
3136 ));
3137 }
3138 let guessed;
3139 let options = if report.guessed {
3140 guessed = OpenOptions {
3141 format_guessed: true,
3142 ..options.clone()
3143 };
3144 &guessed
3145 } else {
3146 options
3147 };
3148 crate::formats::readers::scan(crate::formats::readers::ScanIn {
3149 format,
3150 paths,
3151 options,
3152 report,
3153 formats,
3154 })
3155 }
3156}