Skip to main content

datui_lib/
follow.rs

1//! Following a file that is still being written, as `tail -f` does: `--follow` and
2//! `t` at the table.
3//!
4//! The frame scans the file as any delimited or NDJSON file is scanned, with a slice
5//! right above the scan that bounds it to the rows whose records are complete
6//! ([`bound`]). A watcher thread ([`Follow`]) checks the file's size every interval,
7//! reads only the bytes that arrived, counts the records they complete ([`Tail`]), and
8//! sends what it found as [`AppEvent::Followed`]. The app then moves the slice in every
9//! frame the view holds, so the query, filters and sort run over the new rows with no
10//! frame rebuilt, and reads the window on screen. A partial last line is never in the
11//! bound: it waits for its newline.
12//!
13//! Standard input followed (`datui -f -`) is spooled to a file by a [`Spool`] that goes
14//! on copying after the first rows show; the file is followed like any other.
15//!
16//! An Arrow IPC stream is followed the same way, its record batches counted in place
17//! of records and read by a scan of its own ([`stream`]).
18
19use std::fs::File;
20use std::io::{Read, Seek, SeekFrom, Write};
21use std::path::{Path, PathBuf};
22use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
23use std::sync::mpsc::Sender;
24use std::sync::{Arc, Condvar, Mutex};
25use std::time::{Duration, Instant};
26
27use polars::prelude::*;
28
29use crate::download::TempDownload;
30use crate::unfinished::Writer;
31use crate::{AppEvent, CompressionFormat, FileFormat, OpenOptions};
32
33pub(crate) mod lines;
34#[cfg(target_os = "linux")]
35mod notify;
36pub(crate) mod stream;
37#[doc(hidden)]
38pub use stream::stream_messages;
39
40/// How often the watcher checks the file, or where it hears of changes (Linux) the least
41/// time between two reads, unless `[read] follow_interval` says otherwise. A burst of
42/// appends inside one interval is one refresh.
43pub const DEFAULT_INTERVAL: Duration = Duration::from_millis(250);
44
45/// Bytes read from the file per step while counting records.
46const CHUNK: usize = 1 << 20;
47
48/// A record longer than this is kept only in part: enough to classify it.
49const LONGEST_RECORD: usize = 16 << 20;
50
51/// A row's start is marked once this many rows, or this many bytes, have passed since
52/// the last mark: a window is read from the mark before it, so it costs at most this
53/// much beyond its own rows at any file size.
54pub(crate) const MARK_ROWS: u64 = 8192;
55const MARK_BYTES: u64 = 1 << 20;
56
57/// The most bytes a read from a mark takes in. A view that needs more (a filtered
58/// window far behind the last count, a count after a long pause) is read by Polars
59/// from the start of the file as before.
60const MOST_FROM_A_MARK: u64 = 64 << 20;
61
62/// Why `format` cannot be followed, or `None` when it can. Only text read line by
63/// line, and an Arrow IPC stream (see [`followed_stream`]), can: a file whose footer is
64/// written last (Parquet, an Arrow IPC file, Excel) cannot be read before it is
65/// finished, and a compressed one cannot be read from the middle.
66pub fn refusal(format: Option<FileFormat>, options: &OpenOptions) -> Option<String> {
67    let format = format.unwrap_or(FileFormat::TEXT);
68    if format == FileFormat::Arrow {
69        return Some(
70            "An Arrow IPC file is read once it is finished; an Arrow IPC stream can be followed."
71                .to_string(),
72        );
73    }
74    if !format.follows() {
75        let followed: Vec<&str> = FileFormat::ALL
76            .into_iter()
77            .filter(|f| f.follows())
78            .map(FileFormat::title)
79            .collect();
80        let followed = match followed.split_last() {
81            Some((last, rest)) if !rest.is_empty() => format!("{} and {last}", rest.join(", ")),
82            _ => followed.join(""),
83        };
84        return Some(format!(
85            "Only {followed} and Arrow IPC streams can be followed as they grow; {} is read \
86             once it is finished.",
87            format.title()
88        ));
89    }
90    if options.compression.is_some() {
91        return Some("A compressed file cannot be followed as it grows.".to_string());
92    }
93    if options.header_rows().is_some() || options.skip_tail_rows.is_some() {
94        return Some(
95            "A file read with --header-rows or --footer-rows cannot be followed.".to_string(),
96        );
97    }
98    if options.spec_name.is_some() || options.spec_file.is_some() || options.delimited.is_some() {
99        return Some("A file read through a format spec cannot be followed.".to_string());
100    }
101    None
102}
103
104/// Whether `path`'s format, as `options` say or its name does, is one `--follow` reads.
105pub fn followable_path(path: &Path, options: &OpenOptions) -> bool {
106    let format = options.format.or_else(|| {
107        path.extension()
108            .and_then(|e| e.to_str())
109            .and_then(FileFormat::from_extension)
110    });
111    let compression = options
112        .compression
113        .or_else(|| CompressionFormat::from_extension(path));
114    compression.is_none()
115        && (refusal(format, options).is_none() || followed_stream(path, format, options))
116}
117
118/// Whether `path`, read as `format`, is an Arrow IPC stream `--follow` reads as it grows,
119/// by its contents.
120pub(crate) fn followed_stream(
121    path: &Path,
122    format: Option<FileFormat>,
123    options: &OpenOptions,
124) -> bool {
125    format == Some(FileFormat::Arrow)
126        && options.compression.is_none()
127        && options.spec_name.is_none()
128        && options.spec_file.is_none()
129        && crate::ipc_stream::is_stream_file(path)
130}
131
132/// Why `paths` cannot be followed as `options` ask, before anything is read: only one
133/// local file, or standard input, can be.
134pub fn refuse_paths(paths: &[PathBuf], options: &OpenOptions) -> Option<String> {
135    let [path] = paths else {
136        return Some("Only one file can be followed at a time.".to_string());
137    };
138    if crate::stdin::is_stdin(path) {
139        // What it holds is known once its first bytes are in.
140        return None;
141    }
142    if !matches!(
143        crate::source::input_source(path),
144        crate::source::InputSource::Local(_)
145    ) {
146        return Some("Only a local file or standard input can be followed.".to_string());
147    }
148    if path.is_dir() {
149        return Some("A directory cannot be followed: name a file in it.".to_string());
150    }
151    let format = options.format.or_else(|| {
152        path.extension()
153            .and_then(|e| e.to_str())
154            .and_then(FileFormat::from_extension)
155    });
156    let options = OpenOptions {
157        compression: options
158            .compression
159            .or_else(|| CompressionFormat::from_extension(path)),
160        ..options.clone()
161    };
162    if followed_stream(path, format, &options) {
163        return None;
164    }
165    refusal(format, &options)
166}
167
168/// The format a followed `path` is read as.
169pub(crate) fn format_of(path: &Path, found: Option<FileFormat>) -> FileFormat {
170    found
171        .or_else(|| {
172            path.extension()
173                .and_then(|e| e.to_str())
174                .and_then(FileFormat::from_extension)
175        })
176        .unwrap_or(FileFormat::TEXT)
177}
178
179/// An NDJSON file followed, scanned lazily rather than read whole as an unfollowed one
180/// is: the frame reads more of it as it grows, and only its complete lines
181/// ([`lines::LinesScan`]). The schema comes from the first `infer_schema_length` lines,
182/// or from every line there is when `every_line` (the journal, whose fields vary).
183pub(crate) fn scan_lines(
184    path: &Path,
185    options: &OpenOptions,
186    every_line: bool,
187    read_python: &mut Vec<String>,
188) -> color_eyre::Result<LazyFrame> {
189    let infer = if every_line {
190        None
191    } else {
192        // Polars' own default for an NDJSON scan.
193        Some(
194            options
195                .infer_schema_length
196                .and_then(std::num::NonZeroUsize::new)
197                .unwrap_or(std::num::NonZeroUsize::new(100).expect("not zero")),
198        )
199    };
200    let spool = options.spool.as_ref().map(|handle| handle.spool().clone());
201    let lf = lines::LinesScan::open(path, infer, true, spool)?.lazy()?;
202    crate::widgets::datatable::DataTableState::apply_parse_dates_to_json_lazyframe(
203        lf,
204        options,
205        read_python,
206    )
207}
208
209/// `lf`, the scan of `path` read as `format`, bounded to the rows of the file's complete
210/// records, and the count its watcher reads on from.
211pub(crate) fn bound_to_complete(
212    mut lf: LazyFrame,
213    path: &Path,
214    format: FileFormat,
215    options: &OpenOptions,
216) -> color_eyre::Result<(LazyFrame, Tail)> {
217    let schema = lf.collect_schema()?;
218    let mut tail = Tail::new(format, options, &schema);
219    if options.spool.is_some() {
220        tail = tail.widening();
221    }
222    tail.path = path.to_path_buf();
223    let mut file = File::open(path)?;
224    let len = file.metadata()?.len();
225    tail.read_on(&mut file, len, false)?;
226    bound(&mut lf, path, tail.rows());
227    Ok((lf, tail))
228}
229
230/// What a field of a row has to read as, for a row to fit the schema inferred from the
231/// first rows.
232#[derive(Clone, Copy, Debug, PartialEq)]
233enum Fits {
234    Anything,
235    Integer,
236    Number,
237    Boolean,
238    Text,
239}
240
241impl Fits {
242    fn of(dtype: &DataType) -> Fits {
243        if dtype.is_integer() {
244            Fits::Integer
245        } else if dtype.is_float() {
246            Fits::Number
247        } else if matches!(dtype, DataType::Boolean) {
248            Fits::Boolean
249        } else if matches!(dtype, DataType::String) {
250            Fits::Text
251        } else {
252            Fits::Anything
253        }
254    }
255
256    /// Whether a CSV cell fits. An empty cell or a null value fits anything.
257    fn cell(self, cell: &str, nulls: &[String]) -> bool {
258        let cell = cell.trim();
259        if cell.is_empty() || nulls.iter().any(|n| n == cell) {
260            return true;
261        }
262        match self {
263            Fits::Integer => cell.parse::<i64>().is_ok() || cell.parse::<u64>().is_ok(),
264            Fits::Number => cell.parse::<f64>().is_ok(),
265            Fits::Boolean => {
266                cell.eq_ignore_ascii_case("true") || cell.eq_ignore_ascii_case("false")
267            }
268            Fits::Text | Fits::Anything => true,
269        }
270    }
271
272    /// Whether a JSON value fits.
273    fn value(self, value: &serde_json::Value) -> bool {
274        match self {
275            _ if value.is_null() => true,
276            Fits::Integer => value.is_i64() || value.is_u64(),
277            Fits::Number => value.is_number(),
278            Fits::Boolean => value.is_boolean(),
279            // Polars reads an array into a text column as its JSON text: journalctl's
280            // bytes for a message that is not UTF-8.
281            Fits::Text => value.is_string() || value.is_array(),
282            Fits::Anything => true,
283        }
284    }
285}
286
287/// How the file's records become rows.
288#[derive(Clone, Debug)]
289enum Layout {
290    Delimited {
291        separator: u8,
292        /// Records before the header that are not rows: `--skip-lines` and
293        /// `--skip-rows`.
294        skip: u64,
295        header: bool,
296        comment: Option<Vec<u8>>,
297        nulls: Vec<String>,
298    },
299    Lines,
300    /// An Arrow IPC stream: its record batch messages.
301    Stream,
302    /// Text read as lines: every line a row, blank ones too.
303    Text,
304}
305
306/// Rows a [`Tail`] marked that its follow's [`Marks`] does not have yet, as (row, byte
307/// where its record starts), and the last mark made.
308#[derive(Clone, Debug, Default)]
309struct NewMarks {
310    new: Vec<(u64, u64)>,
311    last: Option<(u64, u64)>,
312}
313
314/// The complete records of a growing delimited or NDJSON file: where they end, and
315/// how many rows they hold. Extended with the bytes that arrive; a partial last record
316/// is not counted until its newline lands.
317#[derive(Clone, Debug)]
318pub struct Tail {
319    /// The file counted.
320    path: PathBuf,
321    layout: Layout,
322    /// Each column's kind, by position for delimited text and by name for NDJSON.
323    columns: Vec<(String, Fits)>,
324    /// Bytes through the end of the last complete record.
325    complete: u64,
326    /// Records counted, rows among them, and the rows that did not fit.
327    records: u64,
328    rows: u64,
329    misfits: u64,
330    /// Fields in the header, which every row should have.
331    fields: Option<usize>,
332    /// Rows marked since the marks were last handed to the follow's [`Marks`].
333    marks: NewMarks,
334    /// How far apart marks are: rows, bytes. Small in tests.
335    mark_every: (u64, u64),
336    /// Standard input's NDJSON: a field the schema does not have is not a misfit but a
337    /// column to come, joined once the stream ends ([`Self::new_fields`]).
338    widens: bool,
339    /// The fields that arrived after the open, in the order they first came, and what
340    /// their values have been.
341    arrived: Vec<(String, Arrived)>,
342}
343
344/// The most fields that can arrive after the open: past this many, a row with another
345/// is a misfit, as in a followed file.
346const MOST_NEW_FIELDS: usize = 4096;
347
348/// What the values of a field that arrived after the open have been, for its column's
349/// type: the widest of them, and text once they disagree.
350#[derive(Clone, Copy, Debug, PartialEq)]
351enum Arrived {
352    Nothing,
353    Integer,
354    Number,
355    Boolean,
356    Text,
357}
358
359impl Arrived {
360    fn of(value: &serde_json::Value) -> Arrived {
361        match value {
362            serde_json::Value::Null => Arrived::Nothing,
363            serde_json::Value::Bool(_) => Arrived::Boolean,
364            serde_json::Value::Number(n) if n.is_i64() || n.is_u64() => Arrived::Integer,
365            serde_json::Value::Number(_) => Arrived::Number,
366            // Polars reads an array or an object into a text column as its JSON text.
367            _ => Arrived::Text,
368        }
369    }
370
371    fn and(self, other: Arrived) -> Arrived {
372        match (self, other) {
373            (a, b) if a == b => a,
374            (Arrived::Nothing, x) | (x, Arrived::Nothing) => x,
375            (Arrived::Integer, Arrived::Number) | (Arrived::Number, Arrived::Integer) => {
376                Arrived::Number
377            }
378            _ => Arrived::Text,
379        }
380    }
381
382    fn dtype(self) -> DataType {
383        match self {
384            Arrived::Integer => DataType::Int64,
385            Arrived::Number => DataType::Float64,
386            Arrived::Boolean => DataType::Boolean,
387            Arrived::Nothing | Arrived::Text => DataType::String,
388        }
389    }
390}
391
392impl Tail {
393    /// The tail of a file read as `format` with `options`, whose frame has `schema`,
394    /// before anything is counted.
395    pub fn new(format: FileFormat, options: &OpenOptions, schema: &Schema) -> Tail {
396        let layout = match format.separator() {
397            _ if format == FileFormat::Arrow => Layout::Stream,
398            Some(separator) => Layout::Delimited {
399                separator: options.separator_or(separator),
400                skip: options.skip_lines.unwrap_or(0) as u64
401                    + options.skip_rows.unwrap_or(0) as u64,
402                header: options.has_header != Some(false),
403                comment: options
404                    .comment_char
405                    .as_ref()
406                    .filter(|c| !c.is_empty())
407                    .map(|c| c.as_bytes().to_vec()),
408                nulls: options
409                    .null_values
410                    .iter()
411                    .flatten()
412                    .filter(|spec| !spec.contains('='))
413                    .cloned()
414                    .collect(),
415            },
416            None if format.is_lines() => Layout::Text,
417            None => Layout::Lines,
418        };
419        let columns = schema
420            .iter()
421            .map(|(name, dtype)| (name.to_string(), Fits::of(dtype)))
422            .collect();
423        Tail {
424            path: PathBuf::new(),
425            layout,
426            columns,
427            complete: 0,
428            records: 0,
429            rows: 0,
430            misfits: 0,
431            fields: None,
432            marks: NewMarks::default(),
433            mark_every: (MARK_ROWS, MARK_BYTES),
434            widens: false,
435            arrived: Vec::new(),
436        }
437    }
438
439    /// Fields the schema does not have are columns to come rather than misfits:
440    /// standard input's NDJSON, read again with them once it ends.
441    pub fn widening(mut self) -> Tail {
442        self.widens = matches!(self.layout, Layout::Lines);
443        self
444    }
445
446    /// The fields that arrived that the schema does not have, typed by their values.
447    pub fn new_fields(&self) -> Vec<Field> {
448        self.arrived
449            .iter()
450            .map(|(name, values)| Field::new(name.as_str().into(), values.dtype()))
451            .collect()
452    }
453
454    /// The file counted.
455    pub fn path(&self) -> &Path {
456        &self.path
457    }
458
459    /// Rows in the complete records.
460    pub fn rows(&self) -> usize {
461        self.rows as usize
462    }
463
464    /// Rows that arrived after the first count and did not fit the schema.
465    pub fn misfits(&self) -> usize {
466        self.misfits as usize
467    }
468
469    /// Bytes through the end of the last complete record.
470    pub fn complete(&self) -> u64 {
471        self.complete
472    }
473
474    /// Forget what was counted, to count the file again from its start.
475    fn restart(&mut self) {
476        self.complete = 0;
477        self.records = 0;
478        self.rows = 0;
479        self.misfits = 0;
480        self.fields = None;
481        self.marks = NewMarks::default();
482        self.arrived.clear();
483    }
484
485    /// Mark where row `row`, whose record starts at byte `start`, is: the first row, and
486    /// then once enough has passed since the last mark. The first is marked so that no
487    /// page is read through a Polars slice with an offset, which counts an NDJSON
488    /// file's blank lines as rows (#672). A blank record is never marked: read first
489    /// from a mark, it could be taken for no row at all.
490    fn mark(marks: &mut NewMarks, every: (u64, u64), row: u64, start: u64, blank: bool) {
491        let (rows, bytes) = every;
492        if blank
493            || marks.last.is_some_and(|(last_row, last_start)| {
494                row - last_row < rows && start - last_start < bytes
495            })
496        {
497            return;
498        }
499        marks.last = Some((row, start));
500        marks.new.push((row, start));
501    }
502
503    /// Count the records `file` completes between what was counted and `len`, checking
504    /// each new row against the schema when `check`.
505    pub fn read_on(&mut self, file: &mut File, len: u64, check: bool) -> std::io::Result<()> {
506        self.read(file, len, check, false)
507    }
508
509    /// As [`Self::read_on`], with nothing more to come: a last record with no newline
510    /// is a record too. Text read as lines and streams wait for their ends as before.
511    pub fn read_to_end(&mut self, file: &mut File, len: u64, check: bool) -> std::io::Result<()> {
512        let last = matches!(self.layout, Layout::Lines | Layout::Delimited { .. });
513        self.read(file, len, check, last)
514    }
515
516    fn read(&mut self, file: &mut File, len: u64, check: bool, last: bool) -> std::io::Result<()> {
517        if len <= self.complete {
518            return Ok(());
519        }
520        if matches!(self.layout, Layout::Stream) {
521            return self.read_messages(file, len);
522        }
523        file.seek(SeekFrom::Start(self.complete))?;
524        let mut reader = file.take(len - self.complete);
525        let quoted = matches!(self.layout, Layout::Delimited { .. });
526        let mut buf = vec![0u8; CHUNK];
527        let mut record: Vec<u8> = Vec::new();
528        let mut oversized = false;
529        let mut in_quotes = false;
530        let mut at = self.complete;
531        loop {
532            let n = match reader.read(&mut buf) {
533                Ok(0) => break,
534                Ok(n) => n,
535                Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
536                Err(e) => return Err(e),
537            };
538            let chunk = &buf[..n];
539            let mut i = 0;
540            while i < n {
541                let rest = &chunk[i..];
542                let next = if quoted {
543                    rest.iter().position(|&b| b == b'\n' || b == b'"')
544                } else {
545                    rest.iter().position(|&b| b == b'\n')
546                };
547                let Some(k) = next else {
548                    keep(&mut record, rest, &mut oversized);
549                    break;
550                };
551                keep(&mut record, &rest[..k], &mut oversized);
552                let byte = rest[k];
553                i += k + 1;
554                if byte == b'"' {
555                    in_quotes = !in_quotes;
556                    keep(&mut record, b"\"", &mut oversized);
557                } else if in_quotes {
558                    keep(&mut record, b"\n", &mut oversized);
559                } else {
560                    self.end_record(&record, oversized, check);
561                    record.clear();
562                    oversized = false;
563                    self.complete = at + i as u64;
564                }
565            }
566            at += n as u64;
567        }
568        if last && at > self.complete {
569            self.end_record(&record, oversized, check);
570            self.complete = at;
571        }
572        Ok(())
573    }
574
575    /// Count the record batches of a stream whose messages `file` completes between what
576    /// was counted and `len`. Only their headers are read.
577    fn read_messages(&mut self, file: &mut File, len: u64) -> std::io::Result<()> {
578        file.seek(SeekFrom::Start(self.complete))?;
579        let mut reader = std::io::BufReader::with_capacity(CHUNK, file);
580        while let Some((message, size)) = stream::next_message(&mut reader, len - self.complete)? {
581            match message {
582                stream::Message::Batch { rows } => {
583                    Self::mark(
584                        &mut self.marks,
585                        self.mark_every,
586                        self.rows,
587                        self.complete,
588                        false,
589                    );
590                    self.rows += rows;
591                }
592                stream::Message::Dictionary => {
593                    return Err(std::io::Error::other(
594                        "a dictionary batch arrived, which a followed stream cannot read",
595                    ));
596                }
597                stream::Message::Schema | stream::Message::End | stream::Message::Other => {}
598            }
599            self.records += 1;
600            self.complete += size;
601        }
602        Ok(())
603    }
604
605    /// One complete record, without its newline. It starts where the records counted
606    /// before it end.
607    fn end_record(&mut self, record: &[u8], oversized: bool, check: bool) {
608        let record = record.strip_suffix(b"\r").unwrap_or(record);
609        let start = self.complete;
610        let index = self.records;
611        self.records += 1;
612        match &self.layout {
613            Layout::Delimited {
614                separator,
615                skip,
616                header,
617                comment,
618                nulls,
619            } => {
620                if index < *skip {
621                    return;
622                }
623                if comment.as_ref().is_some_and(|c| record.starts_with(c)) {
624                    return;
625                }
626                if *header && self.fields.is_none() {
627                    self.fields = Some(split_fields(record, *separator).len());
628                    return;
629                }
630                Self::mark(
631                    &mut self.marks,
632                    self.mark_every,
633                    self.rows,
634                    start,
635                    record.is_empty(),
636                );
637                self.rows += 1;
638                if check && (oversized || !self.cells_fit(record, *separator, nulls)) {
639                    self.misfits += 1;
640                }
641            }
642            Layout::Stream => {}
643            Layout::Text => self.rows += 1,
644            Layout::Lines => {
645                if record.iter().all(u8::is_ascii_whitespace) {
646                    return;
647                }
648                Self::mark(&mut self.marks, self.mark_every, self.rows, start, false);
649                self.rows += 1;
650                // Read at the open too when widening: a field first seen there, past the
651                // lines the schema came from, is a column to come as well.
652                if check || self.widens {
653                    let fits = !oversized && self.object_fits(record);
654                    if check && !fits {
655                        self.misfits += 1;
656                    }
657                }
658            }
659        }
660    }
661
662    fn cells_fit(&self, record: &[u8], separator: u8, nulls: &[String]) -> bool {
663        let cells = split_fields(record, separator);
664        let expected = self.fields.unwrap_or(self.columns.len());
665        // An empty line is a row of nulls, as Polars reads it.
666        if record.is_empty() {
667            return true;
668        }
669        cells.len() == expected
670            && cells.iter().zip(&self.columns).all(|(cell, (_, fits))| {
671                let text = String::from_utf8_lossy(cell);
672                fits.cell(unquote(&text), nulls)
673            })
674    }
675
676    /// Whether an NDJSON record fits the schema. Widening, a field it does not have is
677    /// noted as arrived, and fits.
678    fn object_fits(&mut self, record: &[u8]) -> bool {
679        let Ok(serde_json::Value::Object(object)) = serde_json::from_slice(record) else {
680            return false;
681        };
682        let mut fits = true;
683        for (key, value) in &object {
684            match self.columns.iter().find(|(name, _)| name == key) {
685                Some((_, kind)) => fits &= kind.value(value),
686                None if self.widens => fits &= Self::arrive(&mut self.arrived, key, value),
687                None => fits = false,
688            }
689        }
690        fits
691    }
692
693    /// Note `key`, a field the schema does not have, and its value. False when there
694    /// is no room for another field.
695    fn arrive(arrived: &mut Vec<(String, Arrived)>, key: &str, value: &serde_json::Value) -> bool {
696        let kind = Arrived::of(value);
697        if let Some((_, seen)) = arrived.iter_mut().find(|(name, _)| name == key) {
698            *seen = seen.and(kind);
699            return true;
700        }
701        if arrived.len() >= MOST_NEW_FIELDS {
702            return false;
703        }
704        arrived.push((key.to_string(), kind));
705        true
706    }
707}
708
709/// Add `bytes` to the record being read, up to [`LONGEST_RECORD`].
710fn keep(record: &mut Vec<u8>, bytes: &[u8], oversized: &mut bool) {
711    if record.len() + bytes.len() > LONGEST_RECORD {
712        *oversized = true;
713        return;
714    }
715    record.extend_from_slice(bytes);
716}
717
718/// A delimited record's fields, quotes respected.
719fn split_fields(record: &[u8], separator: u8) -> Vec<&[u8]> {
720    let mut fields = Vec::new();
721    let mut in_quotes = false;
722    let mut start = 0;
723    for (i, &b) in record.iter().enumerate() {
724        if b == b'"' {
725            in_quotes = !in_quotes;
726        } else if b == separator && !in_quotes {
727            fields.push(&record[start..i]);
728            start = i + 1;
729        }
730    }
731    fields.push(&record[start..]);
732    fields
733}
734
735fn unquote(cell: &str) -> &str {
736    let trimmed = cell.trim();
737    trimmed
738        .strip_prefix('"')
739        .and_then(|c| c.strip_suffix('"'))
740        .unwrap_or(trimmed)
741}
742
743/// Whether two spellings of a path name one file. Polars keeps its own spelling of a
744/// scan's path, which on Windows need not match ours character for character.
745fn same_file(a: &str, b: &str) -> bool {
746    if a == b || Path::new(a) == Path::new(b) {
747        return true;
748    }
749    matches!(
750        (std::fs::canonicalize(a), std::fs::canonicalize(b)),
751        (Ok(a), Ok(b)) if a == b
752    )
753}
754
755/// Whether `plan` is a scan of `path` and nothing else.
756fn scans(plan: &polars::lazy::dsl::DslPlan, path: &str) -> bool {
757    use polars::lazy::dsl::DslPlan;
758    match plan {
759        DslPlan::Scan {
760            sources: ScanSources::Paths(paths),
761            ..
762        } if paths.len() == 1 => same_file(paths[0].as_str(), path),
763        DslPlan::Scan { .. } => {
764            stream::StreamScan::of(plan, path).is_some()
765                || lines::LinesScan::of(plan, path).is_some()
766        }
767        DslPlan::IR { dsl, .. } => scans(dsl, path),
768        _ => false,
769    }
770}
771
772/// `lf` reading the file at `path` only up to its first `rows` rows. A scan of it with
773/// no bound gets one right above it; one bounded already has its bound moved. The
774/// frames built on a scan carry the scan in their plans, so moving its bound moves
775/// what every one of them reads.
776pub fn bound(lf: &mut LazyFrame, path: &Path, rows: usize) {
777    let path = path.to_string_lossy();
778    let rows = IdxSize::try_from(rows).unwrap_or(IdxSize::MAX);
779    bound_plan(&mut lf.logical_plan, &path, rows);
780}
781
782fn bound_plan(plan: &mut polars::lazy::dsl::DslPlan, path: &str, rows: IdxSize) {
783    use polars::lazy::dsl::DslPlan;
784    if crate::lines::bound(plan, rows) {
785        return;
786    }
787    match plan {
788        // A plan asked for its schema is wrapped as IR, which would run as converted:
789        // bound the plan it came from, and leave the IR behind.
790        DslPlan::IR { dsl, .. } => {
791            let mut inner = Arc::unwrap_or_clone(dsl.clone());
792            bound_plan(&mut inner, path, rows);
793            *plan = inner;
794            return;
795        }
796        DslPlan::Slice {
797            input,
798            offset: 0,
799            len,
800        } if scans(input, path) => {
801            *len = rows;
802            return;
803        }
804        DslPlan::Scan { .. } if scans(plan, path) => {
805            let scan = std::mem::take(plan);
806            *plan = DslPlan::Slice {
807                input: Arc::new(scan),
808                offset: 0,
809                len: rows,
810            };
811            return;
812        }
813        _ => {}
814    }
815    crate::widgets::datatable::for_each_input(plan, &mut |input| bound_plan(input, path, rows));
816}
817
818/// `lf` reading the file at `path` through `file`, an open handle on it, rather than by
819/// its name: the file was deleted, and the handle still reads what it held.
820pub fn read_through(lf: &mut LazyFrame, path: &Path, file: &File) {
821    let path = path.to_string_lossy();
822    read_through_plan(&mut lf.logical_plan, &path, file);
823}
824
825fn read_through_plan(plan: &mut polars::lazy::dsl::DslPlan, path: &str, file: &File) {
826    use polars::lazy::dsl::DslPlan;
827    match plan {
828        DslPlan::IR { dsl, .. } => {
829            let mut inner = Arc::unwrap_or_clone(dsl.clone());
830            read_through_plan(&mut inner, path, file);
831            *plan = inner;
832            return;
833        }
834        DslPlan::Scan { .. } if scans(plan, path) => {
835            let held: Option<Arc<dyn AnonymousScan>> = match stream::StreamScan::of(plan, path) {
836                Some(scan) => scan
837                    .held(file)
838                    .map(|s| Arc::new(s) as Arc<dyn AnonymousScan>),
839                None => lines::LinesScan::of(plan, path)
840                    .and_then(|scan| scan.held(file))
841                    .map(|s| Arc::new(s) as Arc<dyn AnonymousScan>),
842            };
843            if let DslPlan::Scan {
844                sources,
845                scan_type,
846                cached_ir,
847                ..
848            } = plan
849                && let Ok(handle) = file.try_clone()
850            {
851                match (held, &mut **scan_type) {
852                    (Some(held), polars::lazy::dsl::FileScanDsl::Anonymous { function, .. }) => {
853                        *function = held;
854                    }
855                    _ => *sources = ScanSources::Files(Arc::from([handle])),
856                }
857                // The conversion cached for the path would read the path.
858                *cached_ir = Default::default();
859            }
860            return;
861        }
862        _ => {}
863    }
864    crate::widgets::datatable::for_each_input(plan, &mut |input| {
865        read_through_plan(input, path, file)
866    });
867}
868
869/// Where rows of a followed file start, every so many rows ([`MARK_ROWS`],
870/// [`MARK_BYTES`]), and where its complete records end: a window deep in the file is
871/// read from the mark before it rather than from the file's start. The watcher makes
872/// the marks in the pass that counts the new records, so they cost no read of their own.
873#[derive(Default)]
874pub struct Marks {
875    inner: Mutex<MarksInner>,
876}
877
878#[derive(Default)]
879struct MarksInner {
880    /// (row, byte where its record starts), rows ascending.
881    at: Vec<(u64, u64)>,
882    complete: u64,
883}
884
885/// The bytes holding a run of rows: from the start of the record of `row`, the mark at or
886/// before the run, to `end`.
887#[derive(Clone, Copy, Debug, PartialEq)]
888struct Span {
889    row: u64,
890    start: u64,
891    end: u64,
892}
893
894impl Marks {
895    fn lock(&self) -> std::sync::MutexGuard<'_, MarksInner> {
896        self.inner.lock().unwrap_or_else(|e| e.into_inner())
897    }
898
899    /// Take the marks `tail` made since the last call, and where its records end.
900    fn take_from(&self, tail: &mut Tail) {
901        let mut inner = self.lock();
902        inner.at.append(&mut tail.marks.new);
903        inner.complete = tail.complete;
904    }
905
906    /// The file is read again from its start: the marks so far are of another file.
907    fn clear(&self) {
908        let mut inner = self.lock();
909        inner.at.clear();
910        inner.complete = 0;
911    }
912
913    /// The bytes holding rows `[from, to)`: from the last mark at or before `from` to
914    /// the first at or after `to`, or to the end of the complete records. `None` before
915    /// the first mark (a blank first row) and when the bytes are more than
916    /// [`MOST_FROM_A_MARK`].
917    fn span(&self, from: u64, to: u64) -> Option<Span> {
918        let inner = self.lock();
919        let before = inner.at.partition_point(|&(row, _)| row <= from);
920        let (row, start) = *inner.at.get(before.checked_sub(1)?)?;
921        let after = inner.at.partition_point(|&(row, _)| row < to);
922        let end = inner.at.get(after).map_or(inner.complete, |&(_, at)| at);
923        (end >= start && end - start <= MOST_FROM_A_MARK).then_some(Span { row, start, end })
924    }
925}
926
927/// How a run of a followed file's bytes, starting at a record, becomes rows: as the
928/// file's scan reads them, with no header and nothing skipped.
929#[derive(Clone)]
930enum Parse {
931    Csv(Box<CsvReadOptions>),
932    Lines { ignore_errors: bool },
933    Stream(Arc<stream::StreamSchema>),
934}
935
936/// The scan under `plan`, seen through the wrapper a schema request leaves.
937fn scan_node(plan: &polars::lazy::dsl::DslPlan) -> &polars::lazy::dsl::DslPlan {
938    match plan {
939        polars::lazy::dsl::DslPlan::IR { dsl, .. } => scan_node(dsl),
940        plan => plan,
941    }
942}
943
944impl Parse {
945    /// How `scan` reads its rows, when a run of them can be read the same way: a CSV
946    /// or NDJSON scan of every column, with no row index or path column.
947    fn of(scan: &polars::lazy::dsl::DslPlan, schema: &SchemaRef) -> Option<Parse> {
948        use polars::lazy::dsl::{DslPlan, FileScanDsl};
949        if let Some(stream) = stream::StreamScan::in_plan(scan_node(scan)) {
950            return Some(Parse::Stream(stream.schema().clone()));
951        }
952        if let Some(lines) = lines::LinesScan::in_plan(scan_node(scan)) {
953            return Some(Parse::Lines {
954                ignore_errors: lines.ignore_errors(),
955            });
956        }
957        let DslPlan::Scan {
958            scan_type,
959            unified_scan_args,
960            ..
961        } = scan_node(scan)
962        else {
963            return None;
964        };
965        if unified_scan_args.row_index.is_some() || unified_scan_args.include_file_paths.is_some() {
966            return None;
967        }
968        match &**scan_type {
969            FileScanDsl::Csv { options } => {
970                if options.columns.is_some()
971                    || options.projection.is_some()
972                    || options.row_index.is_some()
973                {
974                    return None;
975                }
976                let mut options = (**options).clone();
977                options.path = None;
978                options.has_header = false;
979                options.skip_rows = 0;
980                options.skip_lines = 0;
981                options.skip_rows_after_header = 0;
982                options.n_rows = None;
983                // The names and types the scan settled on, by position.
984                options.schema = Some(schema.clone());
985                options.schema_overwrite = None;
986                options.dtype_overwrite = None;
987                options.column_names_overwrite = None;
988                options.raise_if_empty = false;
989                Some(Parse::Csv(Box::new(options)))
990            }
991            FileScanDsl::NDJson { options } => Some(Parse::Lines {
992                ignore_errors: options.ignore_errors,
993            }),
994            _ => None,
995        }
996    }
997}
998
999/// Rows `[skip, skip + take)` of the records in `span` of a followed file, read when
1000/// the frame is collected.
1001struct Piece {
1002    path: PathBuf,
1003    span: Span,
1004    skip: usize,
1005    take: usize,
1006    parse: Parse,
1007    schema: SchemaRef,
1008}
1009
1010/// The name a read from a mark carries in a plan.
1011const PIECE_NAME: &str = "FOLLOWED";
1012
1013impl polars::prelude::AnonymousScan for Piece {
1014    fn as_any(&self) -> &dyn std::any::Any {
1015        self
1016    }
1017
1018    fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
1019        Ok(self.schema.clone())
1020    }
1021
1022    fn scan(&self, args: polars::prelude::AnonymousScanArgs) -> PolarsResult<DataFrame> {
1023        let take = args.n_rows.map_or(self.take, |n| n.min(self.take));
1024        let mut file = File::open(&self.path)?;
1025        file.seek(SeekFrom::Start(self.span.start))?;
1026        let mut bytes = Vec::with_capacity((self.span.end - self.span.start) as usize);
1027        file.take(self.span.end - self.span.start)
1028            .read_to_end(&mut bytes)?;
1029        let df = match &self.parse {
1030            Parse::Csv(options) => {
1031                let mut options = (**options).clone();
1032                options.n_rows = Some(self.skip + take);
1033                options
1034                    .into_reader_with_file_handle(std::io::Cursor::new(bytes))
1035                    .finish()?
1036            }
1037            Parse::Lines { ignore_errors } => {
1038                lines::parse_run(&bytes, &self.schema, *ignore_errors)?
1039            }
1040            Parse::Stream(schema) => stream::decode_run(bytes, schema, self.skip + take)?,
1041        };
1042        Ok(df.slice(self.skip as i64, take))
1043    }
1044}
1045
1046/// `lf` with the bounded scan of the followed file at `path` reading only its rows
1047/// `[from, to)` (`to` at most the bound, the bound when `None`), from the mark before
1048/// them: whatever the view does above the scan is done to those rows alone. `None` when
1049/// the marks do not reach them or the plan has no bounded scan of the file.
1050pub(crate) fn from_marks(
1051    lf: &LazyFrame,
1052    path: &Path,
1053    marks: &Marks,
1054    from: usize,
1055    to: Option<usize>,
1056) -> Option<LazyFrame> {
1057    let path_text = path.to_string_lossy();
1058    let mut plan = lf.logical_plan.clone();
1059    let mut replaced = false;
1060    let mut failed = false;
1061    let piece = |scan: &polars::lazy::dsl::DslPlan, bound: usize| {
1062        let to = to.map_or(bound, |to| to.min(bound));
1063        let from = from.min(to);
1064        let span = marks.span(from as u64, to as u64)?;
1065        let schema = LazyFrame::from(scan.clone()).collect_schema().ok()?;
1066        let parse = Parse::of(scan, &schema)?;
1067        let piece = Piece {
1068            path: path.to_path_buf(),
1069            span,
1070            skip: from - span.row as usize,
1071            take: to - from,
1072            parse,
1073            schema: schema.clone(),
1074        };
1075        LazyFrame::anonymous_scan(
1076            Arc::new(piece),
1077            ScanArgsAnonymous {
1078                schema: Some(schema),
1079                name: PIECE_NAME,
1080                ..Default::default()
1081            },
1082        )
1083        .ok()
1084        .map(|lf| lf.logical_plan)
1085    };
1086    replace_bound(
1087        &mut plan,
1088        &path_text,
1089        &mut |scan, bound| match piece(scan, bound) {
1090            Some(plan) => {
1091                replaced = true;
1092                Some(plan)
1093            }
1094            None => {
1095                failed = true;
1096                None
1097            }
1098        },
1099    );
1100    (replaced && !failed).then(|| {
1101        let mut out = lf.clone();
1102        out.logical_plan = plan;
1103        out
1104    })
1105}
1106
1107/// Put `with(scan, bound)` where `plan` reads the file at `path` through its bound.
1108fn replace_bound(
1109    plan: &mut polars::lazy::dsl::DslPlan,
1110    path: &str,
1111    with: &mut dyn FnMut(&polars::lazy::dsl::DslPlan, usize) -> Option<polars::lazy::dsl::DslPlan>,
1112) {
1113    use polars::lazy::dsl::DslPlan;
1114    match plan {
1115        DslPlan::IR { dsl, .. } => {
1116            let mut inner = Arc::unwrap_or_clone(dsl.clone());
1117            replace_bound(&mut inner, path, with);
1118            *plan = inner;
1119            return;
1120        }
1121        DslPlan::Slice {
1122            input,
1123            offset: 0,
1124            len,
1125        } if scans(input, path) => {
1126            if let Some(piece) = with(input, *len as usize) {
1127                *plan = piece;
1128            }
1129            return;
1130        }
1131        _ => {}
1132    }
1133    crate::widgets::datatable::for_each_input(plan, &mut |input| replace_bound(input, path, with));
1134}
1135
1136/// How many rows the frame `lf` reads of the followed file at `path`: its bound.
1137pub(crate) fn bound_of(lf: &LazyFrame, path: &Path) -> Option<usize> {
1138    use polars::lazy::dsl::DslPlan;
1139    let path = path.to_string_lossy();
1140    (&lf.logical_plan).into_iter().find_map(|node| match node {
1141        DslPlan::Slice {
1142            input,
1143            offset: 0,
1144            len,
1145        } if scans(input, &path) => Some(*len as usize),
1146        _ => None,
1147    })
1148}
1149
1150/// `root`, the frame of the followed NDJSON file at `path` read as `format` and bounded
1151/// to `rows`, reading `fields` too, after the columns it has. `None` when it has no
1152/// lines scan of the file, or `fields` brings no column it lacks.
1153pub(crate) fn widen(
1154    root: &LazyFrame,
1155    path: &Path,
1156    format: FileFormat,
1157    fields: &[Field],
1158    rows: usize,
1159) -> Option<LazyFrame> {
1160    let path_text = path.to_string_lossy();
1161    let mut plan = root.logical_plan.clone();
1162    let mut widened = None;
1163    replace_lines_scan(&mut plan, &path_text, &mut |scan| {
1164        let mut schema = (**scan.schema()).clone();
1165        for field in fields {
1166            if !schema.contains(field.name()) {
1167                schema.with_column(field.name().clone(), field.dtype().clone());
1168            }
1169        }
1170        if schema.len() == scan.schema().len() {
1171            return None;
1172        }
1173        let schema = Arc::new(schema);
1174        let lf = scan.with_schema(schema.clone()).lazy().ok()?;
1175        widened = Some((lf.clone(), schema));
1176        Some(lf.logical_plan)
1177    });
1178    let (raw, schema) = widened?;
1179    if format == FileFormat::Journal {
1180        // The journal names every column it shows: built again from the scan.
1181        let mut raw = raw;
1182        bound(&mut raw, path, rows);
1183        return Some(crate::journal::derive(raw, &schema).0);
1184    }
1185    let mut out = root.clone();
1186    out.logical_plan = plan;
1187    Some(out)
1188}
1189
1190/// Put `with(scan)` where `plan` reads the file at `path` through a lines scan.
1191fn replace_lines_scan(
1192    plan: &mut polars::lazy::dsl::DslPlan,
1193    path: &str,
1194    with: &mut dyn FnMut(&lines::LinesScan) -> Option<polars::lazy::dsl::DslPlan>,
1195) {
1196    use polars::lazy::dsl::DslPlan;
1197    if let DslPlan::IR { dsl, .. } = plan {
1198        let mut inner = Arc::unwrap_or_clone(dsl.clone());
1199        replace_lines_scan(&mut inner, path, with);
1200        *plan = inner;
1201        return;
1202    }
1203    if let Some(scan) = lines::LinesScan::of(plan, path) {
1204        if let Some(replaced) = with(scan) {
1205            *plan = replaced;
1206        }
1207        return;
1208    }
1209    crate::widgets::datatable::for_each_input(plan, &mut |input| {
1210        replace_lines_scan(input, path, with)
1211    });
1212}
1213
1214/// The windows of a followed file's view, each read from the mark before it. A view
1215/// of the rows as they are reads its rows straight; one that only filters them reads
1216/// on from `known`, a point where the rows of the view before it are known (view row,
1217/// file row), and slices.
1218pub(crate) struct Window {
1219    pub(crate) lf: LazyFrame,
1220    pub(crate) path: PathBuf,
1221    pub(crate) marks: Arc<Marks>,
1222    pub(crate) known: Option<Vec<(usize, usize)>>,
1223}
1224
1225impl crate::pushdown::Windowed for Window {
1226    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1227        let read = match &self.known {
1228            None => from_marks(&self.lf, &self.path, &self.marks, start, Some(start + len)),
1229            Some(known) => {
1230                let at = known.partition_point(|&(view, _)| view <= start);
1231                known.get(at.wrapping_sub(1)).and_then(|&(view, row)| {
1232                    from_marks(&self.lf, &self.path, &self.marks, row, None)
1233                        .map(|lf| lf.slice((start - view) as i64, len as IdxSize))
1234                })
1235            }
1236        };
1237        // Short of marks, Polars reads from the start of the file.
1238        Ok(read.unwrap_or_else(|| self.lf.clone().slice(start as i64, len as IdxSize)))
1239    }
1240}
1241
1242/// What the watcher found.
1243#[derive(Clone)]
1244pub enum Change {
1245    /// More complete rows: how many there are now, and how many of the rows that came
1246    /// in since the open do not fit the schema.
1247    Grew { rows: usize, misfits: usize },
1248    /// The file shrank or was replaced (truncated, rotated): it is read again from its
1249    /// start, which holds `rows` rows.
1250    Restarted { rows: usize, misfits: usize },
1251    /// The file is gone. `handle` still reads what it held.
1252    Gone { handle: Option<Arc<File>> },
1253    /// Fields that arrived in standard input's NDJSON after the open, which the schema
1254    /// does not have: sent once it has ended, just before [`Change::Ended`].
1255    NewFields(Vec<Field>),
1256    /// Standard input ended: with the reason when it ended in an error.
1257    Ended(Option<String>),
1258    /// The file could not be read.
1259    Failed(String),
1260}
1261
1262/// One report from a watcher, for the follow named `id`.
1263#[derive(Clone)]
1264pub struct News {
1265    pub(crate) id: u64,
1266    pub(crate) change: Change,
1267}
1268
1269/// What stops the watcher, and wakes it to look now.
1270#[derive(Default)]
1271struct Shared {
1272    stop: AtomicBool,
1273    poke: Mutex<bool>,
1274    woken: Condvar,
1275    /// Wakes a watcher waiting on inotify rather than on `woken`.
1276    #[cfg(target_os = "linux")]
1277    bell: notify::Bell,
1278}
1279
1280impl Shared {
1281    /// Wait out `interval`, or until poked or stopped. Whether to go on.
1282    fn wait(&self, interval: Duration) -> bool {
1283        let mut poked = self.poke.lock().unwrap_or_else(|e| e.into_inner());
1284        if !*poked && !self.stop.load(Ordering::Relaxed) {
1285            poked = self
1286                .woken
1287                .wait_timeout(poked, interval)
1288                .unwrap_or_else(|e| e.into_inner())
1289                .0;
1290        }
1291        *poked = false;
1292        !self.stop.load(Ordering::Relaxed)
1293    }
1294
1295    fn wake(&self) {
1296        *self.poke.lock().unwrap_or_else(|e| e.into_inner()) = true;
1297        self.woken.notify_all();
1298        #[cfg(target_os = "linux")]
1299        self.bell.ring();
1300    }
1301
1302    /// Wait until the file changes, as `notify` hears, or until poked or stopped. A
1303    /// change is looked at no sooner than `interval` after the last look, `last`, so a
1304    /// burst of appends is one look. Whether to go on.
1305    #[cfg(target_os = "linux")]
1306    fn wait_for_change(
1307        &self,
1308        notify: &notify::Notify,
1309        interval: Duration,
1310        last: &mut Option<Instant>,
1311    ) -> bool {
1312        loop {
1313            if self.stop.load(Ordering::Relaxed) {
1314                return false;
1315            }
1316            if std::mem::take(&mut *self.poke.lock().unwrap_or_else(|e| e.into_inner())) {
1317                break;
1318            }
1319            if notify.wait(&self.bell, None) == notify::Woke::Changed {
1320                let left = last
1321                    .map(|at| at + interval)
1322                    .and_then(|due| due.checked_duration_since(Instant::now()));
1323                // Returns early when poked.
1324                if let Some(left) = left
1325                    && !self.wait(left)
1326                {
1327                    return false;
1328                }
1329                break;
1330            }
1331        }
1332        // What changed before this look is read by it.
1333        notify.drain();
1334        *last = Some(Instant::now());
1335        !self.stop.load(Ordering::Relaxed)
1336    }
1337}
1338
1339/// How long ago, `elapsed`, as the follow chip says it: in the largest whole unit,
1340/// padded so the chip keeps its width as the number grows.
1341pub fn age(elapsed: Duration) -> String {
1342    let secs = elapsed.as_secs();
1343    match secs {
1344        0..60 => format!("{secs:>2}s ago"),
1345        60..3_600 => format!("{:>2}m ago", secs / 60),
1346        3_600..86_400 => format!("{:>2}h ago", secs / 3_600),
1347        _ => format!("{:>2}d ago", secs / 86_400),
1348    }
1349}
1350
1351/// When the age of an append made at `at` next reads differently.
1352pub fn next_tick(at: Instant) -> Instant {
1353    let secs = at.elapsed().as_secs();
1354    let unit = match secs {
1355        0..60 => 1,
1356        60..3_600 => 60,
1357        3_600..86_400 => 3_600,
1358        _ => 86_400,
1359    };
1360    at + Duration::from_secs((secs / unit + 1) * unit)
1361}
1362
1363/// Where the follow stands, as the control bar says it.
1364#[derive(Clone, Debug, PartialEq)]
1365pub enum Standing {
1366    Following,
1367    Paused,
1368    /// Standard input ended, or the file went: what is on screen is all there is.
1369    Ended,
1370}
1371
1372/// A file being followed: its watcher, and what the view has taken of it. Belongs to
1373/// the dataset; dropping it stops the watcher.
1374pub struct Follow {
1375    id: u64,
1376    /// The file the frame scans.
1377    path: PathBuf,
1378    shared: Arc<Shared>,
1379    spool: Option<Arc<SpoolHandle>>,
1380    /// Rows the watcher has counted, and the misfits among them.
1381    counted: usize,
1382    misfits: usize,
1383    /// The file was read again from its start, and the view has not caught up.
1384    restarted: bool,
1385    /// Rows the frame reads now.
1386    shown: usize,
1387    /// Rows that arrived below the cursor while it was not on the last row.
1388    pub(crate) new_below: usize,
1389    pub(crate) standing: Standing,
1390    pub(crate) last_append: Option<Instant>,
1391    /// The view has not gone to the last row yet: a follow starts there, as `tail -f`
1392    /// does, once the first page is drawn.
1393    pub(crate) settle_at_end: bool,
1394    /// The cursor was on the last row when rows arrived: it goes to the new last row
1395    /// once they are read.
1396    pub(crate) end_pending: bool,
1397    /// Rows were taken while the table was not on screen: it reads its rows again
1398    /// when it is.
1399    pub(crate) stale_view: bool,
1400    /// The handle a deleted file is read through from now on.
1401    held: Option<Arc<File>>,
1402    /// Where its rows start, every so many.
1403    marks: Arc<Marks>,
1404    /// Standard input read as it arrives without `--follow`: the view stays where it
1405    /// is, and the rows so far are a part of what is coming.
1406    pipe: bool,
1407    /// Fields that arrived after the open, for the frames to join once standard input
1408    /// has ended.
1409    new_fields: Vec<Field>,
1410    /// What the format says of the whole stream (the journal's Info tab) has been
1411    /// asked for again, now that it has ended.
1412    pub(crate) described: bool,
1413}
1414
1415static NEXT_ID: AtomicU64 = AtomicU64::new(1);
1416
1417impl Follow {
1418    /// Follow the file `tail` counted, whose first `tail.rows()` rows the frame reads,
1419    /// checking every `interval` and telling `events`. `spool` is standard input being
1420    /// copied to it.
1421    pub fn start(
1422        mut tail: Tail,
1423        interval: Duration,
1424        events: Sender<AppEvent>,
1425        spool: Option<Arc<SpoolHandle>>,
1426    ) -> Follow {
1427        let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
1428        let path = tail.path.clone();
1429        let shared = Arc::new(Shared::default());
1430        let shown = tail.rows();
1431        if let Some(handle) = &spool {
1432            handle.spool.wake_on_end(shared.clone());
1433        }
1434        let marks = Arc::new(Marks::default());
1435        marks.take_from(&mut tail);
1436        let watcher = Watcher {
1437            marks: marks.clone(),
1438            id,
1439            path: path.clone(),
1440            tail,
1441            shared: shared.clone(),
1442            events,
1443            spool: spool.as_ref().map(|handle| handle.spool.clone()),
1444            interval,
1445        };
1446        let _ = std::thread::Builder::new()
1447            .name("datui-follow".to_string())
1448            .spawn(move || watcher.run());
1449        Follow {
1450            id,
1451            path,
1452            shared,
1453            spool,
1454            counted: shown,
1455            misfits: 0,
1456            restarted: false,
1457            shown,
1458            new_below: 0,
1459            standing: Standing::Following,
1460            last_append: None,
1461            settle_at_end: true,
1462            end_pending: false,
1463            stale_view: false,
1464            held: None,
1465            marks,
1466            pipe: false,
1467            new_fields: Vec::new(),
1468            described: false,
1469        }
1470    }
1471
1472    /// This follows standard input read as it arrives rather than `--follow`: the
1473    /// view starts at the top and stays where it is put.
1474    pub fn as_pipe(mut self) -> Follow {
1475        self.pipe = true;
1476        self.settle_at_end = false;
1477        self
1478    }
1479
1480    /// Whether this is standard input read as it arrives ([`Self::as_pipe`]).
1481    pub fn is_pipe(&self) -> bool {
1482        self.pipe
1483    }
1484
1485    /// Whether more rows may still come: followed and not ended.
1486    pub fn live(&self) -> bool {
1487        self.standing != Standing::Ended
1488    }
1489
1490    /// Where the file's rows start, every so many.
1491    pub(crate) fn marks(&self) -> &Arc<Marks> {
1492        &self.marks
1493    }
1494
1495    pub fn id(&self) -> u64 {
1496        self.id
1497    }
1498
1499    pub fn path(&self) -> &Path {
1500        &self.path
1501    }
1502
1503    /// Rows the frame reads.
1504    pub fn shown(&self) -> usize {
1505        self.shown
1506    }
1507
1508    /// Rows counted that the view has not taken yet.
1509    pub fn waiting(&self) -> usize {
1510        self.counted.saturating_sub(self.shown)
1511    }
1512
1513    pub fn misfits(&self) -> usize {
1514        self.misfits
1515    }
1516
1517    pub fn standing(&self) -> &Standing {
1518        &self.standing
1519    }
1520
1521    /// Rows that arrived below the cursor while it was off the last row.
1522    pub fn new_below(&self) -> usize {
1523        self.new_below
1524    }
1525
1526    /// Whether the view has rows or a restart to take: not while paused. The rows
1527    /// counted before standard input ended are taken after it did.
1528    pub fn behind(&self) -> bool {
1529        self.standing != Standing::Paused && (self.restarted || self.counted != self.shown)
1530    }
1531
1532    /// Standard input being copied, if this follows it.
1533    pub fn spool(&self) -> Option<&Arc<Spool>> {
1534        self.spool.as_ref().map(|handle| &handle.spool)
1535    }
1536
1537    /// Look now rather than at the end of the interval.
1538    pub fn check_now(&self) {
1539        self.shared.wake();
1540    }
1541
1542    /// Take a report from this follow's watcher. Returns what the user is told, if
1543    /// anything.
1544    pub fn take(&mut self, change: &Change) -> Option<String> {
1545        match change {
1546            Change::Grew { rows, misfits } => {
1547                if *rows > self.counted {
1548                    self.last_append = Some(Instant::now());
1549                }
1550                self.counted = *rows;
1551                self.misfits = *misfits;
1552                None
1553            }
1554            Change::Restarted { rows, misfits } => {
1555                self.counted = *rows;
1556                self.misfits = *misfits;
1557                self.restarted = true;
1558                self.last_append = Some(Instant::now());
1559                Some("The file was truncated or replaced: reading it from the start".to_string())
1560            }
1561            Change::NewFields(fields) => {
1562                self.new_fields = fields.clone();
1563                None
1564            }
1565            Change::Gone { handle } => {
1566                self.held = handle.clone();
1567                self.end();
1568                Some("The file was deleted: following stopped, the rows read stay".to_string())
1569            }
1570            // A recording's end is said by the recording's own mark.
1571            Change::Ended(_) if self.spool().is_some_and(|s| s.tee().is_some()) => {
1572                self.end();
1573                None
1574            }
1575            Change::Ended(None) => {
1576                self.end();
1577                Some("Standard input ended".to_string())
1578            }
1579            Change::Ended(Some(reason)) | Change::Failed(reason) => {
1580                self.end();
1581                Some(reason.clone())
1582            }
1583        }
1584    }
1585
1586    /// The view takes what was counted: the rows its frame reads from now on, and
1587    /// whether the file was read again from its start.
1588    pub fn catch_up(&mut self) -> (usize, bool) {
1589        self.shown = self.counted;
1590        (self.shown, std::mem::take(&mut self.restarted))
1591    }
1592
1593    /// The fields that arrived after the open, once, for the frames to join.
1594    pub(crate) fn take_new_fields(&mut self) -> Vec<Field> {
1595        std::mem::take(&mut self.new_fields)
1596    }
1597
1598    /// The handle a deleted file is read through, once, for the frames to take.
1599    pub fn take_held(&mut self) -> Option<Arc<File>> {
1600        self.held.take()
1601    }
1602
1603    pub fn pause(&mut self) {
1604        if self.standing == Standing::Following {
1605            self.standing = Standing::Paused;
1606        }
1607    }
1608
1609    pub fn resume(&mut self) {
1610        if self.standing == Standing::Paused {
1611            self.standing = Standing::Following;
1612        }
1613    }
1614
1615    /// Stop watching. What the frame reads stays. Standard input stops being read,
1616    /// unless it is being recorded (`--tee`): that is stopped only when asked.
1617    pub fn end(&mut self) {
1618        self.standing = Standing::Ended;
1619        self.shared.stop.store(true, Ordering::Relaxed);
1620        self.shared.wake();
1621        if let Some(spool) = self.spool.as_ref().filter(|s| s.spool.tee().is_none()) {
1622            spool.spool.stop();
1623        }
1624    }
1625}
1626
1627impl Drop for Follow {
1628    fn drop(&mut self) {
1629        self.shared.stop.store(true, Ordering::Relaxed);
1630        self.shared.wake();
1631    }
1632}
1633
1634/// The thread that watches the file.
1635struct Watcher {
1636    id: u64,
1637    path: PathBuf,
1638    tail: Tail,
1639    shared: Arc<Shared>,
1640    events: Sender<AppEvent>,
1641    spool: Option<Arc<Spool>>,
1642    interval: Duration,
1643    marks: Arc<Marks>,
1644}
1645
1646/// Which file this is, so a replaced file is told from a grown one: (device, inode) on
1647/// Unix, (volume serial, file index) on Windows.
1648type Identity = (u64, u64);
1649
1650#[cfg(unix)]
1651fn identity_of(file: &File) -> Option<Identity> {
1652    use std::os::unix::fs::MetadataExt;
1653    file.metadata().ok().map(|meta| (meta.dev(), meta.ino()))
1654}
1655
1656#[cfg(windows)]
1657fn identity_of(file: &File) -> Option<Identity> {
1658    use std::os::windows::io::AsRawHandle;
1659    use windows_sys::Win32::Storage::FileSystem::{
1660        BY_HANDLE_FILE_INFORMATION, GetFileInformationByHandle,
1661    };
1662    // SAFETY: the handle is `file`'s, open while this runs, and `info` is plain data
1663    // the call fills in; zeroed is a valid value of it.
1664    let mut info: BY_HANDLE_FILE_INFORMATION = unsafe { std::mem::zeroed() };
1665    let ok = unsafe { GetFileInformationByHandle(file.as_raw_handle(), &mut info) };
1666    (ok != 0).then(|| {
1667        (
1668            u64::from(info.dwVolumeSerialNumber),
1669            u64::from(info.nFileIndexHigh) << 32 | u64::from(info.nFileIndexLow),
1670        )
1671    })
1672}
1673
1674#[cfg(not(any(unix, windows)))]
1675fn identity_of(_file: &File) -> Option<Identity> {
1676    None
1677}
1678
1679/// Which file `path` names now, `meta` its metadata. Unix reads it from the metadata;
1680/// Windows has to open the file to ask.
1681#[cfg(unix)]
1682fn identity_at(_path: &Path, meta: &std::fs::Metadata) -> Option<Identity> {
1683    use std::os::unix::fs::MetadataExt;
1684    Some((meta.dev(), meta.ino()))
1685}
1686
1687#[cfg(not(unix))]
1688fn identity_at(path: &Path, _meta: &std::fs::Metadata) -> Option<Identity> {
1689    File::open(path).ok().as_ref().and_then(identity_of)
1690}
1691
1692/// Whether the file a path names is another one than the file followed. Unknown on
1693/// either side is not a replacement: a shrink still tells a truncation.
1694fn replaced(known: Option<Identity>, now: Option<Identity>) -> bool {
1695    matches!((known, now), (Some(known), Some(now)) if known != now)
1696}
1697
1698/// Why following `path` (`None`: standard input) ended: `doing` when `e` stopped it.
1699fn failed_message(path: Option<&Path>, doing: &str, e: &std::io::Error) -> String {
1700    let what = format!(
1701        "{doing}. {}",
1702        crate::error_display::user_message_from_io(e, None)
1703    );
1704    match path {
1705        Some(path) => crate::error_display::file_message(path, &what),
1706        None => crate::error_display::sentence(&format!("standard input: {what}")),
1707    }
1708}
1709
1710impl Watcher {
1711    /// Why following ended, `doing` when `e` stopped it, naming the file followed;
1712    /// standard input's spool is a file the user never named.
1713    fn failed(&self, doing: &str, e: &std::io::Error) -> String {
1714        failed_message(
1715            self.spool.is_none().then_some(self.path.as_path()),
1716            doing,
1717            e,
1718        )
1719    }
1720
1721    fn run(mut self) {
1722        let mut file = match File::open(&self.path) {
1723            Ok(file) => file,
1724            Err(e) => {
1725                self.send(Change::Failed(self.failed("following it stopped", &e)));
1726                return;
1727            }
1728        };
1729        let mut known = identity_of(&file);
1730        let mut sent = (self.tail.rows(), 0usize);
1731        #[cfg(target_os = "linux")]
1732        let notify = notify::Notify::new(&self.path);
1733        #[cfg(target_os = "linux")]
1734        let mut last = None;
1735        loop {
1736            #[cfg(target_os = "linux")]
1737            let go_on = match &notify {
1738                Some(notify) => self
1739                    .shared
1740                    .wait_for_change(notify, self.interval, &mut last),
1741                None => self.shared.wait(self.interval),
1742            };
1743            #[cfg(not(target_os = "linux"))]
1744            let go_on = self.shared.wait(self.interval);
1745            if !go_on {
1746                return;
1747            }
1748            // Read before the file, so nothing the spool wrote before it ended is missed.
1749            let spool_ended = self.spool.as_ref().and_then(|spool| spool.ended());
1750            let meta = match std::fs::metadata(&self.path) {
1751                Ok(meta) => meta,
1752                Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
1753                    self.send(Change::Gone {
1754                        handle: Some(Arc::new(file)),
1755                    });
1756                    return;
1757                }
1758                Err(e) => {
1759                    self.send(Change::Failed(self.failed("following it stopped", &e)));
1760                    return;
1761                }
1762            };
1763            let now = identity_at(&self.path, &meta);
1764            let len = meta.len();
1765            if replaced(known, now) || len < self.tail.complete() {
1766                match File::open(&self.path) {
1767                    Ok(reopened) => file = reopened,
1768                    Err(e) => {
1769                        self.send(Change::Failed(self.failed("following it stopped", &e)));
1770                        return;
1771                    }
1772                }
1773                known = identity_of(&file);
1774                #[cfg(target_os = "linux")]
1775                if let Some(notify) = &notify {
1776                    notify.rewatch(&self.path);
1777                }
1778                self.tail.restart();
1779                self.marks.clear();
1780                if let Err(e) = self.tail.read_on(&mut file, len, false) {
1781                    self.send(Change::Failed(self.failed("reading it stopped", &e)));
1782                    return;
1783                }
1784                self.marks.take_from(&mut self.tail);
1785                sent = (self.tail.rows(), self.tail.misfits());
1786                self.send(Change::Restarted {
1787                    rows: sent.0,
1788                    misfits: sent.1,
1789                });
1790                continue;
1791            }
1792            let read = if spool_ended.is_some() {
1793                self.tail.read_to_end(&mut file, len, true)
1794            } else {
1795                self.tail.read_on(&mut file, len, true)
1796            };
1797            if let Err(e) = read {
1798                self.send(Change::Failed(self.failed("reading it stopped", &e)));
1799                return;
1800            }
1801            // Before the rows are reported, so the view's reads of them find marks.
1802            self.marks.take_from(&mut self.tail);
1803            let now_counted = (self.tail.rows(), self.tail.misfits());
1804            if now_counted != sent {
1805                sent = now_counted;
1806                if !self.send(Change::Grew {
1807                    rows: sent.0,
1808                    misfits: sent.1,
1809                }) {
1810                    return;
1811                }
1812            }
1813            if let Some(reason) = spool_ended {
1814                let fields = self.tail.new_fields();
1815                if !fields.is_empty() && !self.send(Change::NewFields(fields)) {
1816                    return;
1817                }
1818                self.send(Change::Ended(reason));
1819                return;
1820            }
1821        }
1822    }
1823
1824    /// Whether the app is still there to hear it.
1825    fn send(&self, change: Change) -> bool {
1826        self.events
1827            .send(AppEvent::Followed(News {
1828                id: self.id,
1829                change,
1830            }))
1831            .is_ok()
1832    }
1833}
1834
1835/// Standard input being copied to a file while the file is read: the copy goes on
1836/// after the first rows show, until the stream ends or the copy is stopped. The file is
1837/// a temporary one, or the one `--tee` names, which the user keeps; with `--tee -`, a
1838/// temporary one, and the stream is passed on to standard output too.
1839pub struct Spool {
1840    stop: AtomicBool,
1841    bytes: AtomicU64,
1842    state: Mutex<SpoolState>,
1843    changed: Condvar,
1844    /// The file being written. Taken when the copy finishes, so a read still waiting on
1845    /// the producer writes nothing after it.
1846    sink: Mutex<Option<File>>,
1847    /// The file `--tee` named, when it is the one written.
1848    tee: Option<Tee>,
1849    /// Standard output, for `--tee -`. Taken when the copy finishes, which closes it.
1850    pass: Mutex<Option<Box<dyn Write + Send>>>,
1851    started: Instant,
1852}
1853
1854/// The file `--tee` named.
1855#[derive(Clone, Debug)]
1856pub struct Tee {
1857    pub path: PathBuf,
1858    /// `--tee-raw`: the bytes exactly as they came, a WAV header's sizes included.
1859    pub raw: bool,
1860}
1861
1862impl Tee {
1863    /// `--tee -`: the stream is passed on to standard output rather than kept in a file.
1864    pub fn to_stdout(&self) -> bool {
1865        crate::stdin::is_stdin(&self.path)
1866    }
1867
1868    /// Where the stream goes, as a message names it.
1869    pub fn name(&self) -> String {
1870        if self.to_stdout() {
1871            return "standard output".to_string();
1872        }
1873        self.path.file_name().map_or_else(
1874            || self.path.display().to_string(),
1875            |name| name.to_string_lossy().into_owned(),
1876        )
1877    }
1878}
1879
1880#[derive(Default)]
1881struct SpoolState {
1882    /// Newlines copied so far, up to what the open waits for.
1883    lines: usize,
1884    /// The last read took all the producer had: it is slower than the copy.
1885    drained: bool,
1886    /// The stream ended: `Some(reason)` when in an error.
1887    ended: Option<Option<String>>,
1888    /// When the copy finished, and what finishing the file said.
1889    finished: Option<Instant>,
1890    /// Bytes copied by when, a few seconds of them, for the rate.
1891    samples: std::collections::VecDeque<(Instant, u64)>,
1892    /// The watcher following the file, woken when the copy ends: a watcher waiting
1893    /// for the file to change would not hear an end that writes nothing.
1894    watcher: Option<Arc<Shared>>,
1895}
1896
1897/// How far back the rate looks.
1898const RATE_WINDOW: Duration = Duration::from_secs(2);
1899
1900impl Spool {
1901    fn new(sink: File, tee: Option<Tee>, pass: Option<Box<dyn Write + Send>>) -> Spool {
1902        Spool {
1903            stop: AtomicBool::new(false),
1904            bytes: AtomicU64::new(0),
1905            state: Mutex::new(SpoolState::default()),
1906            changed: Condvar::new(),
1907            sink: Mutex::new(Some(sink)),
1908            tee,
1909            pass: Mutex::new(pass),
1910            started: Instant::now(),
1911        }
1912    }
1913
1914    /// Bytes copied so far.
1915    pub fn bytes(&self) -> u64 {
1916        self.bytes.load(Ordering::Relaxed)
1917    }
1918
1919    /// Bytes a second over the last few seconds.
1920    pub fn rate(&self) -> f64 {
1921        let state = self.lock();
1922        match (state.samples.front(), state.samples.back()) {
1923            (Some((t0, b0)), Some((t1, b1))) if t1 > t0 => {
1924                (b1 - b0) as f64 / t1.duration_since(*t0).as_secs_f64()
1925            }
1926            _ => 0.0,
1927        }
1928    }
1929
1930    /// How long the copy ran, or has run.
1931    pub fn duration(&self) -> Duration {
1932        let finished = self.lock().finished;
1933        finished
1934            .unwrap_or_else(Instant::now)
1935            .duration_since(self.started)
1936    }
1937
1938    /// The file `--tee` named, when the copy goes there.
1939    pub fn tee(&self) -> Option<&Tee> {
1940        self.tee.as_ref()
1941    }
1942
1943    /// Stop copying and finish the file: a read waiting on the producer writes
1944    /// nothing more when it returns.
1945    pub fn stop(&self) {
1946        self.stop.store(true, Ordering::Relaxed);
1947        self.finish(None);
1948    }
1949
1950    pub fn stopped(&self) -> bool {
1951        self.stop.load(Ordering::Relaxed)
1952    }
1953
1954    /// Whether the copy ended, and the reason when in an error.
1955    pub fn ended(&self) -> Option<Option<String>> {
1956        self.lock().ended.clone()
1957    }
1958
1959    /// Whether bytes may still arrive: the producer is still sending.
1960    pub fn live(&self) -> bool {
1961        self.lock().ended.is_none()
1962    }
1963
1964    /// Wait until the copy has ended.
1965    pub fn wait(&self) {
1966        let mut state = self.lock();
1967        while state.ended.is_none() {
1968            state = self.changed.wait(state).unwrap_or_else(|e| e.into_inner());
1969        }
1970    }
1971
1972    fn lock(&self) -> std::sync::MutexGuard<'_, SpoolState> {
1973        self.state.lock().unwrap_or_else(|e| e.into_inner())
1974    }
1975
1976    /// Write `bytes` as they came. False once the copy is finished.
1977    fn write(&self, bytes: &[u8]) -> Result<bool, String> {
1978        let mut sink = self.sink.lock().unwrap_or_else(|e| e.into_inner());
1979        let Some(file) = sink.as_mut() else {
1980            return Ok(false);
1981        };
1982        file.write_all(bytes).map_err(|e| {
1983            format!(
1984                "Could not write {}: {e}",
1985                self.tee
1986                    .as_ref()
1987                    .filter(|t| !t.to_stdout())
1988                    .map_or("what came in".to_string(), |t| t.path.display().to_string())
1989            )
1990        })?;
1991        drop(sink);
1992        // Outside the file's lock: a reader downstream that stops reading holds up this
1993        // write, and must not hold up a stop.
1994        if let Some(out) = self.pass.lock().unwrap_or_else(|e| e.into_inner()).as_mut() {
1995            out.write_all(bytes)
1996                .and_then(|()| out.flush())
1997                .map_err(|e| format!("Could not write standard output: {e}"))?;
1998        }
1999        let total =
2000            self.bytes.fetch_add(bytes.len() as u64, Ordering::Relaxed) + bytes.len() as u64;
2001        let now = Instant::now();
2002        let mut state = self.lock();
2003        if state.lines < WANTED_LINES {
2004            state.lines += bytes.iter().filter(|&&b| b == b'\n').count();
2005        }
2006        state.samples.push_back((now, total));
2007        while state
2008            .samples
2009            .front()
2010            .is_some_and(|(at, _)| now.duration_since(*at) > RATE_WINDOW)
2011            && state.samples.len() > 2
2012        {
2013            state.samples.pop_front();
2014        }
2015        drop(state);
2016        self.changed.notify_all();
2017        Ok(true)
2018    }
2019
2020    /// End the copy, `reason` when it ended in an error, and finish the file: a WAV
2021    /// header's sizes filled in for `--tee` (unless `--tee-raw`), and the file synced so
2022    /// that saved means safe to copy. Once; later calls change nothing.
2023    fn finish(&self, reason: Option<String>) {
2024        let file = self.sink.lock().unwrap_or_else(|e| e.into_inner()).take();
2025        // Closed, so the reader downstream sees the stream end. Held by a write a reader
2026        // downstream is not taking, it is closed once that write returns: the copy then
2027        // finds the file finished and finishes again.
2028        if let Ok(mut pass) = self.pass.try_lock() {
2029            pass.take();
2030        }
2031        let mut reason = reason;
2032        if let (Some(mut file), Some(tee)) = (file, self.tee.as_ref().filter(|t| !t.to_stdout())) {
2033            let finished = (if tee.raw {
2034                Ok(())
2035            } else {
2036                crate::tee::fix_wav_sizes(&mut file).map(|_| ())
2037            })
2038            .and_then(|()| file.sync_all());
2039            if let Err(e) = finished
2040                && reason.is_none()
2041            {
2042                reason = Some(format!("Could not finish {}: {e}", tee.path.display()));
2043            }
2044        }
2045        let mut state = self.lock();
2046        if state.ended.is_none() {
2047            state.ended = Some(reason);
2048            state.finished = Some(Instant::now());
2049        }
2050        let watcher = state.watcher.take();
2051        drop(state);
2052        self.changed.notify_all();
2053        if let Some(watcher) = watcher {
2054            watcher.wake();
2055        }
2056    }
2057
2058    /// Wake the watcher `shared` once the copy ends, or now if it has.
2059    fn wake_on_end(&self, shared: Arc<Shared>) {
2060        let mut state = self.lock();
2061        if state.ended.is_some() {
2062            drop(state);
2063            shared.wake();
2064        } else {
2065            state.watcher = Some(shared);
2066        }
2067    }
2068}
2069
2070/// The open's hold on a [`Spool`]: the copy stops when the last holder lets go, so a
2071/// follow put down before its dataset arrived does not go on copying, and quitting
2072/// finishes the file.
2073pub struct SpoolHandle {
2074    spool: Arc<Spool>,
2075}
2076
2077impl SpoolHandle {
2078    pub fn spool(&self) -> &Arc<Spool> {
2079        &self.spool
2080    }
2081}
2082
2083impl Drop for SpoolHandle {
2084    fn drop(&mut self) {
2085        self.spool.stop();
2086    }
2087}
2088
2089/// Copy what `reader` sends into `spool` on a thread of its own, until it ends or the
2090/// spool stops. Each chunk is written whole as it arrives into one buffer, reused:
2091/// nothing is held back, and however long the stream runs the copy holds a megabyte.
2092/// A producer faster than the disk waits on the pipe, not on datui's memory.
2093fn copy_on(mut reader: impl Read + Send + 'static, spool: Arc<Spool>) {
2094    let _ = std::thread::Builder::new()
2095        .name("datui-spool".to_string())
2096        .spawn(move || {
2097            let mut buf = vec![0u8; CHUNK];
2098            let reason = loop {
2099                if spool.stopped() {
2100                    break None;
2101                }
2102                let n = match reader.read(&mut buf) {
2103                    Ok(0) => break None,
2104                    Ok(n) => n,
2105                    Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
2106                    Err(e) => break Some(format!("Standard input failed: {e}")),
2107                };
2108                match spool.write(&buf[..n]) {
2109                    Ok(true) => {}
2110                    Ok(false) => break None,
2111                    Err(reason) => break Some(reason),
2112                }
2113                spool.lock().drained = n < buf.len();
2114            };
2115            spool.finish(reason);
2116        });
2117}
2118
2119/// Lines the open waits for before it reads the spooled file, unless the producer is
2120/// slower than the copy: then two, a header and a row, are enough to show.
2121const WANTED_LINES: usize = 1000;
2122
2123/// What standard input was copied to.
2124pub enum Spooled {
2125    /// A temporary file, removed when its holders let go.
2126    Temp(TempDownload),
2127    /// The file `--tee` named, the user's.
2128    Kept(PathBuf),
2129}
2130
2131/// Copy standard input, from `open`, to the file `--tee` names or else a temporary
2132/// file in the spool directory ([`crate::stdin::spool_dir`]) claimed through `writer`,
2133/// until enough has arrived to show; the copy goes on behind the answer. Says what the
2134/// file holds, as [`crate::stdin::spool`] does, with the copy carried in the options
2135/// for the dataset to hold.
2136///
2137/// Not followed and not recorded, standard input is read as it arrives all the same,
2138/// when what it holds can be (`OpenOptions::pipe`): the rows show as they come and stop
2139/// coming at its end. What cannot be (Parquet, a compressed stream) is copied to its
2140/// end first, as before.
2141pub(crate) fn spool<R: Read + Send + 'static>(
2142    open: impl FnOnce() -> crate::download::Opened<R>,
2143    options: OpenOptions,
2144    writer: &Writer,
2145    read: &AtomicU64,
2146    stdout: Option<Box<dyn Write + Send>>,
2147) -> Result<(Spooled, OpenOptions), String> {
2148    let tee = options.tee.clone().map(|path| Tee {
2149        path,
2150        raw: options.tee_raw,
2151    });
2152    // A pipe cannot be sought back to, to fill in a WAV header.
2153    let tee = tee.map(|tee| Tee {
2154        raw: tee.raw || tee.to_stdout(),
2155        ..tee
2156    });
2157    let pass = match &tee {
2158        Some(tee) if tee.to_stdout() => Some(stdout.ok_or_else(|| {
2159            "--tee - passes the stream on to standard output, which only the datui command has."
2160                .to_string()
2161        })?),
2162        _ => None,
2163    };
2164    let progressive = !options.follow && tee.is_none();
2165    let (reader, _) = open().map_err(|e| format!("Could not read standard input: {e}"))?;
2166    let (spooled, file) = match &tee {
2167        Some(tee) if !tee.to_stdout() => {
2168            let file = crate::tee::create(&tee.path, options.force)?;
2169            (Spooled::Kept(tee.path.clone()), file)
2170        }
2171        _ => {
2172            let dir = crate::stdin::spool_dir(&options);
2173            let Some((named, claim)) = writer
2174                .create(|| TempDownload::create(dir.as_deref(), None))
2175                .map_err(|e| crate::error_display::user_message_from_report(&e, None))?
2176            else {
2177                return Err("Reading standard input was stopped.".to_string());
2178            };
2179            let file = named
2180                .as_file()
2181                .try_clone()
2182                .map_err(|e| format!("Could not write what came in: {e}"))?;
2183            (Spooled::Temp(TempDownload::held(named, Some(claim))), file)
2184        }
2185    };
2186    let spooled_path = match &spooled {
2187        Spooled::Temp(download) => download.path().to_path_buf(),
2188        Spooled::Kept(path) => path.clone(),
2189    };
2190    let spool = Arc::new(Spool::new(file, tee, pass));
2191    let handle = Arc::new(SpoolHandle {
2192        spool: spool.clone(),
2193    });
2194    copy_on(reader, spool.clone());
2195    // A header and a row, at least, before anything is read.
2196    let wanted = if options.has_header == Some(false) {
2197        1
2198    } else {
2199        2
2200    };
2201    let mut state = spool.lock();
2202    loop {
2203        read.store(spool.bytes(), Ordering::Relaxed);
2204        if writer.stopped() {
2205            drop(state);
2206            spool.stop();
2207            return Err("Reading standard input was stopped.".to_string());
2208        }
2209        // An Arrow stream has no lines to count: its schema message is enough.
2210        let enough = (options.follow || progressive)
2211            && (state.lines >= WANTED_LINES
2212                || (state.drained
2213                    && (state.lines >= wanted || stream::begins_with_schema(&spooled_path))));
2214        if enough || state.ended.is_some() {
2215            break;
2216        }
2217        state = spool
2218            .changed
2219            .wait_timeout(state, Duration::from_millis(100))
2220            .unwrap_or_else(|e| e.into_inner())
2221            .0;
2222    }
2223    if let Some(Some(reason)) = &state.ended {
2224        return Err(reason.clone());
2225    }
2226    drop(state);
2227    read.store(spool.bytes(), Ordering::Relaxed);
2228    let path = spooled_path;
2229    let mut head = Vec::new();
2230    File::open(&path)
2231        .and_then(|f| f.take(4096).read_to_end(&mut head))
2232        .map_err(|e| format!("Could not read standard input back: {e}"))?;
2233    if head.is_empty() {
2234        return Err("Nothing came in on standard input.".to_string());
2235    }
2236    let (format, compression, guessed) = crate::stdin::sniff_for(&head, &options);
2237    let asked = options.clone();
2238    let options = OpenOptions {
2239        format_guessed: options.format.is_none() && guessed,
2240        format: options.format.or(Some(format)),
2241        compression: options.compression.or(compression),
2242        ..options
2243    };
2244    if progressive {
2245        let read_on = followed_stream(&path, options.format, &options)
2246            || refusal(options.format, &options).is_none();
2247        if read_on {
2248            return Ok((
2249                spooled,
2250                OpenOptions {
2251                    spool: Some(handle),
2252                    follow: true,
2253                    pipe: true,
2254                    ..options
2255                },
2256            ));
2257        }
2258        // Read once it is finished: the rest is copied first, then looked at whole.
2259        while spool.ended().is_none() {
2260            if writer.stopped() {
2261                spool.stop();
2262                return Err("Reading standard input was stopped.".to_string());
2263            }
2264            read.store(spool.bytes(), Ordering::Relaxed);
2265            let state = spool.lock();
2266            if state.ended.is_none() {
2267                let _ = spool
2268                    .changed
2269                    .wait_timeout(state, Duration::from_millis(100))
2270                    .unwrap_or_else(|e| e.into_inner());
2271            }
2272        }
2273        if let Some(Some(reason)) = spool.ended() {
2274            return Err(reason);
2275        }
2276        read.store(spool.bytes(), Ordering::Relaxed);
2277        let options = crate::stdin::described(&path, asked)?;
2278        return Ok((spooled, options));
2279    }
2280    // A recording goes on whatever it holds; only the view is not followed then.
2281    if options.follow
2282        && options.tee.is_none()
2283        && !followed_stream(&path, options.format, &options)
2284        && let Some(refusal) = refusal(options.format, &options)
2285    {
2286        return Err(refusal);
2287    }
2288    Ok((
2289        spooled,
2290        OpenOptions {
2291            spool: Some(handle),
2292            // The file is read from here on, as any file is.
2293            tee: None,
2294            ..options
2295        },
2296    ))
2297}
2298
2299#[cfg(test)]
2300mod tests {
2301    use super::*;
2302
2303    /// Following that stops on an error names the file followed, in the one shape;
2304    /// standard input is called that.
2305    #[test]
2306    fn errors_name_the_file() {
2307        let path = Path::new("/data/app.log");
2308        for (doing, e) in [
2309            (
2310                "following it stopped",
2311                std::io::Error::from(std::io::ErrorKind::NotFound),
2312            ),
2313            (
2314                "reading it stopped",
2315                std::io::Error::from(std::io::ErrorKind::PermissionDenied),
2316            ),
2317        ] {
2318            let message = failed_message(Some(path), doing, &e);
2319            eprintln!("{message}");
2320            crate::readers::bad_input::assert_shape(&message, path);
2321            assert!(message.contains(&doing[1..]), "{message}");
2322        }
2323        let e = std::io::Error::from(std::io::ErrorKind::UnexpectedEof);
2324        let message = failed_message(None, "reading it stopped", &e);
2325        assert_eq!(
2326            message,
2327            "Standard input: reading it stopped. Unexpected end of file."
2328        );
2329    }
2330
2331    fn tail_of(text: &[u8], format: FileFormat, options: &OpenOptions) -> Tail {
2332        let dir = tempfile::tempdir().unwrap();
2333        let path = dir.path().join("t");
2334        std::fs::write(&path, text).unwrap();
2335        let mut tail = Tail::new(format, options, &Schema::default());
2336        let mut file = File::open(&path).unwrap();
2337        tail.read_on(&mut file, text.len() as u64, false).unwrap();
2338        tail
2339    }
2340
2341    /// Polars' own count of the complete records of `text` read as CSV.
2342    fn polars_rows(text: &[u8], options: &OpenOptions) -> usize {
2343        let complete = text.iter().rposition(|&b| b == b'\n').map_or(0, |i| i + 1);
2344        let mut read = CsvReadOptions::default().with_has_header(options.has_header != Some(false));
2345        read = read.map_parse_options(|p| {
2346            p.with_comment_prefix(
2347                options
2348                    .comment_char
2349                    .as_deref()
2350                    .map(polars::io::csv::read::CommentPrefix::new_from_str),
2351            )
2352        });
2353        if let Some(skip) = options.skip_lines {
2354            read.skip_lines = skip;
2355        }
2356        CsvReader::new(std::io::Cursor::new(text[..complete].to_vec()))
2357            .with_options(read)
2358            .finish()
2359            .map(|df| df.height())
2360            .unwrap_or(0)
2361    }
2362
2363    /// Rows counted as Polars counts them: blank lines are rows of nulls, a quoted
2364    /// newline is not a record's end, comments and skipped lines are not rows, and a
2365    /// partial last line waits.
2366    #[test]
2367    fn records_are_counted_as_polars_reads_them() {
2368        let cases: [(&[u8], OpenOptions); 6] = [
2369            (b"a,b\n1,2\n3,4\n", OpenOptions::default()),
2370            (b"a,b\n1,2\n\n3,4\n5,", OpenOptions::default()),
2371            (b"a,b\n1,\"x\ny\"\n3,4\n", OpenOptions::default()),
2372            (b"a,b\r\n1,2\r\n3,4\r\n", OpenOptions::default()),
2373            (
2374                b"#c\na,b\n#x\n1,2\n3,4\n",
2375                OpenOptions {
2376                    comment_char: Some("#".to_string()),
2377                    ..Default::default()
2378                },
2379            ),
2380            (
2381                b"junk\na,b\n1,2\n",
2382                OpenOptions::default().with_skip_lines(1),
2383            ),
2384        ];
2385        for (text, options) in cases {
2386            let tail = tail_of(text, FileFormat::Csv, &options);
2387            assert_eq!(
2388                tail.rows(),
2389                polars_rows(text, &options),
2390                "{}",
2391                String::from_utf8_lossy(text)
2392            );
2393        }
2394        let lines = tail_of(
2395            b"{\"a\":1}\n\n{\"a\":2}\n{\"a\":",
2396            FileFormat::Jsonl,
2397            &Default::default(),
2398        );
2399        assert_eq!(lines.rows(), 2);
2400        assert_eq!(lines.complete(), 17);
2401    }
2402
2403    /// On Linux the watcher hears an append through inotify: with an interval of an
2404    /// hour, no size check would see it, and nothing pokes it.
2405    #[cfg(target_os = "linux")]
2406    #[test]
2407    fn an_append_is_heard_of_without_a_check() {
2408        let dir = tempfile::tempdir().unwrap();
2409        let path = dir.path().join("log.csv");
2410        std::fs::write(&path, "t\n1\n").unwrap();
2411        let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2412            .finish()
2413            .unwrap();
2414        let (_, tail) =
2415            bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
2416        let (tx, rx) = std::sync::mpsc::channel();
2417        let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
2418        let guard = Duration::from_secs(30);
2419        // The watcher may not be waiting yet: append until it reports, each append a
2420        // change it hears once it is.
2421        let mut file = std::fs::OpenOptions::new()
2422            .append(true)
2423            .open(&path)
2424            .unwrap();
2425        let mut rows = 1;
2426        let deadline = Instant::now() + guard;
2427        let news = loop {
2428            assert!(Instant::now() < deadline, "the watcher never heard");
2429            file.write_all(format!("{}\n", rows + 1).as_bytes())
2430                .unwrap();
2431            rows += 1;
2432            if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
2433                break news;
2434            }
2435        };
2436        assert!(matches!(news.change, Change::Grew { rows: 2.., .. }));
2437        drop(follow);
2438    }
2439
2440    /// The bytes that arrive later are read from where the count stopped, a partial
2441    /// line among them once it completes; a row that does not fit is counted.
2442    #[test]
2443    fn a_tail_reads_on_from_where_it_stopped() {
2444        let dir = tempfile::tempdir().unwrap();
2445        let path = dir.path().join("grow.csv");
2446        let mut out = File::create(&path).unwrap();
2447        out.write_all(b"t,n\n1.5,2\n2.5,").unwrap();
2448        let schema = Schema::from_iter([
2449            Field::new("t".into(), DataType::Float64),
2450            Field::new("n".into(), DataType::Int64),
2451        ]);
2452        let mut tail = Tail::new(FileFormat::Csv, &OpenOptions::default(), &schema);
2453        let mut file = File::open(&path).unwrap();
2454        let len = |p: &Path| std::fs::metadata(p).unwrap().len();
2455        tail.read_on(&mut file, len(&path), true).unwrap();
2456        assert_eq!((tail.rows(), tail.complete()), (1, 10));
2457        out.write_all(b"3\nx,4\n4.5,5,6\n").unwrap();
2458        tail.read_on(&mut file, len(&path), true).unwrap();
2459        assert_eq!(tail.rows(), 4);
2460        assert_eq!(
2461            tail.misfits(),
2462            2,
2463            "a word for a number, and a field too many"
2464        );
2465    }
2466
2467    /// A scan bounded to its complete rows reads more once the bound moves, through
2468    /// the filters built on it, and only as many as the bound says.
2469    #[test]
2470    fn moving_the_bound_reads_the_new_rows_through_the_view() {
2471        let dir = tempfile::tempdir().unwrap();
2472        let path = dir.path().join("grow.csv");
2473        std::fs::write(&path, "a,b\n1,x\n2,y\n3,").unwrap();
2474        let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2475            .with_truncate_ragged_lines(true)
2476            .with_ignore_errors(true)
2477            .finish()
2478            .unwrap();
2479        let mut root = scan.clone();
2480        bound(&mut root, &path, 2);
2481        let mut view = root.clone().filter(col("a").gt(lit(1)));
2482        view.collect_schema().unwrap();
2483        assert_eq!(view.clone().collect().unwrap().height(), 1);
2484        let mut out = std::fs::OpenOptions::new()
2485            .append(true)
2486            .open(&path)
2487            .unwrap();
2488        out.write_all(b"z\n4,w\n5").unwrap();
2489        bound(&mut view, &path, 4);
2490        let df = view.clone().collect().unwrap();
2491        assert_eq!(df.height(), 3, "{df}");
2492        bound(&mut root, &path, 4);
2493        assert_eq!(root.collect().unwrap().height(), 4, "the partial row waits");
2494    }
2495
2496    /// `text` written to a file and scanned as `scan` does, bounded to its complete
2497    /// rows, with a mark every `every` rows.
2498    fn marked(
2499        text: &[u8],
2500        format: FileFormat,
2501        options: &OpenOptions,
2502        every: u64,
2503        scan: impl Fn(&Path) -> LazyFrame,
2504    ) -> (tempfile::TempDir, PathBuf, LazyFrame, Arc<Marks>, usize) {
2505        let dir = tempfile::tempdir().unwrap();
2506        let path = dir.path().join("marked");
2507        std::fs::write(&path, text).unwrap();
2508        let mut lf = scan(&path);
2509        let schema = lf.collect_schema().unwrap();
2510        let mut tail = Tail::new(format, options, &schema);
2511        tail.mark_every = (every, u64::MAX);
2512        tail.read_on(&mut File::open(&path).unwrap(), text.len() as u64, false)
2513            .unwrap();
2514        let marks = Arc::new(Marks::default());
2515        marks.take_from(&mut tail);
2516        bound(&mut lf, &path, tail.rows());
2517        (dir, path, lf, marks, tail.rows())
2518    }
2519
2520    fn csv_scan(path: &Path, options: &OpenOptions) -> LazyFrame {
2521        let mut reader = LazyCsvReader::new(PlRefPath::try_from_path(path).unwrap())
2522            .with_ignore_errors(true)
2523            .with_truncate_ragged_lines(true)
2524            .with_has_header(options.has_header != Some(false))
2525            .with_comment_prefix(options.comment_char.as_deref().map(PlSmallStr::from_str));
2526        if let Some(skip) = options.skip_lines {
2527            reader = reader.with_skip_lines(skip);
2528        }
2529        reader.finish().unwrap()
2530    }
2531
2532    /// Every window read from the marks holds the rows a read from the start of the
2533    /// file gives: quoted newlines, blank lines, comments, skipped lines, carriage
2534    /// returns and NDJSON, at every offset.
2535    #[test]
2536    fn a_window_from_a_mark_reads_what_a_read_from_the_start_does() {
2537        let mut csv = b"skipped\nt,s,n\n".to_vec();
2538        let mut crlf = b"t,s,n\r\n".to_vec();
2539        let mut lines = Vec::new();
2540        for i in 0..120 {
2541            let row = match i % 9 {
2542                0 => format!("{i},\"two\nlines\",{}\n", i * 2),
2543                3 => "\n".to_string(),
2544                5 => "# a comment\n".to_string(),
2545                7 => format!("{i},x,oops\n"),
2546                _ => format!("{i},s{i},{}\n", i * 2),
2547            };
2548            csv.extend(row.as_bytes());
2549            crlf.extend(format!("{i},s{i},{}\r\n", i * 2).as_bytes());
2550            lines.extend(format!("{{\"t\":{i},\"s\":\"s{i}\"}}\n").as_bytes());
2551            if i % 4 == 1 {
2552                lines.extend(b"\n  \n");
2553            }
2554        }
2555        csv.extend(b"999,partial");
2556        let commented = OpenOptions {
2557            comment_char: Some("#".to_string()),
2558            ..OpenOptions::default().with_skip_lines(1)
2559        };
2560        // An Arrow stream of batches of three rows, the last message cut short.
2561        let (schema, batches) = stream_messages(&arrow_rows(0, 100), 3);
2562        let mut arrows = schema;
2563        batches.iter().for_each(|batch| arrows.extend(batch));
2564        arrows.extend(&batches[0][..20]);
2565        let cases: Vec<(&[u8], FileFormat, OpenOptions)> = vec![
2566            (&csv, FileFormat::Csv, commented),
2567            (&crlf, FileFormat::Csv, OpenOptions::default()),
2568            (&lines, FileFormat::Jsonl, OpenOptions::default()),
2569            (&arrows, FileFormat::Arrow, OpenOptions::default()),
2570        ];
2571        for (text, format, options) in cases {
2572            let scan = |path: &Path| match format {
2573                FileFormat::Jsonl => scan_lines(path, &options, false, &mut Vec::new()).unwrap(),
2574                FileFormat::Arrow => stream::scan(path).unwrap(),
2575                _ => csv_scan(path, &options),
2576            };
2577            let (_dir, path, lf, marks, rows) = marked(text, format, &options, 7, scan);
2578            let window = Window {
2579                lf: lf.clone(),
2580                path: path.clone(),
2581                marks: marks.clone(),
2582                known: None,
2583            };
2584            let whole = lf.clone().collect().unwrap();
2585            assert_eq!(whole.height(), rows);
2586            for start in (0..rows + 3).step_by(5) {
2587                for len in [1, 6, 40] {
2588                    let read = crate::pushdown::Windowed::window(&window, start, len)
2589                        .unwrap()
2590                        .collect()
2591                        .unwrap();
2592                    let expected = whole.slice(start as i64, len);
2593                    assert!(
2594                        read.equals_missing(&expected),
2595                        "{format:?} rows {start}+{len}:\n{read:?}\n{expected:?}"
2596                    );
2597                }
2598                if start < rows {
2599                    assert!(
2600                        from_marks(&lf, &path, &marks, start, Some(start + 1)).is_some(),
2601                        "{format:?} row {start} is read from a mark"
2602                    );
2603                }
2604            }
2605        }
2606    }
2607
2608    /// `n` rows from `from`: a number, its text, and a float.
2609    fn arrow_rows(from: i64, n: i64) -> DataFrame {
2610        df!(
2611            "t" => (from..from + n).collect::<Vec<_>>(),
2612            "s" => (from..from + n).map(|i| format!("s{i}")).collect::<Vec<_>>(),
2613            "x" => (from..from + n).map(|i| i as f64 / 2.0).collect::<Vec<_>>(),
2614        )
2615        .unwrap()
2616    }
2617
2618    /// An Arrow stream's batches are counted as their messages complete, a batch cut
2619    /// short waiting for the rest; the stream's own scan reads them, filtered and
2620    /// projected a batch at a time; and a stream with dictionaries is refused.
2621    #[test]
2622    fn an_arrow_stream_is_counted_and_read_by_its_batches() {
2623        let dir = tempfile::tempdir().unwrap();
2624        let path = dir.path().join("live.arrows");
2625        let (schema, batches) = stream_messages(&arrow_rows(0, 50), 4);
2626        let mut head = schema.clone();
2627        head.extend(&batches[0]);
2628        head.extend(&batches[1][..batches[1].len() - 3]);
2629        std::fs::write(&path, &head).unwrap();
2630        let lf = stream::scan(&path).unwrap();
2631        let (mut lf, mut tail) =
2632            bound_to_complete(lf, &path, FileFormat::Arrow, &OpenOptions::default()).unwrap();
2633        assert_eq!(tail.rows(), 4, "the second batch is not all there");
2634        assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 4));
2635
2636        let mut rest = batches[1][batches[1].len() - 3..].to_vec();
2637        batches[2..].iter().for_each(|batch| rest.extend(batch));
2638        let mut file = std::fs::OpenOptions::new()
2639            .append(true)
2640            .open(&path)
2641            .unwrap();
2642        file.write_all(&rest).unwrap();
2643        let size = file.metadata().unwrap().len();
2644        tail.read_on(&mut File::open(&path).unwrap(), size, true)
2645            .unwrap();
2646        assert_eq!((tail.rows(), tail.misfits()), (50, 0));
2647        bound(&mut lf, &path, tail.rows());
2648        assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 50));
2649        let kept = lf
2650            .clone()
2651            .filter(col("t").gt_eq(lit(45)))
2652            .select([col("s")])
2653            .collect()
2654            .unwrap();
2655        assert_eq!(kept, arrow_rows(45, 5).select(["s"]).unwrap());
2656        let count = lf.select([len()]).collect().unwrap();
2657        assert_eq!(count.column("len").unwrap().u32().unwrap().get(0), Some(50));
2658
2659        // Written before Arrow 0.15, with no continuation markers, and ended.
2660        let legacy = dir.path().join("legacy.arrows");
2661        std::fs::write(
2662            &legacy,
2663            crate::ipc_stream::tests::stream(&arrow_rows(0, 10), None, true),
2664        )
2665        .unwrap();
2666        let (lf, tail) = bound_to_complete(
2667            stream::scan(&legacy).unwrap(),
2668            &legacy,
2669            FileFormat::Arrow,
2670            &OpenOptions::default(),
2671        )
2672        .unwrap();
2673        assert_eq!(tail.rows(), 10);
2674        assert_eq!(lf.collect().unwrap(), arrow_rows(0, 10));
2675
2676        // Dictionary-encoded columns.
2677        let cats = df!("c" => ["a", "b", "a"])
2678            .unwrap()
2679            .lazy()
2680            .with_column(col("c").cast(DataType::from_categories(Categories::global())))
2681            .collect()
2682            .unwrap();
2683        let (schema, batches) = stream_messages(&cats, 3);
2684        let dictionary = dir.path().join("dict.arrows");
2685        std::fs::write(&dictionary, [schema, batches.concat()].concat()).unwrap();
2686        assert!(
2687            stream::scan(&dictionary).is_err_and(|e| e.contains("dictionary")),
2688            "refused"
2689        );
2690    }
2691
2692    /// A filtered view reads on from a point where its rows are known, and counts the
2693    /// rows after it alone.
2694    #[test]
2695    fn a_filtered_view_reads_and_counts_on_from_what_is_known() {
2696        let mut text = b"t,n\n".to_vec();
2697        for i in 0..300 {
2698            text.extend(format!("{i},{}\n", i % 5).as_bytes());
2699        }
2700        let options = OpenOptions::default();
2701        let (_dir, path, lf, marks, rows) = marked(&text, FileFormat::Csv, &options, 16, |p| {
2702            csv_scan(p, &options)
2703        });
2704        let view = lf.filter(col("n").eq(lit(3)));
2705        let whole = view.clone().collect().unwrap();
2706        // The view's rows among the first 200 of the file.
2707        let known = (whole.column("t").unwrap().i64().unwrap().to_vec())
2708            .into_iter()
2709            .filter(|t| t.unwrap() < 200)
2710            .count();
2711        let rest = from_marks(&view, &path, &marks, 200, None).unwrap();
2712        let after = rest.collect().unwrap().height();
2713        assert_eq!(known + after, whole.height());
2714        assert_eq!(rows, 300);
2715        let window = Window {
2716            lf: view.clone(),
2717            path,
2718            marks,
2719            known: Some(vec![(0, 0), (known, 200)]),
2720        };
2721        for start in [0, 10, known - 1, known, known + 5, whole.height() - 3] {
2722            let read = crate::pushdown::Windowed::window(&window, start, 4)
2723                .unwrap()
2724                .collect()
2725                .unwrap();
2726            assert!(
2727                read.equals_missing(&whole.slice(start as i64, 4)),
2728                "{start}"
2729            );
2730        }
2731    }
2732
2733    /// A file put in place of the followed one, as big or bigger, is another file;
2734    /// the file grown in place is the same one. Windows reads this from the volume and
2735    /// file index, Unix from the device and inode.
2736    #[test]
2737    fn a_replaced_file_is_told_from_a_grown_one() {
2738        let dir = tempfile::tempdir().unwrap();
2739        let path = dir.path().join("log.csv");
2740        std::fs::write(&path, "t\n1\n").unwrap();
2741        let held = File::open(&path).unwrap();
2742        let known = identity_of(&held);
2743        assert!(cfg!(not(any(unix, windows))) || known.is_some());
2744        let now = |path: &Path| identity_at(path, &std::fs::metadata(path).unwrap());
2745        std::fs::OpenOptions::new()
2746            .append(true)
2747            .open(&path)
2748            .unwrap()
2749            .write_all(b"2\n")
2750            .unwrap();
2751        assert!(!replaced(known, now(&path)), "grown in place");
2752        let other = dir.path().join("next.csv");
2753        std::fs::write(&other, "t\n1\n2\n3\n").unwrap();
2754        std::fs::rename(&other, &path).unwrap();
2755        assert_eq!(
2756            replaced(known, now(&path)),
2757            known.is_some(),
2758            "renamed over it"
2759        );
2760        assert!(!replaced(None, now(&path)), "unknown is no replacement");
2761        drop(held);
2762    }
2763
2764    /// A deleted file is read through the handle held on it.
2765    #[cfg(unix)]
2766    #[test]
2767    fn a_deleted_file_reads_through_its_handle() {
2768        let dir = tempfile::tempdir().unwrap();
2769        let path = dir.path().join("gone.csv");
2770        std::fs::write(&path, "a\n1\n2\n").unwrap();
2771        let mut lf = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2772            .finish()
2773            .unwrap();
2774        bound(&mut lf, &path, 2);
2775        let mut view = lf.filter(col("a").gt(lit(0)));
2776        view.collect_schema().unwrap();
2777        let handle = File::open(&path).unwrap();
2778        std::fs::remove_file(&path).unwrap();
2779        read_through(&mut view, &path, &handle);
2780        assert_eq!(view.collect().unwrap().height(), 2);
2781    }
2782}