1pub(crate) mod counting;
23pub(crate) mod first_rows_trace;
24pub mod follow;
25pub(crate) mod local_glob;
26pub mod measurements;
27pub(crate) mod open_options;
28pub(crate) mod open_scan;
29pub(crate) mod scan;
30pub mod stdin;
31pub mod tee;
32pub(crate) mod unfinished;
33
34use std::path::{Path, PathBuf};
35use std::sync::Arc;
36use std::sync::atomic::{AtomicU64, Ordering};
37
38use polars::prelude::LazyFrame;
39
40use crate::cloud::download::TempDownload;
41use crate::formats::schema_union::FooterProgress;
42use crate::loading::unfinished::{Unfinished, Writer};
43use crate::table::DataTableState;
44use crate::{CompressionFormat, FileFormat, OpenOptions, cloud::source};
45
46use crate::app::jobs::Hold;
47#[cfg(any(feature = "http", feature = "cloud"))]
48use crate::app::jobs::Jobs;
49
50#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
52pub(crate) struct LoadId(u64);
53
54impl LoadId {
55 #[cfg(test)]
57 pub(crate) fn for_tests(n: u64) -> Self {
58 Self(n)
59 }
60}
61
62#[cfg(any(feature = "http", feature = "cloud"))]
64#[derive(Clone)]
65pub(crate) enum PendingDownload {
66 #[cfg(feature = "http")]
67 Http {
68 url: String,
69 size: Option<u64>,
70 options: OpenOptions,
71 },
72 #[cfg(feature = "cloud")]
73 S3 {
74 url: String,
75 size: Option<u64>,
76 options: OpenOptions,
77 },
78 #[cfg(feature = "cloud")]
79 Gcs {
80 url: String,
81 size: Option<u64>,
82 options: OpenOptions,
83 },
84 #[cfg(feature = "cloud")]
85 Azure {
86 url: String,
87 size: Option<u64>,
88 options: OpenOptions,
89 },
90 #[cfg(feature = "cloud")]
94 Arrow {
95 url: String,
96 objects: Vec<crate::cloud::cloud_arrow::Object>,
97 size: Option<u64>,
98 options: OpenOptions,
99 },
100}
101
102#[cfg(any(feature = "http", feature = "cloud"))]
103impl PendingDownload {
104 pub(crate) fn parts(&self) -> (&str, Option<u64>, &OpenOptions) {
107 match self {
108 #[cfg(feature = "http")]
109 PendingDownload::Http { url, size, options } => (url, *size, options),
110 #[cfg(feature = "cloud")]
111 PendingDownload::S3 { url, size, options } => (url, *size, options),
112 #[cfg(feature = "cloud")]
113 PendingDownload::Gcs { url, size, options } => (url, *size, options),
114 #[cfg(feature = "cloud")]
115 PendingDownload::Azure { url, size, options } => (url, *size, options),
116 #[cfg(feature = "cloud")]
117 PendingDownload::Arrow {
118 url, size, options, ..
119 } => (url, *size, options),
120 }
121 }
122
123 pub(crate) fn arrow_files(&self) -> Option<(usize, usize)> {
126 match self {
127 #[cfg(feature = "cloud")]
128 PendingDownload::Arrow { objects, .. } => {
129 let streams = objects.iter().filter(|o| o.stream).count();
130 Some((streams, objects.len() - streams))
131 }
132 _ => None,
133 }
134 }
135
136 pub(crate) fn asked(mut self) -> Self {
139 match &mut self {
140 #[cfg(feature = "http")]
141 PendingDownload::Http { options, .. } => options.download_unasked = None,
142 #[cfg(feature = "cloud")]
143 PendingDownload::S3 { options, .. }
144 | PendingDownload::Gcs { options, .. }
145 | PendingDownload::Azure { options, .. }
146 | PendingDownload::Arrow { options, .. } => options.download_unasked = None,
147 }
148 self
149 }
150
151 pub(crate) fn with_size(mut self, found: Option<u64>) -> Self {
153 match &mut self {
154 #[cfg(feature = "http")]
155 PendingDownload::Http { size, .. } => *size = found,
156 #[cfg(feature = "cloud")]
157 PendingDownload::S3 { size, .. } => *size = found,
158 #[cfg(feature = "cloud")]
159 PendingDownload::Gcs { size, .. } => *size = found,
160 #[cfg(feature = "cloud")]
161 PendingDownload::Azure { size, .. } => *size = found,
162 #[cfg(feature = "cloud")]
163 PendingDownload::Arrow { size, .. } => *size = found,
164 }
165 self
166 }
167}
168
169#[derive(Clone)]
172pub(crate) struct OpenRequest {
173 pub(crate) paths: Vec<PathBuf>,
174 pub(crate) options: OpenOptions,
175 pub(crate) size: u64,
177 pub(crate) recent: Option<PathBuf>,
180 pub(crate) shown: Option<PathBuf>,
183 pub(crate) warn_in_memory_above: Option<u64>,
186}
187
188impl OpenRequest {
189 pub(crate) fn named(
194 mut paths: Vec<PathBuf>,
195 mut options: OpenOptions,
196 formats: &crate::formats::Registry,
197 ) -> Self {
198 let first = paths[0].clone();
199 let piped = stdin::is_stdin(&first);
200 let local = !piped && matches!(source::input_source(&first), source::InputSource::Local(_));
201 let mut table = None;
202 if local
203 && paths.len() == 1
204 && let Some((db, name)) = crate::formats::members::split(&first)
205 {
206 table = Some(first.clone());
207 options.table = Some(name);
208 paths = vec![db];
209 } else if local
210 && first.is_file()
211 && let Some(name) = options.table.as_deref()
212 && crate::formats::members::holder(&first).is_some()
213 {
214 table = Some(crate::formats::members::place(&first, name));
215 } else if local
216 && paths.len() == 1
217 && options.table.is_none()
218 && let Some((dir, split)) = crate::formats::hf_splits::split_place(&first)
219 {
220 table = Some(first.clone());
222 options.table = Some(split);
223 paths = vec![dir];
224 } else if local
225 && first.is_dir()
226 && let Some(split) = options.table.as_deref()
227 && crate::formats::hf_splits::cache_splits(&first)
228 .iter()
229 .any(|s| s == split)
230 {
231 table = Some(first.join(split));
232 } else if local
233 && paths.len() == 1
234 && options.table.is_none()
235 && let Some((file, variant)) = crate::formats::members::split_variant(&first, formats)
236 {
237 table = Some(first.clone());
239 options.table = Some(variant);
240 paths = vec![file];
241 } else if local
242 && let Some(variant) = options.table.as_deref()
243 && first.is_file()
244 && formats.variants_of(&first).is_some()
245 {
246 table = Some(crate::formats::members::place(&first, variant));
247 }
248 let first = &paths[0];
249 let size = if local {
250 std::fs::metadata(first).map(|m| m.len()).unwrap_or(0)
251 } else {
252 0
253 };
254 let recent = table
256 .clone()
257 .or_else(|| (!piped && (!local || first.exists())).then(|| first.clone()));
258 Self {
259 paths,
260 options,
261 size,
262 recent,
263 shown: table,
264 warn_in_memory_above: None,
265 }
266 }
267}
268
269pub(crate) enum Phase {
272 Starting {
275 label: String,
276 percent: u16,
277 },
278 LookingAtPaths,
280 ReadingSpec,
282 LookingAtDirectory,
284 #[cfg(any(feature = "http", feature = "cloud"))]
286 ReadingHeaders,
287 #[cfg(any(feature = "http", feature = "cloud"))]
290 CheckingSize {
291 note: Option<&'static str>,
292 },
293 ConfirmingRead {
296 scan: Box<Scan>,
297 _hold: Option<Hold>,
298 },
299 #[cfg(any(feature = "http", feature = "cloud"))]
302 Confirming {
303 pending: Box<PendingDownload>,
304 note: Option<&'static str>,
305 _hold: Hold,
306 },
307 #[cfg(any(feature = "http", feature = "cloud"))]
308 Downloading,
309 Spooling {
311 read: Arc<AtomicU64>,
312 },
313 Decompressing,
314 DecompressingRecords,
316 ReadingRecords,
318 Converting {
321 what: Conversion,
322 read: Arc<AtomicU64>,
323 total: u64,
324 },
325 ScanningStrings,
327 CountingFooter,
330 Scanning {
332 downloaded: bool,
333 },
334 ReadingLines,
336 ReadingSchema,
338 FirstRows,
341}
342
343impl Phase {
344 pub(crate) fn label(&self) -> (&str, u16) {
347 match self {
348 Phase::Starting { label, percent } => (label, *percent),
349 Phase::LookingAtPaths => ("Scanning input", 10),
350 Phase::ReadingSpec => ("Reading spec", 5),
351 Phase::LookingAtDirectory => (crate::App::LOOKING_AT_A_DIRECTORY, 5),
352 #[cfg(any(feature = "http", feature = "cloud"))]
353 Phase::ReadingHeaders => ("Reading headers", 20),
354 #[cfg(any(feature = "http", feature = "cloud"))]
355 Phase::CheckingSize { .. } | Phase::Confirming { .. } => ("Checking size", 0),
356 #[cfg(any(feature = "http", feature = "cloud"))]
357 Phase::Downloading => ("Downloading", 20),
358 Phase::Spooling { .. } => ("Reading stdin", 5),
359 Phase::ConfirmingRead { .. } => ("Scanning input", 0),
360 Phase::Decompressing | Phase::DecompressingRecords => ("Decompressing", 30),
361 Phase::ReadingRecords => ("Reading records", 35),
362 Phase::Converting { what, read, total } => {
363 let done = read.load(Ordering::Relaxed).min(*total);
364 let share = (done * 20).checked_div(*total).unwrap_or(0);
366 (what.label(), 10 + share as u16)
367 }
368 Phase::ScanningStrings => ("Scanning string columns", 55),
369 Phase::CountingFooter => (COUNTING_FOOTER, 10),
370 Phase::Scanning { downloaded: false } => ("Scanning input", 10),
371 Phase::Scanning { downloaded: true } => ("Scanning", 30),
372 Phase::ReadingLines => (READING_LINES, 10),
373 Phase::ReadingSchema => ("Reading schema", 40),
374 Phase::FirstRows => ("Loading buffer", 70),
375 }
376 }
377
378 fn asks(&self) -> bool {
380 match self {
381 Phase::ConfirmingRead { .. } => true,
382 #[cfg(any(feature = "http", feature = "cloud"))]
383 Phase::Confirming { .. } => true,
384 _ => false,
385 }
386 }
387
388 fn builds_the_dataset(&self) -> bool {
390 match self {
391 Phase::ReadingSchema | Phase::Decompressing => true,
392 #[cfg(any(feature = "http", feature = "cloud"))]
393 Phase::ReadingHeaders => true,
394 _ => false,
395 }
396 }
397
398 fn starting(&self) -> bool {
401 matches!(
402 self,
403 Phase::Starting { .. } | Phase::LookingAtPaths | Phase::LookingAtDirectory
404 )
405 }
406}
407
408#[derive(Clone)]
411struct Fetched {
412 url: PathBuf,
413 file: TempDownload,
414 arrow: Option<KeptArrow>,
417}
418
419#[derive(Clone)]
420struct KeptArrow {
421 parts: Arc<Vec<crate::formats::ipc_stream::Part>>,
422 splits: Option<Arc<crate::formats::hf_splits::Splits>>,
423 table: Option<String>,
424}
425
426impl Fetched {
427 fn serves(&self, options: &OpenOptions) -> bool {
430 self.arrow
431 .as_ref()
432 .is_none_or(|arrow| arrow.table == options.table)
433 }
434}
435
436pub(crate) struct Load {
438 id: LoadId,
439 from_home: bool,
442 phase: Phase,
443 path: Option<PathBuf>,
446 size: u64,
447 paths: Option<Vec<PathBuf>>,
449 recent: Option<PathBuf>,
450 progress: Arc<FooterProgress>,
452 writer: Writer,
454 download: Option<Fetched>,
455 converted: Vec<TempDownload>,
458 warn_in_memory_above: Option<u64>,
460}
461
462pub(crate) struct Scan {
464 paths: Vec<PathBuf>,
465 options: OpenOptions,
466 display: Option<PathBuf>,
467}
468
469#[derive(Debug, Clone, PartialEq, Eq)]
471pub(crate) struct InMemory {
472 pub(crate) bytes: u64,
474 pub(crate) format: FileFormat,
475 pub(crate) files: usize,
476 pub(crate) name: PathBuf,
478}
479
480impl Load {
481 pub(crate) fn phase(&self) -> &Phase {
482 &self.phase
483 }
484
485 fn made(&self) -> Made {
488 Made {
489 download: self.download.as_ref().map(|fetched| fetched.file.clone()),
490 converted: self.converted.clone(),
491 ..Made::default()
492 }
493 }
494
495 pub(crate) fn path(&self) -> Option<&Path> {
497 self.path.as_deref()
498 }
499
500 pub(crate) fn size(&self) -> u64 {
503 match &self.phase {
504 Phase::Spooling { read } => read.load(Ordering::Relaxed),
505 _ => self.size,
506 }
507 }
508}
509
510pub(crate) enum Step {
512 Nothing,
514 Crash(String),
516 #[cfg(any(feature = "http", feature = "cloud"))]
519 ReadHeaders {
520 url: PathBuf,
521 format: FileFormat,
522 options: OpenOptions,
523 writer: Writer,
524 },
525 #[cfg(any(feature = "http", feature = "cloud"))]
527 Probe(PendingDownload),
528 #[cfg(any(feature = "http", feature = "cloud"))]
530 Ask(PendingDownload),
531 AskRead(InMemory),
533 #[cfg(any(feature = "http", feature = "cloud"))]
536 Download {
537 pending: PendingDownload,
538 writer: Writer,
539 },
540 FetchSpec {
542 url: PathBuf,
543 options: OpenOptions,
544 writer: Writer,
545 },
546 Spool {
548 options: OpenOptions,
549 writer: Writer,
550 read: Arc<AtomicU64>,
551 },
552 Decompress {
555 file: PathBuf,
556 path: PathBuf,
557 options: OpenOptions,
558 writer: Writer,
559 download: Option<TempDownload>,
561 },
562 DecompressRecords {
565 file: PathBuf,
566 path: PathBuf,
567 choice: crate::formats::Choice,
568 options: OpenOptions,
569 writer: Writer,
570 },
571 ReadRecords {
574 copy: PathBuf,
575 path: PathBuf,
576 choice: crate::formats::Choice,
577 options: OpenOptions,
578 },
579 Convert {
582 what: Conversion,
583 files: Vec<PathBuf>,
584 path: Option<PathBuf>,
585 options: OpenOptions,
586 writer: Writer,
587 read: Arc<AtomicU64>,
588 },
589 Scan {
592 paths: Vec<PathBuf>,
593 options: OpenOptions,
594 display: Option<PathBuf>,
595 status: &'static str,
596 },
597 ReadSchema {
599 lf: Box<LazyFrame>,
600 path: Option<PathBuf>,
601 options: OpenOptions,
602 progress: Arc<FooterProgress>,
603 made: Made,
605 },
606 Install(Box<Loaded>),
608 Failed(Failed),
610 Tables(Tables),
613 Hex(Hex),
616}
617
618#[derive(Debug)]
620pub(crate) struct Hex {
621 pub(crate) file: PathBuf,
622 pub(crate) from_home: bool,
623 pub(crate) asked: bool,
625 pub(crate) record_size: Option<usize>,
627}
628
629#[derive(Debug)]
631pub(crate) struct Tables {
632 pub(crate) database: PathBuf,
634 pub(crate) from_home: bool,
635}
636
637#[derive(Debug, Clone, Copy, PartialEq, Eq)]
639pub(crate) enum Conversion {
640 Streams,
642 Text(FileFormat),
645}
646
647impl Conversion {
648 pub(crate) fn label(self) -> &'static str {
650 match self {
651 Conversion::Streams => "Converting Arrow stream",
652 Conversion::Text(format) => format.conversion().label,
653 }
654 }
655
656 pub(crate) fn status(self) -> &'static str {
658 match self {
659 Conversion::Streams => "Converting Arrow stream...",
660 Conversion::Text(format) => format.conversion().status,
661 }
662 }
663}
664
665pub(crate) enum Converted {
667 Streams {
669 file: TempDownload,
670 parts: Vec<crate::formats::ipc_stream::Part>,
671 },
672 Frame {
675 files: Vec<TempDownload>,
676 lf: Box<LazyFrame>,
677 notes: Vec<crate::notes::Note>,
678 other_tables: Vec<String>,
679 detail: Option<Arc<crate::formats::text_formats::Detail>>,
681 },
682}
683
684#[derive(Default)]
687pub(crate) struct Made {
688 pub(crate) download: Option<TempDownload>,
690 pub(crate) converted: Vec<TempDownload>,
692 pub(crate) notes: Vec<crate::notes::Note>,
693 pub(crate) other_tables: Vec<String>,
694 pub(crate) detail: Option<Arc<crate::formats::text_formats::Detail>>,
695}
696
697#[derive(Debug)]
699pub(crate) struct Failed {
700 pub(crate) message: String,
701 pub(crate) from_home: bool,
702}
703
704pub(crate) struct Loaded {
706 pub(crate) state: DataTableState,
707 pub(crate) path: Option<PathBuf>,
709 pub(crate) options: OpenOptions,
710 pub(crate) debug_label: Option<String>,
711 pub(crate) paths: Option<Vec<PathBuf>>,
713 pub(crate) recent: Option<PathBuf>,
715 pub(crate) from_home: bool,
716 pub(crate) footers: Arc<FooterProgress>,
718}
719
720pub(crate) enum LoadAnswer {
722 Scanned {
724 lf: Box<LazyFrame>,
725 path: Option<PathBuf>,
726 options: OpenOptions,
727 },
728 Compressed {
730 file: PathBuf,
731 path: Option<PathBuf>,
732 options: OpenOptions,
733 },
734 SpecFetched {
736 spec: Arc<crate::formats::Spec>,
737 options: OpenOptions,
738 },
739 CompressedRecords {
742 file: PathBuf,
743 path: Option<PathBuf>,
744 choice: crate::formats::Choice,
745 options: OpenOptions,
746 },
747 DecompressedRecords {
749 copy: TempDownload,
750 path: PathBuf,
751 choice: crate::formats::Choice,
752 options: OpenOptions,
753 },
754 Convert {
757 what: Conversion,
758 files: Vec<PathBuf>,
759 bytes: u64,
760 path: Option<PathBuf>,
761 options: OpenOptions,
762 },
763 Converted {
765 converted: Converted,
766 path: Option<PathBuf>,
767 options: OpenOptions,
768 },
769 Tables {
772 file: PathBuf,
773 tables: Vec<String>,
774 path: Option<PathBuf>,
775 },
776 Hex {
779 file: PathBuf,
780 asked: bool,
781 record_size: Option<usize>,
782 },
783 SchemaRead {
785 state: Box<DataTableState>,
786 path: Option<PathBuf>,
787 options: OpenOptions,
788 debug_label: Option<String>,
789 },
790 #[cfg(any(feature = "http", feature = "cloud"))]
792 NoRanges { options: OpenOptions },
793 #[cfg(any(feature = "http", feature = "cloud"))]
795 Sized(PendingDownload),
796 #[cfg(any(feature = "http", feature = "cloud"))]
799 PastLimit(PendingDownload),
800 #[cfg(any(feature = "http", feature = "cloud"))]
802 Downloaded {
803 download: TempDownload,
804 options: OpenOptions,
805 },
806 Spooled {
808 download: TempDownload,
809 options: OpenOptions,
810 },
811 Recorded { file: PathBuf, options: OpenOptions },
814}
815
816#[derive(Debug, Clone, Copy)]
818pub(crate) struct Retired {
819 pub(crate) id: LoadId,
820 pub(crate) asking: bool,
822}
823
824#[derive(Default)]
826pub(crate) struct Loader {
827 next_id: u64,
828 load: Option<Load>,
829 kept: Option<Fetched>,
833 unfinished: Unfinished,
836}
837
838impl Loader {
839 pub(crate) fn unfinished(&self) -> &Unfinished {
841 &self.unfinished
842 }
843
844 pub(crate) fn current(&self) -> Option<&Load> {
846 self.load.as_ref()
847 }
848
849 pub(crate) fn id(&self) -> Option<LoadId> {
850 self.load.as_ref().map(|load| load.id)
851 }
852
853 fn current_in(&self, id: LoadId, phase: impl Fn(&Phase) -> bool) -> bool {
854 self.load
855 .as_ref()
856 .is_some_and(|load| load.id == id && phase(&load.phase))
857 }
858
859 pub(crate) fn looking_at_paths(&self, id: LoadId) -> bool {
861 self.current_in(id, |phase| matches!(phase, Phase::LookingAtPaths))
862 }
863
864 pub(crate) fn looking_at_directory(&self, id: LoadId) -> bool {
866 self.current_in(id, |phase| matches!(phase, Phase::LookingAtDirectory))
867 }
868
869 pub(crate) fn awaiting_dataset(&self) -> bool {
872 self.load
873 .as_ref()
874 .is_some_and(|load| !matches!(load.phase, Phase::FirstRows))
875 }
876
877 pub(crate) fn waits(&self) -> bool {
880 self.awaiting_dataset() && !self.asking()
881 }
882
883 pub(crate) fn asking(&self) -> bool {
885 self.load.as_ref().is_some_and(|load| load.phase.asks())
886 }
887
888 pub(crate) fn hold_while_asking(&mut self, hold: Hold) {
890 if let Some(Load {
891 phase: Phase::ConfirmingRead { _hold, .. },
892 ..
893 }) = self.load.as_mut()
894 {
895 *_hold = Some(hold);
896 }
897 }
898
899 #[cfg(any(feature = "http", feature = "cloud"))]
901 pub(crate) fn download_note(&self) -> Option<&'static str> {
902 match self.load.as_ref().map(|load| &load.phase) {
903 Some(Phase::Confirming { note, .. }) => *note,
904 _ => None,
905 }
906 }
907
908 pub(crate) fn progress(&self) -> Option<&Arc<FooterProgress>> {
911 self.load
912 .as_ref()
913 .filter(|load| !matches!(load.phase, Phase::FirstRows))
914 .map(|load| &load.progress)
915 }
916
917 pub(crate) fn retire(&mut self) -> Option<Retired> {
920 let load = self.load.take()?;
921 if !matches!(load.phase, Phase::FirstRows) {
922 load.progress.cancel();
923 }
924 Some(Retired {
925 id: load.id,
926 asking: load.phase.asks(),
927 })
928 }
929
930 pub(crate) fn make_way(&mut self) -> Option<Retired> {
933 if self.load.as_ref().is_some_and(|load| load.phase.starting()) {
934 return None;
935 }
936 self.retire()
937 }
938
939 fn start(&mut self, from_home: bool) -> &mut Load {
941 if self
944 .load
945 .as_ref()
946 .is_some_and(|load| !load.phase.starting())
947 {
948 debug_assert!(false, "an open started without making way");
949 self.retire();
950 }
951 if self.load.is_none() {
952 self.next_id = self.next_id.wrapping_add(1);
953 let progress = Arc::<FooterProgress>::default();
954 self.load = Some(Load {
955 id: LoadId(self.next_id),
956 from_home,
957 phase: Phase::Starting {
958 label: "Loading".to_string(),
959 percent: 0,
960 },
961 path: None,
962 size: 0,
963 paths: None,
964 recent: None,
965 writer: self.unfinished.writer(progress.cancel_flag()),
966 progress,
967 download: None,
968 converted: Vec::new(),
969 warn_in_memory_above: None,
970 });
971 }
972 let load = self.load.as_mut().expect("started just above");
973 load.from_home |= from_home;
974 load
975 }
976
977 pub(crate) fn announce(&mut self, from_home: bool, label: String, percent: u16) -> LoadId {
980 let load = self.start(from_home);
981 load.phase = Phase::Starting { label, percent };
982 load.id
983 }
984
985 #[cfg(test)]
988 pub(crate) fn first_rows_for_tests(&mut self) {
989 self.retire();
990 self.start(false).phase = Phase::FirstRows;
991 }
992
993 #[cfg(test)]
995 pub(crate) fn size_for_tests(&mut self, size: u64) {
996 if let Some(load) = self.load.as_mut() {
997 load.size = size;
998 }
999 }
1000
1001 pub(crate) fn name(&mut self, path: PathBuf) {
1003 if let Some(load) = self
1004 .load
1005 .as_mut()
1006 .filter(|load| !matches!(load.phase, Phase::FirstRows))
1007 {
1008 load.path = Some(stdin::named(&path));
1009 }
1010 }
1011
1012 pub(crate) fn look_at_paths(&mut self) -> LoadId {
1014 let load = self.start(false);
1015 load.phase = Phase::LookingAtPaths;
1016 load.id
1017 }
1018
1019 pub(crate) fn look_at_directory(&mut self, dir: PathBuf) -> LoadId {
1021 let load = self.start(false);
1022 load.phase = Phase::LookingAtDirectory;
1023 load.path = Some(dir);
1024 load.id
1025 }
1026
1027 pub(crate) fn open(&mut self, request: OpenRequest) -> Step {
1030 if !(request.paths.len() == 1
1031 && self
1032 .kept
1033 .as_ref()
1034 .is_some_and(|kept| kept.url == request.paths[0]))
1035 {
1036 self.kept = None;
1037 }
1038 let OpenRequest {
1039 paths,
1040 mut options,
1041 size,
1042 recent,
1043 shown,
1044 warn_in_memory_above,
1045 } = request;
1046 options.arrow_parts = None;
1048 let prepared = options
1049 .prepared
1050 .take()
1051 .and_then(|handoff| handoff.lock().ok()?.take());
1052 let load = self.start(false);
1053 load.path = Some(shown.unwrap_or_else(|| stdin::named(&paths[0])));
1054 load.size = size;
1055 load.recent = recent;
1056 load.paths = Some(paths.clone());
1057 load.warn_in_memory_above = warn_in_memory_above;
1058 match prepared {
1059 Some(prepared) => self.install_prepared(*prepared),
1060 None => self.first_step(paths, options),
1061 }
1062 }
1063
1064 fn install_prepared(&mut self, prepared: crate::home::home_preview::Prepared) -> Step {
1067 let crate::home::home_preview::Prepared {
1068 state,
1069 options,
1070 debug_label,
1071 progress,
1072 } = prepared;
1073 let writer = self.unfinished.writer(progress.cancel_flag());
1074 let load = self.load.as_mut().expect("started by the open");
1075 load.phase = Phase::FirstRows;
1076 load.progress = progress;
1079 load.writer = writer;
1080 Step::Install(Box::new(Loaded {
1081 state: *state,
1082 path: load.path.clone(),
1083 options,
1084 debug_label,
1085 paths: load.paths.clone(),
1086 recent: load.recent.clone(),
1087 from_home: load.from_home,
1088 footers: load.progress.clone(),
1089 }))
1090 }
1091
1092 pub(crate) fn open_frame(&mut self, lf: LazyFrame, options: OpenOptions) -> Step {
1094 let load = self.start(false);
1095 load.path = None;
1096 load.size = 0;
1097 load.paths = None;
1098 load.recent = None;
1099 load.phase = Phase::ReadingSchema;
1100 Step::ReadSchema {
1101 lf: Box::new(lf),
1102 path: None,
1103 options,
1104 progress: load.progress.clone(),
1105 made: Made::default(),
1106 }
1107 }
1108
1109 fn first_step(&mut self, paths: Vec<PathBuf>, options: OpenOptions) -> Step {
1112 let first = paths[0].clone();
1113 let src = source::input_source(&first);
1114 if paths.len() > 1 {
1115 let only_one = match &src {
1116 source::InputSource::S3(_) => {
1117 Some("Only one S3 URL at a time. Open a single s3:// path.")
1118 }
1119 source::InputSource::Gcs(_) => {
1120 Some("Only one GCS URL at a time. Open a single gs:// path.")
1121 }
1122 source::InputSource::Azure(_) => {
1123 Some("Only one Azure URL at a time. Open a single abfss:// path.")
1124 }
1125 source::InputSource::Http(_) => {
1126 Some("Only one HTTP/HTTPS URL at a time. Open a single URL.")
1127 }
1128 source::InputSource::Local(_) => None,
1129 };
1130 if let Some(message) = only_one {
1131 self.load = None;
1132 return Step::Crash(message.to_string());
1133 }
1134 }
1135 if let Some(message) = stdin::refuse(&paths, true) {
1136 self.load = None;
1137 return Step::Crash(message.to_string());
1138 }
1139 if options.follow
1140 && let Some(message) = crate::loading::follow::refuse_paths(&paths, &options)
1141 {
1142 self.load = None;
1143 return Step::Crash(message);
1144 }
1145 if let Some(url) = options
1148 .spec_file
1149 .clone()
1150 .filter(|file| options.spec_fetched.is_none() && source::is_remote_url(file))
1151 {
1152 let load = self.load.as_mut().expect("an open has a load");
1153 load.phase = Phase::ReadingSpec;
1154 return Step::FetchSpec {
1155 url,
1156 options,
1157 writer: load.writer.clone(),
1158 };
1159 }
1160 if stdin::is_stdin(&first) {
1161 if let Some(kept) = self.kept.clone().filter(|kept| {
1163 kept.url == first && kept.file.path().exists() && kept.serves(&options)
1164 }) {
1165 return self.read_download(kept, options);
1166 }
1167 let load = self.load.as_mut().expect("an open has a load");
1168 let read = Arc::<AtomicU64>::default();
1169 load.phase = Phase::Spooling { read: read.clone() };
1170 return Step::Spool {
1171 options,
1172 writer: load.writer.clone(),
1173 read,
1174 };
1175 }
1176 let load = self.load.as_mut().expect("an open has a load");
1177 let compression = options
1178 .compression
1179 .or_else(|| CompressionFormat::from_extension(&first));
1180 let delimited = delimited_format(&first, &options);
1181 if matches!(src, source::InputSource::Local(_))
1182 && paths.len() == 1
1183 && compression.is_some()
1184 && let Some(format) = delimited
1185 {
1186 load.phase = Phase::Decompressing;
1187 return Step::Decompress {
1188 file: first.clone(),
1189 path: first,
1190 options: OpenOptions {
1191 format: Some(format),
1192 ..options
1193 },
1194 writer: load.writer.clone(),
1195 download: None,
1196 };
1197 }
1198 #[cfg(any(feature = "http", feature = "cloud"))]
1200 if let Some(kept) = self
1201 .kept
1202 .clone()
1203 .filter(|kept| paths.len() == 1 && kept.file.path().exists() && kept.serves(&options))
1204 {
1205 return self.read_download(kept, options);
1206 }
1207 #[cfg(any(feature = "http", feature = "cloud"))]
1209 if paths.len() == 1
1210 && let Some(format) = crate::cloud::remote_model::model_format(&first, options.format)
1211 {
1212 let load = self.load.as_mut().expect("an open has a load");
1213 load.phase = Phase::ReadingHeaders;
1214 return Step::ReadHeaders {
1215 url: first,
1216 format,
1217 options,
1218 writer: load.writer.clone(),
1219 };
1220 }
1221 #[cfg(any(feature = "http", feature = "cloud"))]
1222 if let Some(pending) = remote_download(&src, &options) {
1223 let load = self.load.as_mut().expect("an open has a load");
1224 load.phase = Phase::CheckingSize { note: None };
1225 return Step::Probe(pending);
1226 }
1227 let load = self.load.as_mut().expect("an open has a load");
1228 if counts_footer(&paths, &options) {
1229 load.phase = Phase::CountingFooter;
1230 let display = load.path.clone().filter(|shown| *shown != paths[0]);
1231 return Step::Scan {
1232 paths,
1233 options,
1234 display,
1235 status: COUNTING_FOOTER_STATUS,
1236 };
1237 }
1238 if paths.len() == 1 && delimited.is_some() && options.parse_strings.is_some() {
1239 load.phase = Phase::ScanningStrings;
1240 return Step::Scan {
1241 paths,
1242 options,
1243 display: None,
1244 status: "Scanning string columns...",
1245 };
1246 }
1247 let display = load.path.clone().filter(|shown| *shown != paths[0]);
1249 if let Some(limit) = load.warn_in_memory_above
1251 && let Some(read) = in_memory(&paths, &options)
1252 && read.bytes > limit
1253 {
1254 load.phase = Phase::ConfirmingRead {
1255 scan: Box::new(Scan {
1256 paths,
1257 options,
1258 display,
1259 }),
1260 _hold: None,
1261 };
1262 return Step::AskRead(read);
1263 }
1264 if reads_lines(&paths, &options) {
1265 load.phase = Phase::ReadingLines;
1266 return Step::Scan {
1267 paths,
1268 options,
1269 display,
1270 status: READING_LINES_STATUS,
1271 };
1272 }
1273 load.phase = Phase::Scanning { downloaded: false };
1274 Step::Scan {
1275 paths,
1276 options,
1277 display,
1278 status: "Scanning input...",
1279 }
1280 }
1281
1282 fn read_download(&mut self, fetched: Fetched, mut options: OpenOptions) -> Step {
1285 if let Some(arrow) = &fetched.arrow {
1286 options.format = Some(FileFormat::Arrow);
1287 options.hive = false;
1288 options.arrow_parts = Some(arrow.parts.clone());
1289 options.splits = arrow.splits.clone();
1290 }
1291 let load = self.load.as_mut().expect("a download read has a load");
1292 let file = fetched.file.path().to_path_buf();
1293 let url = stdin::named(&fetched.url);
1294 let download = fetched.file.clone();
1295 load.download = Some(fetched);
1296 let compressed = options
1298 .compression
1299 .or_else(|| CompressionFormat::from_extension(&file))
1300 .is_some();
1301 if compressed && let Some(format) = delimited_format(&file, &options) {
1302 load.phase = Phase::Decompressing;
1303 return Step::Decompress {
1304 file,
1305 path: url,
1306 options: OpenOptions {
1307 format: Some(format),
1308 ..options
1309 },
1310 writer: load.writer.clone(),
1311 download: Some(download),
1312 };
1313 }
1314 let paths = vec![file];
1315 let status = if counts_footer(&paths, &options) {
1316 load.phase = Phase::CountingFooter;
1317 COUNTING_FOOTER_STATUS
1318 } else if reads_lines(&paths, &options) {
1319 load.phase = Phase::ReadingLines;
1320 READING_LINES_STATUS
1321 } else {
1322 load.phase = Phase::Scanning { downloaded: true };
1323 "Scanning..."
1324 };
1325 Step::Scan {
1326 paths,
1327 options,
1328 display: Some(url),
1329 status,
1330 }
1331 }
1332
1333 pub(crate) fn answered(
1336 &mut self,
1337 id: LoadId,
1338 answer: LoadAnswer,
1339 #[cfg(any(feature = "http", feature = "cloud"))] jobs: &Jobs,
1340 ) -> Step {
1341 let Some(load) = self.load.as_mut().filter(|load| load.id == id) else {
1342 return Step::Nothing;
1343 };
1344 match (answer, &load.phase) {
1345 (
1346 LoadAnswer::Scanned { lf, path, options },
1347 Phase::Scanning { .. }
1348 | Phase::ScanningStrings
1349 | Phase::CountingFooter
1350 | Phase::ReadingRecords
1351 | Phase::ReadingLines,
1352 ) => {
1353 load.phase = Phase::ReadingSchema;
1354 Step::ReadSchema {
1355 lf,
1356 path,
1357 options,
1358 progress: load.progress.clone(),
1359 made: load.made(),
1360 }
1361 }
1362 (
1363 LoadAnswer::Convert {
1364 what,
1365 files,
1366 bytes,
1367 path,
1368 options,
1369 },
1370 Phase::Scanning { .. }
1371 | Phase::ScanningStrings
1372 | Phase::CountingFooter
1373 | Phase::ReadingLines,
1374 ) if load.converted.is_empty() => {
1375 let read = Arc::<AtomicU64>::default();
1376 load.phase = Phase::Converting {
1377 what,
1378 read: read.clone(),
1379 total: bytes,
1380 };
1381 Step::Convert {
1382 what,
1383 files,
1384 path,
1385 options,
1386 writer: load.writer.clone(),
1387 read,
1388 }
1389 }
1390 (
1391 LoadAnswer::Converted {
1392 converted: Converted::Streams { file, parts },
1393 path,
1394 options,
1395 },
1396 Phase::Converting { .. },
1397 ) => {
1398 let paths = vec![file.path().to_path_buf()];
1399 if let Some(fetched) = load.download.as_mut() {
1402 fetched.file = file.clone();
1403 }
1404 load.converted = vec![file];
1405 load.phase = Phase::Scanning { downloaded: true };
1406 Step::Scan {
1407 paths,
1408 options: OpenOptions {
1409 format: Some(FileFormat::Arrow),
1410 hive: false,
1411 arrow_parts: Some(Arc::new(parts)),
1412 ..options
1413 },
1414 display: path,
1415 status: "Scanning...",
1416 }
1417 }
1418 (
1419 LoadAnswer::Converted {
1420 converted:
1421 Converted::Frame {
1422 files,
1423 lf,
1424 notes,
1425 other_tables,
1426 detail,
1427 },
1428 path,
1429 options,
1430 },
1431 Phase::Converting { .. },
1432 ) => {
1433 load.converted = files;
1434 load.phase = Phase::ReadingSchema;
1435 Step::ReadSchema {
1436 lf,
1437 path,
1438 options,
1439 progress: load.progress.clone(),
1440 made: Made {
1441 notes,
1442 other_tables,
1443 detail,
1444 ..load.made()
1445 },
1446 }
1447 }
1448 (
1449 LoadAnswer::Compressed {
1450 file,
1451 path,
1452 options,
1453 },
1454 Phase::Scanning { .. }
1455 | Phase::ScanningStrings
1456 | Phase::CountingFooter
1457 | Phase::ReadingLines,
1458 ) => {
1459 load.phase = Phase::Decompressing;
1460 Step::Decompress {
1461 path: path.unwrap_or_else(|| file.clone()),
1462 file,
1463 options,
1464 writer: load.writer.clone(),
1465 download: load.download.as_ref().map(|fetched| fetched.file.clone()),
1466 }
1467 }
1468 (LoadAnswer::SpecFetched { spec, options }, Phase::ReadingSpec) => {
1469 let paths = load.paths.clone().unwrap_or_default();
1470 if paths.is_empty() {
1471 return Step::Nothing;
1472 }
1473 self.first_step(
1474 paths,
1475 OpenOptions {
1476 spec_fetched: Some(spec),
1477 ..options
1478 },
1479 )
1480 }
1481 (
1482 LoadAnswer::CompressedRecords {
1483 file,
1484 path,
1485 choice,
1486 options,
1487 },
1488 Phase::Scanning { .. } | Phase::ScanningStrings | Phase::ReadingLines,
1489 ) => {
1490 load.phase = Phase::DecompressingRecords;
1491 Step::DecompressRecords {
1492 path: path.unwrap_or_else(|| file.clone()),
1493 file,
1494 choice,
1495 options,
1496 writer: load.writer.clone(),
1497 }
1498 }
1499 (
1500 LoadAnswer::DecompressedRecords {
1501 copy,
1502 path,
1503 choice,
1504 options,
1505 },
1506 Phase::DecompressingRecords,
1507 ) => {
1508 let file = copy.path().to_path_buf();
1510 load.converted = vec![copy];
1511 load.phase = Phase::ReadingRecords;
1512 Step::ReadRecords {
1513 copy: file,
1514 path,
1515 choice,
1516 options,
1517 }
1518 }
1519 (
1520 LoadAnswer::Hex {
1521 file,
1522 asked,
1523 record_size,
1524 },
1525 Phase::Scanning { .. }
1526 | Phase::ScanningStrings
1527 | Phase::CountingFooter
1528 | Phase::ReadingLines,
1529 ) => {
1530 let from_home = load.from_home;
1531 let fetched = load.download.is_some();
1534 let named = load.path.clone();
1535 self.retire();
1536 if fetched {
1537 let message = match named {
1538 Some(path) => crate::error_display::file_message(&path, crate::UNSUPPORTED),
1539 None => crate::UNSUPPORTED.to_string(),
1540 };
1541 return Step::Failed(Failed { message, from_home });
1542 }
1543 Step::Hex(Hex {
1544 file,
1545 from_home,
1546 asked,
1547 record_size,
1548 })
1549 }
1550 (
1551 LoadAnswer::Tables { file, tables, path },
1552 Phase::Scanning { .. }
1553 | Phase::ScanningStrings
1554 | Phase::CountingFooter
1555 | Phase::ReadingLines,
1556 ) => {
1557 let from_home = load.from_home;
1558 let database = path.unwrap_or(file);
1559 let fetched = load.download.is_some();
1562 self.retire();
1563 if fetched {
1564 let shown = tables.iter().take(20).cloned().collect::<Vec<_>>();
1565 let more = tables.len().saturating_sub(shown.len());
1566 let more = match more {
1567 0 => String::new(),
1568 n => format!(" and {n} more"),
1569 };
1570 return Step::Failed(Failed {
1571 message: format!(
1572 "{} holds {} tables: {}{more}. Open one with --table NAME.",
1573 database.display(),
1574 tables.len(),
1575 shown.join(", ")
1576 ),
1577 from_home,
1578 });
1579 }
1580 Step::Tables(Tables {
1581 database,
1582 from_home,
1583 })
1584 }
1585 (
1586 LoadAnswer::SchemaRead {
1587 state,
1588 path,
1589 options,
1590 debug_label,
1591 },
1592 phase,
1593 ) if phase.builds_the_dataset() => {
1594 load.phase = Phase::FirstRows;
1595 if let Some(fetched) = load.download.take() {
1598 self.kept = Some(fetched);
1599 }
1600 let state = *state;
1601 let load = self.load.as_ref().expect("installing its load");
1602 Step::Install(Box::new(Loaded {
1603 state,
1604 path,
1605 options,
1606 debug_label,
1607 paths: load.paths.clone(),
1608 recent: load.recent.clone(),
1609 from_home: load.from_home,
1610 footers: load.progress.clone(),
1611 }))
1612 }
1613 #[cfg(feature = "cloud")]
1615 (
1616 LoadAnswer::Sized(PendingDownload::Arrow {
1617 url,
1618 objects,
1619 options,
1620 ..
1621 }),
1622 Phase::CheckingSize { .. },
1623 ) if crate::cloud::cloud_arrow::in_place(&objects).is_some() => {
1624 load.phase = Phase::Scanning { downloaded: false };
1625 let parts = crate::cloud::cloud_arrow::in_place(&objects).unwrap_or_default();
1626 Step::Scan {
1627 paths: vec![PathBuf::from(url)],
1628 options: OpenOptions {
1629 format: Some(FileFormat::Arrow),
1630 hive: false,
1631 arrow_parts: Some(Arc::new(parts)),
1632 ..options
1633 },
1634 display: None,
1635 status: "Scanning...",
1636 }
1637 }
1638 #[cfg(any(feature = "http", feature = "cloud"))]
1639 (LoadAnswer::NoRanges { options }, Phase::ReadingHeaders) => {
1640 let first = load.paths.as_ref().and_then(|paths| paths.first().cloned());
1641 match first
1642 .and_then(|first| remote_download(&source::input_source(&first), &options))
1643 {
1644 Some(pending) => {
1645 load.phase = Phase::CheckingSize {
1646 note: Some(NO_RANGES),
1647 };
1648 Step::Probe(pending)
1649 }
1650 None => self.failed(id, NO_RANGES),
1651 }
1652 }
1653 #[cfg(any(feature = "http", feature = "cloud"))]
1656 (LoadAnswer::Sized(pending), Phase::CheckingSize { note: None })
1657 if pending
1658 .parts()
1659 .2
1660 .download_unasked
1661 .is_some_and(|unasked| unasked.covers(pending.parts().1)) =>
1662 {
1663 load.phase = Phase::Downloading;
1664 Step::Download {
1665 pending,
1666 writer: load.writer.clone(),
1667 }
1668 }
1669 #[cfg(any(feature = "http", feature = "cloud"))]
1670 (LoadAnswer::Sized(pending), Phase::CheckingSize { note }) => {
1671 load.phase = Phase::Confirming {
1672 pending: Box::new(pending.clone()),
1673 note: *note,
1674 _hold: jobs.hold(),
1675 };
1676 Step::Ask(pending)
1677 }
1678 #[cfg(any(feature = "http", feature = "cloud"))]
1681 (LoadAnswer::PastLimit(pending), Phase::Downloading) => {
1682 let pending = pending.with_size(None);
1683 load.phase = Phase::Confirming {
1684 pending: Box::new(pending.clone()),
1685 note: Some(PAST_LIMIT),
1686 _hold: jobs.hold(),
1687 };
1688 Step::Ask(pending)
1689 }
1690 #[cfg(any(feature = "http", feature = "cloud"))]
1691 (LoadAnswer::Downloaded { download, options }, Phase::Downloading) => {
1692 let fetched = Fetched {
1693 url: load.path.clone().unwrap_or_default(),
1695 file: download,
1696 arrow: options.arrow_parts.clone().map(|parts| KeptArrow {
1697 parts,
1698 splits: options.splits.clone(),
1699 table: options.table.clone(),
1700 }),
1701 };
1702 self.read_download(fetched, options)
1703 }
1704 (LoadAnswer::Recorded { file, options }, Phase::Spooling { read }) => {
1705 load.size = read.load(Ordering::Relaxed);
1706 self.kept = None;
1708 load.path = Some(file.clone());
1709 load.paths = Some(vec![file.clone()]);
1710 load.recent = Some(file.clone());
1711 let paths = vec![file];
1712 let status = if counts_footer(&paths, &options) {
1713 load.phase = Phase::CountingFooter;
1714 COUNTING_FOOTER_STATUS
1715 } else {
1716 load.phase = Phase::Scanning { downloaded: false };
1717 "Scanning input..."
1718 };
1719 Step::Scan {
1720 paths,
1721 options,
1722 display: None,
1723 status,
1724 }
1725 }
1726 (LoadAnswer::Spooled { download, options }, Phase::Spooling { read }) => {
1727 load.size = read.load(Ordering::Relaxed);
1728 let fetched = Fetched {
1729 url: PathBuf::from(stdin::PATH),
1730 file: download,
1731 arrow: None,
1732 };
1733 self.read_download(fetched, options)
1734 }
1735 _ => Step::Nothing,
1736 }
1737 }
1738
1739 pub(crate) fn failed(&mut self, id: LoadId, message: &str) -> Step {
1742 let Some(load) = self.load.as_ref().filter(|load| load.id == id) else {
1743 return Step::Nothing;
1744 };
1745 if matches!(load.phase, Phase::FirstRows) {
1746 return Step::Nothing;
1747 }
1748 let from_home = load.from_home;
1749 let mut message = match &load.download {
1752 Some(fetched) => crate::error_display::named_by_source(
1753 message,
1754 fetched.file.path(),
1755 &stdin::named(&fetched.url),
1756 ),
1757 None => message.to_string(),
1758 };
1759 if let Some(path) = &load.path {
1761 for converted in &load.converted {
1762 message = crate::error_display::named_by_source(&message, converted.path(), path);
1763 }
1764 }
1765 self.retire();
1766 Step::Failed(Failed { message, from_home })
1767 }
1768
1769 pub(crate) fn confirmed(&mut self) -> Step {
1772 let Some(load) = self.load.as_mut().filter(|load| load.phase.asks()) else {
1773 return Step::Nothing;
1774 };
1775 match std::mem::replace(&mut load.phase, Phase::Scanning { downloaded: false }) {
1778 Phase::ConfirmingRead { scan, .. } => {
1779 let Scan {
1780 paths,
1781 options,
1782 display,
1783 } = *scan;
1784 Step::Scan {
1785 paths,
1786 options,
1787 display,
1788 status: "Scanning input...",
1789 }
1790 }
1791 #[cfg(any(feature = "http", feature = "cloud"))]
1792 Phase::Confirming { pending, .. } => {
1793 load.phase = Phase::Downloading;
1794 Step::Download {
1795 pending: pending.asked(),
1796 writer: load.writer.clone(),
1797 }
1798 }
1799 _ => unreachable!("a phase that asks"),
1800 }
1801 }
1802
1803 pub(crate) fn first_rows_settled(&mut self) {
1806 if self
1807 .load
1808 .as_ref()
1809 .is_some_and(|load| matches!(load.phase, Phase::FirstRows))
1810 {
1811 self.load = None;
1812 }
1813 }
1814}
1815
1816impl Drop for Loader {
1817 fn drop(&mut self) {
1819 self.retire();
1820 }
1821}
1822
1823const READING_LINES: &str = "Reading as lines";
1825const READING_LINES_STATUS: &str = "Reading as lines...";
1826
1827fn reads_lines(paths: &[PathBuf], options: &OpenOptions) -> bool {
1830 options
1831 .format
1832 .or_else(|| match paths {
1833 [one] => FileFormat::from_path(one),
1834 _ => None,
1835 })
1836 .is_some_and(FileFormat::is_lines)
1837}
1838
1839pub(crate) fn in_memory(paths: &[PathBuf], options: &OpenOptions) -> Option<InMemory> {
1843 let mut found: Option<InMemory> = None;
1844 for path in paths {
1845 if !matches!(source::input_source(path), source::InputSource::Local(_)) {
1846 continue;
1847 }
1848 let compression = options
1849 .compression
1850 .or_else(|| CompressionFormat::from_extension(path));
1851 let format = options.format.or_else(|| match compression {
1854 Some(_) => path
1855 .file_stem()
1856 .and_then(|stem| FileFormat::from_path(Path::new(stem))),
1857 None => match FileFormat::from_path(path) {
1858 Some(named) => Some(crate::formats::readers::refined(path, named).unwrap_or(named)),
1859 None => crate::formats::readers::sniff_open(path, None),
1860 },
1861 });
1862 let Some(format) = format else {
1863 continue;
1864 };
1865 let stored = match compression {
1866 Some(_) => crate::Stored::Compressed {
1867 in_memory: options.decompress_in_memory,
1868 },
1869 None => crate::Stored::Plain,
1870 };
1871 if format.read_mode(stored) != Some(crate::ReadMode::InMemory)
1874 || format.http_file() == crate::RemoteRead::InPlace
1875 {
1876 continue;
1877 }
1878 let Some(bytes) = std::fs::metadata(path)
1879 .ok()
1880 .filter(|m| m.is_file())
1881 .map(|m| m.len())
1882 else {
1883 continue;
1884 };
1885 match &mut found {
1886 Some(read) => {
1887 read.bytes += bytes;
1888 read.files += 1;
1889 }
1890 None => {
1891 found = Some(InMemory {
1892 bytes,
1893 format,
1894 files: 1,
1895 name: path.clone(),
1896 })
1897 }
1898 }
1899 }
1900 found
1901}
1902
1903const COUNTING_FOOTER: &str = "Counting rows to skip the footer";
1905const COUNTING_FOOTER_STATUS: &str = "Counting rows to skip the footer...";
1907
1908fn counts_footer(paths: &[PathBuf], options: &OpenOptions) -> bool {
1911 options.skip_tail_rows.is_some_and(|n| n > 0)
1912 && paths.iter().all(|p| delimited_format(p, options).is_some())
1913}
1914
1915pub(crate) fn delimited_format(path: &Path, options: &OpenOptions) -> Option<FileFormat> {
1918 let format = options.format.or_else(|| {
1919 FileFormat::from_path(path).or_else(|| {
1920 CompressionFormat::from_extension(path)?;
1921 FileFormat::from_path(Path::new(path.file_stem()?))
1922 })
1923 })?;
1924 format.decompressed_once().then_some(format)
1925}
1926
1927#[cfg(any(feature = "http", feature = "cloud"))]
1929pub(crate) const PAST_LIMIT: &str =
1930 "This catalog file passed 50 MiB, more than it may download without asking.";
1931
1932#[cfg(any(feature = "http", feature = "cloud"))]
1934pub(crate) const NO_RANGES: &str = "The server does not send byte ranges, so the model's header cannot be read without downloading the whole file.";
1935
1936#[cfg(any(feature = "http", feature = "cloud"))]
1939fn remote_download(src: &source::InputSource, options: &OpenOptions) -> Option<PendingDownload> {
1940 #[cfg(feature = "cloud")]
1941 let spec = options.spec_file.is_some() || options.spec_name.is_some();
1942 #[cfg(feature = "cloud")]
1943 let should_download =
1944 |url: &str| should_download(url) || (spec && !source::is_prefix_or_glob(url));
1945 #[cfg(feature = "cloud")]
1946 if !spec && let Some(arrow) = cloud_arrow(src, options) {
1947 return Some(arrow);
1948 }
1949 let options = options.clone();
1950 match src {
1951 #[cfg(feature = "http")]
1952 source::InputSource::Http(url) => Some(PendingDownload::Http {
1953 url: url.clone(),
1954 size: None,
1955 options,
1956 }),
1957 #[cfg(feature = "cloud")]
1958 source::InputSource::S3(url) => {
1959 let full = format!("s3://{url}");
1960 should_download(&full).then_some(PendingDownload::S3 {
1961 url: full,
1962 size: None,
1963 options,
1964 })
1965 }
1966 #[cfg(feature = "cloud")]
1967 source::InputSource::Gcs(url) => {
1968 let full = format!("gs://{url}");
1969 should_download(&full).then_some(PendingDownload::Gcs {
1970 url: full,
1971 size: None,
1972 options,
1973 })
1974 }
1975 #[cfg(feature = "cloud")]
1976 source::InputSource::Azure(url) => should_download(url).then(|| PendingDownload::Azure {
1977 url: url.clone(),
1978 size: None,
1979 options,
1980 }),
1981 _ => None,
1982 }
1983}
1984
1985#[cfg(feature = "cloud")]
1988fn cloud_arrow(src: &source::InputSource, options: &OpenOptions) -> Option<PendingDownload> {
1989 let url = match src {
1990 source::InputSource::S3(url) => format!("s3://{url}"),
1991 source::InputSource::Gcs(url) => format!("gs://{url}"),
1992 source::InputSource::Azure(url) => url.clone(),
1993 _ => return None,
1994 };
1995 if url.contains('*') {
1996 return None;
1997 }
1998 let arrow = if url.ends_with('/') {
1999 options.format == Some(FileFormat::Arrow)
2000 } else {
2001 options.compression.is_none()
2003 && options
2004 .format
2005 .or_else(|| FileFormat::from_path(Path::new(&url)))
2006 == Some(FileFormat::Arrow)
2007 };
2008 arrow.then(|| PendingDownload::Arrow {
2009 url,
2010 objects: Vec::new(),
2011 size: None,
2012 options: options.clone(),
2013 })
2014}
2015
2016#[cfg(feature = "cloud")]
2017fn should_download(url: &str) -> bool {
2018 let (_, ext) = source::url_path_extension(url);
2019 source::cloud_path_should_download(ext.as_deref(), source::is_prefix_or_glob(url))
2020}
2021
2022#[cfg(test)]
2023mod tests;