Skip to main content

datui_lib/formats/
lines.rs

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