Skip to main content

datui_lib/loading/
mod.rs

1//! Opening a dataset, from the request to its first rows, and its one owner.
2//!
3//! [`Loader`] holds the open in flight ([`Load`]): its origin, paths, phase, loading
4//! screen text, and what it holds (stop flag, footer counter, download, converted IPC
5//! file, the generation hold while a download is asked about). The app reports events
6//! (a request, a worker's answer or failure, the user's download answer) and carries
7//! out the [`Step`] it returns. Nothing else keeps the open's state.
8//!
9//! - **Identity.** Each load has a [`LoadId`], carried by its phases' jobs. An answer
10//!   is taken only for the load in flight, in the phase that asked; others are dropped
11//!   with their payload.
12//! - **Retirement.** Abandoning, replacing or failing a load raises its stop flag (a
13//!   download or conversion stops at its next chunk and removes its file; a footer
14//!   pass stops reading), releases its files and its hold.
15//! - **Handover.** The dataset is built holding the load's download or converted file
16//!   and everything found ([`crate::table::OpenFacts`]), and takes the footer counter on
17//!   install; what remains is the first-rows read ([`Phase::FirstRows`]), whose
18//!   abandonment stops neither.
19//!
20//! Home's looks, analyses and charts are not loads.
21
22pub(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/// Names one open, from the moment it is asked for until it is done.
51#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
52pub(crate) struct LoadId(u64);
53
54impl LoadId {
55    /// A load's name for a job a test starts by hand.
56    #[cfg(test)]
57    pub(crate) fn for_tests(n: u64) -> Self {
58        Self(n)
59    }
60}
61
62/// A remote file to download once the user agrees; its size is what the probe found.
63#[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    /// Arrow in S3, GCS or Azure, one object or a prefix: `objects` are what the probe
91    /// chose, `size` its streams' bytes (downloaded and converted); IPC files are scanned
92    /// in place (`cloud_arrow`).
93    #[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    /// The url, the size the probe found, and the open options — the same three
105    /// fields whichever store this came from.
106    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    /// How many Arrow streams the download is, and how many IPC files are read in
124    /// place beside them, for the question about it.
125    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    /// The download once the user has agreed to it: it never asks again, so it has
137    /// no limit to stop at.
138    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    /// Replace the placeholder size with what the probe actually found.
152    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/// Paths asked to be opened, and what the screen and the recents need of them. Built
170/// once, when the open is asked for, and not changed after.
171#[derive(Clone)]
172pub(crate) struct OpenRequest {
173    pub(crate) paths: Vec<PathBuf>,
174    pub(crate) options: OpenOptions,
175    /// The first path's size on disk, for the loading screen; 0 for a URL.
176    pub(crate) size: u64,
177    /// The path recorded as a recent if the dataset installs: the first, as named,
178    /// unless it is a local path that is not there.
179    pub(crate) recent: Option<PathBuf>,
180    /// What the loading screen names in place of the first path: a table inside a
181    /// database.
182    pub(crate) shown: Option<PathBuf>,
183    /// Ask before reading more than this many bytes whole into memory
184    /// (`[read] memory_warning`); `None` never asks.
185    pub(crate) warn_in_memory_above: Option<u64>,
186}
187
188impl OpenRequest {
189    /// The request for `paths`, asking the filesystem for the first one's size and
190    /// existence. Every open records a recent (command-line ones especially; URLs
191    /// verbatim), under the name given, not a download's temp copy. A table inside a file
192    /// (`app.db/users`) opens as the file with `--table` and is recorded as the table.
193    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            // A split of a Hugging Face cache, as the home screen lists it.
221            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            // A variant of a file a spec reads as several, as the home screen lists it.
238            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        // Standard input cannot be opened again from a list.
255        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
269/// Where an open stands. Each phase but the first rows is waiting on one worker, or on
270/// the user.
271pub(crate) enum Phase {
272    /// Asked for, before its first step: the frame between a key and the open it asks
273    /// for, or a startup open before its paths are looked at. Says what its caller says.
274    Starting {
275        label: String,
276        percent: u16,
277    },
278    /// Whether the paths named on the command line are there, and which is a directory.
279    LookingAtPaths,
280    /// The spec a `--format` URL names, fetched before anything is read with it.
281    ReadingSpec,
282    /// What a directory named on the command line holds, before it is opened.
283    LookingAtDirectory,
284    /// A remote model's headers, read by range rather than downloaded.
285    #[cfg(any(feature = "http", feature = "cloud"))]
286    ReadingHeaders,
287    /// The size of a remote file, to put the download to the user; `note` says why it
288    /// is downloaded when it would not have been.
289    #[cfg(any(feature = "http", feature = "cloud"))]
290    CheckingSize {
291        note: Option<&'static str>,
292    },
293    /// Waiting on the user to agree to read files whole into memory, past the size
294    /// that asks first. Holds the generation meanwhile, once the app gives it a hold.
295    ConfirmingRead {
296        scan: Box<Scan>,
297        _hold: Option<Hold>,
298    },
299    /// Waiting on the user to agree to the download. Holds the generation meanwhile:
300    /// nothing is running, and the open is very much unfinished.
301    #[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    /// Standard input being read to a file; `read` counts its bytes.
310    Spooling {
311        read: Arc<AtomicU64>,
312    },
313    Decompressing,
314    /// A compressed file a format spec reads, being decompressed to a copy.
315    DecompressingRecords,
316    /// The decompressed copy's records, read through the spec.
317    ReadingRecords,
318    /// Files being converted to ones the dataset scans, `read` of their `total` bytes:
319    /// Arrow IPC streams to one IPC file, or GPS logs to their table.
320    Converting {
321        what: Conversion,
322        read: Arc<AtomicU64>,
323        total: u64,
324    },
325    /// A CSV read with its string columns parsed.
326    ScanningStrings,
327    /// A CSV whose footer rows are dropped (`--footer-rows`): the scan counts every
328    /// row of the file first, which is the wait.
329    CountingFooter,
330    /// The scan; `downloaded` when it reads a download rather than what was named.
331    Scanning {
332        downloaded: bool,
333    },
334    /// The scan of text read as lines, which indexes them.
335    ReadingLines,
336    /// The schema, and whatever the dataset needs before its first rows.
337    ReadingSchema,
338    /// Installed: its first rows are being read. The dataset is the one on screen, and
339    /// what the load held is its now.
340    FirstRows,
341}
342
343impl Phase {
344    /// What the loading screen and the footer call this phase, and the flat
345    /// percentage the bar shows beside it (0 for none).
346    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                // Up to the scan or the schema read that follows.
365                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    /// Whether the open waits on the user's answer to a question.
379    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    /// Whether this phase's worker answers with the dataset itself.
389    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    /// Whether the open has not started any work of its own yet, so an open asked for
399    /// now carries it on rather than replacing it: the `Open` a key or a look returns.
400    fn starting(&self) -> bool {
401        matches!(
402            self,
403            Phase::Starting { .. } | Phase::LookingAtPaths | Phase::LookingAtDirectory
404        )
405    }
406}
407
408/// A download a load fetched, and the URL it was fetched from: standard input's `-`
409/// for what was piped in.
410#[derive(Clone)]
411struct Fetched {
412    url: PathBuf,
413    file: TempDownload,
414    /// What the listing of a store's Arrow chose, read again with the file: each
415    /// input's place, the split and the `--table` that chose it.
416    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    /// Whether an open with `options` reads what this download holds: not when it
428    /// asks for another split.
429    fn serves(&self, options: &OpenOptions) -> bool {
430        self.arrow
431            .as_ref()
432            .is_none_or(|arrow| arrow.table == options.table)
433    }
434}
435
436/// The open in flight.
437pub(crate) struct Load {
438    id: LoadId,
439    /// Chosen on the home screen, which is where its failure is reported: the dataset
440    /// left over from before is not what the user was looking at when they chose.
441    from_home: bool,
442    phase: Phase,
443    /// What the screen names: the path as asked for (a URL, not the temporary file it
444    /// landed in), and its size on disk.
445    path: Option<PathBuf>,
446    size: u64,
447    /// The paths asked for, which `H` opens again; `None` for a frame handed over.
448    paths: Option<Vec<PathBuf>>,
449    recent: Option<PathBuf>,
450    /// This load's footer counter, and through it its stop flag.
451    progress: Arc<FooterProgress>,
452    /// The stop flag again, for the workers that write files, and where they claim them.
453    writer: Writer,
454    download: Option<Fetched>,
455    /// The IPC files its Arrow streams or GPS logs were converted to, which the
456    /// dataset scans.
457    converted: Vec<TempDownload>,
458    /// See [`OpenRequest::warn_in_memory_above`].
459    warn_in_memory_above: Option<u64>,
460}
461
462/// A scan held while the user is asked about it.
463pub(crate) struct Scan {
464    paths: Vec<PathBuf>,
465    options: OpenOptions,
466    display: Option<PathBuf>,
467}
468
469/// What a read whole into memory would take, put to the user before it starts.
470#[derive(Debug, Clone, PartialEq, Eq)]
471pub(crate) struct InMemory {
472    /// The files' bytes on disk.
473    pub(crate) bytes: u64,
474    pub(crate) format: FileFormat,
475    pub(crate) files: usize,
476    /// The first file, as the loading screen names it.
477    pub(crate) name: PathBuf,
478}
479
480impl Load {
481    pub(crate) fn phase(&self) -> &Phase {
482        &self.phase
483    }
484
485    /// The files the frame about to have its schema read is made from: the download,
486    /// and what a conversion wrote.
487    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    /// The path the screen names, if it names one.
496    pub(crate) fn path(&self) -> Option<&Path> {
497        self.path.as_deref()
498    }
499
500    /// The size the screen shows beside the path: while standard input is read, what
501    /// has come in so far.
502    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
510/// What the app does next for the open.
511pub(crate) enum Step {
512    /// Nothing: the answer was for an open nobody is waiting on, or the phase waits.
513    Nothing,
514    /// The open cannot be read at all, and the session ends saying why.
515    Crash(String),
516    /// Read the headers of the remote model at `url` as `format`, by range; `writer`
517    /// carries the load's stop flag.
518    #[cfg(any(feature = "http", feature = "cloud"))]
519    ReadHeaders {
520        url: PathBuf,
521        format: FileFormat,
522        options: OpenOptions,
523        writer: Writer,
524    },
525    /// Find the remote file's size.
526    #[cfg(any(feature = "http", feature = "cloud"))]
527    Probe(PendingDownload),
528    /// Ask the user whether to download it.
529    #[cfg(any(feature = "http", feature = "cloud"))]
530    Ask(PendingDownload),
531    /// Ask the user whether to read files this large whole into memory.
532    AskRead(InMemory),
533    /// Download it, writing through `writer`: the load's stop flag, and its claim on
534    /// the file for quitting to find.
535    #[cfg(any(feature = "http", feature = "cloud"))]
536    Download {
537        pending: PendingDownload,
538        writer: Writer,
539    },
540    /// Fetch the spec at `url`, asking `writer`'s stop flag before the request.
541    FetchSpec {
542        url: PathBuf,
543        options: OpenOptions,
544        writer: Writer,
545    },
546    /// Read standard input to a file through `writer`, counting its bytes in `read`.
547    Spool {
548        options: OpenOptions,
549        writer: Writer,
550        read: Arc<AtomicU64>,
551    },
552    /// Decompress the CSV in `file`, writing through `writer`; `path` names it on screen
553    /// and in errors.
554    Decompress {
555        file: PathBuf,
556        path: PathBuf,
557        options: OpenOptions,
558        writer: Writer,
559        /// The download `file` is, given to the dataset built from it.
560        download: Option<TempDownload>,
561    },
562    /// Decompress `file`, which `choice`'s spec reads, to a copy written through
563    /// `writer`; `path` names it on screen and in errors.
564    DecompressRecords {
565        file: PathBuf,
566        path: PathBuf,
567        choice: crate::formats::Choice,
568        options: OpenOptions,
569        writer: Writer,
570    },
571    /// Read the records of `copy`, the decompressed file `path` names, with
572    /// `choice`'s spec.
573    ReadRecords {
574        copy: PathBuf,
575        path: PathBuf,
576        choice: crate::formats::Choice,
577        options: OpenOptions,
578    },
579    /// Convert `files` as `what` says into temp IPC files via `writer`, counting bytes read
580    /// in `read`; `path` names them on screen and in errors.
581    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 `paths`, saying `status` on the footer; `display` names the dataset when
590    /// what is scanned is a download.
591    Scan {
592        paths: Vec<PathBuf>,
593        options: OpenOptions,
594        display: Option<PathBuf>,
595        status: &'static str,
596    },
597    /// Read the scan's schema, reporting footers to `progress`.
598    ReadSchema {
599        lf: Box<LazyFrame>,
600        path: Option<PathBuf>,
601        options: OpenOptions,
602        progress: Arc<FooterProgress>,
603        /// What the open made the frame from, given to the dataset built from it.
604        made: Made,
605    },
606    /// Install the dataset, then read its first rows.
607    Install(Box<Loaded>),
608    /// The open failed. Its load is retired.
609    Failed(Failed),
610    /// The open found a SQLite database of several tables and was asked for none: the
611    /// home screen lists them. Its load is retired.
612    Tables(Tables),
613    /// The open found a local file no reader and no spec takes, or was asked for its
614    /// bytes: the hex view shows it. Its load is retired.
615    Hex(Hex),
616}
617
618/// A local file to show in the hex view.
619#[derive(Debug)]
620pub(crate) struct Hex {
621    pub(crate) file: PathBuf,
622    pub(crate) from_home: bool,
623    /// Asked for (`--hex`) rather than fallen back to.
624    pub(crate) asked: bool,
625    /// `--hex-width`: the bytes a row holds.
626    pub(crate) record_size: Option<usize>,
627}
628
629/// A file of several tables, to be listed on the home screen.
630#[derive(Debug)]
631pub(crate) struct Tables {
632    /// The database file, as the user named it.
633    pub(crate) database: PathBuf,
634    pub(crate) from_home: bool,
635}
636
637/// What a conversion turns into files the dataset scans.
638#[derive(Debug, Clone, Copy, PartialEq, Eq)]
639pub(crate) enum Conversion {
640    /// Arrow IPC streams, into one IPC file; the IPC files among them stay put.
641    Streams,
642    /// GPS logs of this format, each into IPC files of its own, read as one table; or
643    /// a VCD dump, FIX log or SDF file into its own.
644    Text(FileFormat),
645}
646
647impl Conversion {
648    /// What the loading screen calls it.
649    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    /// What the footer says while it runs.
657    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
665/// What a conversion wrote. Dropped unused, it removes its files.
666pub(crate) enum Converted {
667    /// The streams in one IPC file, and where each input's rows are: scanned next.
668    Streams {
669        file: TempDownload,
670        parts: Vec<crate::formats::ipc_stream::Part>,
671    },
672    /// The logs' IPC files and the frame over them, which only needs its schema read,
673    /// and what reading them noticed.
674    Frame {
675        files: Vec<TempDownload>,
676        lf: Box<LazyFrame>,
677        notes: Vec<crate::notes::Note>,
678        other_tables: Vec<String>,
679        /// What the file says besides its rows, for the Info panel.
680        detail: Option<Arc<crate::formats::text_formats::Detail>>,
681    },
682}
683
684/// The files an open made or fetched for the frame it scans, and what making them
685/// found: handed to the dataset as it is built, which holds the files from then on.
686#[derive(Default)]
687pub(crate) struct Made {
688    /// The download the frame reads.
689    pub(crate) download: Option<TempDownload>,
690    /// The files a conversion wrote.
691    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/// A failed open: why, and whether it was chosen on the home screen.
698#[derive(Debug)]
699pub(crate) struct Failed {
700    pub(crate) message: String,
701    pub(crate) from_home: bool,
702}
703
704/// A dataset read and ready to install, with everything the load hands over to it.
705pub(crate) struct Loaded {
706    pub(crate) state: DataTableState,
707    /// What names the dataset: the path or URL asked for.
708    pub(crate) path: Option<PathBuf>,
709    pub(crate) options: OpenOptions,
710    pub(crate) debug_label: Option<String>,
711    /// The paths asked for, which `H` opens again; `None` for a frame handed over.
712    pub(crate) paths: Option<Vec<PathBuf>>,
713    /// Recorded as a recent once installed.
714    pub(crate) recent: Option<PathBuf>,
715    pub(crate) from_home: bool,
716    /// The footer counter the dataset's own pass reports to from now on.
717    pub(crate) footers: Arc<FooterProgress>,
718}
719
720/// What a phase's worker found.
721pub(crate) enum LoadAnswer {
722    /// The scan's frame; `path` names the dataset.
723    Scanned {
724        lf: Box<LazyFrame>,
725        path: Option<PathBuf>,
726        options: OpenOptions,
727    },
728    /// The scan found one compressed delimited file, `file`, to decompress first.
729    Compressed {
730        file: PathBuf,
731        path: Option<PathBuf>,
732        options: OpenOptions,
733    },
734    /// The spec a `--format` URL names, and the open's options to carry it in.
735    SpecFetched {
736        spec: Arc<crate::formats::Spec>,
737        options: OpenOptions,
738    },
739    /// The scan found a compressed file, `file`, that `choice`'s spec reads once it is
740    /// decompressed.
741    CompressedRecords {
742        file: PathBuf,
743        path: Option<PathBuf>,
744        choice: crate::formats::Choice,
745        options: OpenOptions,
746    },
747    /// The decompressed copy of a file a spec reads. Dropped unused, it removes itself.
748    DecompressedRecords {
749        copy: TempDownload,
750        path: PathBuf,
751        choice: crate::formats::Choice,
752        options: OpenOptions,
753    },
754    /// The scan found `files`, `bytes` in all as stored, which have to be converted as
755    /// `what` says before they can be read: Arrow IPC streams, or GPS logs.
756    Convert {
757        what: Conversion,
758        files: Vec<PathBuf>,
759        bytes: u64,
760        path: Option<PathBuf>,
761        options: OpenOptions,
762    },
763    /// What the conversion wrote. Dropped unused, it removes its files.
764    Converted {
765        converted: Converted,
766        path: Option<PathBuf>,
767        options: OpenOptions,
768    },
769    /// The scan found a file of several tables (a SQLite database, a NumPy archive),
770    /// `file`, its tables named `tables`, and no `--table` to say which.
771    Tables {
772        file: PathBuf,
773        tables: Vec<String>,
774        path: Option<PathBuf>,
775    },
776    /// The scan found a local file no reader and no spec takes, or `--hex` asked for
777    /// its bytes.
778    Hex {
779        file: PathBuf,
780        asked: bool,
781        record_size: Option<usize>,
782    },
783    /// The dataset, its schema read.
784    SchemaRead {
785        state: Box<DataTableState>,
786        path: Option<PathBuf>,
787        options: OpenOptions,
788        debug_label: Option<String>,
789    },
790    /// The remote model's server sends whole files, not ranges: it is downloaded.
791    #[cfg(any(feature = "http", feature = "cloud"))]
792    NoRanges { options: OpenOptions },
793    /// The remote file's size.
794    #[cfg(any(feature = "http", feature = "cloud"))]
795    Sized(PendingDownload),
796    /// A download started without asking passed its limit and was stopped, its
797    /// partial file removed: the user is asked before it is fetched whole.
798    #[cfg(any(feature = "http", feature = "cloud"))]
799    PastLimit(PendingDownload),
800    /// The remote file, downloaded. Dropped unused, it removes the file.
801    #[cfg(any(feature = "http", feature = "cloud"))]
802    Downloaded {
803        download: TempDownload,
804        options: OpenOptions,
805    },
806    /// Standard input, read to a file, and `options` with the format it holds.
807    Spooled {
808        download: TempDownload,
809        options: OpenOptions,
810    },
811    /// Standard input, recorded to the file `--tee` named, which is read from here on
812    /// as any file is.
813    Recorded { file: PathBuf, options: OpenOptions },
814}
815
816/// A load put down before it finished: which, and what of it the app has to put down.
817#[derive(Debug, Clone, Copy)]
818pub(crate) struct Retired {
819    pub(crate) id: LoadId,
820    /// It was asking the user about a download: the question goes with it.
821    pub(crate) asking: bool,
822}
823
824/// The owner of the open in flight, and of the last download, kept to be read again.
825#[derive(Default)]
826pub(crate) struct Loader {
827    next_id: u64,
828    load: Option<Load>,
829    /// The last remote file an installed open downloaded, kept so reopening the URL (`H`)
830    /// reads it instead of downloading again; released when other data opens, removed
831    /// once the scanning dataset also lets go. Stdin, read once, is kept the same way.
832    kept: Option<Fetched>,
833    /// The files every load's workers have written and not yet let go, for quitting to
834    /// remove. See [`crate::loading::unfinished`].
835    unfinished: Unfinished,
836}
837
838impl Loader {
839    /// The files this loader's opens are writing, swept when the app exits.
840    pub(crate) fn unfinished(&self) -> &Unfinished {
841        &self.unfinished
842    }
843
844    /// The open in flight, if there is one.
845    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    /// Whether `id` is still looking at the paths named on the command line.
860    pub(crate) fn looking_at_paths(&self, id: LoadId) -> bool {
861        self.current_in(id, |phase| matches!(phase, Phase::LookingAtPaths))
862    }
863
864    /// Whether `id` is still looking at a directory named on the command line.
865    pub(crate) fn looking_at_directory(&self, id: LoadId) -> bool {
866        self.current_in(id, |phase| matches!(phase, Phase::LookingAtDirectory))
867    }
868
869    /// Whether an open is on its way and its dataset is not installed yet. Whatever
870    /// table is up meanwhile belongs to the dataset being replaced.
871    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    /// Whether the user waits on the open (keys held): not while it asks about a download
878    /// (the question takes the keys), nor once the dataset is up.
879    pub(crate) fn waits(&self) -> bool {
880        self.awaiting_dataset() && !self.asking()
881    }
882
883    /// Whether the open is waiting on the user to agree to a download.
884    pub(crate) fn asking(&self) -> bool {
885        self.load.as_ref().is_some_and(|load| load.phase.asks())
886    }
887
888    /// Give the question being asked the hold that keeps the generation while it waits.
889    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    /// Why the download being asked about is a download, when it would not have been.
900    #[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    /// The footer counter of the open that has not installed yet: what the loading
909    /// screen counts.
910    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    /// Put down the open in flight, if there is one, and what it holds. Its stop flag is
918    /// raised unless its dataset is installed, when the counter is the dataset's.
919    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    /// Make way for an open being asked for: the load in flight is retired, unless it
931    /// has not started work of its own yet, in which case the open carries it on.
932    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    /// The load an open being asked for belongs to: the starting one, or a new one.
940    fn start(&mut self, from_home: bool) -> &mut Load {
941        // A caller that did not make way would leave a load doing work behind with its
942        // stop flag down; retired here so its download and its footer pass stop.
943        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    /// An open is on its way, asked for from the home screen when `from_home`: the
978    /// screen is handed over to it now, saying `label`.
979    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    /// An open whose dataset is up and whose first rows are being read, with nothing
986    /// running: for tests of what ends that wait.
987    #[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    /// Give the starting load's path a size, as an open asking the filesystem does.
994    #[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    /// Put `path` on the loading screen, so a wait says what it is waiting for.
1002    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    /// The paths named on the command line are looked at before they are opened.
1013    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    /// A directory named on the command line is looked at before it is opened.
1020    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    /// Open `request`: carry on the starting load, or begin one. The last download is
1028    /// let go unless this opens it again.
1029    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        // What an earlier load found of its Arrow is not this one's to read.
1047        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    /// Install the dataset home's preview built: scan, schema and first page are read, so
1065    /// the open goes straight to its (on-hand) first rows.
1066    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        // The dataset was built reporting to the preview's counter, which is the one
1077        // its own footer pass goes on with.
1078        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    /// Open a frame handed over (the Python binding's): only its schema is read.
1093    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    /// What the open of `paths` does first: read standard input, decompress, download,
1110    /// or scan.
1111    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        // A remote spec is fetched once, before any phase reads with it; the open
1146        // then starts again from here with it in hand.
1147        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            // Read once: opened again (`H`), the copy on hand is read.
1162            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        // Opened again, and downloaded already: read the copy on hand.
1199        #[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        // A remote model's headers are all it needs: read by range, not downloaded.
1208        #[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        // A table inside a database goes by its path there, on screen and once open.
1248        let display = load.path.clone().filter(|shown| *shown != paths[0]);
1249        // A large file read whole is put to the user before the read starts.
1250        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    /// Read a download: decompress compressed CSV, TSV or PSV first, else scan it, named by
1283    /// the URL (or `stdin`), not the temp file.
1284    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        // Compressed delimited text must be decompressed before scanning, as from disk.
1297        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    /// A worker of `id` answered. Taken only while `id` is in flight and in the phase
1334    /// that asked; otherwise the answer, and whatever it carries, is dropped.
1335    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                // A converted download or pipe spool is kept as its copy, to reread without
1400                // converting; the stream is released.
1401                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                // The load holds the copy until the dataset built from it does.
1509                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                // A download is a temporary file the load owns; it has no bytes to show
1532                // once the load is put down.
1533                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                // A download or standard input has no place on the home screen to list
1560                // the tables at: the table is named on the command line instead.
1561                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                // The dataset was built holding its download (`Step::ReadSchema`); the
1596                // loader keeps it too, to be read again.
1597                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            // Arrow in a store with no stream in it: nothing to download or ask about.
1614            #[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            // A small file of the built-in catalog: its row said what it is and what it
1654            // weighs, so it is fetched without a question.
1655            #[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            // More arrived than the catalog said: asked once, as any large download is.
1679            // Its size is unknown now, whatever the server said.
1680            #[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                    // The URL the open was asked for, which is what opening it again names.
1694                    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                // The user's file, not a copy: `H` reads it again, and it is a recent.
1707                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    /// A worker of `id` failed. The open ends there, with its reason, unless it is no
1740    /// longer the one in flight or its dataset is already up.
1741    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        // A download is read from a temp file the user never typed: the reason names
1750        // the URL they did, or `stdin`.
1751        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        // So are the IPC files a conversion wrote.
1760        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    /// The user agreed to what the open asked: let go of the hold, and fetch the
1770    /// download or start the read.
1771    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        // The hold goes as the phase changes: the job the caller starts next holds the
1776        // generation before anything else can look at it.
1777        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    /// The first rows of the installed dataset are on screen, or will not be read: the
1804    /// open is done.
1805    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    /// A download still running stops at its next chunk and removes its partial file.
1818    fn drop(&mut self) {
1819        self.retire();
1820    }
1821}
1822
1823/// What the loading screen and the footer say while text is read as lines.
1824const READING_LINES: &str = "Reading as lines";
1825const READING_LINES_STATUS: &str = "Reading as lines...";
1826
1827/// Whether `paths` are read as lines per `--format` or their names, for the loading line
1828/// (bytes may still say otherwise, e.g. a candump `.log`).
1829fn 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
1839/// What reading `paths` would read whole into memory, by names and options only (not
1840/// bytes, on the event thread): local files of an in-memory format
1841/// ([`FileFormat::read_mode`]). `None` when none is.
1842pub(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        // By name, or by the first bytes the open will judge it by: journal JSON
1852        // piped to a file with no name, or in a `.json` one.
1853        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        // A model file's table comes from its header, which is what a URL of one is
1872        // read in place for: small however large the file.
1873        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
1903/// What the loading screen says while a CSV is counted to drop its footer rows.
1904const COUNTING_FOOTER: &str = "Counting rows to skip the footer";
1905/// The footer's line for the same wait.
1906const COUNTING_FOOTER_STATUS: &str = "Counting rows to skip the footer...";
1907
1908/// Whether the scan of `paths` counts every row first: delimited text whose footer
1909/// rows are dropped (`--footer-rows`).
1910fn 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
1915/// The delimited format (CSV, TSV, PSV) or text `path` is read as, if any: `--format`,
1916/// else the extension past any compression suffix (`x.tsv.gz` is TSV).
1917pub(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/// Why a catalog file that started without a question asks partway.
1928#[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/// Why a remote model is downloaded rather than read by its headers.
1933#[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/// The download a remote source needs first, if any: always for HTTP, and for a store
1937/// object that cannot be scanned in place, or one a format spec reads.
1938#[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/// Arrow in a store (object or listed prefix, not a glob): the probe lists it, its
1986/// streams are downloaded, its IPC files scanned in place.
1987#[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        // Compressed, it is a download like any other compressed object.
2002        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;