Skip to main content

datui_lib/loading/
follow.rs

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