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