Skip to main content

datui_lib/
lines.rs

1//! Text read as it stands: a row per line, in a `line` column.
2//!
3//! One pass records where each line starts ([`LineIndex`]); a line is then read from a
4//! map of the file where it is shown, so only the rows a view reaches are decoded.
5//! Every line is a row, blank ones included, and each row carries its place in the
6//! file ([`crate::row_index::INDEX`]), so `#` numbers it as `less -N` does through any
7//! sort or filter. A `\r\n` ending is a line ending; bytes that are not UTF-8 are shown
8//! as `�` and counted in a note. Control characters stay in the value: the screen
9//! escapes them when it draws ([`crate::sanitize`]).
10//!
11//! A large file shows its first rows once its first [`FIRST_BYTES`] are indexed; the
12//! rest is indexed behind them ([`Lines::index_more`]), and the frame's height moves
13//! as it goes ([`bound`]).
14//!
15//! A followed file's lines are counted by the watcher ([`crate::follow`]), which moves
16//! the frame's height ([`bound`]); the index reads on from its last whole line when a
17//! row past it is asked for.
18//!
19//! [`guess`] says what text no format's signature claims is: JSON, CSV or TSV only on
20//! real evidence, lines otherwise.
21
22use std::path::{Path, PathBuf};
23use std::sync::{Arc, RwLock};
24
25use color_eyre::Result;
26use polars::prelude::*;
27
28use crate::FileFormat;
29use crate::fixed_records::Bytes;
30use crate::indexed::Offsets;
31
32/// The columns: the file a line is from (several files only), and the line.
33pub const FILE: &str = "file";
34pub const LINE: &str = "line";
35
36/// Bytes an open indexes before it shows the first rows. Past this the rest of a file
37/// is indexed in the background.
38pub const FIRST_BYTES: usize = 8 << 20;
39
40/// Bytes checked as UTF-8 at once, as the newlines in them are found: a chunk is
41/// still in cache when its lines are indexed, and the file is read once.
42const CHUNK: usize = 4 << 20;
43
44/// What datui does with text read as lines: see [`crate::readers`].
45pub(crate) const READER: crate::readers::Reader = crate::readers::Reader {
46    scan,
47    preview: Some(crate::readers::Preview::Scan),
48    python: Some(crate::python_script::Python {
49        call: "pl.LazyFrame",
50        eager: false,
51        glob_flag: false,
52        arguments: Some(crate::python_script::lines_arguments),
53    }),
54    export: Some(crate::export_modal::ExportFormat::Csv),
55    ..crate::readers::BASE
56};
57
58/// The columns of one file's lines, or several files' with the file named.
59pub fn schema(several: bool) -> Schema {
60    let mut fields = Vec::with_capacity(2);
61    if several {
62        fields.push(Field::new(FILE.into(), DataType::String));
63    }
64    fields.push(Field::new(LINE.into(), DataType::String));
65    Schema::from_iter(fields)
66}
67
68/// Where each line of some bytes starts, from one pass. Read on from where it left
69/// off as the bytes grow: a last line with no newline yet is read again.
70#[derive(Debug, Clone, Default)]
71pub struct LineIndex {
72    offsets: Offsets,
73    /// Bytes through the last newline.
74    end: usize,
75    /// The last line has no newline: whether it was indexed, and whether it was
76    /// counted as not UTF-8, so a read on from `end` takes it back first.
77    partial: Option<(bool, bool)>,
78    /// Lines that are not valid UTF-8.
79    pub invalid: usize,
80    /// Lines past [`crate::indexed::MAX_RECORDS`], left out.
81    pub past_limit: usize,
82    /// Lines indexed before this index begins: a step taken apart from the index it
83    /// goes on ([`Self::step_after`]) counts them toward the limit.
84    before: usize,
85}
86
87impl LineIndex {
88    /// The lines of `bytes`, all of them, the last whether or not it ends in a newline.
89    pub fn of(bytes: &[u8]) -> Self {
90        let mut index = Self {
91            offsets: Offsets::for_file(bytes.len()),
92            ..Default::default()
93        };
94        index.extend(bytes);
95        index
96    }
97
98    /// Lines indexed.
99    pub fn lines(&self) -> usize {
100        self.offsets.len()
101    }
102
103    /// Lines that end in a newline.
104    pub fn complete(&self) -> usize {
105        self.lines() - usize::from(self.partial.is_some_and(|(indexed, _)| indexed))
106    }
107
108    /// Bytes through the last newline.
109    pub fn end(&self) -> usize {
110        self.end
111    }
112
113    /// Whether every byte of `bytes` is indexed.
114    pub fn whole(&self, bytes: &[u8]) -> bool {
115        self.end >= bytes.len() || self.partial.is_some()
116    }
117
118    /// Index what `bytes`, the same bytes as before and maybe more, hold past the
119    /// last newline.
120    pub fn extend(&mut self, bytes: &[u8]) {
121        self.extend_by(bytes, usize::MAX);
122    }
123
124    /// [`Self::extend`] through `budget` more bytes at least: to the end of the line
125    /// that passes it. Returns whether every byte is indexed.
126    pub fn extend_by(&mut self, bytes: &[u8], budget: usize) -> bool {
127        if let Some((indexed, invalid)) = self.partial.take() {
128            if indexed {
129                self.offsets.truncate(self.offsets.len() - 1);
130            } else {
131                self.past_limit -= 1;
132            }
133            self.invalid -= usize::from(invalid);
134        }
135        if bytes.len() <= self.end {
136            return true;
137        }
138        self.offsets.widen_for(bytes.len());
139        let stop = self.end.saturating_add(budget).min(bytes.len());
140        while self.end < stop && self.partial.is_none() {
141            // A chunk ends at a line's end, so each line is checked whole.
142            let cut = (self.end + CHUNK).min(stop);
143            let to = if cut < bytes.len() {
144                memchr::memchr(b'\n', &bytes[cut - 1..]).map_or(bytes.len(), |i| cut + i)
145            } else {
146                cut
147            };
148            self.index_chunk(bytes, to);
149        }
150        self.whole(bytes)
151    }
152
153    /// Index the lines from `end` to `to`: the end of `bytes`, or just past a newline.
154    fn index_chunk(&mut self, bytes: &[u8], to: usize) {
155        // Checked once for the chunk; line by line only when that fails.
156        let all_utf8 = std::str::from_utf8(&bytes[self.end..to]).is_ok();
157        let mut at = self.end;
158        while at < to {
159            let newline = memchr::memchr(b'\n', &bytes[at..to]).map(|i| at + i);
160            let end = newline.unwrap_or(to);
161            let indexed = self.before + self.offsets.len() < crate::indexed::MAX_RECORDS;
162            if indexed {
163                self.offsets.push(at);
164            } else {
165                self.past_limit += 1;
166            }
167            let invalid = !all_utf8 && std::str::from_utf8(&bytes[at..end]).is_err();
168            self.invalid += usize::from(invalid);
169            match newline {
170                Some(newline) => {
171                    at = newline + 1;
172                    self.end = at;
173                }
174                None => {
175                    self.partial = Some((indexed, invalid));
176                    at = to;
177                }
178            }
179        }
180    }
181
182    /// An empty index that goes on from this one's last whole line, for a step taken
183    /// without holding this one: [`Self::take_step`] appends it. Not for an index with
184    /// a last line still open, which only a whole read leaves.
185    fn step_after(&self, len: usize) -> Self {
186        Self {
187            offsets: Offsets::for_file(len),
188            end: self.end,
189            before: self.before + self.offsets.len(),
190            ..Default::default()
191        }
192    }
193
194    /// Append `step`, taken from [`Self::step_after`] over the same bytes.
195    fn take_step(&mut self, step: LineIndex) {
196        debug_assert_eq!(step.before, self.offsets.len());
197        self.offsets.append(&step.offsets);
198        self.end = step.end;
199        self.partial = step.partial;
200        self.invalid += step.invalid;
201        self.past_limit += step.past_limit;
202    }
203
204    /// Line `i` of `bytes`, without its newline or a `\r` before it.
205    fn line<'a>(&self, bytes: &'a [u8], i: usize) -> &'a [u8] {
206        let start = self.offsets.get(i);
207        let end = if i + 1 < self.offsets.len() {
208            self.offsets.get(i + 1) - 1
209        } else {
210            memchr::memchr(b'\n', &bytes[start..]).map_or(bytes.len(), |n| start + n)
211        };
212        let line = &bytes[start..end.min(bytes.len())];
213        line.strip_suffix(b"\r").unwrap_or(line)
214    }
215}
216
217/// One file's bytes and their lines.
218struct Mapped {
219    bytes: Arc<Bytes>,
220    index: LineIndex,
221}
222
223struct LineFile {
224    /// The name shown in the `file` column.
225    name: String,
226    /// The file, when it is followed and its lines are read on as it grows.
227    grows: Option<PathBuf>,
228    mapped: RwLock<Mapped>,
229}
230
231impl LineFile {
232    fn rows(&self) -> usize {
233        self.mapped
234            .read()
235            .unwrap_or_else(|e| e.into_inner())
236            .index
237            .lines()
238    }
239
240    /// Map the file again and index on from the last newline, or from the start when
241    /// it is shorter than what was indexed: it was truncated or replaced.
242    fn grow(&self) -> PolarsResult<()> {
243        let Some(path) = &self.grows else {
244            return Ok(());
245        };
246        let mut mapped = self.mapped.write().unwrap_or_else(|e| e.into_inner());
247        // By name while there is one; a deleted file is read through the handle on it.
248        let bytes = match Bytes::map(path) {
249            Ok(bytes) => bytes,
250            Err(_) => match mapped.bytes.as_ref() {
251                Bytes::Mapped(_, file) => remap(file)?,
252                Bytes::Owned(_) => return Ok(()),
253            },
254        };
255        if bytes.len() < mapped.index.end() {
256            mapped.index = LineIndex::default();
257        }
258        mapped.index.extend(bytes.as_slice());
259        mapped.bytes = Arc::new(bytes);
260        Ok(())
261    }
262}
263
264fn remap(file: &std::fs::File) -> PolarsResult<Bytes> {
265    let file = file.try_clone()?;
266    if file.metadata()?.len() == 0 {
267        return Ok(Bytes::Owned(Vec::new()));
268    }
269    // SAFETY: as `Bytes::map`: read-only, and a shorter file is indexed again.
270    let map = unsafe { memmap2::Mmap::map(&file)? };
271    Ok(Bytes::Mapped(map, file))
272}
273
274/// The lines of one file or several, read where they are shown.
275pub struct Lines {
276    files: Vec<LineFile>,
277    schema: SchemaRef,
278    /// The one file's lines are still being indexed, behind the first rows, and what
279    /// a read of every line waits on until they are.
280    indexing: (std::sync::Mutex<bool>, std::sync::Condvar),
281    /// One indexing step at a time: a paused indexing taken up again may start its
282    /// thread while the last one finishes its step.
283    stepping: std::sync::Mutex<()>,
284    /// The file came back shorter than it was mapped while it was indexed.
285    shrank: std::sync::atomic::AtomicBool,
286}
287
288/// What a read of every line says when the file shrank as it was indexed.
289pub const SHRANK: &str = "the file shrank while it was indexed; open it again to read it";
290
291impl Lines {
292    /// `bytes` read as lines; `name` is what the `file` column says when there are
293    /// several.
294    pub fn from_bytes(parts: Vec<(String, Arc<Bytes>)>) -> Self {
295        Self::from_bytes_first(parts, usize::MAX)
296    }
297
298    /// [`Self::from_bytes`], one file's first `first` bytes indexed and the rest left
299    /// to [`Self::index_more`]. Several files are indexed whole: a row's place in one
300    /// would move as the file before it was indexed.
301    pub fn from_bytes_first(parts: Vec<(String, Arc<Bytes>)>, first: usize) -> Self {
302        let several = parts.len() > 1;
303        let mut indexing = false;
304        let files = parts
305            .into_iter()
306            .map(|(name, bytes)| {
307                let mut index = LineIndex {
308                    offsets: Offsets::for_file(bytes.len()),
309                    ..Default::default()
310                };
311                let budget = if several { usize::MAX } else { first };
312                indexing |= !index.extend_by(bytes.as_slice(), budget);
313                LineFile {
314                    name,
315                    grows: None,
316                    mapped: RwLock::new(Mapped { index, bytes }),
317                }
318            })
319            .collect();
320        Self {
321            files,
322            schema: Arc::new(schema(several)),
323            indexing: (indexing.into(), Default::default()),
324            stepping: Default::default(),
325            shrank: Default::default(),
326        }
327    }
328
329    /// Whether every line of every file is indexed.
330    pub fn whole(&self) -> bool {
331        self.files.iter().all(|f| {
332            let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
333            m.index.whole(m.bytes.as_slice())
334        })
335    }
336
337    /// Whether the file came back shorter than it was mapped while it was indexed.
338    pub fn shrank(&self) -> bool {
339        self.shrank.load(std::sync::atomic::Ordering::Relaxed)
340    }
341
342    /// Indexing set aside ([`Self::stop_indexing`] was not called: a pause leaves
343    /// the reads waiting) is taken up again: whether there is more to index.
344    pub fn resume_indexing(&self) -> bool {
345        let more = !self.whole() && !self.shrank();
346        *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner()) = more;
347        more
348    }
349
350    /// Which line of its own file the row at `place` is, from 0: the row's place in
351    /// the table, for one file.
352    pub fn line_in_file(&self, place: usize) -> Option<usize> {
353        let mut start = 0;
354        for f in &self.files {
355            let rows = f.rows();
356            if place < start + rows {
357                return Some(place - start);
358            }
359            start += rows;
360        }
361        None
362    }
363
364    /// Whether the rows are of several files, each numbered on its own.
365    pub fn several(&self) -> bool {
366        self.files.len() > 1
367    }
368
369    /// Whether lines are still being indexed behind the first rows.
370    pub fn indexing(&self) -> bool {
371        *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner())
372    }
373
374    /// No more lines will be indexed: all of them are, or the indexing stopped (the
375    /// dataset went). Whatever waits on them goes on with what there is.
376    pub fn stop_indexing(&self) {
377        *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner()) = false;
378        self.indexing.1.notify_all();
379    }
380
381    /// Wait until every line is indexed, on a worker, never the UI thread: `Err` when
382    /// they will not all be (the indexing stopped, or the file shrank), or the job
383    /// waiting was superseded and its answer is not wanted.
384    fn wait_indexed(&self) -> PolarsResult<()> {
385        let mut indexing = self.indexing.0.lock().unwrap_or_else(|e| e.into_inner());
386        while *indexing {
387            polars_ensure!(
388                !crate::jobs::superseded(),
389                ComputeError: "the read is no longer wanted"
390            );
391            indexing = self
392                .indexing
393                .1
394                .wait_timeout(indexing, std::time::Duration::from_millis(100))
395                .unwrap_or_else(|e| e.into_inner())
396                .0;
397        }
398        drop(indexing);
399        polars_ensure!(!self.shrank(), ComputeError: "{SHRANK}");
400        polars_ensure!(
401            self.whole(),
402            ComputeError: "the file's lines were not all indexed; open it again"
403        );
404        Ok(())
405    }
406
407    /// Index at least `budget` more bytes of a file indexed in part. Returns whether
408    /// every line is indexed. Holds the lines' write lock for the step only, so rows
409    /// are read between steps.
410    pub fn index_more(&self, budget: usize) -> bool {
411        if !self.indexing() {
412            return true;
413        }
414        let _step = self.stepping.lock().unwrap_or_else(|e| e.into_inner());
415        let whole = self.files.iter().all(|f| {
416            // Read with no lock held, so rows go on being read meanwhile, then taken
417            // in under the write lock: this is the only writer of a file not followed.
418            let (bytes, mut step) = {
419                let mapped = f.mapped.read().unwrap_or_else(|e| e.into_inner());
420                if mapped.index.whole(mapped.bytes.as_slice()) {
421                    return true;
422                }
423                // A file cut short meanwhile is not read past its end, and the lines
424                // so far are not taken for all of them. Checked before every step: a
425                // truncation inside one step can still fault the map (SIGBUS), which
426                // only a copy of the file would rule out.
427                if mapped.bytes.still_whole().is_err() {
428                    self.shrank
429                        .store(true, std::sync::atomic::Ordering::Relaxed);
430                    return true;
431                }
432                let bytes = mapped.bytes.clone();
433                let step = mapped.index.step_after(bytes.len());
434                (bytes, step)
435            };
436            let whole = step.extend_by(bytes.as_slice(), budget);
437            f.mapped
438                .write()
439                .unwrap_or_else(|e| e.into_inner())
440                .index
441                .take_step(step);
442            whole
443        });
444        if whole {
445            self.stop_indexing();
446        }
447        whole
448    }
449
450    /// Bytes indexed, and the bytes there are.
451    pub fn indexed_bytes(&self) -> (u64, u64) {
452        self.files.iter().fold((0, 0), |(done, all), f| {
453            let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
454            let len = m.bytes.as_slice().len();
455            let indexed = if m.index.whole(m.bytes.as_slice()) {
456                len
457            } else {
458                m.index.end()
459            };
460            (done + indexed as u64, all + len as u64)
461        })
462    }
463
464    /// The files at `paths`, mapped and indexed. `follow` reads one on as it grows.
465    pub fn open(paths: &[PathBuf], follow: bool) -> Result<Self> {
466        Self::open_first(paths, follow, usize::MAX)
467    }
468
469    /// [`Self::open`], a file not followed indexed through its first `first` bytes.
470    pub fn open_first(paths: &[PathBuf], follow: bool, first: usize) -> Result<Self> {
471        let parts = paths
472            .iter()
473            .map(|path| {
474                let bytes =
475                    Bytes::map(path).map_err(|e| crate::error_display::in_file(path, e.into()))?;
476                Ok((file_name(path), Arc::new(bytes)))
477            })
478            .collect::<Result<Vec<_>>>()?;
479        let first = if follow { usize::MAX } else { first };
480        let mut lines = Self::from_bytes_first(parts, first);
481        if follow && let ([file], [path]) = (lines.files.as_mut_slice(), paths) {
482            file.grows = Some(path.clone());
483        }
484        Ok(lines)
485    }
486
487    /// Whether lines are read on as the file grows.
488    pub fn grows(&self) -> bool {
489        self.files.iter().any(|f| f.grows.is_some())
490    }
491
492    /// Rows on hand.
493    pub fn rows(&self) -> usize {
494        self.files
495            .iter()
496            .map(LineFile::rows)
497            .sum::<usize>()
498            .min(crate::row_index::MAX_ROWS)
499    }
500
501    /// Rows of lines that end in a newline.
502    fn complete_rows(&self) -> usize {
503        self.files
504            .iter()
505            .map(|f| {
506                f.mapped
507                    .read()
508                    .unwrap_or_else(|e| e.into_inner())
509                    .index
510                    .complete()
511            })
512            .sum()
513    }
514
515    /// What the files hold that the rows do not say as they are: lines not UTF-8,
516    /// lines left out.
517    fn counts(&self) -> (usize, usize) {
518        self.files.iter().fold((0, 0), |(invalid, past), f| {
519            let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
520            (invalid + m.index.invalid, past + m.index.past_limit)
521        })
522    }
523
524    /// The frame: a row index decoded a column at a time. A followed file's frame has
525    /// the rows of its complete lines, moved as it grows ([`bound`]).
526    pub fn lazy(self: &Arc<Self>) -> LazyFrame {
527        if self.indexing() {
528            // Every line, once they are all indexed: the frame's height is taken when
529            // it runs, on a worker, which waits for the indexing. The first rows come
530            // from the window ([`opened`]), which does not wait.
531            let lines = self.clone();
532            let height = DataFrame::empty_with_height(0).lazy().map(
533                move |_| {
534                    lines.wait_indexed()?;
535                    Ok(DataFrame::empty_with_height(lines.rows()))
536                },
537                AllowedOptimizations::empty(),
538                None,
539                Some("every line"),
540            );
541            return crate::row_index::lazy_numbered_over(self, height);
542        }
543        let height = if self.grows() {
544            self.complete_rows()
545        } else {
546            self.rows()
547        };
548        crate::row_index::lazy_numbered(self, height)
549    }
550
551    /// For each row, its file and the line in it, in order.
552    fn place(&self, rows: &[IdxSize]) -> PolarsResult<Vec<(usize, usize)>> {
553        let mut starts = Vec::with_capacity(self.files.len());
554        let mut total = 0usize;
555        for f in &self.files {
556            starts.push(total);
557            total += f.rows();
558        }
559        rows.iter()
560            .map(|&r| {
561                let r = r as usize;
562                polars_ensure!(r < total, OutOfBounds: "row {r} is past the {total} lines on hand");
563                let file = starts.partition_point(|&s| s <= r) - 1;
564                Ok((file, r - starts[file]))
565            })
566            .collect()
567    }
568
569    fn column(&self, column: usize, rows: &[IdxSize]) -> PolarsResult<Column> {
570        // A followed file is read on when a row is past its last whole line.
571        if self.grows()
572            && let Some(max) = rows.iter().max()
573            && *max as usize >= self.complete_rows()
574        {
575            for f in &self.files {
576                f.grow()?;
577            }
578        }
579        let placed = self.place(rows)?;
580        let several = self.files.len() > 1;
581        let name = self
582            .schema
583            .get_at_index(column)
584            .expect("a column")
585            .0
586            .clone();
587        let column = if several { column } else { column + 1 };
588        let series = match column {
589            0 => placed
590                .iter()
591                .map(|&(f, _)| Some(self.files[f].name.as_str()))
592                .collect::<StringChunked>()
593                .into_series(),
594            _ => {
595                let held: Vec<_> = self
596                    .files
597                    .iter()
598                    .map(|f| f.mapped.read().unwrap_or_else(|e| e.into_inner()))
599                    .collect();
600                if !self.grows() {
601                    for m in &held {
602                        m.bytes.still_whole()?;
603                    }
604                }
605                let mut builder = StringChunkedBuilder::new(name.clone(), placed.len());
606                for &(f, i) in &placed {
607                    let m = &held[f];
608                    let bytes = m.bytes.as_slice();
609                    match std::str::from_utf8(m.index.line(bytes, i)) {
610                        Ok(text) => builder.append_value(text),
611                        Err(_) => builder
612                            .append_value(String::from_utf8_lossy(m.index.line(bytes, i)).as_ref()),
613                    }
614                }
615                builder.finish().into_series()
616            }
617        };
618        Ok(series.with_name(name).into_column())
619    }
620
621    /// Rows `[start, start + len)`, decoded now.
622    pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
623        let rows = self.rows();
624        let start = start.min(rows);
625        let len = len.min(rows - start);
626        let index: Vec<IdxSize> = (start..start + len).map(|r| r as IdxSize).collect();
627        let mut columns = (0..self.schema.len())
628            .map(|c| self.column(c, &index))
629            .collect::<PolarsResult<Vec<_>>>()?;
630        // Each row's place, as the frame carries it.
631        columns.push(IdxCa::from_vec(crate::row_index::INDEX.into(), index).into_column());
632        DataFrame::new(len, columns)
633    }
634}
635
636impl crate::row_index::RowSource for Lines {
637    fn height(&self) -> usize {
638        self.rows()
639    }
640
641    fn schema(&self) -> SchemaRef {
642        self.schema.clone()
643    }
644
645    fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
646        polars_ensure!(
647            index.null_count() == 0,
648            ComputeError: "a row index has a missing row"
649        );
650        let rows: Vec<IdxSize> = index.into_no_null_iter().collect();
651        self.column(column, &rows)
652    }
653}
654
655impl crate::pushdown::Windowed for Lines {
656    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
657        Ok(self.collect_window(start, len)?.lazy())
658    }
659}
660
661fn file_name(path: &Path) -> String {
662    path.file_name().map_or_else(
663        || path.display().to_string(),
664        |n| n.to_string_lossy().into_owned(),
665    )
666}
667
668/// Moves the height of a followed file's lines in `plan` to `rows`, and says whether
669/// it found them. The frame is a row index over a frame of no columns
670/// ([`crate::row_index::lazy_with_height`]), so its height is the rows it has.
671pub(crate) fn bound(plan: &mut polars::lazy::dsl::DslPlan, rows: IdxSize) -> bool {
672    use polars::lazy::dsl::DslPlan;
673    match plan {
674        DslPlan::DataFrameScan { df, .. } if df.width() == 0 => {
675            *df = Arc::new(DataFrame::empty_with_height(rows as usize));
676            true
677        }
678        _ => false,
679    }
680}
681
682/// The lines of `input`, for its reader.
683fn scan(input: crate::readers::ScanIn<'_>) -> Result<crate::scan::Scan> {
684    let lines = Arc::new(Lines::open_first(
685        input.paths,
686        input.options.follow,
687        FIRST_BYTES,
688    )?);
689    let lf = lines.lazy();
690    input.report.opened = Some(Arc::new(opened(&lines, input.options)));
691    Ok(lf.into())
692}
693
694/// What the Info panel says of lines read, the window a view reads straight from them
695/// when it can, and the lines themselves while they are still being indexed.
696pub(crate) fn opened(lines: &Arc<Lines>, options: &crate::OpenOptions) -> crate::members::Opened {
697    crate::members::Opened {
698        // A followed file's rows are the watcher's to count.
699        window: (!lines.grows()).then(|| {
700            (
701                lines.clone() as Arc<dyn crate::pushdown::Windowed>,
702                lines.rows(),
703            )
704        }),
705        notes: notes(lines, options.format_guessed),
706        indexing: lines.indexing().then(|| lines.clone()),
707        numbering: lines.several().then(|| lines.clone()),
708        ..Default::default()
709    }
710}
711
712/// How the note that the format was guessed begins.
713pub(crate) const GUESSED: &str = "no format detected";
714
715/// What the Info panel says of lines read, as far as they are indexed.
716pub(crate) fn notes(lines: &Lines, format_guessed: bool) -> Vec<crate::notes::Note> {
717    let (invalid, past_limit) = lines.counts();
718    let scope = match lines.files.len() {
719        1 => "the file".to_string(),
720        n => format!("the {n} files"),
721    };
722    let middot = crate::glyphs::get().middot;
723    let mut notes = vec![format!(
724        "read as lines {middot} a row per line, blank lines included"
725    )];
726    if format_guessed {
727        notes.push(format!("{GUESSED} {middot} --format csv reads it as CSV"));
728    }
729    if invalid > 0 {
730        notes.push(format!(
731            "{} with invalid UTF-8 {middot} shown as \u{fffd}",
732            crate::text_formats::count(invalid as u64, "line", "lines")
733        ));
734    }
735    if past_limit > 0 {
736        notes.push(format!(
737            "{} left out: past the first {}",
738            crate::text_formats::count(past_limit as u64, "line", "lines"),
739            crate::numfmt::group_chrome(crate::indexed::MAX_RECORDS)
740        ));
741    }
742    notes
743        .into_iter()
744        .map(|n| crate::text_formats::note(n, scope.clone()))
745        .collect()
746}
747
748// --- What text is -----------------------------------------------------------------
749
750/// Complete records looked at for a field count.
751const MOST_RECORDS: usize = 50;
752
753/// What text no format's signature claims is, from its first bytes `head`; `whole`
754/// when they are all there is, so a last line without a newline is complete. `None`
755/// for bytes that are not text, which a local file shows in the hex view.
756///
757/// JSON only when the bytes parse as JSON so far; NDJSON when the first line is an
758/// object. CSV or TSV only when several complete records have one field count, two or
759/// more, with quotes where CSV allows them. Anything else is lines.
760pub fn guess(head: &[u8], whole: bool) -> Option<FileFormat> {
761    let text = head.strip_prefix(b"\xef\xbb\xbf").unwrap_or(head);
762    if !is_text(text, whole) {
763        return None;
764    }
765    let start = text
766        .iter()
767        .position(|b| !b.is_ascii_whitespace())
768        .unwrap_or(text.len());
769    let trimmed = &text[start..];
770    match trimmed.first() {
771        Some(b'[') if json_so_far(trimmed, whole) => return Some(FileFormat::Json),
772        Some(b'{') => {
773            let first = trimmed.split(|&b| b == b'\n').next().unwrap_or_default();
774            let line_whole = first.len() < trimmed.len() || whole;
775            if line_whole
776                && matches!(
777                    serde_json::from_slice::<serde_json::Value>(first),
778                    Ok(serde_json::Value::Object(_))
779                )
780            {
781                return Some(FileFormat::Jsonl);
782            }
783            if json_so_far(trimmed, whole) {
784                return Some(FileFormat::Json);
785            }
786        }
787        _ => {}
788    }
789    let mut best: Option<(FileFormat, usize)> = None;
790    for (separator, format) in [(b',', FileFormat::Csv), (b'\t', FileFormat::Tsv)] {
791        if let Some(fields) = consistent_fields(text, separator, whole)
792            && best.is_none_or(|(_, most)| fields > most)
793        {
794            best = Some((format, fields));
795        }
796    }
797    Some(best.map_or(FileFormat::Text, |(format, _)| format))
798}
799
800/// A guess of lines turned to CSV when the user gave a delimited dialect
801/// (`--delimiter`, `--comment`, `--header-rows`, `--skip-initial-space`): they
802/// said the text is delimited, in a shape the guess does not try.
803pub fn as_asked(format: FileFormat, options: &crate::OpenOptions) -> FileFormat {
804    let delimited = options.delimiter.is_some()
805        || options.comment_char.is_some()
806        || options.header_rows().is_some()
807        || options.skip_initial_space;
808    match format {
809        FileFormat::Text if delimited => FileFormat::Csv,
810        other => other,
811    }
812}
813
814/// [`guess`] for the local file at `path`, through its compression. What is inside a
815/// compressed file is read through it only as delimited text or lines.
816pub fn guess_file(
817    path: &Path,
818    compression: Option<crate::CompressionFormat>,
819) -> Option<FileFormat> {
820    let compression = compression.or_else(|| crate::CompressionFormat::from_extension(path));
821    let reach = crate::readers::HEAD as u64;
822    let head = crate::formats::head_of(path, compression, reach)?;
823    // An empty file holds no lines to show; the hex view says it is empty.
824    if head.is_empty() {
825        return None;
826    }
827    let whole = match compression {
828        Some(_) => (head.len() as u64) < reach,
829        None => std::fs::metadata(path).is_ok_and(|m| m.len() <= head.len() as u64),
830    };
831    guess(&head, whole).map(|f| match compression {
832        Some(_) if !f.decompressed_once() => FileFormat::Text,
833        _ => f,
834    })
835}
836
837/// Whether `bytes` read as text: no NUL, and at most one byte in ten a control
838/// character or not UTF-8. A character cut off at the end of a head is not counted.
839fn is_text(bytes: &[u8], whole: bool) -> bool {
840    if bytes.contains(&0) {
841        return false;
842    }
843    let mut odd = 0usize;
844    let mut chunks = bytes.utf8_chunks().peekable();
845    while let Some(chunk) = chunks.next() {
846        odd += chunk
847            .valid()
848            .chars()
849            .filter(|c| {
850                c.is_control() && !matches!(c, '\t' | '\n' | '\r' | '\x0c' | '\x1b' | '\x08')
851            })
852            .count();
853        let cut_off = !whole && chunks.peek().is_none();
854        if !cut_off {
855            odd += chunk.invalid().len();
856        }
857    }
858    odd * 10 <= bytes.len()
859}
860
861/// Whether `bytes` are JSON, or the start of it when they are not `whole`.
862fn json_so_far(bytes: &[u8], whole: bool) -> bool {
863    match serde_json::from_slice::<serde::de::IgnoredAny>(bytes) {
864        Ok(_) => true,
865        Err(e) => !whole && e.is_eof(),
866    }
867}
868
869/// The field count every complete record of `text` has when split on `separator`,
870/// quotes respected, if they agree, have two fields or more, and are enough records to
871/// say so: three, or two when that is all there is. Two fields split at a comma that
872/// is always followed by a space is prose (`ok, then`), not evidence.
873fn consistent_fields(text: &[u8], separator: u8, whole: bool) -> Option<usize> {
874    let (mut counts, last, spaced) = field_counts(text, separator)?;
875    // A last line with no newline is a record of a whole input when it agrees: a
876    // stream caught mid-line is not refused for it.
877    if whole
878        && let Some(last) = last
879        && counts.first().is_none_or(|&first| first == last)
880    {
881        counts.push(last);
882    }
883    let first = *counts.first()?;
884    let enough = if whole {
885        counts.len() >= 2
886    } else {
887        counts.len() >= 3
888    };
889    let prose = separator == b',' && first == 2 && spaced;
890    (enough && first >= 2 && !prose && counts.iter().all(|&n| n == first)).then_some(first)
891}
892
893/// Each complete, non-blank record's field count, the last record's when it has no
894/// newline, and whether every separator is followed by a space; `None` when a quote is
895/// where CSV does not allow one.
896fn field_counts(text: &[u8], separator: u8) -> Option<(Vec<usize>, Option<usize>, bool)> {
897    let mut counts = Vec::new();
898    let mut spaced = true;
899    let mut fields = 1usize;
900    let mut blank = true;
901    let mut field_start = true;
902    let mut in_quotes = false;
903    let mut i = 0;
904    while i < text.len() && counts.len() < MOST_RECORDS {
905        let b = text[i];
906        i += 1;
907        if in_quotes {
908            if b == b'"' {
909                if text.get(i) == Some(&b'"') {
910                    i += 1;
911                } else {
912                    in_quotes = false;
913                    // A closing quote ends its field.
914                    match text.get(i) {
915                        Some(&next) if next == separator || next == b'\n' || next == b'\r' => {}
916                        None => {}
917                        Some(_) => return None,
918                    }
919                }
920            }
921            continue;
922        }
923        match b {
924            b'"' if field_start => {
925                in_quotes = true;
926                field_start = false;
927                blank = false;
928            }
929            b'\n' => {
930                if !blank {
931                    counts.push(fields);
932                }
933                fields = 1;
934                blank = true;
935                field_start = true;
936            }
937            b'\r' if text.get(i) == Some(&b'\n') => {}
938            _ if b == separator => {
939                spaced &= text.get(i) == Some(&b' ');
940                fields += 1;
941                blank = false;
942                field_start = true;
943            }
944            _ => {
945                if !b.is_ascii_whitespace() {
946                    blank = false;
947                }
948                field_start = false;
949            }
950        }
951    }
952    // A quoted field still open is cut off by the head, or never closed: not a record.
953    let last = (!in_quotes && !blank && counts.len() < MOST_RECORDS).then_some(fields);
954    Some((counts, last, spaced))
955}
956
957#[cfg(test)]
958mod tests {
959    use super::*;
960
961    fn lines_of(bytes: &[u8]) -> Vec<String> {
962        let lines = Arc::new(Lines::from_bytes(vec![(
963            "a.log".to_string(),
964            Arc::new(Bytes::Owned(bytes.to_vec())),
965        )]));
966        let df = lines.lazy().collect().unwrap();
967        let places: Vec<u32> = df
968            .column(crate::row_index::INDEX)
969            .unwrap()
970            .u32()
971            .unwrap()
972            .into_no_null_iter()
973            .collect();
974        assert_eq!(places, (0..df.height() as u32).collect::<Vec<_>>());
975        df.column(LINE)
976            .unwrap()
977            .str()
978            .unwrap()
979            .iter()
980            .map(|v| v.unwrap().to_string())
981            .collect()
982    }
983
984    /// Every line as written: blank lines kept, `\r\n` an ending, a lone `\r` and
985    /// separators, quotes and control characters kept, bytes not UTF-8 shown lossily.
986    #[test]
987    fn every_line_as_written() {
988        assert_eq!(lines_of(b"a\n\nb\n\n\nc\n"), ["a", "", "b", "", "", "c"]);
989        assert_eq!(lines_of(b"a\r\nb\r\n\r\nc"), ["a", "b", "", "c"]);
990        assert_eq!(
991            lines_of(b"x\x1fy,\"z\tw\n#c\x00\x1b[1m\r\nq\rr\n"),
992            ["x\x1fy,\"z\tw", "#c\0\x1b[1m", "q\rr"]
993        );
994        assert_eq!(
995            lines_of(b"ok\n\xff\xfe bad\n"),
996            ["ok", "\u{fffd}\u{fffd} bad"]
997        );
998        assert_eq!(lines_of(b"\n\n"), ["", ""]);
999        assert!(lines_of(b"").is_empty());
1000    }
1001
1002    /// An index read on as its bytes grow matches one read whole, the partial last
1003    /// line taken back first.
1004    #[test]
1005    fn an_index_reads_on_as_bytes_grow() {
1006        let all = b"one\ntw\xffo\npartial line\n\nlast";
1007        for cut in 0..all.len() {
1008            let mut index = LineIndex::of(&all[..cut]);
1009            index.extend(all);
1010            let whole = LineIndex::of(all);
1011            assert_eq!(index.lines(), whole.lines(), "{cut}");
1012            assert_eq!(index.invalid, whole.invalid, "{cut}");
1013            assert_eq!(index.complete(), 4);
1014            for i in 0..whole.lines() {
1015                assert_eq!(index.line(all, i), whole.line(all, i));
1016            }
1017        }
1018    }
1019
1020    #[test]
1021    fn several_files_name_theirs() {
1022        let lines = Arc::new(Lines::from_bytes(vec![
1023            ("a.log".into(), Arc::new(Bytes::Owned(b"1\n2\n".to_vec()))),
1024            ("b.log".into(), Arc::new(Bytes::Owned(b"3\n".to_vec()))),
1025        ]));
1026        let df = lines.lazy().collect().unwrap();
1027        assert_eq!(df.get_column_names(), [FILE, LINE, crate::row_index::INDEX]);
1028        let file: Vec<Option<&str>> = df.column(FILE).unwrap().str().unwrap().iter().collect();
1029        assert_eq!(file, [Some("a.log"), Some("a.log"), Some("b.log")]);
1030        let w = lines.collect_window(1, 5).unwrap();
1031        assert_eq!(w.height(), 2);
1032        assert_eq!(w.get_column_names(), df.get_column_names());
1033        let places: Vec<u32> = w
1034            .column(crate::row_index::INDEX)
1035            .unwrap()
1036            .u32()
1037            .unwrap()
1038            .into_no_null_iter()
1039            .collect();
1040        assert_eq!(places, [1, 2]);
1041    }
1042
1043    /// An index taken a chunk at a time, in steps of any size, is the one taken whole:
1044    /// the same lines, the same count not UTF-8, the last line with no newline kept.
1045    #[test]
1046    fn an_index_in_steps_is_the_index_whole() {
1047        let mut bytes = Vec::new();
1048        let mut bad = 0;
1049        for i in 0..200_000u32 {
1050            match i % 7 {
1051                0 => bytes.extend_from_slice(b"\r\n"),
1052                3 => {
1053                    bytes.extend_from_slice(b"bad \xff\xfe line\n");
1054                    bad += 1;
1055                }
1056                _ => bytes.extend_from_slice(
1057                    format!("line {i} {}\n", "x".repeat((i % 50) as usize)).as_bytes(),
1058                ),
1059            }
1060        }
1061        bytes.extend_from_slice(b"no newline at the end \xff");
1062        let whole = LineIndex::of(&bytes);
1063        assert_eq!(whole.lines(), 200_001);
1064        assert_eq!(whole.invalid, bad + 1);
1065        for step in [1, 4096, CHUNK - 3, CHUNK * 2 + 17] {
1066            let mut index = LineIndex {
1067                offsets: Offsets::for_file(bytes.len()),
1068                ..Default::default()
1069            };
1070            let mut steps = 0;
1071            while !index.extend_by(&bytes, step) {
1072                steps += 1;
1073                assert!(index.lines() < whole.lines(), "{step}");
1074            }
1075            assert!(steps > 0 || step > bytes.len(), "{step}");
1076            assert_eq!(index.lines(), whole.lines(), "{step}");
1077            assert_eq!(index.invalid, whole.invalid, "{step}");
1078            assert_eq!(index.complete(), whole.complete(), "{step}");
1079            for i in [0, 1, 3, 1000, whole.lines() - 1] {
1080                assert_eq!(index.line(&bytes, i), whole.line(&bytes, i), "{step} {i}");
1081            }
1082        }
1083    }
1084
1085    /// Indexing stopped before every line is in (the dataset went) gives a read
1086    /// waiting on it an error, not the lines so far as all of them; paused and taken up
1087    /// again, it reads on.
1088    #[test]
1089    fn a_stopped_index_is_no_count_and_a_paused_one_reads_on() {
1090        let bytes: Vec<u8> = (0..50_000u32)
1091            .flat_map(|i| format!("{i}\n").into_bytes())
1092            .collect();
1093        let lines = Arc::new(Lines::from_bytes_first(
1094            vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1095            1000,
1096        ));
1097        let lf = lines.lazy();
1098        let waiting = {
1099            let lf = lf.clone();
1100            std::thread::spawn(move || lf.collect())
1101        };
1102        // Paused: the flag stays, the read waits; taken up again, it reads on.
1103        assert!(!lines.index_more(1000));
1104        assert!(lines.resume_indexing());
1105        while !lines.index_more(100_000) {}
1106        assert_eq!(waiting.join().unwrap().unwrap().height(), 50_000);
1107
1108        let bytes: Vec<u8> = (0..50_000u32)
1109            .flat_map(|i| format!("{i}\n").into_bytes())
1110            .collect();
1111        let lines = Arc::new(Lines::from_bytes_first(
1112            vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1113            1000,
1114        ));
1115        let lf = lines.lazy();
1116        let waiting = std::thread::spawn(move || lf.collect());
1117        lines.stop_indexing();
1118        let error = waiting.join().unwrap().unwrap_err().to_string();
1119        assert!(error.contains("not all indexed"), "{error}");
1120        assert!(!lines.resume_indexing() || !lines.whole());
1121    }
1122
1123    /// A file that shrinks while it is indexed stops the indexing where it is, and a
1124    /// read of every line says so rather than taking the lines so far for all. Unix
1125    /// only: Windows refuses to shorten a file another handle has mapped (os error
1126    /// 1224), so a rotation there fails in the rotating process, not here.
1127    #[test]
1128    #[cfg(unix)]
1129    fn a_file_that_shrinks_while_indexed_has_no_count() {
1130        let dir = tempfile::tempdir().unwrap();
1131        let path = dir.path().join("rotated.log");
1132        let bytes: Vec<u8> = (0..50_000u32)
1133            .flat_map(|i| format!("{i}\n").into_bytes())
1134            .collect();
1135        std::fs::write(&path, &bytes).unwrap();
1136        let lines = Arc::new(Lines::open_first(std::slice::from_ref(&path), false, 1000).unwrap());
1137        assert!(lines.indexing());
1138        // copytruncate, between two steps.
1139        std::fs::OpenOptions::new()
1140            .write(true)
1141            .open(&path)
1142            .unwrap()
1143            .set_len(10)
1144            .unwrap();
1145        assert!(lines.index_more(100_000), "the indexing stops");
1146        assert!(lines.shrank());
1147        let error = lines.lazy().collect().unwrap_err().to_string();
1148        // Whichever reads first says the file is shorter than it was.
1149        assert!(
1150            error.contains(SHRANK) || error.contains("shorter"),
1151            "{error}"
1152        );
1153    }
1154
1155    /// Several files' rows are numbered by their line in their own file.
1156    #[test]
1157    fn several_files_number_their_own_lines() {
1158        let lines = Lines::from_bytes(vec![
1159            ("a.log".into(), Arc::new(Bytes::Owned(b"1\n2\n".to_vec()))),
1160            (
1161                "b.log".into(),
1162                Arc::new(Bytes::Owned(b"3\n4\n5\n".to_vec())),
1163            ),
1164        ]);
1165        assert!(lines.several());
1166        let numbered: Vec<Option<usize>> = (0..6).map(|p| lines.line_in_file(p)).collect();
1167        assert_eq!(
1168            numbered,
1169            [Some(0), Some(1), Some(0), Some(1), Some(2), None]
1170        );
1171    }
1172
1173    /// A large file shows its first lines indexed, and indexes the rest in steps; its
1174    /// frame reads every line once they are in, numbered by their place, on either
1175    /// engine, and a read of it started meanwhile waits for them.
1176    #[test]
1177    fn a_large_file_indexes_behind_its_first_rows() {
1178        let mut bytes = Vec::new();
1179        for i in 0..100_000u32 {
1180            bytes.extend_from_slice(format!("{i}\n").as_bytes());
1181        }
1182        let lines = Arc::new(Lines::from_bytes_first(
1183            vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1184            1000,
1185        ));
1186        assert!(lines.indexing());
1187        let first = lines.rows();
1188        assert!(first > 0 && first < 100_000, "{first}");
1189        assert_eq!(lines.collect_window(0, 200_000).unwrap().height(), first);
1190        let lf = lines.lazy();
1191        let waiting = {
1192            let lf = lf.clone().filter(col(LINE).eq(lit("99999")));
1193            std::thread::spawn(move || lf.collect().unwrap())
1194        };
1195        while !lines.index_more(50_000) {}
1196        assert!(!lines.indexing());
1197        assert_eq!(lines.rows(), 100_000);
1198        let (done, all) = lines.indexed_bytes();
1199        assert_eq!(done, all);
1200        let df = waiting.join().unwrap();
1201        assert_eq!(df.height(), 1);
1202        assert_eq!(
1203            df.column(crate::row_index::INDEX)
1204                .unwrap()
1205                .u32()
1206                .unwrap()
1207                .get(0),
1208            Some(99_999)
1209        );
1210        for streaming in [false, true] {
1211            let df = crate::statistics::collect_lazy(lf.clone(), streaming).unwrap();
1212            assert_eq!(df.height(), 100_000, "streaming {streaming}");
1213        }
1214    }
1215
1216    #[test]
1217    fn delimited_only_on_evidence() {
1218        let cases: [(&[u8], bool, Option<FileFormat>); 18] = [
1219            (b"a,b,c\n1,2,3\n4,5,6\n", true, Some(FileFormat::Csv)),
1220            (b"a,b\n1,2\n", true, Some(FileFormat::Csv)),
1221            (b"a\tb\n1\t2\n3\t4\n", true, Some(FileFormat::Tsv)),
1222            (
1223                b"name,note\n\"x\",\"a, b\nc\"\ny,z\n",
1224                true,
1225                Some(FileFormat::Csv),
1226            ),
1227            // A cut line is not a record.
1228            (b"a,b,c\n1,2,3\n4,5,6\n7,8", false, Some(FileFormat::Csv)),
1229            (b"a,b,c\n1,2,3\n", false, Some(FileFormat::Text)),
1230            // Logs: commas and tabs that do not line up.
1231            (
1232                b"2024-01-01 INFO started, ok\n2024-01-01 WARN slow, very, slow\nINFO done\n",
1233                true,
1234                Some(FileFormat::Text),
1235            ),
1236            (
1237                b"[2024-01-01 12:00] INFO hi\n[2024-01-01 12:01] INFO there\n",
1238                true,
1239                Some(FileFormat::Text),
1240            ),
1241            (b"[INFO] a\n", true, Some(FileFormat::Text)),
1242            (
1243                b"say \"hi\", she said\nok, then\nno, way\n",
1244                true,
1245                Some(FileFormat::Text),
1246            ),
1247            (b"one line, with a comma\n", true, Some(FileFormat::Text)),
1248            (b"[1, 2,", false, Some(FileFormat::Json)),
1249            (b"[1, 2]", true, Some(FileFormat::Json)),
1250            (b"{\"a\": 1}\n{\"a\": 2}\n", true, Some(FileFormat::Jsonl)),
1251            (b"{\n  \"a\": 1\n}\n", true, Some(FileFormat::Json)),
1252            (b"{not json}\n", true, Some(FileFormat::Text)),
1253            (b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR", false, None),
1254            (b"\xef\xbb\xbfid,v\n1,2\n", true, Some(FileFormat::Csv)),
1255        ];
1256        for (head, whole, format) in cases {
1257            assert_eq!(
1258                guess(head, whole),
1259                format,
1260                "{}",
1261                String::from_utf8_lossy(head)
1262            );
1263        }
1264        // Latin-1 is still text; a run of random bytes is not.
1265        assert_eq!(guess(b"caf\xe9 au lait\n", true), Some(FileFormat::Text));
1266        let noise: Vec<u8> = (0..512u32).map(|i| (i * 97 % 255 + 1) as u8).collect();
1267        assert_eq!(guess(&noise, true), None);
1268    }
1269}