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 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
952pub(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 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
975struct Piece {
978 path: PathBuf,
979 span: Span,
980 skip: usize,
981 take: usize,
982 parse: Parse,
983 schema: SchemaRef,
984}
985
986const 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
1022pub(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
1082fn 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
1111pub(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
1125pub(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 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
1164fn 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
1186pub(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 Ok(read.unwrap_or_else(|| self.lf.clone().slice(start as i64, len as IdxSize)))
1210 }
1211}
1212
1213#[derive(Clone)]
1215pub enum Change {
1216 Grew { rows: usize, misfits: usize },
1219 Restarted { rows: usize, misfits: usize },
1222 Gone { handle: Option<Arc<File>> },
1224 NewFields(Vec<Field>),
1227 Ended(Option<String>),
1229 Failed(String),
1231}
1232
1233#[derive(Clone)]
1235pub struct News {
1236 pub(crate) id: u64,
1237 pub(crate) change: Change,
1238}
1239
1240#[derive(Default)]
1242struct Shared {
1243 stop: AtomicBool,
1244 poke: Mutex<bool>,
1245 woken: Condvar,
1246 #[cfg(target_os = "linux")]
1248 bell: notify::Bell,
1249}
1250
1251impl Shared {
1252 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 #[cfg(target_os = "linux")]
1276 fn wait_for_change(
1277 &self,
1278 notify: ¬ify::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 if let Some(left) = left
1295 && !self.wait(left)
1296 {
1297 return false;
1298 }
1299 break;
1300 }
1301 }
1302 notify.drain();
1304 *last = Some(Instant::now());
1305 !self.stop.load(Ordering::Relaxed)
1306 }
1307}
1308
1309pub 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
1321pub 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#[derive(Clone, Debug, PartialEq)]
1335pub enum Standing {
1336 Following,
1337 Paused,
1338 Ended,
1340}
1341
1342pub struct Follow {
1345 id: u64,
1346 path: PathBuf,
1348 shared: Arc<Shared>,
1349 spool: Option<Arc<SpoolHandle>>,
1350 counted: usize,
1352 misfits: usize,
1353 restarted: bool,
1355 shown: usize,
1357 pub(crate) new_below: usize,
1359 pub(crate) standing: Standing,
1360 pub(crate) last_append: Option<Instant>,
1361 pub(crate) settle_at_end: bool,
1364 pub(crate) end_pending: bool,
1367 pub(crate) stale_view: bool,
1370 held: Option<Arc<File>>,
1372 marks: Arc<Marks>,
1374 pipe: bool,
1377 new_fields: Vec<Field>,
1380 pub(crate) described: bool,
1383}
1384
1385static NEXT_ID: AtomicU64 = AtomicU64::new(1);
1386
1387impl Follow {
1388 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 pub fn as_pipe(mut self) -> Follow {
1444 self.pipe = true;
1445 self.settle_at_end = false;
1446 self
1447 }
1448
1449 pub fn is_pipe(&self) -> bool {
1451 self.pipe
1452 }
1453
1454 pub fn live(&self) -> bool {
1456 self.standing != Standing::Ended
1457 }
1458
1459 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 pub fn shown(&self) -> usize {
1474 self.shown
1475 }
1476
1477 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 pub fn new_below(&self) -> usize {
1492 self.new_below
1493 }
1494
1495 pub fn behind(&self) -> bool {
1498 self.standing != Standing::Paused && (self.restarted || self.counted != self.shown)
1499 }
1500
1501 pub fn spool(&self) -> Option<&Arc<Spool>> {
1503 self.spool.as_ref().map(|handle| &handle.spool)
1504 }
1505
1506 pub fn check_now(&self) {
1508 self.shared.wake();
1509 }
1510
1511 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 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 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 pub(crate) fn take_new_fields(&mut self) -> Vec<Field> {
1564 std::mem::take(&mut self.new_fields)
1565 }
1566
1567 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 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
1603struct 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
1615type 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 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#[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
1661fn replaced(known: Option<Identity>, now: Option<Identity>) -> bool {
1664 matches!((known, now), (Some(known), Some(now)) if known != now)
1665}
1666
1667fn 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 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 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 ¬ify {
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 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) = ¬ify {
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 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 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
1806pub struct Spool {
1810 stop: AtomicBool,
1811 bytes: AtomicU64,
1812 state: Mutex<SpoolState>,
1813 changed: Condvar,
1814 sink: Mutex<Option<File>>,
1817 tee: Option<Tee>,
1819 pass: Mutex<Option<Box<dyn Write + Send>>>,
1821 started: Instant,
1822}
1823
1824#[derive(Clone, Debug)]
1826pub struct Tee {
1827 pub path: PathBuf,
1828 pub raw: bool,
1830}
1831
1832impl Tee {
1833 pub fn to_stdout(&self) -> bool {
1835 crate::loading::stdin::is_stdin(&self.path)
1836 }
1837
1838 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 lines: usize,
1854 drained: bool,
1856 ended: Option<Option<String>>,
1858 finished: Option<Instant>,
1860 samples: std::collections::VecDeque<(Instant, u64)>,
1862 watcher: Option<Arc<Shared>>,
1865}
1866
1867const 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 pub fn bytes(&self) -> u64 {
1886 self.bytes.load(Ordering::Relaxed)
1887 }
1888
1889 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 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 pub fn tee(&self) -> Option<&Tee> {
1910 self.tee.as_ref()
1911 }
1912
1913 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 pub fn ended(&self) -> Option<Option<String>> {
1926 self.lock().ended.clone()
1927 }
1928
1929 pub fn live(&self) -> bool {
1931 self.lock().ended.is_none()
1932 }
1933
1934 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 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 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 fn finish(&self, reason: Option<String>) {
1993 let file = self.sink.lock().unwrap_or_else(|e| e.into_inner()).take();
1994 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 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
2038pub 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
2056fn 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
2085const WANTED_LINES: usize = 1000;
2088
2089pub enum Spooled {
2091 Temp(TempDownload),
2093 Kept(PathBuf),
2095}
2096
2097pub(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 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 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 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 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 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 tee: None,
2255 ..options
2256 },
2257 ))
2258}
2259
2260#[cfg(test)]
2261mod tests {
2262 use super::*;
2263
2264 #[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 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 #[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 #[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 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 #[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 #[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 #[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 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 #[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 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 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 #[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 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 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 #[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 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 #[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 #[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}