1use 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
35pub const DEFAULT_INTERVAL: Duration = Duration::from_millis(250);
39
40const CHUNK: usize = 1 << 20;
42
43const LONGEST_RECORD: usize = 16 << 20;
45
46pub(crate) const MARK_ROWS: u64 = 8192;
49const MARK_BYTES: u64 = 1 << 20;
50
51const MOST_FROM_A_MARK: u64 = 64 << 20;
54
55pub 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
96pub(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
110pub 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 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
146pub(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
157pub(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 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
183pub(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#[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 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 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 Fits::Text => value.is_string() || value.is_array(),
257 Fits::Anything => true,
258 }
259 }
260}
261
262#[derive(Clone, Debug)]
264enum Layout {
265 Delimited {
266 separator: u8,
267 skip: u64,
270 header: bool,
271 comment: Option<Vec<u8>>,
272 nulls: Vec<String>,
273 },
274 Lines,
275 Stream,
277 Text,
279}
280
281#[derive(Clone, Debug, Default)]
284struct NewMarks {
285 new: Vec<(u64, u64)>,
286 last: Option<(u64, u64)>,
287}
288
289#[derive(Clone, Debug)]
292pub struct Tail {
293 path: PathBuf,
295 identity: Option<Identity>,
298 layout: Layout,
299 columns: Vec<(String, Fits)>,
301 complete: u64,
303 records: u64,
305 rows: u64,
306 misfits: u64,
307 fields: Option<usize>,
309 marks: NewMarks,
311 mark_every: (u64, u64),
313 widens: bool,
316 arrived: Vec<(String, Arrived)>,
319}
320
321const MOST_NEW_FIELDS: usize = 4096;
324
325#[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 _ => 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 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 pub fn widening(mut self) -> Tail {
420 self.widens = matches!(self.layout, Layout::Lines);
421 self
422 }
423
424 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 pub fn path(&self) -> &Path {
434 &self.path
435 }
436
437 pub fn rows(&self) -> usize {
439 self.rows as usize
440 }
441
442 pub fn misfits(&self) -> usize {
444 self.misfits as usize
445 }
446
447 pub fn complete(&self) -> u64 {
449 self.complete
450 }
451
452 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 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 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 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 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 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 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 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 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 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
686fn 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
695fn 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
720fn 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
732fn 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
749pub 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 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
794pub 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 *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#[derive(Default)]
847pub struct Marks {
848 inner: Mutex<MarksInner>,
849}
850
851#[derive(Default)]
852struct MarksInner {
853 at: Vec<(u64, u64)>,
855 complete: u64,
856}
857
858#[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 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 fn clear(&self) {
881 let mut inner = self.lock();
882 inner.at.clear();
883 inner.complete = 0;
884 }
885
886 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#[derive(Clone)]
902enum Parse {
903 Csv(Box<CsvReadOptions>),
904 Lines { ignore_errors: bool },
905 Stream(Arc<stream::StreamSchema>),
906}
907
908fn 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 fn of(scan: &polars::lazy::dsl::DslPlan, schema: &SchemaRef) -> Option<Parse> {
920 use polars::lazy::dsl::{DslPlan, FileScanDsl};
921 if let Some(stream) = stream::StreamScan::in_plan(scan_node(scan)) {
922 return Some(Parse::Stream(stream.schema().clone()));
923 }
924 if let Some(lines) = lines::LinesScan::in_plan(scan_node(scan)) {
925 return Some(Parse::Lines {
926 ignore_errors: lines.ignore_errors(),
927 });
928 }
929 let DslPlan::Scan {
930 scan_type,
931 unified_scan_args,
932 ..
933 } = scan_node(scan)
934 else {
935 return None;
936 };
937 if unified_scan_args.row_index.is_some() || unified_scan_args.include_file_paths.is_some() {
938 return None;
939 }
940 match &**scan_type {
941 FileScanDsl::Csv { options } => {
942 if options.columns.is_some()
943 || options.projection.is_some()
944 || options.row_index.is_some()
945 {
946 return None;
947 }
948 let mut options = (**options).clone();
949 options.path = None;
950 options.has_header = false;
951 options.skip_rows = 0;
952 options.skip_lines = 0;
953 options.skip_rows_after_header = 0;
954 options.n_rows = None;
955 options.schema = Some(schema.clone());
957 options.schema_overwrite = None;
958 options.dtype_overwrite = None;
959 options.column_names_overwrite = None;
960 options.raise_if_empty = false;
961 Some(Parse::Csv(Box::new(options)))
962 }
963 FileScanDsl::NDJson { options } => Some(Parse::Lines {
964 ignore_errors: options.ignore_errors,
965 }),
966 _ => None,
967 }
968 }
969}
970
971struct Piece {
974 path: PathBuf,
975 span: Span,
976 skip: usize,
977 take: usize,
978 parse: Parse,
979 schema: SchemaRef,
980}
981
982const PIECE_NAME: &str = "FOLLOWED";
984
985impl polars::prelude::AnonymousScan for Piece {
986 fn as_any(&self) -> &dyn std::any::Any {
987 self
988 }
989
990 fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
991 Ok(self.schema.clone())
992 }
993
994 fn scan(&self, args: polars::prelude::AnonymousScanArgs) -> PolarsResult<DataFrame> {
995 let take = args.n_rows.map_or(self.take, |n| n.min(self.take));
996 let mut file = File::open(&self.path)?;
997 file.seek(SeekFrom::Start(self.span.start))?;
998 let mut bytes = Vec::with_capacity((self.span.end - self.span.start) as usize);
999 file.take(self.span.end - self.span.start)
1000 .read_to_end(&mut bytes)?;
1001 let df = match &self.parse {
1002 Parse::Csv(options) => {
1003 let mut options = (**options).clone();
1004 options.n_rows = Some(self.skip + take);
1005 options
1006 .into_reader_with_file_handle(std::io::Cursor::new(bytes))
1007 .finish()?
1008 }
1009 Parse::Lines { ignore_errors } => {
1010 lines::parse_run(&bytes, &self.schema, *ignore_errors)?
1011 }
1012 Parse::Stream(schema) => stream::decode_run(bytes, schema, self.skip + take)?,
1013 };
1014 Ok(df.slice(self.skip as i64, take))
1015 }
1016}
1017
1018pub(crate) fn from_marks(
1022 lf: &LazyFrame,
1023 path: &Path,
1024 marks: &Marks,
1025 from: usize,
1026 to: Option<usize>,
1027) -> Option<LazyFrame> {
1028 let path_text = path.to_string_lossy();
1029 let mut plan = lf.logical_plan.clone();
1030 let mut replaced = false;
1031 let mut failed = false;
1032 let piece = |scan: &polars::lazy::dsl::DslPlan, bound: usize| {
1033 let to = to.map_or(bound, |to| to.min(bound));
1034 let from = from.min(to);
1035 let span = marks.span(from as u64, to as u64)?;
1036 let schema = LazyFrame::from(scan.clone()).collect_schema().ok()?;
1037 let parse = Parse::of(scan, &schema)?;
1038 let piece = Piece {
1039 path: path.to_path_buf(),
1040 span,
1041 skip: from - span.row as usize,
1042 take: to - from,
1043 parse,
1044 schema: schema.clone(),
1045 };
1046 LazyFrame::anonymous_scan(
1047 Arc::new(piece),
1048 ScanArgsAnonymous {
1049 schema: Some(schema),
1050 name: PIECE_NAME,
1051 ..Default::default()
1052 },
1053 )
1054 .ok()
1055 .map(|lf| lf.logical_plan)
1056 };
1057 replace_bound(
1058 &mut plan,
1059 &path_text,
1060 &mut |scan, bound| match piece(scan, bound) {
1061 Some(plan) => {
1062 replaced = true;
1063 Some(plan)
1064 }
1065 None => {
1066 failed = true;
1067 None
1068 }
1069 },
1070 );
1071 (replaced && !failed).then(|| {
1072 let mut out = lf.clone();
1073 out.logical_plan = plan;
1074 out
1075 })
1076}
1077
1078fn replace_bound(
1080 plan: &mut polars::lazy::dsl::DslPlan,
1081 path: &str,
1082 with: &mut dyn FnMut(&polars::lazy::dsl::DslPlan, usize) -> Option<polars::lazy::dsl::DslPlan>,
1083) {
1084 use polars::lazy::dsl::DslPlan;
1085 match plan {
1086 DslPlan::IR { dsl, .. } => {
1087 let mut inner = Arc::unwrap_or_clone(dsl.clone());
1088 replace_bound(&mut inner, path, with);
1089 *plan = inner;
1090 return;
1091 }
1092 DslPlan::Slice {
1093 input,
1094 offset: 0,
1095 len,
1096 } if scans(input, path) => {
1097 if let Some(piece) = with(input, *len as usize) {
1098 *plan = piece;
1099 }
1100 return;
1101 }
1102 _ => {}
1103 }
1104 crate::table::for_each_input(plan, &mut |input| replace_bound(input, path, with));
1105}
1106
1107pub(crate) fn bound_of(lf: &LazyFrame, path: &Path) -> Option<usize> {
1109 use polars::lazy::dsl::DslPlan;
1110 let path = path.to_string_lossy();
1111 (&lf.logical_plan).into_iter().find_map(|node| match node {
1112 DslPlan::Slice {
1113 input,
1114 offset: 0,
1115 len,
1116 } if scans(input, &path) => Some(*len as usize),
1117 _ => None,
1118 })
1119}
1120
1121pub(crate) fn widen(
1124 root: &LazyFrame,
1125 path: &Path,
1126 format: FileFormat,
1127 fields: &[Field],
1128 rows: usize,
1129) -> Option<LazyFrame> {
1130 let path_text = path.to_string_lossy();
1131 let mut plan = root.logical_plan.clone();
1132 let mut widened = None;
1133 replace_lines_scan(&mut plan, &path_text, &mut |scan| {
1134 let mut schema = (**scan.schema()).clone();
1135 for field in fields {
1136 if !schema.contains(field.name()) {
1137 schema.with_column(field.name().clone(), field.dtype().clone());
1138 }
1139 }
1140 if schema.len() == scan.schema().len() {
1141 return None;
1142 }
1143 let schema = Arc::new(schema);
1144 let lf = scan.with_schema(schema.clone()).lazy().ok()?;
1145 widened = Some((lf.clone(), schema));
1146 Some(lf.logical_plan)
1147 });
1148 let (raw, schema) = widened?;
1149 if format == FileFormat::Journal {
1150 let mut raw = raw;
1152 bound(&mut raw, path, rows);
1153 return Some(crate::formats::journal::derive(raw, &schema).0);
1154 }
1155 let mut out = root.clone();
1156 out.logical_plan = plan;
1157 Some(out)
1158}
1159
1160fn replace_lines_scan(
1162 plan: &mut polars::lazy::dsl::DslPlan,
1163 path: &str,
1164 with: &mut dyn FnMut(&lines::LinesScan) -> Option<polars::lazy::dsl::DslPlan>,
1165) {
1166 use polars::lazy::dsl::DslPlan;
1167 if let DslPlan::IR { dsl, .. } = plan {
1168 let mut inner = Arc::unwrap_or_clone(dsl.clone());
1169 replace_lines_scan(&mut inner, path, with);
1170 *plan = inner;
1171 return;
1172 }
1173 if let Some(scan) = lines::LinesScan::of(plan, path) {
1174 if let Some(replaced) = with(scan) {
1175 *plan = replaced;
1176 }
1177 return;
1178 }
1179 crate::table::for_each_input(plan, &mut |input| replace_lines_scan(input, path, with));
1180}
1181
1182pub(crate) struct Window {
1186 pub(crate) lf: LazyFrame,
1187 pub(crate) path: PathBuf,
1188 pub(crate) marks: Arc<Marks>,
1189 pub(crate) known: Option<Vec<(usize, usize)>>,
1190}
1191
1192impl crate::formats::pushdown::Windowed for Window {
1193 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
1194 let read = match &self.known {
1195 None => from_marks(&self.lf, &self.path, &self.marks, start, Some(start + len)),
1196 Some(known) => {
1197 let at = known.partition_point(|&(view, _)| view <= start);
1198 known.get(at.wrapping_sub(1)).and_then(|&(view, row)| {
1199 from_marks(&self.lf, &self.path, &self.marks, row, None)
1200 .map(|lf| lf.slice((start - view) as i64, len as IdxSize))
1201 })
1202 }
1203 };
1204 Ok(read.unwrap_or_else(|| self.lf.clone().slice(start as i64, len as IdxSize)))
1206 }
1207}
1208
1209#[derive(Clone)]
1211pub enum Change {
1212 Grew { rows: usize, misfits: usize },
1215 Restarted { rows: usize, misfits: usize },
1218 Gone { handle: Option<Arc<File>> },
1220 NewFields(Vec<Field>),
1223 Ended(Option<String>),
1225 Failed(String),
1227}
1228
1229#[derive(Clone)]
1231pub struct News {
1232 pub(crate) id: u64,
1233 pub(crate) change: Change,
1234}
1235
1236#[derive(Default)]
1238struct Shared {
1239 stop: AtomicBool,
1240 poke: Mutex<bool>,
1241 woken: Condvar,
1242 #[cfg(target_os = "linux")]
1244 bell: notify::Bell,
1245}
1246
1247impl Shared {
1248 fn wait(&self, interval: Duration) -> bool {
1250 let mut poked = self.poke.lock().unwrap_or_else(|e| e.into_inner());
1251 if !*poked && !self.stop.load(Ordering::Relaxed) {
1252 poked = self
1253 .woken
1254 .wait_timeout(poked, interval)
1255 .unwrap_or_else(|e| e.into_inner())
1256 .0;
1257 }
1258 *poked = false;
1259 !self.stop.load(Ordering::Relaxed)
1260 }
1261
1262 fn wake(&self) {
1263 *self.poke.lock().unwrap_or_else(|e| e.into_inner()) = true;
1264 self.woken.notify_all();
1265 #[cfg(target_os = "linux")]
1266 self.bell.ring();
1267 }
1268
1269 #[cfg(target_os = "linux")]
1272 fn wait_for_change(
1273 &self,
1274 notify: ¬ify::Notify,
1275 interval: Duration,
1276 last: &mut Option<Instant>,
1277 ) -> bool {
1278 loop {
1279 if self.stop.load(Ordering::Relaxed) {
1280 return false;
1281 }
1282 if std::mem::take(&mut *self.poke.lock().unwrap_or_else(|e| e.into_inner())) {
1283 break;
1284 }
1285 if notify.wait(&self.bell, None) == notify::Woke::Changed {
1286 let left = last
1287 .map(|at| at + interval)
1288 .and_then(|due| due.checked_duration_since(Instant::now()));
1289 if let Some(left) = left
1291 && !self.wait(left)
1292 {
1293 return false;
1294 }
1295 break;
1296 }
1297 }
1298 notify.drain();
1300 *last = Some(Instant::now());
1301 !self.stop.load(Ordering::Relaxed)
1302 }
1303}
1304
1305pub fn age(elapsed: Duration) -> String {
1308 let secs = elapsed.as_secs();
1309 match secs {
1310 0..60 => format!("{secs:>2}s ago"),
1311 60..3_600 => format!("{:>2}m ago", secs / 60),
1312 3_600..86_400 => format!("{:>2}h ago", secs / 3_600),
1313 _ => format!("{:>2}d ago", secs / 86_400),
1314 }
1315}
1316
1317pub fn next_tick(at: Instant) -> Instant {
1319 let secs = at.elapsed().as_secs();
1320 let unit = match secs {
1321 0..60 => 1,
1322 60..3_600 => 60,
1323 3_600..86_400 => 3_600,
1324 _ => 86_400,
1325 };
1326 at + Duration::from_secs((secs / unit + 1) * unit)
1327}
1328
1329#[derive(Clone, Debug, PartialEq)]
1331pub enum Standing {
1332 Following,
1333 Paused,
1334 Ended,
1336}
1337
1338pub struct Follow {
1341 id: u64,
1342 path: PathBuf,
1344 shared: Arc<Shared>,
1345 spool: Option<Arc<SpoolHandle>>,
1346 counted: usize,
1348 misfits: usize,
1349 restarted: bool,
1351 shown: usize,
1353 pub(crate) new_below: usize,
1355 pub(crate) standing: Standing,
1356 pub(crate) last_append: Option<Instant>,
1357 pub(crate) settle_at_end: bool,
1360 pub(crate) end_pending: bool,
1363 pub(crate) stale_view: bool,
1366 held: Option<Arc<File>>,
1368 marks: Arc<Marks>,
1370 pipe: bool,
1373 new_fields: Vec<Field>,
1376 pub(crate) described: bool,
1379}
1380
1381static NEXT_ID: AtomicU64 = AtomicU64::new(1);
1382
1383impl Follow {
1384 pub fn start(
1387 mut tail: Tail,
1388 interval: Duration,
1389 events: Sender<AppEvent>,
1390 spool: Option<Arc<SpoolHandle>>,
1391 ) -> Follow {
1392 let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
1393 let path = tail.path.clone();
1394 let shared = Arc::new(Shared::default());
1395 let shown = tail.rows();
1396 if let Some(handle) = &spool {
1397 handle.spool.wake_on_end(shared.clone());
1398 }
1399 let marks = Arc::new(Marks::default());
1400 marks.take_from(&mut tail);
1401 let watcher = Watcher {
1402 marks: marks.clone(),
1403 id,
1404 path: path.clone(),
1405 tail,
1406 shared: shared.clone(),
1407 events,
1408 spool: spool.as_ref().map(|handle| handle.spool.clone()),
1409 interval,
1410 };
1411 let _ = std::thread::Builder::new()
1412 .name("datui-follow".to_string())
1413 .spawn(move || watcher.run());
1414 Follow {
1415 id,
1416 path,
1417 shared,
1418 spool,
1419 counted: shown,
1420 misfits: 0,
1421 restarted: false,
1422 shown,
1423 new_below: 0,
1424 standing: Standing::Following,
1425 last_append: None,
1426 settle_at_end: true,
1427 end_pending: false,
1428 stale_view: false,
1429 held: None,
1430 marks,
1431 pipe: false,
1432 new_fields: Vec::new(),
1433 described: false,
1434 }
1435 }
1436
1437 pub fn as_pipe(mut self) -> Follow {
1440 self.pipe = true;
1441 self.settle_at_end = false;
1442 self
1443 }
1444
1445 pub fn is_pipe(&self) -> bool {
1447 self.pipe
1448 }
1449
1450 pub fn live(&self) -> bool {
1452 self.standing != Standing::Ended
1453 }
1454
1455 pub(crate) fn marks(&self) -> &Arc<Marks> {
1457 &self.marks
1458 }
1459
1460 pub fn id(&self) -> u64 {
1461 self.id
1462 }
1463
1464 pub fn path(&self) -> &Path {
1465 &self.path
1466 }
1467
1468 pub fn shown(&self) -> usize {
1470 self.shown
1471 }
1472
1473 pub fn waiting(&self) -> usize {
1475 self.counted.saturating_sub(self.shown)
1476 }
1477
1478 pub fn misfits(&self) -> usize {
1479 self.misfits
1480 }
1481
1482 pub fn standing(&self) -> &Standing {
1483 &self.standing
1484 }
1485
1486 pub fn new_below(&self) -> usize {
1488 self.new_below
1489 }
1490
1491 pub fn behind(&self) -> bool {
1494 self.standing != Standing::Paused && (self.restarted || self.counted != self.shown)
1495 }
1496
1497 pub fn spool(&self) -> Option<&Arc<Spool>> {
1499 self.spool.as_ref().map(|handle| &handle.spool)
1500 }
1501
1502 pub fn check_now(&self) {
1504 self.shared.wake();
1505 }
1506
1507 pub fn take(&mut self, change: &Change) -> Option<String> {
1510 match change {
1511 Change::Grew { rows, misfits } => {
1512 if *rows > self.counted {
1513 self.last_append = Some(Instant::now());
1514 }
1515 self.counted = *rows;
1516 self.misfits = *misfits;
1517 None
1518 }
1519 Change::Restarted { rows, misfits } => {
1520 self.counted = *rows;
1521 self.misfits = *misfits;
1522 self.restarted = true;
1523 self.last_append = Some(Instant::now());
1524 Some("The file was truncated or replaced: reading it from the start".to_string())
1525 }
1526 Change::NewFields(fields) => {
1527 self.new_fields = fields.clone();
1528 None
1529 }
1530 Change::Gone { handle } => {
1531 self.held = handle.clone();
1532 self.end();
1533 Some("The file was deleted: following stopped, the rows read stay".to_string())
1534 }
1535 Change::Ended(_) if self.spool().is_some_and(|s| s.tee().is_some()) => {
1537 self.end();
1538 None
1539 }
1540 Change::Ended(None) => {
1541 self.end();
1542 Some("Standard input ended".to_string())
1543 }
1544 Change::Ended(Some(reason)) | Change::Failed(reason) => {
1545 self.end();
1546 Some(reason.clone())
1547 }
1548 }
1549 }
1550
1551 pub fn catch_up(&mut self) -> (usize, bool) {
1554 self.shown = self.counted;
1555 (self.shown, std::mem::take(&mut self.restarted))
1556 }
1557
1558 pub(crate) fn take_new_fields(&mut self) -> Vec<Field> {
1560 std::mem::take(&mut self.new_fields)
1561 }
1562
1563 pub fn take_held(&mut self) -> Option<Arc<File>> {
1565 self.held.take()
1566 }
1567
1568 pub fn pause(&mut self) {
1569 if self.standing == Standing::Following {
1570 self.standing = Standing::Paused;
1571 }
1572 }
1573
1574 pub fn resume(&mut self) {
1575 if self.standing == Standing::Paused {
1576 self.standing = Standing::Following;
1577 }
1578 }
1579
1580 pub fn end(&mut self) {
1583 self.standing = Standing::Ended;
1584 self.shared.stop.store(true, Ordering::Relaxed);
1585 self.shared.wake();
1586 if let Some(spool) = self.spool.as_ref().filter(|s| s.spool.tee().is_none()) {
1587 spool.spool.stop();
1588 }
1589 }
1590}
1591
1592impl Drop for Follow {
1593 fn drop(&mut self) {
1594 self.shared.stop.store(true, Ordering::Relaxed);
1595 self.shared.wake();
1596 }
1597}
1598
1599struct Watcher {
1601 id: u64,
1602 path: PathBuf,
1603 tail: Tail,
1604 shared: Arc<Shared>,
1605 events: Sender<AppEvent>,
1606 spool: Option<Arc<Spool>>,
1607 interval: Duration,
1608 marks: Arc<Marks>,
1609}
1610
1611type Identity = (u64, u64);
1614
1615#[cfg(unix)]
1616fn identity_of(file: &File) -> Option<Identity> {
1617 use std::os::unix::fs::MetadataExt;
1618 file.metadata().ok().map(|meta| (meta.dev(), meta.ino()))
1619}
1620
1621#[cfg(windows)]
1622fn identity_of(file: &File) -> Option<Identity> {
1623 use std::os::windows::io::AsRawHandle;
1624 use windows_sys::Win32::Storage::FileSystem::{
1625 BY_HANDLE_FILE_INFORMATION, GetFileInformationByHandle,
1626 };
1627 let mut info: BY_HANDLE_FILE_INFORMATION = unsafe { std::mem::zeroed() };
1630 let ok = unsafe { GetFileInformationByHandle(file.as_raw_handle(), &mut info) };
1631 (ok != 0).then(|| {
1632 (
1633 u64::from(info.dwVolumeSerialNumber),
1634 u64::from(info.nFileIndexHigh) << 32 | u64::from(info.nFileIndexLow),
1635 )
1636 })
1637}
1638
1639#[cfg(not(any(unix, windows)))]
1640fn identity_of(_file: &File) -> Option<Identity> {
1641 None
1642}
1643
1644#[cfg(unix)]
1647fn identity_at(_path: &Path, meta: &std::fs::Metadata) -> Option<Identity> {
1648 use std::os::unix::fs::MetadataExt;
1649 Some((meta.dev(), meta.ino()))
1650}
1651
1652#[cfg(not(unix))]
1653fn identity_at(path: &Path, _meta: &std::fs::Metadata) -> Option<Identity> {
1654 File::open(path).ok().as_ref().and_then(identity_of)
1655}
1656
1657fn replaced(known: Option<Identity>, now: Option<Identity>) -> bool {
1660 matches!((known, now), (Some(known), Some(now)) if known != now)
1661}
1662
1663fn failed_message(path: Option<&Path>, doing: &str, e: &std::io::Error) -> String {
1665 let what = format!(
1666 "{doing}. {}",
1667 crate::error_display::user_message_from_io(e, None)
1668 );
1669 match path {
1670 Some(path) => crate::error_display::file_message(path, &what),
1671 None => crate::error_display::sentence(&format!("standard input: {what}")),
1672 }
1673}
1674
1675impl Watcher {
1676 fn failed(&self, doing: &str, e: &std::io::Error) -> String {
1679 failed_message(
1680 self.spool.is_none().then_some(self.path.as_path()),
1681 doing,
1682 e,
1683 )
1684 }
1685
1686 fn run(mut self) {
1687 let mut file = match File::open(&self.path) {
1688 Ok(file) => file,
1689 Err(e) => {
1690 self.send(Change::Failed(self.failed("following it stopped", &e)));
1691 return;
1692 }
1693 };
1694 let mut known = self.tail.identity.or_else(|| identity_of(&file));
1697 let mut sent = (self.tail.rows(), 0usize);
1698 #[cfg(target_os = "linux")]
1699 let notify = notify::Notify::new(&self.path);
1700 #[cfg(target_os = "linux")]
1701 let mut last = None;
1702 loop {
1703 #[cfg(target_os = "linux")]
1704 let go_on = match ¬ify {
1705 Some(notify) => self
1706 .shared
1707 .wait_for_change(notify, self.interval, &mut last),
1708 None => self.shared.wait(self.interval),
1709 };
1710 #[cfg(not(target_os = "linux"))]
1711 let go_on = self.shared.wait(self.interval);
1712 if !go_on {
1713 return;
1714 }
1715 let spool_ended = self.spool.as_ref().and_then(|spool| spool.ended());
1717 let meta = match std::fs::metadata(&self.path) {
1718 Ok(meta) => meta,
1719 Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
1720 self.send(Change::Gone {
1721 handle: Some(Arc::new(file)),
1722 });
1723 return;
1724 }
1725 Err(e) => {
1726 self.send(Change::Failed(self.failed("following it stopped", &e)));
1727 return;
1728 }
1729 };
1730 let now = identity_at(&self.path, &meta);
1731 let len = meta.len();
1732 if replaced(known, now) || len < self.tail.complete() {
1733 match File::open(&self.path) {
1734 Ok(reopened) => file = reopened,
1735 Err(e) => {
1736 self.send(Change::Failed(self.failed("following it stopped", &e)));
1737 return;
1738 }
1739 }
1740 known = identity_of(&file);
1741 #[cfg(target_os = "linux")]
1742 if let Some(notify) = ¬ify {
1743 notify.rewatch(&self.path);
1744 }
1745 self.tail.restart();
1746 self.marks.clear();
1747 if let Err(e) = self.tail.read_on(&mut file, len, false) {
1748 self.send(Change::Failed(self.failed("reading it stopped", &e)));
1749 return;
1750 }
1751 self.marks.take_from(&mut self.tail);
1752 sent = (self.tail.rows(), self.tail.misfits());
1753 self.send(Change::Restarted {
1754 rows: sent.0,
1755 misfits: sent.1,
1756 });
1757 continue;
1758 }
1759 let read = if spool_ended.is_some() {
1760 self.tail.read_to_end(&mut file, len, true)
1761 } else {
1762 self.tail.read_on(&mut file, len, true)
1763 };
1764 if let Err(e) = read {
1765 self.send(Change::Failed(self.failed("reading it stopped", &e)));
1766 return;
1767 }
1768 self.marks.take_from(&mut self.tail);
1770 let now_counted = (self.tail.rows(), self.tail.misfits());
1771 if now_counted != sent {
1772 sent = now_counted;
1773 if !self.send(Change::Grew {
1774 rows: sent.0,
1775 misfits: sent.1,
1776 }) {
1777 return;
1778 }
1779 }
1780 if let Some(reason) = spool_ended {
1781 let fields = self.tail.new_fields();
1782 if !fields.is_empty() && !self.send(Change::NewFields(fields)) {
1783 return;
1784 }
1785 self.send(Change::Ended(reason));
1786 return;
1787 }
1788 }
1789 }
1790
1791 fn send(&self, change: Change) -> bool {
1793 self.events
1794 .send(AppEvent::Followed(News {
1795 id: self.id,
1796 change,
1797 }))
1798 .is_ok()
1799 }
1800}
1801
1802pub struct Spool {
1806 stop: AtomicBool,
1807 bytes: AtomicU64,
1808 state: Mutex<SpoolState>,
1809 changed: Condvar,
1810 sink: Mutex<Option<File>>,
1813 tee: Option<Tee>,
1815 pass: Mutex<Option<Box<dyn Write + Send>>>,
1817 started: Instant,
1818}
1819
1820#[derive(Clone, Debug)]
1822pub struct Tee {
1823 pub path: PathBuf,
1824 pub raw: bool,
1826}
1827
1828impl Tee {
1829 pub fn to_stdout(&self) -> bool {
1831 crate::loading::stdin::is_stdin(&self.path)
1832 }
1833
1834 pub fn name(&self) -> String {
1836 if self.to_stdout() {
1837 return "standard output".to_string();
1838 }
1839 self.path.file_name().map_or_else(
1840 || self.path.display().to_string(),
1841 |name| name.to_string_lossy().into_owned(),
1842 )
1843 }
1844}
1845
1846#[derive(Default)]
1847struct SpoolState {
1848 lines: usize,
1850 drained: bool,
1852 ended: Option<Option<String>>,
1854 finished: Option<Instant>,
1856 samples: std::collections::VecDeque<(Instant, u64)>,
1858 watcher: Option<Arc<Shared>>,
1861}
1862
1863const RATE_WINDOW: Duration = Duration::from_secs(2);
1865
1866impl Spool {
1867 fn new(sink: File, tee: Option<Tee>, pass: Option<Box<dyn Write + Send>>) -> Spool {
1868 Spool {
1869 stop: AtomicBool::new(false),
1870 bytes: AtomicU64::new(0),
1871 state: Mutex::new(SpoolState::default()),
1872 changed: Condvar::new(),
1873 sink: Mutex::new(Some(sink)),
1874 tee,
1875 pass: Mutex::new(pass),
1876 started: Instant::now(),
1877 }
1878 }
1879
1880 pub fn bytes(&self) -> u64 {
1882 self.bytes.load(Ordering::Relaxed)
1883 }
1884
1885 pub fn rate(&self) -> f64 {
1887 let state = self.lock();
1888 match (state.samples.front(), state.samples.back()) {
1889 (Some((t0, b0)), Some((t1, b1))) if t1 > t0 => {
1890 (b1 - b0) as f64 / t1.duration_since(*t0).as_secs_f64()
1891 }
1892 _ => 0.0,
1893 }
1894 }
1895
1896 pub fn duration(&self) -> Duration {
1898 let finished = self.lock().finished;
1899 finished
1900 .unwrap_or_else(Instant::now)
1901 .duration_since(self.started)
1902 }
1903
1904 pub fn tee(&self) -> Option<&Tee> {
1906 self.tee.as_ref()
1907 }
1908
1909 pub fn stop(&self) {
1912 self.stop.store(true, Ordering::Relaxed);
1913 self.finish(None);
1914 }
1915
1916 pub fn stopped(&self) -> bool {
1917 self.stop.load(Ordering::Relaxed)
1918 }
1919
1920 pub fn ended(&self) -> Option<Option<String>> {
1922 self.lock().ended.clone()
1923 }
1924
1925 pub fn live(&self) -> bool {
1927 self.lock().ended.is_none()
1928 }
1929
1930 pub fn wait(&self) {
1932 let mut state = self.lock();
1933 while state.ended.is_none() {
1934 state = self.changed.wait(state).unwrap_or_else(|e| e.into_inner());
1935 }
1936 }
1937
1938 fn lock(&self) -> std::sync::MutexGuard<'_, SpoolState> {
1939 self.state.lock().unwrap_or_else(|e| e.into_inner())
1940 }
1941
1942 fn write(&self, bytes: &[u8]) -> Result<bool, String> {
1944 let mut sink = self.sink.lock().unwrap_or_else(|e| e.into_inner());
1945 let Some(file) = sink.as_mut() else {
1946 return Ok(false);
1947 };
1948 file.write_all(bytes).map_err(|e| {
1949 format!(
1950 "Could not write {}: {e}",
1951 self.tee
1952 .as_ref()
1953 .filter(|t| !t.to_stdout())
1954 .map_or("what came in".to_string(), |t| t.path.display().to_string())
1955 )
1956 })?;
1957 drop(sink);
1958 if let Some(out) = self.pass.lock().unwrap_or_else(|e| e.into_inner()).as_mut() {
1961 out.write_all(bytes)
1962 .and_then(|()| out.flush())
1963 .map_err(|e| format!("Could not write standard output: {e}"))?;
1964 }
1965 let total =
1966 self.bytes.fetch_add(bytes.len() as u64, Ordering::Relaxed) + bytes.len() as u64;
1967 let now = Instant::now();
1968 let mut state = self.lock();
1969 if state.lines < WANTED_LINES {
1970 state.lines += bytes.iter().filter(|&&b| b == b'\n').count();
1971 }
1972 state.samples.push_back((now, total));
1973 while state
1974 .samples
1975 .front()
1976 .is_some_and(|(at, _)| now.duration_since(*at) > RATE_WINDOW)
1977 && state.samples.len() > 2
1978 {
1979 state.samples.pop_front();
1980 }
1981 drop(state);
1982 self.changed.notify_all();
1983 Ok(true)
1984 }
1985
1986 fn finish(&self, reason: Option<String>) {
1989 let file = self.sink.lock().unwrap_or_else(|e| e.into_inner()).take();
1990 if let Ok(mut pass) = self.pass.try_lock() {
1993 pass.take();
1994 }
1995 let mut reason = reason;
1996 if let (Some(mut file), Some(tee)) = (file, self.tee.as_ref().filter(|t| !t.to_stdout())) {
1997 let finished = (if tee.raw {
1998 Ok(())
1999 } else {
2000 crate::loading::tee::fix_wav_sizes(&mut file).map(|_| ())
2001 })
2002 .and_then(|()| file.sync_all());
2003 if let Err(e) = finished
2004 && reason.is_none()
2005 {
2006 reason = Some(format!("Could not finish {}: {e}", tee.path.display()));
2007 }
2008 }
2009 let mut state = self.lock();
2010 if state.ended.is_none() {
2011 state.ended = Some(reason);
2012 state.finished = Some(Instant::now());
2013 }
2014 let watcher = state.watcher.take();
2015 drop(state);
2016 self.changed.notify_all();
2017 if let Some(watcher) = watcher {
2018 watcher.wake();
2019 }
2020 }
2021
2022 fn wake_on_end(&self, shared: Arc<Shared>) {
2024 let mut state = self.lock();
2025 if state.ended.is_some() {
2026 drop(state);
2027 shared.wake();
2028 } else {
2029 state.watcher = Some(shared);
2030 }
2031 }
2032}
2033
2034pub struct SpoolHandle {
2037 spool: Arc<Spool>,
2038}
2039
2040impl SpoolHandle {
2041 pub fn spool(&self) -> &Arc<Spool> {
2042 &self.spool
2043 }
2044}
2045
2046impl Drop for SpoolHandle {
2047 fn drop(&mut self) {
2048 self.spool.stop();
2049 }
2050}
2051
2052fn copy_on(mut reader: impl Read + Send + 'static, spool: Arc<Spool>) {
2056 let _ = std::thread::Builder::new()
2057 .name("datui-spool".to_string())
2058 .spawn(move || {
2059 let mut buf = vec![0u8; CHUNK];
2060 let reason = loop {
2061 if spool.stopped() {
2062 break None;
2063 }
2064 let n = match reader.read(&mut buf) {
2065 Ok(0) => break None,
2066 Ok(n) => n,
2067 Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
2068 Err(e) => break Some(format!("Standard input failed: {e}")),
2069 };
2070 match spool.write(&buf[..n]) {
2071 Ok(true) => {}
2072 Ok(false) => break None,
2073 Err(reason) => break Some(reason),
2074 }
2075 spool.lock().drained = n < buf.len();
2076 };
2077 spool.finish(reason);
2078 });
2079}
2080
2081const WANTED_LINES: usize = 1000;
2084
2085pub enum Spooled {
2087 Temp(TempDownload),
2089 Kept(PathBuf),
2091}
2092
2093pub(crate) fn spool<R: Read + Send + 'static>(
2099 open: impl FnOnce() -> crate::cloud::download::Opened<R>,
2100 options: OpenOptions,
2101 writer: &Writer,
2102 read: &AtomicU64,
2103 stdout: Option<Box<dyn Write + Send>>,
2104) -> Result<(Spooled, OpenOptions), String> {
2105 let tee = options.tee.clone().map(|path| Tee {
2106 path,
2107 raw: options.tee_raw,
2108 });
2109 let tee = tee.map(|tee| Tee {
2111 raw: tee.raw || tee.to_stdout(),
2112 ..tee
2113 });
2114 let pass = match &tee {
2115 Some(tee) if tee.to_stdout() => Some(stdout.ok_or_else(|| {
2116 "--tee - passes the stream on to standard output, which only the datui command has."
2117 .to_string()
2118 })?),
2119 _ => None,
2120 };
2121 let progressive = !options.follow && tee.is_none();
2122 let (reader, _) = open().map_err(|e| format!("Could not read standard input: {e}"))?;
2123 let (spooled, file) = match &tee {
2124 Some(tee) if !tee.to_stdout() => {
2125 let file = crate::loading::tee::create(&tee.path, options.force)?;
2126 (Spooled::Kept(tee.path.clone()), file)
2127 }
2128 _ => {
2129 let dir = crate::loading::stdin::spool_dir(&options);
2130 let Some((named, claim)) = writer
2131 .create(|| TempDownload::create(dir.as_deref(), None))
2132 .map_err(|e| crate::error_display::user_message_from_report(&e, None))?
2133 else {
2134 return Err("Reading standard input was stopped.".to_string());
2135 };
2136 let file = named
2137 .as_file()
2138 .try_clone()
2139 .map_err(|e| format!("Could not write what came in: {e}"))?;
2140 (Spooled::Temp(TempDownload::held(named, Some(claim))), file)
2141 }
2142 };
2143 let spooled_path = match &spooled {
2144 Spooled::Temp(download) => download.path().to_path_buf(),
2145 Spooled::Kept(path) => path.clone(),
2146 };
2147 let spool = Arc::new(Spool::new(file, tee, pass));
2148 let handle = Arc::new(SpoolHandle {
2149 spool: spool.clone(),
2150 });
2151 copy_on(reader, spool.clone());
2152 let wanted = if options.has_header == Some(false) {
2154 1
2155 } else {
2156 2
2157 };
2158 let mut state = spool.lock();
2159 loop {
2160 read.store(spool.bytes(), Ordering::Relaxed);
2161 if writer.stopped() {
2162 drop(state);
2163 spool.stop();
2164 return Err("Reading standard input was stopped.".to_string());
2165 }
2166 let enough = (options.follow || progressive)
2168 && (state.lines >= WANTED_LINES
2169 || (state.drained
2170 && (state.lines >= wanted || stream::begins_with_schema(&spooled_path))));
2171 if enough || state.ended.is_some() {
2172 break;
2173 }
2174 state = spool
2175 .changed
2176 .wait_timeout(state, Duration::from_millis(100))
2177 .unwrap_or_else(|e| e.into_inner())
2178 .0;
2179 }
2180 if let Some(Some(reason)) = &state.ended {
2181 return Err(reason.clone());
2182 }
2183 drop(state);
2184 read.store(spool.bytes(), Ordering::Relaxed);
2185 let path = spooled_path;
2186 let mut head = Vec::new();
2187 File::open(&path)
2188 .and_then(|f| f.take(4096).read_to_end(&mut head))
2189 .map_err(|e| format!("Could not read standard input back: {e}"))?;
2190 if head.is_empty() {
2191 return Err("Nothing came in on standard input.".to_string());
2192 }
2193 let (format, compression, guessed) = crate::loading::stdin::sniff_for(&head, &options);
2194 let asked = options.clone();
2195 let options = OpenOptions {
2196 format_guessed: options.format.is_none() && guessed,
2197 format: options.format.or(Some(format)),
2198 compression: options.compression.or(compression),
2199 ..options
2200 };
2201 if progressive {
2202 let read_on = followed_stream(&path, options.format, &options)
2203 || refusal(options.format, &options).is_none();
2204 if read_on {
2205 return Ok((
2206 spooled,
2207 OpenOptions {
2208 spool: Some(handle),
2209 follow: true,
2210 pipe: true,
2211 ..options
2212 },
2213 ));
2214 }
2215 while spool.ended().is_none() {
2217 if writer.stopped() {
2218 spool.stop();
2219 return Err("Reading standard input was stopped.".to_string());
2220 }
2221 read.store(spool.bytes(), Ordering::Relaxed);
2222 let state = spool.lock();
2223 if state.ended.is_none() {
2224 let _ = spool
2225 .changed
2226 .wait_timeout(state, Duration::from_millis(100))
2227 .unwrap_or_else(|e| e.into_inner());
2228 }
2229 }
2230 if let Some(Some(reason)) = spool.ended() {
2231 return Err(reason);
2232 }
2233 read.store(spool.bytes(), Ordering::Relaxed);
2234 let options = crate::loading::stdin::described(&path, asked)?;
2235 return Ok((spooled, options));
2236 }
2237 if options.follow
2239 && options.tee.is_none()
2240 && !followed_stream(&path, options.format, &options)
2241 && let Some(refusal) = refusal(options.format, &options)
2242 {
2243 return Err(refusal);
2244 }
2245 Ok((
2246 spooled,
2247 OpenOptions {
2248 spool: Some(handle),
2249 tee: None,
2251 ..options
2252 },
2253 ))
2254}
2255
2256#[cfg(test)]
2257mod tests {
2258 use super::*;
2259
2260 #[test]
2263 fn errors_name_the_file() {
2264 let path = Path::new("/data/app.log");
2265 for (doing, e) in [
2266 (
2267 "following it stopped",
2268 std::io::Error::from(std::io::ErrorKind::NotFound),
2269 ),
2270 (
2271 "reading it stopped",
2272 std::io::Error::from(std::io::ErrorKind::PermissionDenied),
2273 ),
2274 ] {
2275 let message = failed_message(Some(path), doing, &e);
2276 eprintln!("{message}");
2277 crate::formats::readers::bad_input::assert_shape(&message, path);
2278 assert!(message.contains(&doing[1..]), "{message}");
2279 }
2280 let e = std::io::Error::from(std::io::ErrorKind::UnexpectedEof);
2281 let message = failed_message(None, "reading it stopped", &e);
2282 assert_eq!(
2283 message,
2284 "Standard input: reading it stopped. Unexpected end of file."
2285 );
2286 }
2287
2288 fn tail_of(text: &[u8], format: FileFormat, options: &OpenOptions) -> Tail {
2289 let dir = tempfile::tempdir().unwrap();
2290 let path = dir.path().join("t");
2291 std::fs::write(&path, text).unwrap();
2292 let mut tail = Tail::new(format, options, &Schema::default());
2293 let mut file = File::open(&path).unwrap();
2294 tail.read_on(&mut file, text.len() as u64, false).unwrap();
2295 tail
2296 }
2297
2298 fn polars_rows(text: &[u8], options: &OpenOptions) -> usize {
2300 let complete = text.iter().rposition(|&b| b == b'\n').map_or(0, |i| i + 1);
2301 let mut read = CsvReadOptions::default().with_has_header(options.has_header != Some(false));
2302 read = read.map_parse_options(|p| {
2303 p.with_comment_prefix(
2304 options
2305 .comment_char
2306 .as_deref()
2307 .map(polars::io::csv::read::CommentPrefix::new_from_str),
2308 )
2309 });
2310 if let Some(skip) = options.skip_lines {
2311 read.skip_lines = skip;
2312 }
2313 CsvReader::new(std::io::Cursor::new(text[..complete].to_vec()))
2314 .with_options(read)
2315 .finish()
2316 .map(|df| df.height())
2317 .unwrap_or(0)
2318 }
2319
2320 #[test]
2324 fn records_are_counted_as_polars_reads_them() {
2325 let cases: [(&[u8], OpenOptions); 6] = [
2326 (b"a,b\n1,2\n3,4\n", OpenOptions::default()),
2327 (b"a,b\n1,2\n\n3,4\n5,", OpenOptions::default()),
2328 (b"a,b\n1,\"x\ny\"\n3,4\n", OpenOptions::default()),
2329 (b"a,b\r\n1,2\r\n3,4\r\n", OpenOptions::default()),
2330 (
2331 b"#c\na,b\n#x\n1,2\n3,4\n",
2332 OpenOptions {
2333 comment_char: Some("#".to_string()),
2334 ..Default::default()
2335 },
2336 ),
2337 (
2338 b"junk\na,b\n1,2\n",
2339 OpenOptions::default().with_skip_lines(1),
2340 ),
2341 ];
2342 for (text, options) in cases {
2343 let tail = tail_of(text, FileFormat::Csv, &options);
2344 assert_eq!(
2345 tail.rows(),
2346 polars_rows(text, &options),
2347 "{}",
2348 String::from_utf8_lossy(text)
2349 );
2350 }
2351 let lines = tail_of(
2352 b"{\"a\":1}\n\n{\"a\":2}\n{\"a\":",
2353 FileFormat::Jsonl,
2354 &Default::default(),
2355 );
2356 assert_eq!(lines.rows(), 2);
2357 assert_eq!(lines.complete(), 17);
2358 }
2359
2360 #[cfg(target_os = "linux")]
2363 #[test]
2364 fn an_append_is_heard_of_without_a_check() {
2365 let dir = tempfile::tempdir().unwrap();
2366 let path = dir.path().join("log.csv");
2367 std::fs::write(&path, "t\n1\n").unwrap();
2368 let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2369 .finish()
2370 .unwrap();
2371 let (_, tail) =
2372 bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
2373 let (tx, rx) = std::sync::mpsc::channel();
2374 let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
2375 let guard = Duration::from_secs(30);
2376 let mut file = std::fs::OpenOptions::new()
2379 .append(true)
2380 .open(&path)
2381 .unwrap();
2382 let mut rows = 1;
2383 let deadline = Instant::now() + guard;
2384 let news = loop {
2385 assert!(Instant::now() < deadline, "the watcher never heard");
2386 file.write_all(format!("{}\n", rows + 1).as_bytes())
2387 .unwrap();
2388 rows += 1;
2389 if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
2390 break news;
2391 }
2392 };
2393 assert!(matches!(news.change, Change::Grew { rows: 2.., .. }));
2394 drop(follow);
2395 }
2396
2397 #[test]
2402 fn a_file_replaced_before_the_watcher_opens_it_is_read_again() {
2403 let dir = tempfile::tempdir().unwrap();
2404 let path = dir.path().join("rotated.csv");
2405 std::fs::write(&path, "t\n1\n2\n").unwrap();
2406 let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2407 .finish()
2408 .unwrap();
2409 let (_, tail) =
2410 bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
2411 let next = dir.path().join("next.csv");
2412 std::fs::write(&next, "t\n7\n8\n9\n").unwrap();
2413 std::fs::rename(&next, &path).unwrap();
2414 let (tx, rx) = std::sync::mpsc::channel();
2415 let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
2416 let deadline = Instant::now() + Duration::from_secs(30);
2417 let news = loop {
2418 assert!(Instant::now() < deadline, "the watcher never reported");
2419 follow.check_now();
2420 if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
2421 break news;
2422 }
2423 };
2424 assert!(
2425 matches!(news.change, Change::Restarted { rows: 3, .. }),
2426 "read on, not again"
2427 );
2428 }
2429
2430 #[test]
2433 fn a_tail_reads_on_from_where_it_stopped() {
2434 let dir = tempfile::tempdir().unwrap();
2435 let path = dir.path().join("grow.csv");
2436 let mut out = File::create(&path).unwrap();
2437 out.write_all(b"t,n\n1.5,2\n2.5,").unwrap();
2438 let schema = Schema::from_iter([
2439 Field::new("t".into(), DataType::Float64),
2440 Field::new("n".into(), DataType::Int64),
2441 ]);
2442 let mut tail = Tail::new(FileFormat::Csv, &OpenOptions::default(), &schema);
2443 let mut file = File::open(&path).unwrap();
2444 let len = |p: &Path| std::fs::metadata(p).unwrap().len();
2445 tail.read_on(&mut file, len(&path), true).unwrap();
2446 assert_eq!((tail.rows(), tail.complete()), (1, 10));
2447 out.write_all(b"3\nx,4\n4.5,5,6\n").unwrap();
2448 tail.read_on(&mut file, len(&path), true).unwrap();
2449 assert_eq!(tail.rows(), 4);
2450 assert_eq!(
2451 tail.misfits(),
2452 2,
2453 "a word for a number, and a field too many"
2454 );
2455 }
2456
2457 #[test]
2460 fn moving_the_bound_reads_the_new_rows_through_the_view() {
2461 let dir = tempfile::tempdir().unwrap();
2462 let path = dir.path().join("grow.csv");
2463 std::fs::write(&path, "a,b\n1,x\n2,y\n3,").unwrap();
2464 let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2465 .with_truncate_ragged_lines(true)
2466 .with_ignore_errors(true)
2467 .finish()
2468 .unwrap();
2469 let mut root = scan.clone();
2470 bound(&mut root, &path, 2);
2471 let mut view = root.clone().filter(col("a").gt(lit(1)));
2472 view.collect_schema().unwrap();
2473 assert_eq!(view.clone().collect().unwrap().height(), 1);
2474 let mut out = std::fs::OpenOptions::new()
2475 .append(true)
2476 .open(&path)
2477 .unwrap();
2478 out.write_all(b"z\n4,w\n5").unwrap();
2479 bound(&mut view, &path, 4);
2480 let df = view.clone().collect().unwrap();
2481 assert_eq!(df.height(), 3, "{df}");
2482 bound(&mut root, &path, 4);
2483 assert_eq!(root.collect().unwrap().height(), 4, "the partial row waits");
2484 }
2485
2486 fn marked(
2489 text: &[u8],
2490 format: FileFormat,
2491 options: &OpenOptions,
2492 every: u64,
2493 scan: impl Fn(&Path) -> LazyFrame,
2494 ) -> (tempfile::TempDir, PathBuf, LazyFrame, Arc<Marks>, usize) {
2495 let dir = tempfile::tempdir().unwrap();
2496 let path = dir.path().join("marked");
2497 std::fs::write(&path, text).unwrap();
2498 let mut lf = scan(&path);
2499 let schema = lf.collect_schema().unwrap();
2500 let mut tail = Tail::new(format, options, &schema);
2501 tail.mark_every = (every, u64::MAX);
2502 tail.read_on(&mut File::open(&path).unwrap(), text.len() as u64, false)
2503 .unwrap();
2504 let marks = Arc::new(Marks::default());
2505 marks.take_from(&mut tail);
2506 bound(&mut lf, &path, tail.rows());
2507 (dir, path, lf, marks, tail.rows())
2508 }
2509
2510 fn csv_scan(path: &Path, options: &OpenOptions) -> LazyFrame {
2511 let mut reader = LazyCsvReader::new(PlRefPath::try_from_path(path).unwrap())
2512 .with_ignore_errors(true)
2513 .with_truncate_ragged_lines(true)
2514 .with_has_header(options.has_header != Some(false))
2515 .with_comment_prefix(options.comment_char.as_deref().map(PlSmallStr::from_str));
2516 if let Some(skip) = options.skip_lines {
2517 reader = reader.with_skip_lines(skip);
2518 }
2519 reader.finish().unwrap()
2520 }
2521
2522 #[test]
2526 fn a_window_from_a_mark_reads_what_a_read_from_the_start_does() {
2527 let mut csv = b"skipped\nt,s,n\n".to_vec();
2528 let mut crlf = b"t,s,n\r\n".to_vec();
2529 let mut lines = Vec::new();
2530 for i in 0..120 {
2531 let row = match i % 9 {
2532 0 => format!("{i},\"two\nlines\",{}\n", i * 2),
2533 3 => "\n".to_string(),
2534 5 => "# a comment\n".to_string(),
2535 7 => format!("{i},x,oops\n"),
2536 _ => format!("{i},s{i},{}\n", i * 2),
2537 };
2538 csv.extend(row.as_bytes());
2539 crlf.extend(format!("{i},s{i},{}\r\n", i * 2).as_bytes());
2540 lines.extend(format!("{{\"t\":{i},\"s\":\"s{i}\"}}\n").as_bytes());
2541 if i % 4 == 1 {
2542 lines.extend(b"\n \n");
2543 }
2544 }
2545 csv.extend(b"999,partial");
2546 let commented = OpenOptions {
2547 comment_char: Some("#".to_string()),
2548 ..OpenOptions::default().with_skip_lines(1)
2549 };
2550 let (schema, batches) = stream_messages(&arrow_rows(0, 100), 3);
2552 let mut arrows = schema;
2553 batches.iter().for_each(|batch| arrows.extend(batch));
2554 arrows.extend(&batches[0][..20]);
2555 let cases: Vec<(&[u8], FileFormat, OpenOptions)> = vec![
2556 (&csv, FileFormat::Csv, commented),
2557 (&crlf, FileFormat::Csv, OpenOptions::default()),
2558 (&lines, FileFormat::Jsonl, OpenOptions::default()),
2559 (&arrows, FileFormat::Arrow, OpenOptions::default()),
2560 ];
2561 for (text, format, options) in cases {
2562 let scan = |path: &Path| match format {
2563 FileFormat::Jsonl => scan_lines(path, &options, false, &mut Vec::new()).unwrap(),
2564 FileFormat::Arrow => stream::scan(path).unwrap(),
2565 _ => csv_scan(path, &options),
2566 };
2567 let (_dir, path, lf, marks, rows) = marked(text, format, &options, 7, scan);
2568 let window = Window {
2569 lf: lf.clone(),
2570 path: path.clone(),
2571 marks: marks.clone(),
2572 known: None,
2573 };
2574 let whole = lf.clone().collect().unwrap();
2575 assert_eq!(whole.height(), rows);
2576 for start in (0..rows + 3).step_by(5) {
2577 for len in [1, 6, 40] {
2578 let read = crate::formats::pushdown::Windowed::window(&window, start, len)
2579 .unwrap()
2580 .collect()
2581 .unwrap();
2582 let expected = whole.slice(start as i64, len);
2583 assert!(
2584 read.equals_missing(&expected),
2585 "{format:?} rows {start}+{len}:\n{read:?}\n{expected:?}"
2586 );
2587 }
2588 if start < rows {
2589 assert!(
2590 from_marks(&lf, &path, &marks, start, Some(start + 1)).is_some(),
2591 "{format:?} row {start} is read from a mark"
2592 );
2593 }
2594 }
2595 }
2596 }
2597
2598 fn arrow_rows(from: i64, n: i64) -> DataFrame {
2600 df!(
2601 "t" => (from..from + n).collect::<Vec<_>>(),
2602 "s" => (from..from + n).map(|i| format!("s{i}")).collect::<Vec<_>>(),
2603 "x" => (from..from + n).map(|i| i as f64 / 2.0).collect::<Vec<_>>(),
2604 )
2605 .unwrap()
2606 }
2607
2608 #[test]
2612 fn an_arrow_stream_is_counted_and_read_by_its_batches() {
2613 let dir = tempfile::tempdir().unwrap();
2614 let path = dir.path().join("live.arrows");
2615 let (schema, batches) = stream_messages(&arrow_rows(0, 50), 4);
2616 let mut head = schema.clone();
2617 head.extend(&batches[0]);
2618 head.extend(&batches[1][..batches[1].len() - 3]);
2619 std::fs::write(&path, &head).unwrap();
2620 let lf = stream::scan(&path).unwrap();
2621 let (mut lf, mut tail) =
2622 bound_to_complete(lf, &path, FileFormat::Arrow, &OpenOptions::default()).unwrap();
2623 assert_eq!(tail.rows(), 4, "the second batch is not all there");
2624 assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 4));
2625
2626 let mut rest = batches[1][batches[1].len() - 3..].to_vec();
2627 batches[2..].iter().for_each(|batch| rest.extend(batch));
2628 let mut file = std::fs::OpenOptions::new()
2629 .append(true)
2630 .open(&path)
2631 .unwrap();
2632 file.write_all(&rest).unwrap();
2633 let size = file.metadata().unwrap().len();
2634 tail.read_on(&mut File::open(&path).unwrap(), size, true)
2635 .unwrap();
2636 assert_eq!((tail.rows(), tail.misfits()), (50, 0));
2637 bound(&mut lf, &path, tail.rows());
2638 assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 50));
2639 let kept = lf
2640 .clone()
2641 .filter(col("t").gt_eq(lit(45)))
2642 .select([col("s")])
2643 .collect()
2644 .unwrap();
2645 assert_eq!(kept, arrow_rows(45, 5).select(["s"]).unwrap());
2646 let count = lf.select([len()]).collect().unwrap();
2647 assert_eq!(count.column("len").unwrap().u32().unwrap().get(0), Some(50));
2648
2649 let legacy = dir.path().join("legacy.arrows");
2651 std::fs::write(
2652 &legacy,
2653 crate::formats::ipc_stream::tests::stream(&arrow_rows(0, 10), None, true),
2654 )
2655 .unwrap();
2656 let (lf, tail) = bound_to_complete(
2657 stream::scan(&legacy).unwrap(),
2658 &legacy,
2659 FileFormat::Arrow,
2660 &OpenOptions::default(),
2661 )
2662 .unwrap();
2663 assert_eq!(tail.rows(), 10);
2664 assert_eq!(lf.collect().unwrap(), arrow_rows(0, 10));
2665
2666 let cats = df!("c" => ["a", "b", "a"])
2668 .unwrap()
2669 .lazy()
2670 .with_column(col("c").cast(DataType::from_categories(Categories::global())))
2671 .collect()
2672 .unwrap();
2673 let (schema, batches) = stream_messages(&cats, 3);
2674 let dictionary = dir.path().join("dict.arrows");
2675 std::fs::write(&dictionary, [schema, batches.concat()].concat()).unwrap();
2676 assert!(
2677 stream::scan(&dictionary).is_err_and(|e| e.contains("dictionary")),
2678 "refused"
2679 );
2680 }
2681
2682 #[test]
2685 fn a_filtered_view_reads_and_counts_on_from_what_is_known() {
2686 let mut text = b"t,n\n".to_vec();
2687 for i in 0..300 {
2688 text.extend(format!("{i},{}\n", i % 5).as_bytes());
2689 }
2690 let options = OpenOptions::default();
2691 let (_dir, path, lf, marks, rows) = marked(&text, FileFormat::Csv, &options, 16, |p| {
2692 csv_scan(p, &options)
2693 });
2694 let view = lf.filter(col("n").eq(lit(3)));
2695 let whole = view.clone().collect().unwrap();
2696 let known = (whole.column("t").unwrap().i64().unwrap().to_vec())
2698 .into_iter()
2699 .filter(|t| t.unwrap() < 200)
2700 .count();
2701 let rest = from_marks(&view, &path, &marks, 200, None).unwrap();
2702 let after = rest.collect().unwrap().height();
2703 assert_eq!(known + after, whole.height());
2704 assert_eq!(rows, 300);
2705 let window = Window {
2706 lf: view.clone(),
2707 path,
2708 marks,
2709 known: Some(vec![(0, 0), (known, 200)]),
2710 };
2711 for start in [0, 10, known - 1, known, known + 5, whole.height() - 3] {
2712 let read = crate::formats::pushdown::Windowed::window(&window, start, 4)
2713 .unwrap()
2714 .collect()
2715 .unwrap();
2716 assert!(
2717 read.equals_missing(&whole.slice(start as i64, 4)),
2718 "{start}"
2719 );
2720 }
2721 }
2722
2723 #[test]
2727 fn a_replaced_file_is_told_from_a_grown_one() {
2728 let dir = tempfile::tempdir().unwrap();
2729 let path = dir.path().join("log.csv");
2730 std::fs::write(&path, "t\n1\n").unwrap();
2731 let held = File::open(&path).unwrap();
2732 let known = identity_of(&held);
2733 assert!(cfg!(not(any(unix, windows))) || known.is_some());
2734 let now = |path: &Path| identity_at(path, &std::fs::metadata(path).unwrap());
2735 std::fs::OpenOptions::new()
2736 .append(true)
2737 .open(&path)
2738 .unwrap()
2739 .write_all(b"2\n")
2740 .unwrap();
2741 assert!(!replaced(known, now(&path)), "grown in place");
2742 let other = dir.path().join("next.csv");
2743 std::fs::write(&other, "t\n1\n2\n3\n").unwrap();
2744 std::fs::rename(&other, &path).unwrap();
2745 assert_eq!(
2746 replaced(known, now(&path)),
2747 known.is_some(),
2748 "renamed over it"
2749 );
2750 assert!(!replaced(None, now(&path)), "unknown is no replacement");
2751 drop(held);
2752 }
2753
2754 #[cfg(unix)]
2756 #[test]
2757 fn a_deleted_file_reads_through_its_handle() {
2758 let dir = tempfile::tempdir().unwrap();
2759 let path = dir.path().join("gone.csv");
2760 std::fs::write(&path, "a\n1\n2\n").unwrap();
2761 let mut lf = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
2762 .finish()
2763 .unwrap();
2764 bound(&mut lf, &path, 2);
2765 let mut view = lf.filter(col("a").gt(lit(0)));
2766 view.collect_schema().unwrap();
2767 let handle = File::open(&path).unwrap();
2768 std::fs::remove_file(&path).unwrap();
2769 read_through(&mut view, &path, &handle);
2770 assert_eq!(view.collect().unwrap().height(), 2);
2771 }
2772}