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                run_options(options, schema).map(|options| Parse::Csv(Box::new(options)))
943            }
944            FileScanDsl::NDJson { options } => Some(Parse::Lines {
945                ignore_errors: options.ignore_errors,
946            }),
947            _ => None,
948        }
949    }
950}
951
952/// How a CSV scan with `options` reads a run of its rows that starts at a record: no
953/// header, nothing skipped, the names and types the scan settled on (`schema`). `None`
954/// when the scan picks columns or numbers rows itself.
955pub(crate) fn run_options(options: &CsvReadOptions, schema: &SchemaRef) -> Option<CsvReadOptions> {
956    if options.columns.is_some() || options.projection.is_some() || options.row_index.is_some() {
957        return None;
958    }
959    let mut options = options.clone();
960    options.path = None;
961    options.has_header = false;
962    options.skip_rows = 0;
963    options.skip_lines = 0;
964    options.skip_rows_after_header = 0;
965    options.n_rows = None;
966    // The names and types the scan settled on, by position.
967    options.schema = Some(schema.clone());
968    options.schema_overwrite = None;
969    options.dtype_overwrite = None;
970    options.column_names_overwrite = None;
971    options.raise_if_empty = false;
972    Some(options)
973}
974
975/// Rows `[skip, skip + take)` of the records in `span` of a followed file, read when
976/// the frame is collected.
977struct Piece {
978    path: PathBuf,
979    span: Span,
980    skip: usize,
981    take: usize,
982    parse: Parse,
983    schema: SchemaRef,
984}
985
986/// The name a read from a mark carries in a plan.
987const PIECE_NAME: &str = "FOLLOWED";
988
989impl polars::prelude::AnonymousScan for Piece {
990    fn as_any(&self) -> &dyn std::any::Any {
991        self
992    }
993
994    fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
995        Ok(self.schema.clone())
996    }
997
998    fn scan(&self, args: polars::prelude::AnonymousScanArgs) -> PolarsResult<DataFrame> {
999        let take = args.n_rows.map_or(self.take, |n| n.min(self.take));
1000        let mut file = File::open(&self.path)?;
1001        file.seek(SeekFrom::Start(self.span.start))?;
1002        let mut bytes = Vec::with_capacity((self.span.end - self.span.start) as usize);
1003        file.take(self.span.end - self.span.start)
1004            .read_to_end(&mut bytes)?;
1005        let df = match &self.parse {
1006            Parse::Csv(options) => {
1007                let mut options = (**options).clone();
1008                options.n_rows = Some(self.skip + take);
1009                options
1010                    .into_reader_with_file_handle(std::io::Cursor::new(bytes))
1011                    .finish()?
1012            }
1013            Parse::Lines { ignore_errors } => {
1014                lines::parse_run(&bytes, &self.schema, *ignore_errors)?
1015            }
1016            Parse::Stream(schema) => stream::decode_run(bytes, schema, self.skip + take)?,
1017        };
1018        Ok(df.slice(self.skip as i64, take))
1019    }
1020}
1021
1022/// `lf` with its bounded scan of `path` reading only rows `[from, to)` (`to` capped at
1023/// the bound), from the preceding mark, so the view's work covers those rows alone.
1024/// `None` if the marks do not reach them or there is no bounded scan.
1025pub(crate) fn from_marks(
1026    lf: &LazyFrame,
1027    path: &Path,
1028    marks: &Marks,
1029    from: usize,
1030    to: Option<usize>,
1031) -> Option<LazyFrame> {
1032    let path_text = path.to_string_lossy();
1033    let mut plan = lf.logical_plan.clone();
1034    let mut replaced = false;
1035    let mut failed = false;
1036    let piece = |scan: &polars::lazy::dsl::DslPlan, bound: usize| {
1037        let to = to.map_or(bound, |to| to.min(bound));
1038        let from = from.min(to);
1039        let span = marks.span(from as u64, to as u64)?;
1040        let schema = LazyFrame::from(scan.clone()).collect_schema().ok()?;
1041        let parse = Parse::of(scan, &schema)?;
1042        let piece = Piece {
1043            path: path.to_path_buf(),
1044            span,
1045            skip: from - span.row as usize,
1046            take: to - from,
1047            parse,
1048            schema: schema.clone(),
1049        };
1050        LazyFrame::anonymous_scan(
1051            Arc::new(piece),
1052            ScanArgsAnonymous {
1053                schema: Some(schema),
1054                name: PIECE_NAME,
1055                ..Default::default()
1056            },
1057        )
1058        .ok()
1059        .map(|lf| lf.logical_plan)
1060    };
1061    replace_bound(
1062        &mut plan,
1063        &path_text,
1064        &mut |scan, bound| match piece(scan, bound) {
1065            Some(plan) => {
1066                replaced = true;
1067                Some(plan)
1068            }
1069            None => {
1070                failed = true;
1071                None
1072            }
1073        },
1074    );
1075    (replaced && !failed).then(|| {
1076        let mut out = lf.clone();
1077        out.logical_plan = plan;
1078        out
1079    })
1080}
1081
1082/// Put `with(scan, bound)` where `plan` reads the file at `path` through its bound.
1083fn replace_bound(
1084    plan: &mut polars::lazy::dsl::DslPlan,
1085    path: &str,
1086    with: &mut dyn FnMut(&polars::lazy::dsl::DslPlan, usize) -> Option<polars::lazy::dsl::DslPlan>,
1087) {
1088    use polars::lazy::dsl::DslPlan;
1089    match plan {
1090        DslPlan::IR { dsl, .. } => {
1091            let mut inner = Arc::unwrap_or_clone(dsl.clone());
1092            replace_bound(&mut inner, path, with);
1093            *plan = inner;
1094            return;
1095        }
1096        DslPlan::Slice {
1097            input,
1098            offset: 0,
1099            len,
1100        } if scans(input, path) => {
1101            if let Some(piece) = with(input, *len as usize) {
1102                *plan = piece;
1103            }
1104            return;
1105        }
1106        _ => {}
1107    }
1108    crate::table::for_each_input(plan, &mut |input| replace_bound(input, path, with));
1109}
1110
1111/// How many rows the frame `lf` reads of the followed file at `path`: its bound.
1112pub(crate) fn bound_of(lf: &LazyFrame, path: &Path) -> Option<usize> {
1113    use polars::lazy::dsl::DslPlan;
1114    let path = path.to_string_lossy();
1115    (&lf.logical_plan).into_iter().find_map(|node| match node {
1116        DslPlan::Slice {
1117            input,
1118            offset: 0,
1119            len,
1120        } if scans(input, &path) => Some(*len as usize),
1121        _ => None,
1122    })
1123}
1124
1125/// `root`, the bounded NDJSON frame of `path` as `format`, also reading `fields` after
1126/// its columns. `None` without a lines scan of the file, or if `fields` adds nothing.
1127pub(crate) fn widen(
1128    root: &LazyFrame,
1129    path: &Path,
1130    format: FileFormat,
1131    fields: &[Field],
1132    rows: usize,
1133) -> Option<LazyFrame> {
1134    let path_text = path.to_string_lossy();
1135    let mut plan = root.logical_plan.clone();
1136    let mut widened = None;
1137    replace_lines_scan(&mut plan, &path_text, &mut |scan| {
1138        let mut schema = (**scan.schema()).clone();
1139        for field in fields {
1140            if !schema.contains(field.name()) {
1141                schema.with_column(field.name().clone(), field.dtype().clone());
1142            }
1143        }
1144        if schema.len() == scan.schema().len() {
1145            return None;
1146        }
1147        let schema = Arc::new(schema);
1148        let lf = scan.with_schema(schema.clone()).lazy().ok()?;
1149        widened = Some((lf.clone(), schema));
1150        Some(lf.logical_plan)
1151    });
1152    let (raw, schema) = widened?;
1153    if format == FileFormat::Journal {
1154        // The journal names every column it shows: built again from the scan.
1155        let mut raw = raw;
1156        bound(&mut raw, path, rows);
1157        return Some(crate::formats::journal::derive(raw, &schema).0);
1158    }
1159    let mut out = root.clone();
1160    out.logical_plan = plan;
1161    Some(out)
1162}
1163
1164/// Put `with(scan)` where `plan` reads the file at `path` through a lines scan.
1165fn replace_lines_scan(
1166    plan: &mut polars::lazy::dsl::DslPlan,
1167    path: &str,
1168    with: &mut dyn FnMut(&lines::LinesScan) -> Option<polars::lazy::dsl::DslPlan>,
1169) {
1170    use polars::lazy::dsl::DslPlan;
1171    if let DslPlan::IR { dsl, .. } = plan {
1172        let mut inner = Arc::unwrap_or_clone(dsl.clone());
1173        replace_lines_scan(&mut inner, path, with);
1174        *plan = inner;
1175        return;
1176    }
1177    if let Some(scan) = lines::LinesScan::of(plan, path) {
1178        if let Some(replaced) = with(scan) {
1179            *plan = replaced;
1180        }
1181        return;
1182    }
1183    crate::table::for_each_input(plan, &mut |input| replace_lines_scan(input, path, with));
1184}
1185
1186/// The windows of a followed view, each read from the preceding mark. An unfiltered
1187/// view reads its rows straight; a filtering one reads on from `known` (view row, file
1188/// row) and slices.
1189pub(crate) struct Window {
1190    pub(crate) lf: LazyFrame,
1191    pub(crate) path: PathBuf,
1192    pub(crate) marks: Arc<Marks>,
1193    pub(crate) known: Option<Vec<(usize, usize)>>,
1194}
1195
1196impl crate::formats::pushdown::Windowed for Window {
1197    fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1198        let read = match &self.known {
1199            None => from_marks(&self.lf, &self.path, &self.marks, start, Some(start + len)),
1200            Some(known) => {
1201                let at = known.partition_point(|&(view, _)| view <= start);
1202                known.get(at.wrapping_sub(1)).and_then(|&(view, row)| {
1203                    from_marks(&self.lf, &self.path, &self.marks, row, None)
1204                        .map(|lf| lf.slice((start - view) as i64, len as IdxSize))
1205                })
1206            }
1207        };
1208        // Short of marks, Polars reads from the start of the file.
1209        Ok(read.unwrap_or_else(|| self.lf.clone().slice(start as i64, len as IdxSize)))
1210    }
1211}
1212
1213/// What the watcher found.
1214#[derive(Clone)]
1215pub enum Change {
1216    /// More complete rows: how many there are now, and how many of the rows that came
1217    /// in since the open do not fit the schema.
1218    Grew { rows: usize, misfits: usize },
1219    /// The file shrank or was replaced (truncated, rotated): it is read again from its
1220    /// start, which holds `rows` rows.
1221    Restarted { rows: usize, misfits: usize },
1222    /// The file is gone. `handle` still reads what it held.
1223    Gone { handle: Option<Arc<File>> },
1224    /// Fields that arrived in standard input's NDJSON after the open, which the schema
1225    /// does not have: sent once it has ended, just before [`Change::Ended`].
1226    NewFields(Vec<Field>),
1227    /// Standard input ended: with the reason when it ended in an error.
1228    Ended(Option<String>),
1229    /// The file could not be read.
1230    Failed(String),
1231}
1232
1233/// One report from a watcher, for the follow named `id`.
1234#[derive(Clone)]
1235pub struct News {
1236    pub(crate) id: u64,
1237    pub(crate) change: Change,
1238}
1239
1240/// What stops the watcher, and wakes it to look now.
1241#[derive(Default)]
1242struct Shared {
1243    stop: AtomicBool,
1244    poke: Mutex<bool>,
1245    woken: Condvar,
1246    /// Wakes a watcher waiting on inotify rather than on `woken`.
1247    #[cfg(target_os = "linux")]
1248    bell: notify::Bell,
1249}
1250
1251impl Shared {
1252    /// Wait out `interval`, or until poked or stopped. Whether to go on.
1253    fn wait(&self, interval: Duration) -> bool {
1254        let mut poked = self.poke.lock().unwrap_or_else(|e| e.into_inner());
1255        if !*poked && !self.stop.load(Ordering::Relaxed) {
1256            poked = self
1257                .woken
1258                .wait_timeout(poked, interval)
1259                .unwrap_or_else(|e| e.into_inner())
1260                .0;
1261        }
1262        *poked = false;
1263        !self.stop.load(Ordering::Relaxed)
1264    }
1265
1266    fn wake(&self) {
1267        *self.poke.lock().unwrap_or_else(|e| e.into_inner()) = true;
1268        self.woken.notify_all();
1269        #[cfg(target_os = "linux")]
1270        self.bell.ring();
1271    }
1272
1273    /// Wait for a change (`notify`), a poke or a stop, looking no sooner than `interval`
1274    /// after `last` so a burst is one look. Returns whether to go on.
1275    #[cfg(target_os = "linux")]
1276    fn wait_for_change(
1277        &self,
1278        notify: &notify::Notify,
1279        interval: Duration,
1280        last: &mut Option<Instant>,
1281    ) -> bool {
1282        loop {
1283            if self.stop.load(Ordering::Relaxed) {
1284                return false;
1285            }
1286            if std::mem::take(&mut *self.poke.lock().unwrap_or_else(|e| e.into_inner())) {
1287                break;
1288            }
1289            if notify.wait(&self.bell, None) == notify::Woke::Changed {
1290                let left = last
1291                    .map(|at| at + interval)
1292                    .and_then(|due| due.checked_duration_since(Instant::now()));
1293                // Returns early when poked.
1294                if let Some(left) = left
1295                    && !self.wait(left)
1296                {
1297                    return false;
1298                }
1299                break;
1300            }
1301        }
1302        // What changed before this look is read by it.
1303        notify.drain();
1304        *last = Some(Instant::now());
1305        !self.stop.load(Ordering::Relaxed)
1306    }
1307}
1308
1309/// How long ago, `elapsed`, as the follow chip says it: in the largest whole unit,
1310/// padded so the chip keeps its width as the number grows.
1311pub fn age(elapsed: Duration) -> String {
1312    let secs = elapsed.as_secs();
1313    match secs {
1314        0..60 => format!("{secs:>2}s ago"),
1315        60..3_600 => format!("{:>2}m ago", secs / 60),
1316        3_600..86_400 => format!("{:>2}h ago", secs / 3_600),
1317        _ => format!("{:>2}d ago", secs / 86_400),
1318    }
1319}
1320
1321/// When the age of an append made at `at` next reads differently.
1322pub fn next_tick(at: Instant) -> Instant {
1323    let secs = at.elapsed().as_secs();
1324    let unit = match secs {
1325        0..60 => 1,
1326        60..3_600 => 60,
1327        3_600..86_400 => 3_600,
1328        _ => 86_400,
1329    };
1330    at + Duration::from_secs((secs / unit + 1) * unit)
1331}
1332
1333/// Where the follow stands, as the footer says it.
1334#[derive(Clone, Debug, PartialEq)]
1335pub enum Standing {
1336    Following,
1337    Paused,
1338    /// Standard input ended, or the file went: what is on screen is all there is.
1339    Ended,
1340}
1341
1342/// A file being followed: its watcher, and what the view has taken of it. Belongs to
1343/// the dataset; dropping it stops the watcher.
1344pub struct Follow {
1345    id: u64,
1346    /// The file the frame scans.
1347    path: PathBuf,
1348    shared: Arc<Shared>,
1349    spool: Option<Arc<SpoolHandle>>,
1350    /// Rows the watcher has counted, and the misfits among them.
1351    counted: usize,
1352    misfits: usize,
1353    /// The file was read again from its start, and the view has not caught up.
1354    restarted: bool,
1355    /// Rows the frame reads now.
1356    shown: usize,
1357    /// Rows that arrived below the cursor while it was not on the last row.
1358    pub(crate) new_below: usize,
1359    pub(crate) standing: Standing,
1360    pub(crate) last_append: Option<Instant>,
1361    /// The view has not gone to the last row yet: a follow starts there, as `tail -f`
1362    /// does, once the first page is drawn.
1363    pub(crate) settle_at_end: bool,
1364    /// The cursor was on the last row when rows arrived: it goes to the new last row
1365    /// once they are read.
1366    pub(crate) end_pending: bool,
1367    /// Rows were taken while the table was not on screen: it reads its rows again
1368    /// when it is.
1369    pub(crate) stale_view: bool,
1370    /// The handle a deleted file is read through from now on.
1371    held: Option<Arc<File>>,
1372    /// Where its rows start, every so many.
1373    marks: Arc<Marks>,
1374    /// Standard input read as it arrives without `--follow`: the view stays where it
1375    /// is, and the rows so far are a part of what is coming.
1376    pipe: bool,
1377    /// Fields that arrived after the open, for the frames to join once standard input
1378    /// has ended.
1379    new_fields: Vec<Field>,
1380    /// What the format says of the whole stream (the journal's Info tab) has been
1381    /// asked for again, now that it has ended.
1382    pub(crate) described: bool,
1383}
1384
1385static NEXT_ID: AtomicU64 = AtomicU64::new(1);
1386
1387impl Follow {
1388    /// Follow the file `tail` counted (the frame reads its first `tail.rows()` rows),
1389    /// checking every `interval` and telling `events`. `spool` is stdin being copied to it.
1390    pub fn start(
1391        mut tail: Tail,
1392        interval: Duration,
1393        events: Sender<AppEvent>,
1394        spool: Option<Arc<SpoolHandle>>,
1395    ) -> Follow {
1396        let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
1397        let path = tail.path.clone();
1398        let shared = Arc::new(Shared::default());
1399        let shown = tail.rows();
1400        if let Some(handle) = &spool {
1401            handle.spool.wake_on_end(shared.clone());
1402        }
1403        let marks = Arc::new(Marks::default());
1404        marks.take_from(&mut tail);
1405        let watcher = Watcher {
1406            marks: marks.clone(),
1407            id,
1408            path: path.clone(),
1409            tail,
1410            shared: shared.clone(),
1411            events,
1412            spool: spool.as_ref().map(|handle| handle.spool.clone()),
1413            interval,
1414        };
1415        let _ = std::thread::Builder::new()
1416            .name("datui-follow".to_string())
1417            .spawn(move || watcher.run());
1418        Follow {
1419            id,
1420            path,
1421            shared,
1422            spool,
1423            counted: shown,
1424            misfits: 0,
1425            restarted: false,
1426            shown,
1427            new_below: 0,
1428            standing: Standing::Following,
1429            last_append: None,
1430            settle_at_end: true,
1431            end_pending: false,
1432            stale_view: false,
1433            held: None,
1434            marks,
1435            pipe: false,
1436            new_fields: Vec::new(),
1437            described: false,
1438        }
1439    }
1440
1441    /// This follows standard input read as it arrives rather than `--follow`: the
1442    /// view starts at the top and stays where it is put.
1443    pub fn as_pipe(mut self) -> Follow {
1444        self.pipe = true;
1445        self.settle_at_end = false;
1446        self
1447    }
1448
1449    /// Whether this is standard input read as it arrives ([`Self::as_pipe`]).
1450    pub fn is_pipe(&self) -> bool {
1451        self.pipe
1452    }
1453
1454    /// Whether more rows may still come: followed and not ended.
1455    pub fn live(&self) -> bool {
1456        self.standing != Standing::Ended
1457    }
1458
1459    /// Where the file's rows start, every so many.
1460    pub(crate) fn marks(&self) -> &Arc<Marks> {
1461        &self.marks
1462    }
1463
1464    pub fn id(&self) -> u64 {
1465        self.id
1466    }
1467
1468    pub fn path(&self) -> &Path {
1469        &self.path
1470    }
1471
1472    /// Rows the frame reads.
1473    pub fn shown(&self) -> usize {
1474        self.shown
1475    }
1476
1477    /// Rows counted that the view has not taken yet.
1478    pub fn waiting(&self) -> usize {
1479        self.counted.saturating_sub(self.shown)
1480    }
1481
1482    pub fn misfits(&self) -> usize {
1483        self.misfits
1484    }
1485
1486    pub fn standing(&self) -> &Standing {
1487        &self.standing
1488    }
1489
1490    /// Rows that arrived below the cursor while it was off the last row.
1491    pub fn new_below(&self) -> usize {
1492        self.new_below
1493    }
1494
1495    /// Whether the view has rows or a restart to take: not while paused. The rows
1496    /// counted before standard input ended are taken after it did.
1497    pub fn behind(&self) -> bool {
1498        self.standing != Standing::Paused && (self.restarted || self.counted != self.shown)
1499    }
1500
1501    /// Standard input being copied, if this follows it.
1502    pub fn spool(&self) -> Option<&Arc<Spool>> {
1503        self.spool.as_ref().map(|handle| &handle.spool)
1504    }
1505
1506    /// Look now rather than at the end of the interval.
1507    pub fn check_now(&self) {
1508        self.shared.wake();
1509    }
1510
1511    /// Take a report from this follow's watcher. Returns what the user is told, if
1512    /// anything.
1513    pub fn take(&mut self, change: &Change) -> Option<String> {
1514        match change {
1515            Change::Grew { rows, misfits } => {
1516                if *rows > self.counted {
1517                    self.last_append = Some(Instant::now());
1518                }
1519                self.counted = *rows;
1520                self.misfits = *misfits;
1521                None
1522            }
1523            Change::Restarted { rows, misfits } => {
1524                self.counted = *rows;
1525                self.misfits = *misfits;
1526                self.restarted = true;
1527                self.last_append = Some(Instant::now());
1528                Some("The file was truncated or replaced: reading it from the start".to_string())
1529            }
1530            Change::NewFields(fields) => {
1531                self.new_fields = fields.clone();
1532                None
1533            }
1534            Change::Gone { handle } => {
1535                self.held = handle.clone();
1536                self.end();
1537                Some("The file was deleted: following stopped, the rows read stay".to_string())
1538            }
1539            // A recording's end is said by the recording's own mark.
1540            Change::Ended(_) if self.spool().is_some_and(|s| s.tee().is_some()) => {
1541                self.end();
1542                None
1543            }
1544            Change::Ended(None) => {
1545                self.end();
1546                Some("Standard input ended".to_string())
1547            }
1548            Change::Ended(Some(reason)) | Change::Failed(reason) => {
1549                self.end();
1550                Some(reason.clone())
1551            }
1552        }
1553    }
1554
1555    /// The view takes what was counted: the rows its frame reads from now on, and
1556    /// whether the file was read again from its start.
1557    pub fn catch_up(&mut self) -> (usize, bool) {
1558        self.shown = self.counted;
1559        (self.shown, std::mem::take(&mut self.restarted))
1560    }
1561
1562    /// The fields that arrived after the open, once, for the frames to join.
1563    pub(crate) fn take_new_fields(&mut self) -> Vec<Field> {
1564        std::mem::take(&mut self.new_fields)
1565    }
1566
1567    /// The handle a deleted file is read through, once, for the frames to take.
1568    pub fn take_held(&mut self) -> Option<Arc<File>> {
1569        self.held.take()
1570    }
1571
1572    pub fn pause(&mut self) {
1573        if self.standing == Standing::Following {
1574            self.standing = Standing::Paused;
1575        }
1576    }
1577
1578    pub fn resume(&mut self) {
1579        if self.standing == Standing::Paused {
1580            self.standing = Standing::Following;
1581        }
1582    }
1583
1584    /// Stop watching. What the frame reads stays. Standard input stops being read,
1585    /// unless it is being recorded (`--tee`): that is stopped only when asked.
1586    pub fn end(&mut self) {
1587        self.standing = Standing::Ended;
1588        self.shared.stop.store(true, Ordering::Relaxed);
1589        self.shared.wake();
1590        if let Some(spool) = self.spool.as_ref().filter(|s| s.spool.tee().is_none()) {
1591            spool.spool.stop();
1592        }
1593    }
1594}
1595
1596impl Drop for Follow {
1597    fn drop(&mut self) {
1598        self.shared.stop.store(true, Ordering::Relaxed);
1599        self.shared.wake();
1600    }
1601}
1602
1603/// The thread that watches the file.
1604struct Watcher {
1605    id: u64,
1606    path: PathBuf,
1607    tail: Tail,
1608    shared: Arc<Shared>,
1609    events: Sender<AppEvent>,
1610    spool: Option<Arc<Spool>>,
1611    interval: Duration,
1612    marks: Arc<Marks>,
1613}
1614
1615/// Which file this is, so a replaced file is told from a grown one: (device, inode) on
1616/// Unix, (volume serial, file index) on Windows.
1617type Identity = (u64, u64);
1618
1619#[cfg(unix)]
1620fn identity_of(file: &File) -> Option<Identity> {
1621    use std::os::unix::fs::MetadataExt;
1622    file.metadata().ok().map(|meta| (meta.dev(), meta.ino()))
1623}
1624
1625#[cfg(windows)]
1626fn identity_of(file: &File) -> Option<Identity> {
1627    use std::os::windows::io::AsRawHandle;
1628    use windows_sys::Win32::Storage::FileSystem::{
1629        BY_HANDLE_FILE_INFORMATION, GetFileInformationByHandle,
1630    };
1631    // SAFETY: the handle is `file`'s, open while this runs, and `info` is plain data
1632    // the call fills in; zeroed is a valid value of it.
1633    let mut info: BY_HANDLE_FILE_INFORMATION = unsafe { std::mem::zeroed() };
1634    let ok = unsafe { GetFileInformationByHandle(file.as_raw_handle(), &mut info) };
1635    (ok != 0).then(|| {
1636        (
1637            u64::from(info.dwVolumeSerialNumber),
1638            u64::from(info.nFileIndexHigh) << 32 | u64::from(info.nFileIndexLow),
1639        )
1640    })
1641}
1642
1643#[cfg(not(any(unix, windows)))]
1644fn identity_of(_file: &File) -> Option<Identity> {
1645    None
1646}
1647
1648/// Which file `path` names now, `meta` its metadata. Unix reads it from the metadata;
1649/// Windows has to open the file to ask.
1650#[cfg(unix)]
1651fn identity_at(_path: &Path, meta: &std::fs::Metadata) -> Option<Identity> {
1652    use std::os::unix::fs::MetadataExt;
1653    Some((meta.dev(), meta.ino()))
1654}
1655
1656#[cfg(not(unix))]
1657fn identity_at(path: &Path, _meta: &std::fs::Metadata) -> Option<Identity> {
1658    File::open(path).ok().as_ref().and_then(identity_of)
1659}
1660
1661/// Whether the file a path names is another one than the file followed. Unknown on
1662/// either side is not a replacement: a shrink still tells a truncation.
1663fn replaced(known: Option<Identity>, now: Option<Identity>) -> bool {
1664    matches!((known, now), (Some(known), Some(now)) if known != now)
1665}
1666
1667/// Why following `path` (`None`: standard input) ended: `doing` when `e` stopped it.
1668fn failed_message(path: Option<&Path>, doing: &str, e: &std::io::Error) -> String {
1669    let what = format!(
1670        "{doing}. {}",
1671        crate::error_display::user_message_from_io(e, None)
1672    );
1673    match path {
1674        Some(path) => crate::error_display::file_message(path, &what),
1675        None => crate::error_display::sentence(&format!("standard input: {what}")),
1676    }
1677}
1678
1679impl Watcher {
1680    /// Why following ended, `doing` when `e` stopped it, naming the file followed;
1681    /// standard input's spool is a file the user never named.
1682    fn failed(&self, doing: &str, e: &std::io::Error) -> String {
1683        failed_message(
1684            self.spool.is_none().then_some(self.path.as_path()),
1685            doing,
1686            e,
1687        )
1688    }
1689
1690    fn run(mut self) {
1691        let mut file = match File::open(&self.path) {
1692            Ok(file) => file,
1693            Err(e) => {
1694                self.send(Change::Failed(self.failed("following it stopped", &e)));
1695                return;
1696            }
1697        };
1698        // The file the open counted, not the one opened here: a file put in its place
1699        // meanwhile is a replacement, and the next look reads it from the start.
1700        let mut known = self.tail.identity.or_else(|| identity_of(&file));
1701        let mut sent = (self.tail.rows(), 0usize);
1702        #[cfg(target_os = "linux")]
1703        let notify = notify::Notify::new(&self.path);
1704        #[cfg(target_os = "linux")]
1705        let mut last = None;
1706        loop {
1707            #[cfg(target_os = "linux")]
1708            let go_on = match &notify {
1709                Some(notify) => self
1710                    .shared
1711                    .wait_for_change(notify, self.interval, &mut last),
1712                None => self.shared.wait(self.interval),
1713            };
1714            #[cfg(not(target_os = "linux"))]
1715            let go_on = self.shared.wait(self.interval);
1716            if !go_on {
1717                return;
1718            }
1719            // Read before the file, so nothing the spool wrote before it ended is missed.
1720            let spool_ended = self.spool.as_ref().and_then(|spool| spool.ended());
1721            let meta = match std::fs::metadata(&self.path) {
1722                Ok(meta) => meta,
1723                Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
1724                    self.send(Change::Gone {
1725                        handle: Some(Arc::new(file)),
1726                    });
1727                    return;
1728                }
1729                Err(e) => {
1730                    self.send(Change::Failed(self.failed("following it stopped", &e)));
1731                    return;
1732                }
1733            };
1734            let now = identity_at(&self.path, &meta);
1735            let len = meta.len();
1736            if replaced(known, now) || len < self.tail.complete() {
1737                match File::open(&self.path) {
1738                    Ok(reopened) => file = reopened,
1739                    Err(e) => {
1740                        self.send(Change::Failed(self.failed("following it stopped", &e)));
1741                        return;
1742                    }
1743                }
1744                known = identity_of(&file);
1745                #[cfg(target_os = "linux")]
1746                if let Some(notify) = &notify {
1747                    notify.rewatch(&self.path);
1748                }
1749                self.tail.restart();
1750                self.marks.clear();
1751                if let Err(e) = self.tail.read_on(&mut file, len, false) {
1752                    self.send(Change::Failed(self.failed("reading it stopped", &e)));
1753                    return;
1754                }
1755                self.marks.take_from(&mut self.tail);
1756                sent = (self.tail.rows(), self.tail.misfits());
1757                self.send(Change::Restarted {
1758                    rows: sent.0,
1759                    misfits: sent.1,
1760                });
1761                continue;
1762            }
1763            let read = if spool_ended.is_some() {
1764                self.tail.read_to_end(&mut file, len, true)
1765            } else {
1766                self.tail.read_on(&mut file, len, true)
1767            };
1768            if let Err(e) = read {
1769                self.send(Change::Failed(self.failed("reading it stopped", &e)));
1770                return;
1771            }
1772            // Before the rows are reported, so the view's reads of them find marks.
1773            self.marks.take_from(&mut self.tail);
1774            let now_counted = (self.tail.rows(), self.tail.misfits());
1775            if now_counted != sent {
1776                sent = now_counted;
1777                if !self.send(Change::Grew {
1778                    rows: sent.0,
1779                    misfits: sent.1,
1780                }) {
1781                    return;
1782                }
1783            }
1784            if let Some(reason) = spool_ended {
1785                let fields = self.tail.new_fields();
1786                if !fields.is_empty() && !self.send(Change::NewFields(fields)) {
1787                    return;
1788                }
1789                self.send(Change::Ended(reason));
1790                return;
1791            }
1792        }
1793    }
1794
1795    /// Whether the app is still there to hear it.
1796    fn send(&self, change: Change) -> bool {
1797        self.events
1798            .send(AppEvent::Followed(News {
1799                id: self.id,
1800                change,
1801            }))
1802            .is_ok()
1803    }
1804}
1805
1806/// Stdin copied to a file while it is read, the copy continuing after the first rows
1807/// show until the stream ends or is stopped. The file is temporary or the one `--tee`
1808/// names (kept); `--tee -` also passes the stream to stdout.
1809pub struct Spool {
1810    stop: AtomicBool,
1811    bytes: AtomicU64,
1812    state: Mutex<SpoolState>,
1813    changed: Condvar,
1814    /// The file being written. Taken when the copy finishes, so a read still waiting on
1815    /// the producer writes nothing after it.
1816    sink: Mutex<Option<File>>,
1817    /// The file `--tee` named, when it is the one written.
1818    tee: Option<Tee>,
1819    /// Standard output, for `--tee -`. Taken when the copy finishes, which closes it.
1820    pass: Mutex<Option<Box<dyn Write + Send>>>,
1821    started: Instant,
1822}
1823
1824/// The file `--tee` named.
1825#[derive(Clone, Debug)]
1826pub struct Tee {
1827    pub path: PathBuf,
1828    /// `--tee-raw`: the bytes exactly as they came, a WAV header's sizes included.
1829    pub raw: bool,
1830}
1831
1832impl Tee {
1833    /// `--tee -`: the stream is passed on to standard output rather than kept in a file.
1834    pub fn to_stdout(&self) -> bool {
1835        crate::loading::stdin::is_stdin(&self.path)
1836    }
1837
1838    /// Where the stream goes, as a message names it.
1839    pub fn name(&self) -> String {
1840        if self.to_stdout() {
1841            return "standard output".to_string();
1842        }
1843        self.path.file_name().map_or_else(
1844            || self.path.display().to_string(),
1845            |name| name.to_string_lossy().into_owned(),
1846        )
1847    }
1848}
1849
1850#[derive(Default)]
1851struct SpoolState {
1852    /// Newlines copied so far, up to what the open waits for.
1853    lines: usize,
1854    /// The last read took all the producer had: it is slower than the copy.
1855    drained: bool,
1856    /// The stream ended: `Some(reason)` when in an error.
1857    ended: Option<Option<String>>,
1858    /// When the copy finished, and what finishing the file said.
1859    finished: Option<Instant>,
1860    /// Bytes copied by when, a few seconds of them, for the rate.
1861    samples: std::collections::VecDeque<(Instant, u64)>,
1862    /// The watcher following the file, woken when the copy ends: a watcher waiting
1863    /// for the file to change would not hear an end that writes nothing.
1864    watcher: Option<Arc<Shared>>,
1865}
1866
1867/// How far back the rate looks.
1868const RATE_WINDOW: Duration = Duration::from_secs(2);
1869
1870impl Spool {
1871    fn new(sink: File, tee: Option<Tee>, pass: Option<Box<dyn Write + Send>>) -> Spool {
1872        Spool {
1873            stop: AtomicBool::new(false),
1874            bytes: AtomicU64::new(0),
1875            state: Mutex::new(SpoolState::default()),
1876            changed: Condvar::new(),
1877            sink: Mutex::new(Some(sink)),
1878            tee,
1879            pass: Mutex::new(pass),
1880            started: Instant::now(),
1881        }
1882    }
1883
1884    /// Bytes copied so far.
1885    pub fn bytes(&self) -> u64 {
1886        self.bytes.load(Ordering::Relaxed)
1887    }
1888
1889    /// Bytes a second over the last few seconds.
1890    pub fn rate(&self) -> f64 {
1891        let state = self.lock();
1892        match (state.samples.front(), state.samples.back()) {
1893            (Some((t0, b0)), Some((t1, b1))) if t1 > t0 => {
1894                (b1 - b0) as f64 / t1.duration_since(*t0).as_secs_f64()
1895            }
1896            _ => 0.0,
1897        }
1898    }
1899
1900    /// How long the copy ran, or has run.
1901    pub fn duration(&self) -> Duration {
1902        let finished = self.lock().finished;
1903        finished
1904            .unwrap_or_else(Instant::now)
1905            .duration_since(self.started)
1906    }
1907
1908    /// The file `--tee` named, when the copy goes there.
1909    pub fn tee(&self) -> Option<&Tee> {
1910        self.tee.as_ref()
1911    }
1912
1913    /// Stop copying and finish the file: a read waiting on the producer writes
1914    /// nothing more when it returns.
1915    pub fn stop(&self) {
1916        self.stop.store(true, Ordering::Relaxed);
1917        self.finish(None);
1918    }
1919
1920    pub fn stopped(&self) -> bool {
1921        self.stop.load(Ordering::Relaxed)
1922    }
1923
1924    /// Whether the copy ended, and the reason when in an error.
1925    pub fn ended(&self) -> Option<Option<String>> {
1926        self.lock().ended.clone()
1927    }
1928
1929    /// Whether bytes may still arrive: the producer is still sending.
1930    pub fn live(&self) -> bool {
1931        self.lock().ended.is_none()
1932    }
1933
1934    /// Wait until the copy has ended.
1935    pub fn wait(&self) {
1936        let mut state = self.lock();
1937        while state.ended.is_none() {
1938            state = self.changed.wait(state).unwrap_or_else(|e| e.into_inner());
1939        }
1940    }
1941
1942    fn lock(&self) -> std::sync::MutexGuard<'_, SpoolState> {
1943        self.state.lock().unwrap_or_else(|e| e.into_inner())
1944    }
1945
1946    /// Write `bytes` as they came. False once the copy is finished.
1947    fn write(&self, bytes: &[u8]) -> Result<bool, String> {
1948        let mut sink = self.sink.lock().unwrap_or_else(|e| e.into_inner());
1949        let Some(file) = sink.as_mut() else {
1950            return Ok(false);
1951        };
1952        file.write_all(bytes).map_err(|e| {
1953            format!(
1954                "Could not write {}: {e}",
1955                self.tee
1956                    .as_ref()
1957                    .filter(|t| !t.to_stdout())
1958                    .map_or("what came in".to_string(), |t| t.path.display().to_string())
1959            )
1960        })?;
1961        drop(sink);
1962        // Outside the file's lock: a reader downstream that stops reading holds up this
1963        // write, and must not hold up a stop.
1964        if let Some(out) = self.pass.lock().unwrap_or_else(|e| e.into_inner()).as_mut() {
1965            out.write_all(bytes)
1966                .and_then(|()| out.flush())
1967                .map_err(|e| format!("Could not write standard output: {e}"))?;
1968        }
1969        let total =
1970            self.bytes.fetch_add(bytes.len() as u64, Ordering::Relaxed) + bytes.len() as u64;
1971        let now = Instant::now();
1972        let mut state = self.lock();
1973        if state.lines < WANTED_LINES {
1974            state.lines += bytes.iter().filter(|&&b| b == b'\n').count();
1975        }
1976        state.samples.push_back((now, total));
1977        while state
1978            .samples
1979            .front()
1980            .is_some_and(|(at, _)| now.duration_since(*at) > RATE_WINDOW)
1981            && state.samples.len() > 2
1982        {
1983            state.samples.pop_front();
1984        }
1985        drop(state);
1986        self.changed.notify_all();
1987        Ok(true)
1988    }
1989
1990    /// End the copy (`reason` if by error) and finish the file: a WAV header's sizes filled
1991    /// for `--tee` (unless `--tee-raw`), synced so saved means safe. Once.
1992    fn finish(&self, reason: Option<String>) {
1993        let file = self.sink.lock().unwrap_or_else(|e| e.into_inner()).take();
1994        // Close the passthrough so the downstream reader sees the end; if a blocked write
1995        // holds it, it closes when that returns and the copy finishes again.
1996        if let Ok(mut pass) = self.pass.try_lock() {
1997            pass.take();
1998        }
1999        let mut reason = reason;
2000        if let (Some(mut file), Some(tee)) = (file, self.tee.as_ref().filter(|t| !t.to_stdout())) {
2001            let finished = (if tee.raw {
2002                Ok(())
2003            } else {
2004                crate::loading::tee::fix_wav_sizes(&mut file).map(|_| ())
2005            })
2006            .and_then(|()| file.sync_all());
2007            if let Err(e) = finished
2008                && reason.is_none()
2009            {
2010                reason = Some(format!("Could not finish {}: {e}", tee.path.display()));
2011            }
2012        }
2013        let mut state = self.lock();
2014        if state.ended.is_none() {
2015            state.ended = Some(reason);
2016            state.finished = Some(Instant::now());
2017        }
2018        let watcher = state.watcher.take();
2019        drop(state);
2020        self.changed.notify_all();
2021        if let Some(watcher) = watcher {
2022            watcher.wake();
2023        }
2024    }
2025
2026    /// Wake the watcher `shared` once the copy ends, or now if it has.
2027    fn wake_on_end(&self, shared: Arc<Shared>) {
2028        let mut state = self.lock();
2029        if state.ended.is_some() {
2030            drop(state);
2031            shared.wake();
2032        } else {
2033            state.watcher = Some(shared);
2034        }
2035    }
2036}
2037
2038/// The open's hold on a [`Spool`]: the copy stops when the last holder lets go (a
2039/// follow put down early stops copying), and quitting finishes the file.
2040pub struct SpoolHandle {
2041    spool: Arc<Spool>,
2042}
2043
2044impl SpoolHandle {
2045    pub fn spool(&self) -> &Arc<Spool> {
2046        &self.spool
2047    }
2048}
2049
2050impl Drop for SpoolHandle {
2051    fn drop(&mut self) {
2052        self.spool.stop();
2053    }
2054}
2055
2056/// Copy `reader` into `spool` on its own thread until it ends or the spool stops, each
2057/// chunk written as it arrives through one reused buffer (about a megabyte however
2058/// long the stream). A producer faster than the disk waits on the pipe.
2059fn copy_on(mut reader: impl Read + Send + 'static, spool: Arc<Spool>) {
2060    let _ = std::thread::Builder::new()
2061        .name("datui-spool".to_string())
2062        .spawn(move || {
2063            let mut buf = vec![0u8; CHUNK];
2064            let reason = loop {
2065                if spool.stopped() {
2066                    break None;
2067                }
2068                let n = match reader.read(&mut buf) {
2069                    Ok(0) => break None,
2070                    Ok(n) => n,
2071                    Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
2072                    Err(e) => break Some(format!("Standard input failed: {e}")),
2073                };
2074                match spool.write(&buf[..n]) {
2075                    Ok(true) => {}
2076                    Ok(false) => break None,
2077                    Err(reason) => break Some(reason),
2078                }
2079                spool.lock().drained = n < buf.len();
2080            };
2081            spool.finish(reason);
2082        });
2083}
2084
2085/// Lines the open waits for before it reads the spooled file, unless the producer is
2086/// slower than the copy: then two, a header and a row, are enough to show.
2087const WANTED_LINES: usize = 1000;
2088
2089/// What standard input was copied to.
2090pub enum Spooled {
2091    /// A temporary file, removed when its holders let go.
2092    Temp(TempDownload),
2093    /// The file `--tee` named, the user's.
2094    Kept(PathBuf),
2095}
2096
2097/// Copy stdin from `open` to the `--tee` file or a temp file in
2098/// [`crate::loading::stdin::spool_dir`] (claimed via `writer`) until enough has arrived to show;
2099/// the copy continues behind. Reports what the file holds, as [`crate::loading::stdin::spool`]
2100/// does. Unfollowed stdin is still read as it arrives when its format allows
2101/// (`OpenOptions::pipe`); otherwise (Parquet, compressed) it is copied to the end first.
2102pub(crate) fn spool<R: Read + Send + 'static>(
2103    open: impl FnOnce() -> crate::cloud::download::Opened<R>,
2104    options: OpenOptions,
2105    writer: &Writer,
2106    read: &AtomicU64,
2107    stdout: Option<Box<dyn Write + Send>>,
2108) -> Result<(Spooled, OpenOptions), String> {
2109    let tee = options.tee.clone().map(|path| Tee {
2110        path,
2111        raw: options.tee_raw,
2112    });
2113    // A pipe cannot be sought back to, to fill in a WAV header.
2114    let tee = tee.map(|tee| Tee {
2115        raw: tee.raw || tee.to_stdout(),
2116        ..tee
2117    });
2118    let pass = match &tee {
2119        Some(tee) if tee.to_stdout() => Some(stdout.ok_or_else(|| {
2120            "--tee - passes the stream on to standard output, which only the datui command has."
2121                .to_string()
2122        })?),
2123        _ => None,
2124    };
2125    let progressive = !options.follow && tee.is_none();
2126    let (reader, _) = open().map_err(|e| format!("Could not read standard input: {e}"))?;
2127    let (spooled, file) = match &tee {
2128        Some(tee) if !tee.to_stdout() => {
2129            let file = crate::loading::tee::create(&tee.path, options.force)?;
2130            (Spooled::Kept(tee.path.clone()), file)
2131        }
2132        _ => {
2133            let dir = crate::loading::stdin::spool_dir(&options);
2134            let Some((named, claim)) = writer
2135                .create(|| TempDownload::create(dir.as_deref(), None))
2136                .map_err(|e| crate::error_display::user_message_from_report(&e, None))?
2137            else {
2138                return Err("Reading standard input was stopped.".to_string());
2139            };
2140            let file = named
2141                .as_file()
2142                .try_clone()
2143                .map_err(|e| format!("Could not write what came in: {e}"))?;
2144            (Spooled::Temp(TempDownload::held(named, Some(claim))), file)
2145        }
2146    };
2147    let spooled_path = match &spooled {
2148        Spooled::Temp(download) => download.path().to_path_buf(),
2149        Spooled::Kept(path) => path.clone(),
2150    };
2151    let spool = Arc::new(Spool::new(file, tee, pass));
2152    let handle = Arc::new(SpoolHandle {
2153        spool: spool.clone(),
2154    });
2155    copy_on(reader, spool.clone());
2156    // A header and a row, at least, before anything is read.
2157    let wanted = if options.has_header == Some(false) {
2158        1
2159    } else {
2160        2
2161    };
2162    let mut state = spool.lock();
2163    loop {
2164        read.store(spool.bytes(), Ordering::Relaxed);
2165        if writer.stopped() {
2166            drop(state);
2167            spool.stop();
2168            return Err("Reading standard input was stopped.".to_string());
2169        }
2170        // An Arrow stream has no lines to count: its schema message is enough.
2171        let enough = (options.follow || progressive)
2172            && (state.lines >= WANTED_LINES
2173                || (state.drained
2174                    && (state.lines >= wanted || stream::begins_with_schema(&spooled_path))));
2175        if enough || state.ended.is_some() {
2176            break;
2177        }
2178        state = spool
2179            .changed
2180            .wait_timeout(state, Duration::from_millis(100))
2181            .unwrap_or_else(|e| e.into_inner())
2182            .0;
2183    }
2184    if let Some(Some(reason)) = &state.ended {
2185        return Err(reason.clone());
2186    }
2187    drop(state);
2188    read.store(spool.bytes(), Ordering::Relaxed);
2189    let path = spooled_path;
2190    let mut head = Vec::new();
2191    File::open(&path)
2192        .and_then(|f| f.take(4096).read_to_end(&mut head))
2193        .map_err(|e| format!("Could not read standard input back: {e}"))?;
2194    if head.is_empty() {
2195        return Err("Nothing came in on standard input.".to_string());
2196    }
2197    let (format, compression, guessed) = crate::loading::stdin::sniff_for(&head, &options);
2198    let asked = options.clone();
2199    let options = OpenOptions {
2200        format_guessed: options.format.is_none() && guessed,
2201        format: options.format.or(Some(format)),
2202        compression: options.compression.or(compression),
2203        ..options
2204    };
2205    if progressive {
2206        let read_on = followed_stream(&path, options.format, &options)
2207            || refusal(options.format, &options).is_none();
2208        if read_on {
2209            return Ok((
2210                spooled,
2211                OpenOptions {
2212                    spool: Some(handle),
2213                    follow: true,
2214                    pipe: true,
2215                    ..options
2216                },
2217            ));
2218        }
2219        // Read once it is finished: the rest is copied first, then looked at whole.
2220        while spool.ended().is_none() {
2221            if writer.stopped() {
2222                spool.stop();
2223                return Err("Reading standard input was stopped.".to_string());
2224            }
2225            read.store(spool.bytes(), Ordering::Relaxed);
2226            let state = spool.lock();
2227            if state.ended.is_none() {
2228                let _ = spool
2229                    .changed
2230                    .wait_timeout(state, Duration::from_millis(100))
2231                    .unwrap_or_else(|e| e.into_inner());
2232            }
2233        }
2234        if let Some(Some(reason)) = spool.ended() {
2235            return Err(reason);
2236        }
2237        read.store(spool.bytes(), Ordering::Relaxed);
2238        let options = crate::loading::stdin::described(&path, asked)?;
2239        return Ok((spooled, options));
2240    }
2241    // A recording goes on whatever it holds; only the view is not followed then.
2242    if options.follow
2243        && options.tee.is_none()
2244        && !followed_stream(&path, options.format, &options)
2245        && let Some(refusal) = refusal(options.format, &options)
2246    {
2247        return Err(refusal);
2248    }
2249    Ok((
2250        spooled,
2251        OpenOptions {
2252            spool: Some(handle),
2253            // The file is read from here on, as any file is.
2254            tee: None,
2255            ..options
2256        },
2257    ))
2258}
2259
2260#[cfg(test)]
2261mod tests {
2262    use super::*;
2263
2264    /// Following that stops on an error names the file followed, in the one shape;
2265    /// standard input is called that.
2266    #[test]
2267    fn errors_name_the_file() {
2268        let path = Path::new("/data/app.log");
2269        for (doing, e) in [
2270            (
2271                "following it stopped",
2272                std::io::Error::from(std::io::ErrorKind::NotFound),
2273            ),
2274            (
2275                "reading it stopped",
2276                std::io::Error::from(std::io::ErrorKind::PermissionDenied),
2277            ),
2278        ] {
2279            let message = failed_message(Some(path), doing, &e);
2280            eprintln!("{message}");
2281            crate::formats::readers::bad_input::assert_shape(&message, path);
2282            assert!(message.contains(&doing[1..]), "{message}");
2283        }
2284        let e = std::io::Error::from(std::io::ErrorKind::UnexpectedEof);
2285        let message = failed_message(None, "reading it stopped", &e);
2286        assert_eq!(
2287            message,
2288            "Standard input: reading it stopped. Unexpected end of file."
2289        );
2290    }
2291
2292    fn tail_of(text: &[u8], format: FileFormat, options: &OpenOptions) -> Tail {
2293        let dir = tempfile::tempdir().unwrap();
2294        let path = dir.path().join("t");
2295        std::fs::write(&path, text).unwrap();
2296        let mut tail = Tail::new(format, options, &Schema::default());
2297        let mut file = File::open(&path).unwrap();
2298        tail.read_on(&mut file, text.len() as u64, false).unwrap();
2299        tail
2300    }
2301
2302    /// Polars' own count of the complete records of `text` read as CSV.
2303    fn polars_rows(text: &[u8], options: &OpenOptions) -> usize {
2304        let complete = text.iter().rposition(|&b| b == b'\n').map_or(0, |i| i + 1);
2305        let mut read = CsvReadOptions::default().with_has_header(options.has_header != Some(false));
2306        read = read.map_parse_options(|p| {
2307            p.with_comment_prefix(
2308                options
2309                    .comment_char
2310                    .as_deref()
2311                    .map(polars::io::csv::read::CommentPrefix::new_from_str),
2312            )
2313        });
2314        if let Some(skip) = options.skip_lines {
2315            read.skip_lines = skip;
2316        }
2317        CsvReader::new(std::io::Cursor::new(text[..complete].to_vec()))
2318            .with_options(read)
2319            .finish()
2320            .map(|df| df.height())
2321            .unwrap_or(0)
2322    }
2323
2324    /// Rows counted as Polars counts them: blank lines are rows of nulls, a quoted
2325    /// newline is not a record's end, comments and skipped lines are not rows, and a
2326    /// partial last line waits.
2327    #[test]
2328    fn records_are_counted_as_polars_reads_them() {
2329        let cases: [(&[u8], OpenOptions); 6] = [
2330            (b"a,b\n1,2\n3,4\n", OpenOptions::default()),
2331            (b"a,b\n1,2\n\n3,4\n5,", OpenOptions::default()),
2332            (b"a,b\n1,\"x\ny\"\n3,4\n", OpenOptions::default()),
2333            (b"a,b\r\n1,2\r\n3,4\r\n", OpenOptions::default()),
2334            (
2335                b"#c\na,b\n#x\n1,2\n3,4\n",
2336                OpenOptions {
2337                    comment_char: Some("#".to_string()),
2338                    ..Default::default()
2339                },
2340            ),
2341            (
2342                b"junk\na,b\n1,2\n",
2343                OpenOptions::default().with_skip_lines(1),
2344            ),
2345        ];
2346        for (text, options) in cases {
2347            let tail = tail_of(text, FileFormat::Csv, &options);
2348            assert_eq!(
2349                tail.rows(),
2350                polars_rows(text, &options),
2351                "{}",
2352                String::from_utf8_lossy(text)
2353            );
2354        }
2355        let lines = tail_of(
2356            b"{\"a\":1}\n\n{\"a\":2}\n{\"a\":",
2357            FileFormat::Jsonl,
2358            &Default::default(),
2359        );
2360        assert_eq!(lines.rows(), 2);
2361        assert_eq!(lines.complete(), 17);
2362    }
2363
2364    /// On Linux the watcher hears an append through inotify: with an interval of an
2365    /// hour, no size check would see it, and nothing pokes it.
2366    #[cfg(target_os = "linux")]
2367    #[test]
2368    fn an_append_is_heard_of_without_a_check() {
2369        let dir = tempfile::tempdir().unwrap();
2370        let path = dir.path().join("log.csv");
2371        std::fs::write(&path, "t\n1\n").unwrap();
2372        let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2373            .finish()
2374            .unwrap();
2375        let (_, tail) =
2376            bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
2377        let (tx, rx) = std::sync::mpsc::channel();
2378        let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
2379        let guard = Duration::from_secs(30);
2380        // The watcher may not be waiting yet: append until it reports, each append a
2381        // change it hears once it is.
2382        let mut file = std::fs::OpenOptions::new()
2383            .append(true)
2384            .open(&path)
2385            .unwrap();
2386        let mut rows = 1;
2387        let deadline = Instant::now() + guard;
2388        let news = loop {
2389            assert!(Instant::now() < deadline, "the watcher never heard");
2390            file.write_all(format!("{}\n", rows + 1).as_bytes())
2391                .unwrap();
2392            rows += 1;
2393            if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
2394                break news;
2395            }
2396        };
2397        assert!(matches!(news.change, Change::Grew { rows: 2.., .. }));
2398        drop(follow);
2399    }
2400
2401    /// A file put in place of the one the open counted, before the watcher has opened
2402    /// the path, is read from its start: the watcher opens the new file, and a size
2403    /// as big or bigger would otherwise pass for growth. Replacing it before the
2404    /// follow starts is the watcher's thread starting late.
2405    #[test]
2406    fn a_file_replaced_before_the_watcher_opens_it_is_read_again() {
2407        let dir = tempfile::tempdir().unwrap();
2408        let path = dir.path().join("rotated.csv");
2409        std::fs::write(&path, "t\n1\n2\n").unwrap();
2410        let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2411            .finish()
2412            .unwrap();
2413        let (_, tail) =
2414            bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
2415        let next = dir.path().join("next.csv");
2416        std::fs::write(&next, "t\n7\n8\n9\n").unwrap();
2417        std::fs::rename(&next, &path).unwrap();
2418        let (tx, rx) = std::sync::mpsc::channel();
2419        let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
2420        let deadline = Instant::now() + Duration::from_secs(30);
2421        let news = loop {
2422            assert!(Instant::now() < deadline, "the watcher never reported");
2423            follow.check_now();
2424            if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
2425                break news;
2426            }
2427        };
2428        assert!(
2429            matches!(news.change, Change::Restarted { rows: 3, .. }),
2430            "read on, not again"
2431        );
2432    }
2433
2434    /// The bytes that arrive later are read from where the count stopped, a partial
2435    /// line among them once it completes; a row that does not fit is counted.
2436    #[test]
2437    fn a_tail_reads_on_from_where_it_stopped() {
2438        let dir = tempfile::tempdir().unwrap();
2439        let path = dir.path().join("grow.csv");
2440        let mut out = File::create(&path).unwrap();
2441        out.write_all(b"t,n\n1.5,2\n2.5,").unwrap();
2442        let schema = Schema::from_iter([
2443            Field::new("t".into(), DataType::Float64),
2444            Field::new("n".into(), DataType::Int64),
2445        ]);
2446        let mut tail = Tail::new(FileFormat::Csv, &OpenOptions::default(), &schema);
2447        let mut file = File::open(&path).unwrap();
2448        let len = |p: &Path| std::fs::metadata(p).unwrap().len();
2449        tail.read_on(&mut file, len(&path), true).unwrap();
2450        assert_eq!((tail.rows(), tail.complete()), (1, 10));
2451        out.write_all(b"3\nx,4\n4.5,5,6\n").unwrap();
2452        tail.read_on(&mut file, len(&path), true).unwrap();
2453        assert_eq!(tail.rows(), 4);
2454        assert_eq!(
2455            tail.misfits(),
2456            2,
2457            "a word for a number, and a field too many"
2458        );
2459    }
2460
2461    /// A scan bounded to its complete rows reads more once the bound moves, through
2462    /// the filters built on it, and only as many as the bound says.
2463    #[test]
2464    fn moving_the_bound_reads_the_new_rows_through_the_view() {
2465        let dir = tempfile::tempdir().unwrap();
2466        let path = dir.path().join("grow.csv");
2467        std::fs::write(&path, "a,b\n1,x\n2,y\n3,").unwrap();
2468        let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2469            .with_truncate_ragged_lines(true)
2470            .with_ignore_errors(true)
2471            .finish()
2472            .unwrap();
2473        let mut root = scan.clone();
2474        bound(&mut root, &path, 2);
2475        let mut view = root.clone().filter(col("a").gt(lit(1)));
2476        view.collect_schema().unwrap();
2477        assert_eq!(view.clone().collect().unwrap().height(), 1);
2478        let mut out = std::fs::OpenOptions::new()
2479            .append(true)
2480            .open(&path)
2481            .unwrap();
2482        out.write_all(b"z\n4,w\n5").unwrap();
2483        bound(&mut view, &path, 4);
2484        let df = view.clone().collect().unwrap();
2485        assert_eq!(df.height(), 3, "{df}");
2486        bound(&mut root, &path, 4);
2487        assert_eq!(root.collect().unwrap().height(), 4, "the partial row waits");
2488    }
2489
2490    /// `text` written to a file and scanned as `scan` does, bounded to its complete
2491    /// rows, with a mark every `every` rows.
2492    fn marked(
2493        text: &[u8],
2494        format: FileFormat,
2495        options: &OpenOptions,
2496        every: u64,
2497        scan: impl Fn(&Path) -> LazyFrame,
2498    ) -> (tempfile::TempDir, PathBuf, LazyFrame, Arc<Marks>, usize) {
2499        let dir = tempfile::tempdir().unwrap();
2500        let path = dir.path().join("marked");
2501        std::fs::write(&path, text).unwrap();
2502        let mut lf = scan(&path);
2503        let schema = lf.collect_schema().unwrap();
2504        let mut tail = Tail::new(format, options, &schema);
2505        tail.mark_every = (every, u64::MAX);
2506        tail.read_on(&mut File::open(&path).unwrap(), text.len() as u64, false)
2507            .unwrap();
2508        let marks = Arc::new(Marks::default());
2509        marks.take_from(&mut tail);
2510        bound(&mut lf, &path, tail.rows());
2511        (dir, path, lf, marks, tail.rows())
2512    }
2513
2514    fn csv_scan(path: &Path, options: &OpenOptions) -> LazyFrame {
2515        let mut reader = LazyCsvReader::new(PlRefPath::try_from_path(path).unwrap())
2516            .with_ignore_errors(true)
2517            .with_truncate_ragged_lines(true)
2518            .with_has_header(options.has_header != Some(false))
2519            .with_comment_prefix(options.comment_char.as_deref().map(PlSmallStr::from_str));
2520        if let Some(skip) = options.skip_lines {
2521            reader = reader.with_skip_lines(skip);
2522        }
2523        reader.finish().unwrap()
2524    }
2525
2526    /// Every window read from the marks holds the rows a read from the start of the
2527    /// file gives: quoted newlines, blank lines, comments, skipped lines, carriage
2528    /// returns and NDJSON, at every offset.
2529    #[test]
2530    fn a_window_from_a_mark_reads_what_a_read_from_the_start_does() {
2531        let mut csv = b"skipped\nt,s,n\n".to_vec();
2532        let mut crlf = b"t,s,n\r\n".to_vec();
2533        let mut lines = Vec::new();
2534        for i in 0..120 {
2535            let row = match i % 9 {
2536                0 => format!("{i},\"two\nlines\",{}\n", i * 2),
2537                3 => "\n".to_string(),
2538                5 => "# a comment\n".to_string(),
2539                7 => format!("{i},x,oops\n"),
2540                _ => format!("{i},s{i},{}\n", i * 2),
2541            };
2542            csv.extend(row.as_bytes());
2543            crlf.extend(format!("{i},s{i},{}\r\n", i * 2).as_bytes());
2544            lines.extend(format!("{{\"t\":{i},\"s\":\"s{i}\"}}\n").as_bytes());
2545            if i % 4 == 1 {
2546                lines.extend(b"\n  \n");
2547            }
2548        }
2549        csv.extend(b"999,partial");
2550        let commented = OpenOptions {
2551            comment_char: Some("#".to_string()),
2552            ..OpenOptions::default().with_skip_lines(1)
2553        };
2554        // An Arrow stream of batches of three rows, the last message cut short.
2555        let (schema, batches) = stream_messages(&arrow_rows(0, 100), 3);
2556        let mut arrows = schema;
2557        batches.iter().for_each(|batch| arrows.extend(batch));
2558        arrows.extend(&batches[0][..20]);
2559        let cases: Vec<(&[u8], FileFormat, OpenOptions)> = vec![
2560            (&csv, FileFormat::Csv, commented),
2561            (&crlf, FileFormat::Csv, OpenOptions::default()),
2562            (&lines, FileFormat::Jsonl, OpenOptions::default()),
2563            (&arrows, FileFormat::Arrow, OpenOptions::default()),
2564        ];
2565        for (text, format, options) in cases {
2566            let scan = |path: &Path| match format {
2567                FileFormat::Jsonl => scan_lines(path, &options, false, &mut Vec::new()).unwrap(),
2568                FileFormat::Arrow => stream::scan(path).unwrap(),
2569                _ => csv_scan(path, &options),
2570            };
2571            let (_dir, path, lf, marks, rows) = marked(text, format, &options, 7, scan);
2572            let window = Window {
2573                lf: lf.clone(),
2574                path: path.clone(),
2575                marks: marks.clone(),
2576                known: None,
2577            };
2578            let whole = lf.clone().collect().unwrap();
2579            assert_eq!(whole.height(), rows);
2580            for start in (0..rows + 3).step_by(5) {
2581                for len in [1, 6, 40] {
2582                    let read = crate::formats::pushdown::Windowed::window(&window, start, len)
2583                        .unwrap()
2584                        .collect()
2585                        .unwrap();
2586                    let expected = whole.slice(start as i64, len);
2587                    assert!(
2588                        read.equals_missing(&expected),
2589                        "{format:?} rows {start}+{len}:\n{read:?}\n{expected:?}"
2590                    );
2591                }
2592                if start < rows {
2593                    assert!(
2594                        from_marks(&lf, &path, &marks, start, Some(start + 1)).is_some(),
2595                        "{format:?} row {start} is read from a mark"
2596                    );
2597                }
2598            }
2599        }
2600    }
2601
2602    /// `n` rows from `from`: a number, its text, and a float.
2603    fn arrow_rows(from: i64, n: i64) -> DataFrame {
2604        df!(
2605            "t" => (from..from + n).collect::<Vec<_>>(),
2606            "s" => (from..from + n).map(|i| format!("s{i}")).collect::<Vec<_>>(),
2607            "x" => (from..from + n).map(|i| i as f64 / 2.0).collect::<Vec<_>>(),
2608        )
2609        .unwrap()
2610    }
2611
2612    /// An Arrow stream's batches are counted as their messages complete, a batch cut
2613    /// short waiting for the rest; the stream's own scan reads them, filtered and
2614    /// projected a batch at a time; and a stream with dictionaries is refused.
2615    #[test]
2616    fn an_arrow_stream_is_counted_and_read_by_its_batches() {
2617        let dir = tempfile::tempdir().unwrap();
2618        let path = dir.path().join("live.arrows");
2619        let (schema, batches) = stream_messages(&arrow_rows(0, 50), 4);
2620        let mut head = schema.clone();
2621        head.extend(&batches[0]);
2622        head.extend(&batches[1][..batches[1].len() - 3]);
2623        std::fs::write(&path, &head).unwrap();
2624        let lf = stream::scan(&path).unwrap();
2625        let (mut lf, mut tail) =
2626            bound_to_complete(lf, &path, FileFormat::Arrow, &OpenOptions::default()).unwrap();
2627        assert_eq!(tail.rows(), 4, "the second batch is not all there");
2628        assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 4));
2629
2630        let mut rest = batches[1][batches[1].len() - 3..].to_vec();
2631        batches[2..].iter().for_each(|batch| rest.extend(batch));
2632        let mut file = std::fs::OpenOptions::new()
2633            .append(true)
2634            .open(&path)
2635            .unwrap();
2636        file.write_all(&rest).unwrap();
2637        let size = file.metadata().unwrap().len();
2638        tail.read_on(&mut File::open(&path).unwrap(), size, true)
2639            .unwrap();
2640        assert_eq!((tail.rows(), tail.misfits()), (50, 0));
2641        bound(&mut lf, &path, tail.rows());
2642        assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 50));
2643        let kept = lf
2644            .clone()
2645            .filter(col("t").gt_eq(lit(45)))
2646            .select([col("s")])
2647            .collect()
2648            .unwrap();
2649        assert_eq!(kept, arrow_rows(45, 5).select(["s"]).unwrap());
2650        let count = lf.select([len()]).collect().unwrap();
2651        assert_eq!(count.column("len").unwrap().u32().unwrap().get(0), Some(50));
2652
2653        // Written before Arrow 0.15, with no continuation markers, and ended.
2654        let legacy = dir.path().join("legacy.arrows");
2655        std::fs::write(
2656            &legacy,
2657            crate::formats::ipc_stream::tests::stream(&arrow_rows(0, 10), None, true),
2658        )
2659        .unwrap();
2660        let (lf, tail) = bound_to_complete(
2661            stream::scan(&legacy).unwrap(),
2662            &legacy,
2663            FileFormat::Arrow,
2664            &OpenOptions::default(),
2665        )
2666        .unwrap();
2667        assert_eq!(tail.rows(), 10);
2668        assert_eq!(lf.collect().unwrap(), arrow_rows(0, 10));
2669
2670        // Dictionary-encoded columns.
2671        let cats = df!("c" => ["a", "b", "a"])
2672            .unwrap()
2673            .lazy()
2674            .with_column(col("c").cast(DataType::from_categories(Categories::global())))
2675            .collect()
2676            .unwrap();
2677        let (schema, batches) = stream_messages(&cats, 3);
2678        let dictionary = dir.path().join("dict.arrows");
2679        std::fs::write(&dictionary, [schema, batches.concat()].concat()).unwrap();
2680        assert!(
2681            stream::scan(&dictionary).is_err_and(|e| e.contains("dictionary")),
2682            "refused"
2683        );
2684    }
2685
2686    /// A filtered view reads on from a point where its rows are known, and counts the
2687    /// rows after it alone.
2688    #[test]
2689    fn a_filtered_view_reads_and_counts_on_from_what_is_known() {
2690        let mut text = b"t,n\n".to_vec();
2691        for i in 0..300 {
2692            text.extend(format!("{i},{}\n", i % 5).as_bytes());
2693        }
2694        let options = OpenOptions::default();
2695        let (_dir, path, lf, marks, rows) = marked(&text, FileFormat::Csv, &options, 16, |p| {
2696            csv_scan(p, &options)
2697        });
2698        let view = lf.filter(col("n").eq(lit(3)));
2699        let whole = view.clone().collect().unwrap();
2700        // The view's rows among the first 200 of the file.
2701        let known = (whole.column("t").unwrap().i64().unwrap().to_vec())
2702            .into_iter()
2703            .filter(|t| t.unwrap() < 200)
2704            .count();
2705        let rest = from_marks(&view, &path, &marks, 200, None).unwrap();
2706        let after = rest.collect().unwrap().height();
2707        assert_eq!(known + after, whole.height());
2708        assert_eq!(rows, 300);
2709        let window = Window {
2710            lf: view.clone(),
2711            path,
2712            marks,
2713            known: Some(vec![(0, 0), (known, 200)]),
2714        };
2715        for start in [0, 10, known - 1, known, known + 5, whole.height() - 3] {
2716            let read = crate::formats::pushdown::Windowed::window(&window, start, 4)
2717                .unwrap()
2718                .collect()
2719                .unwrap();
2720            assert!(
2721                read.equals_missing(&whole.slice(start as i64, 4)),
2722                "{start}"
2723            );
2724        }
2725    }
2726
2727    /// A file put in place of the followed one, as big or bigger, is another file;
2728    /// the file grown in place is the same one. Windows reads this from the volume and
2729    /// file index, Unix from the device and inode.
2730    #[test]
2731    fn a_replaced_file_is_told_from_a_grown_one() {
2732        let dir = tempfile::tempdir().unwrap();
2733        let path = dir.path().join("log.csv");
2734        std::fs::write(&path, "t\n1\n").unwrap();
2735        let held = File::open(&path).unwrap();
2736        let known = identity_of(&held);
2737        assert!(cfg!(not(any(unix, windows))) || known.is_some());
2738        let now = |path: &Path| identity_at(path, &std::fs::metadata(path).unwrap());
2739        std::fs::OpenOptions::new()
2740            .append(true)
2741            .open(&path)
2742            .unwrap()
2743            .write_all(b"2\n")
2744            .unwrap();
2745        assert!(!replaced(known, now(&path)), "grown in place");
2746        let other = dir.path().join("next.csv");
2747        std::fs::write(&other, "t\n1\n2\n3\n").unwrap();
2748        std::fs::rename(&other, &path).unwrap();
2749        assert_eq!(
2750            replaced(known, now(&path)),
2751            known.is_some(),
2752            "renamed over it"
2753        );
2754        assert!(!replaced(None, now(&path)), "unknown is no replacement");
2755        drop(held);
2756    }
2757
2758    /// A deleted file is read through the handle held on it.
2759    #[cfg(unix)]
2760    #[test]
2761    fn a_deleted_file_reads_through_its_handle() {
2762        let dir = tempfile::tempdir().unwrap();
2763        let path = dir.path().join("gone.csv");
2764        std::fs::write(&path, "a\n1\n2\n").unwrap();
2765        let mut lf = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2766            .finish()
2767            .unwrap();
2768        bound(&mut lf, &path, 2);
2769        let mut view = lf.filter(col("a").gt(lit(0)));
2770        view.collect_schema().unwrap();
2771        let handle = File::open(&path).unwrap();
2772        std::fs::remove_file(&path).unwrap();
2773        read_through(&mut view, &path, &handle);
2774        assert_eq!(view.collect().unwrap().height(), 2);
2775    }
2776}