use std::fs::File;
use std::io::{Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
use std::sync::mpsc::Sender;
use std::sync::{Arc, Condvar, Mutex};
use std::time::{Duration, Instant};
use polars::prelude::*;
use crate::download::TempDownload;
use crate::unfinished::Writer;
use crate::{AppEvent, CompressionFormat, FileFormat, OpenOptions};
pub(crate) mod lines;
#[cfg(target_os = "linux")]
mod notify;
pub(crate) mod stream;
#[doc(hidden)]
pub use stream::stream_messages;
pub const DEFAULT_INTERVAL: Duration = Duration::from_millis(250);
const CHUNK: usize = 1 << 20;
const LONGEST_RECORD: usize = 16 << 20;
pub(crate) const MARK_ROWS: u64 = 8192;
const MARK_BYTES: u64 = 1 << 20;
const MOST_FROM_A_MARK: u64 = 64 << 20;
pub fn refusal(format: Option<FileFormat>, options: &OpenOptions) -> Option<String> {
let format = format.unwrap_or(FileFormat::TEXT);
if format == FileFormat::Arrow {
return Some(
"An Arrow IPC file is read once it is finished; an Arrow IPC stream can be followed."
.to_string(),
);
}
if !format.follows() {
let followed: Vec<&str> = FileFormat::ALL
.into_iter()
.filter(|f| f.follows())
.map(FileFormat::title)
.collect();
let followed = match followed.split_last() {
Some((last, rest)) if !rest.is_empty() => format!("{} and {last}", rest.join(", ")),
_ => followed.join(""),
};
return Some(format!(
"Only {followed} and Arrow IPC streams can be followed as they grow; {} is read \
once it is finished.",
format.title()
));
}
if options.compression.is_some() {
return Some("A compressed file cannot be followed as it grows.".to_string());
}
if options.header_rows().is_some() || options.skip_tail_rows.is_some() {
return Some(
"A file read with --header-rows or --footer-rows cannot be followed.".to_string(),
);
}
if options.spec_name.is_some() || options.spec_file.is_some() || options.delimited.is_some() {
return Some("A file read through a format spec cannot be followed.".to_string());
}
None
}
pub fn followable_path(path: &Path, options: &OpenOptions) -> bool {
let format = options.format.or_else(|| {
path.extension()
.and_then(|e| e.to_str())
.and_then(FileFormat::from_extension)
});
let compression = options
.compression
.or_else(|| CompressionFormat::from_extension(path));
compression.is_none()
&& (refusal(format, options).is_none() || followed_stream(path, format, options))
}
pub(crate) fn followed_stream(
path: &Path,
format: Option<FileFormat>,
options: &OpenOptions,
) -> bool {
format == Some(FileFormat::Arrow)
&& options.compression.is_none()
&& options.spec_name.is_none()
&& options.spec_file.is_none()
&& crate::ipc_stream::is_stream_file(path)
}
pub fn refuse_paths(paths: &[PathBuf], options: &OpenOptions) -> Option<String> {
let [path] = paths else {
return Some("Only one file can be followed at a time.".to_string());
};
if crate::stdin::is_stdin(path) {
return None;
}
if !matches!(
crate::source::input_source(path),
crate::source::InputSource::Local(_)
) {
return Some("Only a local file or standard input can be followed.".to_string());
}
if path.is_dir() {
return Some("A directory cannot be followed: name a file in it.".to_string());
}
let format = options.format.or_else(|| {
path.extension()
.and_then(|e| e.to_str())
.and_then(FileFormat::from_extension)
});
let options = OpenOptions {
compression: options
.compression
.or_else(|| CompressionFormat::from_extension(path)),
..options.clone()
};
if followed_stream(path, format, &options) {
return None;
}
refusal(format, &options)
}
pub(crate) fn format_of(path: &Path, found: Option<FileFormat>) -> FileFormat {
found
.or_else(|| {
path.extension()
.and_then(|e| e.to_str())
.and_then(FileFormat::from_extension)
})
.unwrap_or(FileFormat::TEXT)
}
pub(crate) fn scan_lines(
path: &Path,
options: &OpenOptions,
every_line: bool,
read_python: &mut Vec<String>,
) -> color_eyre::Result<LazyFrame> {
let infer = if every_line {
None
} else {
Some(
options
.infer_schema_length
.and_then(std::num::NonZeroUsize::new)
.unwrap_or(std::num::NonZeroUsize::new(100).expect("not zero")),
)
};
let spool = options.spool.as_ref().map(|handle| handle.spool().clone());
let lf = lines::LinesScan::open(path, infer, true, spool)?.lazy()?;
crate::widgets::datatable::DataTableState::apply_parse_dates_to_json_lazyframe(
lf,
options,
read_python,
)
}
pub(crate) fn bound_to_complete(
mut lf: LazyFrame,
path: &Path,
format: FileFormat,
options: &OpenOptions,
) -> color_eyre::Result<(LazyFrame, Tail)> {
let schema = lf.collect_schema()?;
let mut tail = Tail::new(format, options, &schema);
if options.spool.is_some() {
tail = tail.widening();
}
tail.path = path.to_path_buf();
let mut file = File::open(path)?;
let len = file.metadata()?.len();
tail.read_on(&mut file, len, false)?;
bound(&mut lf, path, tail.rows());
Ok((lf, tail))
}
#[derive(Clone, Copy, Debug, PartialEq)]
enum Fits {
Anything,
Integer,
Number,
Boolean,
Text,
}
impl Fits {
fn of(dtype: &DataType) -> Fits {
if dtype.is_integer() {
Fits::Integer
} else if dtype.is_float() {
Fits::Number
} else if matches!(dtype, DataType::Boolean) {
Fits::Boolean
} else if matches!(dtype, DataType::String) {
Fits::Text
} else {
Fits::Anything
}
}
fn cell(self, cell: &str, nulls: &[String]) -> bool {
let cell = cell.trim();
if cell.is_empty() || nulls.iter().any(|n| n == cell) {
return true;
}
match self {
Fits::Integer => cell.parse::<i64>().is_ok() || cell.parse::<u64>().is_ok(),
Fits::Number => cell.parse::<f64>().is_ok(),
Fits::Boolean => {
cell.eq_ignore_ascii_case("true") || cell.eq_ignore_ascii_case("false")
}
Fits::Text | Fits::Anything => true,
}
}
fn value(self, value: &serde_json::Value) -> bool {
match self {
_ if value.is_null() => true,
Fits::Integer => value.is_i64() || value.is_u64(),
Fits::Number => value.is_number(),
Fits::Boolean => value.is_boolean(),
Fits::Text => value.is_string() || value.is_array(),
Fits::Anything => true,
}
}
}
#[derive(Clone, Debug)]
enum Layout {
Delimited {
separator: u8,
skip: u64,
header: bool,
comment: Option<Vec<u8>>,
nulls: Vec<String>,
},
Lines,
Stream,
Text,
}
#[derive(Clone, Debug, Default)]
struct NewMarks {
new: Vec<(u64, u64)>,
last: Option<(u64, u64)>,
}
#[derive(Clone, Debug)]
pub struct Tail {
path: PathBuf,
layout: Layout,
columns: Vec<(String, Fits)>,
complete: u64,
records: u64,
rows: u64,
misfits: u64,
fields: Option<usize>,
marks: NewMarks,
mark_every: (u64, u64),
widens: bool,
arrived: Vec<(String, Arrived)>,
}
const MOST_NEW_FIELDS: usize = 4096;
#[derive(Clone, Copy, Debug, PartialEq)]
enum Arrived {
Nothing,
Integer,
Number,
Boolean,
Text,
}
impl Arrived {
fn of(value: &serde_json::Value) -> Arrived {
match value {
serde_json::Value::Null => Arrived::Nothing,
serde_json::Value::Bool(_) => Arrived::Boolean,
serde_json::Value::Number(n) if n.is_i64() || n.is_u64() => Arrived::Integer,
serde_json::Value::Number(_) => Arrived::Number,
_ => Arrived::Text,
}
}
fn and(self, other: Arrived) -> Arrived {
match (self, other) {
(a, b) if a == b => a,
(Arrived::Nothing, x) | (x, Arrived::Nothing) => x,
(Arrived::Integer, Arrived::Number) | (Arrived::Number, Arrived::Integer) => {
Arrived::Number
}
_ => Arrived::Text,
}
}
fn dtype(self) -> DataType {
match self {
Arrived::Integer => DataType::Int64,
Arrived::Number => DataType::Float64,
Arrived::Boolean => DataType::Boolean,
Arrived::Nothing | Arrived::Text => DataType::String,
}
}
}
impl Tail {
pub fn new(format: FileFormat, options: &OpenOptions, schema: &Schema) -> Tail {
let layout = match format.separator() {
_ if format == FileFormat::Arrow => Layout::Stream,
Some(separator) => Layout::Delimited {
separator: options.separator_or(separator),
skip: options.skip_lines.unwrap_or(0) as u64
+ options.skip_rows.unwrap_or(0) as u64,
header: options.has_header != Some(false),
comment: options
.comment_char
.as_ref()
.filter(|c| !c.is_empty())
.map(|c| c.as_bytes().to_vec()),
nulls: options
.null_values
.iter()
.flatten()
.filter(|spec| !spec.contains('='))
.cloned()
.collect(),
},
None if format.is_lines() => Layout::Text,
None => Layout::Lines,
};
let columns = schema
.iter()
.map(|(name, dtype)| (name.to_string(), Fits::of(dtype)))
.collect();
Tail {
path: PathBuf::new(),
layout,
columns,
complete: 0,
records: 0,
rows: 0,
misfits: 0,
fields: None,
marks: NewMarks::default(),
mark_every: (MARK_ROWS, MARK_BYTES),
widens: false,
arrived: Vec::new(),
}
}
pub fn widening(mut self) -> Tail {
self.widens = matches!(self.layout, Layout::Lines);
self
}
pub fn new_fields(&self) -> Vec<Field> {
self.arrived
.iter()
.map(|(name, values)| Field::new(name.as_str().into(), values.dtype()))
.collect()
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn rows(&self) -> usize {
self.rows as usize
}
pub fn misfits(&self) -> usize {
self.misfits as usize
}
pub fn complete(&self) -> u64 {
self.complete
}
fn restart(&mut self) {
self.complete = 0;
self.records = 0;
self.rows = 0;
self.misfits = 0;
self.fields = None;
self.marks = NewMarks::default();
self.arrived.clear();
}
fn mark(marks: &mut NewMarks, every: (u64, u64), row: u64, start: u64, blank: bool) {
let (rows, bytes) = every;
if blank
|| marks.last.is_some_and(|(last_row, last_start)| {
row - last_row < rows && start - last_start < bytes
})
{
return;
}
marks.last = Some((row, start));
marks.new.push((row, start));
}
pub fn read_on(&mut self, file: &mut File, len: u64, check: bool) -> std::io::Result<()> {
self.read(file, len, check, false)
}
pub fn read_to_end(&mut self, file: &mut File, len: u64, check: bool) -> std::io::Result<()> {
let last = matches!(self.layout, Layout::Lines | Layout::Delimited { .. });
self.read(file, len, check, last)
}
fn read(&mut self, file: &mut File, len: u64, check: bool, last: bool) -> std::io::Result<()> {
if len <= self.complete {
return Ok(());
}
if matches!(self.layout, Layout::Stream) {
return self.read_messages(file, len);
}
file.seek(SeekFrom::Start(self.complete))?;
let mut reader = file.take(len - self.complete);
let quoted = matches!(self.layout, Layout::Delimited { .. });
let mut buf = vec![0u8; CHUNK];
let mut record: Vec<u8> = Vec::new();
let mut oversized = false;
let mut in_quotes = false;
let mut at = self.complete;
loop {
let n = match reader.read(&mut buf) {
Ok(0) => break,
Ok(n) => n,
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(e) => return Err(e),
};
let chunk = &buf[..n];
let mut i = 0;
while i < n {
let rest = &chunk[i..];
let next = if quoted {
rest.iter().position(|&b| b == b'\n' || b == b'"')
} else {
rest.iter().position(|&b| b == b'\n')
};
let Some(k) = next else {
keep(&mut record, rest, &mut oversized);
break;
};
keep(&mut record, &rest[..k], &mut oversized);
let byte = rest[k];
i += k + 1;
if byte == b'"' {
in_quotes = !in_quotes;
keep(&mut record, b"\"", &mut oversized);
} else if in_quotes {
keep(&mut record, b"\n", &mut oversized);
} else {
self.end_record(&record, oversized, check);
record.clear();
oversized = false;
self.complete = at + i as u64;
}
}
at += n as u64;
}
if last && at > self.complete {
self.end_record(&record, oversized, check);
self.complete = at;
}
Ok(())
}
fn read_messages(&mut self, file: &mut File, len: u64) -> std::io::Result<()> {
file.seek(SeekFrom::Start(self.complete))?;
let mut reader = std::io::BufReader::with_capacity(CHUNK, file);
while let Some((message, size)) = stream::next_message(&mut reader, len - self.complete)? {
match message {
stream::Message::Batch { rows } => {
Self::mark(
&mut self.marks,
self.mark_every,
self.rows,
self.complete,
false,
);
self.rows += rows;
}
stream::Message::Dictionary => {
return Err(std::io::Error::other(
"a dictionary batch arrived, which a followed stream cannot read",
));
}
stream::Message::Schema | stream::Message::End | stream::Message::Other => {}
}
self.records += 1;
self.complete += size;
}
Ok(())
}
fn end_record(&mut self, record: &[u8], oversized: bool, check: bool) {
let record = record.strip_suffix(b"\r").unwrap_or(record);
let start = self.complete;
let index = self.records;
self.records += 1;
match &self.layout {
Layout::Delimited {
separator,
skip,
header,
comment,
nulls,
} => {
if index < *skip {
return;
}
if comment.as_ref().is_some_and(|c| record.starts_with(c)) {
return;
}
if *header && self.fields.is_none() {
self.fields = Some(split_fields(record, *separator).len());
return;
}
Self::mark(
&mut self.marks,
self.mark_every,
self.rows,
start,
record.is_empty(),
);
self.rows += 1;
if check && (oversized || !self.cells_fit(record, *separator, nulls)) {
self.misfits += 1;
}
}
Layout::Stream => {}
Layout::Text => self.rows += 1,
Layout::Lines => {
if record.iter().all(u8::is_ascii_whitespace) {
return;
}
Self::mark(&mut self.marks, self.mark_every, self.rows, start, false);
self.rows += 1;
if check || self.widens {
let fits = !oversized && self.object_fits(record);
if check && !fits {
self.misfits += 1;
}
}
}
}
}
fn cells_fit(&self, record: &[u8], separator: u8, nulls: &[String]) -> bool {
let cells = split_fields(record, separator);
let expected = self.fields.unwrap_or(self.columns.len());
if record.is_empty() {
return true;
}
cells.len() == expected
&& cells.iter().zip(&self.columns).all(|(cell, (_, fits))| {
let text = String::from_utf8_lossy(cell);
fits.cell(unquote(&text), nulls)
})
}
fn object_fits(&mut self, record: &[u8]) -> bool {
let Ok(serde_json::Value::Object(object)) = serde_json::from_slice(record) else {
return false;
};
let mut fits = true;
for (key, value) in &object {
match self.columns.iter().find(|(name, _)| name == key) {
Some((_, kind)) => fits &= kind.value(value),
None if self.widens => fits &= Self::arrive(&mut self.arrived, key, value),
None => fits = false,
}
}
fits
}
fn arrive(arrived: &mut Vec<(String, Arrived)>, key: &str, value: &serde_json::Value) -> bool {
let kind = Arrived::of(value);
if let Some((_, seen)) = arrived.iter_mut().find(|(name, _)| name == key) {
*seen = seen.and(kind);
return true;
}
if arrived.len() >= MOST_NEW_FIELDS {
return false;
}
arrived.push((key.to_string(), kind));
true
}
}
fn keep(record: &mut Vec<u8>, bytes: &[u8], oversized: &mut bool) {
if record.len() + bytes.len() > LONGEST_RECORD {
*oversized = true;
return;
}
record.extend_from_slice(bytes);
}
fn split_fields(record: &[u8], separator: u8) -> Vec<&[u8]> {
let mut fields = Vec::new();
let mut in_quotes = false;
let mut start = 0;
for (i, &b) in record.iter().enumerate() {
if b == b'"' {
in_quotes = !in_quotes;
} else if b == separator && !in_quotes {
fields.push(&record[start..i]);
start = i + 1;
}
}
fields.push(&record[start..]);
fields
}
fn unquote(cell: &str) -> &str {
let trimmed = cell.trim();
trimmed
.strip_prefix('"')
.and_then(|c| c.strip_suffix('"'))
.unwrap_or(trimmed)
}
fn same_file(a: &str, b: &str) -> bool {
if a == b || Path::new(a) == Path::new(b) {
return true;
}
matches!(
(std::fs::canonicalize(a), std::fs::canonicalize(b)),
(Ok(a), Ok(b)) if a == b
)
}
fn scans(plan: &polars::lazy::dsl::DslPlan, path: &str) -> bool {
use polars::lazy::dsl::DslPlan;
match plan {
DslPlan::Scan {
sources: ScanSources::Paths(paths),
..
} if paths.len() == 1 => same_file(paths[0].as_str(), path),
DslPlan::Scan { .. } => {
stream::StreamScan::of(plan, path).is_some()
|| lines::LinesScan::of(plan, path).is_some()
}
DslPlan::IR { dsl, .. } => scans(dsl, path),
_ => false,
}
}
pub fn bound(lf: &mut LazyFrame, path: &Path, rows: usize) {
let path = path.to_string_lossy();
let rows = IdxSize::try_from(rows).unwrap_or(IdxSize::MAX);
bound_plan(&mut lf.logical_plan, &path, rows);
}
fn bound_plan(plan: &mut polars::lazy::dsl::DslPlan, path: &str, rows: IdxSize) {
use polars::lazy::dsl::DslPlan;
if crate::lines::bound(plan, rows) {
return;
}
match plan {
DslPlan::IR { dsl, .. } => {
let mut inner = Arc::unwrap_or_clone(dsl.clone());
bound_plan(&mut inner, path, rows);
*plan = inner;
return;
}
DslPlan::Slice {
input,
offset: 0,
len,
} if scans(input, path) => {
*len = rows;
return;
}
DslPlan::Scan { .. } if scans(plan, path) => {
let scan = std::mem::take(plan);
*plan = DslPlan::Slice {
input: Arc::new(scan),
offset: 0,
len: rows,
};
return;
}
_ => {}
}
crate::widgets::datatable::for_each_input(plan, &mut |input| bound_plan(input, path, rows));
}
pub fn read_through(lf: &mut LazyFrame, path: &Path, file: &File) {
let path = path.to_string_lossy();
read_through_plan(&mut lf.logical_plan, &path, file);
}
fn read_through_plan(plan: &mut polars::lazy::dsl::DslPlan, path: &str, file: &File) {
use polars::lazy::dsl::DslPlan;
match plan {
DslPlan::IR { dsl, .. } => {
let mut inner = Arc::unwrap_or_clone(dsl.clone());
read_through_plan(&mut inner, path, file);
*plan = inner;
return;
}
DslPlan::Scan { .. } if scans(plan, path) => {
let held: Option<Arc<dyn AnonymousScan>> = match stream::StreamScan::of(plan, path) {
Some(scan) => scan
.held(file)
.map(|s| Arc::new(s) as Arc<dyn AnonymousScan>),
None => lines::LinesScan::of(plan, path)
.and_then(|scan| scan.held(file))
.map(|s| Arc::new(s) as Arc<dyn AnonymousScan>),
};
if let DslPlan::Scan {
sources,
scan_type,
cached_ir,
..
} = plan
&& let Ok(handle) = file.try_clone()
{
match (held, &mut **scan_type) {
(Some(held), polars::lazy::dsl::FileScanDsl::Anonymous { function, .. }) => {
*function = held;
}
_ => *sources = ScanSources::Files(Arc::from([handle])),
}
*cached_ir = Default::default();
}
return;
}
_ => {}
}
crate::widgets::datatable::for_each_input(plan, &mut |input| {
read_through_plan(input, path, file)
});
}
#[derive(Default)]
pub struct Marks {
inner: Mutex<MarksInner>,
}
#[derive(Default)]
struct MarksInner {
at: Vec<(u64, u64)>,
complete: u64,
}
#[derive(Clone, Copy, Debug, PartialEq)]
struct Span {
row: u64,
start: u64,
end: u64,
}
impl Marks {
fn lock(&self) -> std::sync::MutexGuard<'_, MarksInner> {
self.inner.lock().unwrap_or_else(|e| e.into_inner())
}
fn take_from(&self, tail: &mut Tail) {
let mut inner = self.lock();
inner.at.append(&mut tail.marks.new);
inner.complete = tail.complete;
}
fn clear(&self) {
let mut inner = self.lock();
inner.at.clear();
inner.complete = 0;
}
fn span(&self, from: u64, to: u64) -> Option<Span> {
let inner = self.lock();
let before = inner.at.partition_point(|&(row, _)| row <= from);
let (row, start) = *inner.at.get(before.checked_sub(1)?)?;
let after = inner.at.partition_point(|&(row, _)| row < to);
let end = inner.at.get(after).map_or(inner.complete, |&(_, at)| at);
(end >= start && end - start <= MOST_FROM_A_MARK).then_some(Span { row, start, end })
}
}
#[derive(Clone)]
enum Parse {
Csv(Box<CsvReadOptions>),
Lines { ignore_errors: bool },
Stream(Arc<stream::StreamSchema>),
}
fn scan_node(plan: &polars::lazy::dsl::DslPlan) -> &polars::lazy::dsl::DslPlan {
match plan {
polars::lazy::dsl::DslPlan::IR { dsl, .. } => scan_node(dsl),
plan => plan,
}
}
impl Parse {
fn of(scan: &polars::lazy::dsl::DslPlan, schema: &SchemaRef) -> Option<Parse> {
use polars::lazy::dsl::{DslPlan, FileScanDsl};
if let Some(stream) = stream::StreamScan::in_plan(scan_node(scan)) {
return Some(Parse::Stream(stream.schema().clone()));
}
if let Some(lines) = lines::LinesScan::in_plan(scan_node(scan)) {
return Some(Parse::Lines {
ignore_errors: lines.ignore_errors(),
});
}
let DslPlan::Scan {
scan_type,
unified_scan_args,
..
} = scan_node(scan)
else {
return None;
};
if unified_scan_args.row_index.is_some() || unified_scan_args.include_file_paths.is_some() {
return None;
}
match &**scan_type {
FileScanDsl::Csv { options } => {
if options.columns.is_some()
|| options.projection.is_some()
|| options.row_index.is_some()
{
return None;
}
let mut options = (**options).clone();
options.path = None;
options.has_header = false;
options.skip_rows = 0;
options.skip_lines = 0;
options.skip_rows_after_header = 0;
options.n_rows = None;
options.schema = Some(schema.clone());
options.schema_overwrite = None;
options.dtype_overwrite = None;
options.column_names_overwrite = None;
options.raise_if_empty = false;
Some(Parse::Csv(Box::new(options)))
}
FileScanDsl::NDJson { options } => Some(Parse::Lines {
ignore_errors: options.ignore_errors,
}),
_ => None,
}
}
}
struct Piece {
path: PathBuf,
span: Span,
skip: usize,
take: usize,
parse: Parse,
schema: SchemaRef,
}
const PIECE_NAME: &str = "FOLLOWED";
impl polars::prelude::AnonymousScan for Piece {
fn as_any(&self) -> &dyn std::any::Any {
self
}
fn schema(&self, _infer_schema_length: Option<usize>) -> PolarsResult<SchemaRef> {
Ok(self.schema.clone())
}
fn scan(&self, args: polars::prelude::AnonymousScanArgs) -> PolarsResult<DataFrame> {
let take = args.n_rows.map_or(self.take, |n| n.min(self.take));
let mut file = File::open(&self.path)?;
file.seek(SeekFrom::Start(self.span.start))?;
let mut bytes = Vec::with_capacity((self.span.end - self.span.start) as usize);
file.take(self.span.end - self.span.start)
.read_to_end(&mut bytes)?;
let df = match &self.parse {
Parse::Csv(options) => {
let mut options = (**options).clone();
options.n_rows = Some(self.skip + take);
options
.into_reader_with_file_handle(std::io::Cursor::new(bytes))
.finish()?
}
Parse::Lines { ignore_errors } => {
lines::parse_run(&bytes, &self.schema, *ignore_errors)?
}
Parse::Stream(schema) => stream::decode_run(bytes, schema, self.skip + take)?,
};
Ok(df.slice(self.skip as i64, take))
}
}
pub(crate) fn from_marks(
lf: &LazyFrame,
path: &Path,
marks: &Marks,
from: usize,
to: Option<usize>,
) -> Option<LazyFrame> {
let path_text = path.to_string_lossy();
let mut plan = lf.logical_plan.clone();
let mut replaced = false;
let mut failed = false;
let piece = |scan: &polars::lazy::dsl::DslPlan, bound: usize| {
let to = to.map_or(bound, |to| to.min(bound));
let from = from.min(to);
let span = marks.span(from as u64, to as u64)?;
let schema = LazyFrame::from(scan.clone()).collect_schema().ok()?;
let parse = Parse::of(scan, &schema)?;
let piece = Piece {
path: path.to_path_buf(),
span,
skip: from - span.row as usize,
take: to - from,
parse,
schema: schema.clone(),
};
LazyFrame::anonymous_scan(
Arc::new(piece),
ScanArgsAnonymous {
schema: Some(schema),
name: PIECE_NAME,
..Default::default()
},
)
.ok()
.map(|lf| lf.logical_plan)
};
replace_bound(
&mut plan,
&path_text,
&mut |scan, bound| match piece(scan, bound) {
Some(plan) => {
replaced = true;
Some(plan)
}
None => {
failed = true;
None
}
},
);
(replaced && !failed).then(|| {
let mut out = lf.clone();
out.logical_plan = plan;
out
})
}
fn replace_bound(
plan: &mut polars::lazy::dsl::DslPlan,
path: &str,
with: &mut dyn FnMut(&polars::lazy::dsl::DslPlan, usize) -> Option<polars::lazy::dsl::DslPlan>,
) {
use polars::lazy::dsl::DslPlan;
match plan {
DslPlan::IR { dsl, .. } => {
let mut inner = Arc::unwrap_or_clone(dsl.clone());
replace_bound(&mut inner, path, with);
*plan = inner;
return;
}
DslPlan::Slice {
input,
offset: 0,
len,
} if scans(input, path) => {
if let Some(piece) = with(input, *len as usize) {
*plan = piece;
}
return;
}
_ => {}
}
crate::widgets::datatable::for_each_input(plan, &mut |input| replace_bound(input, path, with));
}
pub(crate) fn bound_of(lf: &LazyFrame, path: &Path) -> Option<usize> {
use polars::lazy::dsl::DslPlan;
let path = path.to_string_lossy();
(&lf.logical_plan).into_iter().find_map(|node| match node {
DslPlan::Slice {
input,
offset: 0,
len,
} if scans(input, &path) => Some(*len as usize),
_ => None,
})
}
pub(crate) fn widen(
root: &LazyFrame,
path: &Path,
format: FileFormat,
fields: &[Field],
rows: usize,
) -> Option<LazyFrame> {
let path_text = path.to_string_lossy();
let mut plan = root.logical_plan.clone();
let mut widened = None;
replace_lines_scan(&mut plan, &path_text, &mut |scan| {
let mut schema = (**scan.schema()).clone();
for field in fields {
if !schema.contains(field.name()) {
schema.with_column(field.name().clone(), field.dtype().clone());
}
}
if schema.len() == scan.schema().len() {
return None;
}
let schema = Arc::new(schema);
let lf = scan.with_schema(schema.clone()).lazy().ok()?;
widened = Some((lf.clone(), schema));
Some(lf.logical_plan)
});
let (raw, schema) = widened?;
if format == FileFormat::Journal {
let mut raw = raw;
bound(&mut raw, path, rows);
return Some(crate::journal::derive(raw, &schema).0);
}
let mut out = root.clone();
out.logical_plan = plan;
Some(out)
}
fn replace_lines_scan(
plan: &mut polars::lazy::dsl::DslPlan,
path: &str,
with: &mut dyn FnMut(&lines::LinesScan) -> Option<polars::lazy::dsl::DslPlan>,
) {
use polars::lazy::dsl::DslPlan;
if let DslPlan::IR { dsl, .. } = plan {
let mut inner = Arc::unwrap_or_clone(dsl.clone());
replace_lines_scan(&mut inner, path, with);
*plan = inner;
return;
}
if let Some(scan) = lines::LinesScan::of(plan, path) {
if let Some(replaced) = with(scan) {
*plan = replaced;
}
return;
}
crate::widgets::datatable::for_each_input(plan, &mut |input| {
replace_lines_scan(input, path, with)
});
}
pub(crate) struct Window {
pub(crate) lf: LazyFrame,
pub(crate) path: PathBuf,
pub(crate) marks: Arc<Marks>,
pub(crate) known: Option<Vec<(usize, usize)>>,
}
impl crate::pushdown::Windowed for Window {
fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
let read = match &self.known {
None => from_marks(&self.lf, &self.path, &self.marks, start, Some(start + len)),
Some(known) => {
let at = known.partition_point(|&(view, _)| view <= start);
known.get(at.wrapping_sub(1)).and_then(|&(view, row)| {
from_marks(&self.lf, &self.path, &self.marks, row, None)
.map(|lf| lf.slice((start - view) as i64, len as IdxSize))
})
}
};
Ok(read.unwrap_or_else(|| self.lf.clone().slice(start as i64, len as IdxSize)))
}
}
#[derive(Clone)]
pub enum Change {
Grew { rows: usize, misfits: usize },
Restarted { rows: usize, misfits: usize },
Gone { handle: Option<Arc<File>> },
NewFields(Vec<Field>),
Ended(Option<String>),
Failed(String),
}
#[derive(Clone)]
pub struct News {
pub(crate) id: u64,
pub(crate) change: Change,
}
#[derive(Default)]
struct Shared {
stop: AtomicBool,
poke: Mutex<bool>,
woken: Condvar,
#[cfg(target_os = "linux")]
bell: notify::Bell,
}
impl Shared {
fn wait(&self, interval: Duration) -> bool {
let mut poked = self.poke.lock().unwrap_or_else(|e| e.into_inner());
if !*poked && !self.stop.load(Ordering::Relaxed) {
poked = self
.woken
.wait_timeout(poked, interval)
.unwrap_or_else(|e| e.into_inner())
.0;
}
*poked = false;
!self.stop.load(Ordering::Relaxed)
}
fn wake(&self) {
*self.poke.lock().unwrap_or_else(|e| e.into_inner()) = true;
self.woken.notify_all();
#[cfg(target_os = "linux")]
self.bell.ring();
}
#[cfg(target_os = "linux")]
fn wait_for_change(
&self,
notify: ¬ify::Notify,
interval: Duration,
last: &mut Option<Instant>,
) -> bool {
loop {
if self.stop.load(Ordering::Relaxed) {
return false;
}
if std::mem::take(&mut *self.poke.lock().unwrap_or_else(|e| e.into_inner())) {
break;
}
if notify.wait(&self.bell, None) == notify::Woke::Changed {
let left = last
.map(|at| at + interval)
.and_then(|due| due.checked_duration_since(Instant::now()));
if let Some(left) = left
&& !self.wait(left)
{
return false;
}
break;
}
}
notify.drain();
*last = Some(Instant::now());
!self.stop.load(Ordering::Relaxed)
}
}
pub fn age(elapsed: Duration) -> String {
let secs = elapsed.as_secs();
match secs {
0..60 => format!("{secs:>2}s ago"),
60..3_600 => format!("{:>2}m ago", secs / 60),
3_600..86_400 => format!("{:>2}h ago", secs / 3_600),
_ => format!("{:>2}d ago", secs / 86_400),
}
}
pub fn next_tick(at: Instant) -> Instant {
let secs = at.elapsed().as_secs();
let unit = match secs {
0..60 => 1,
60..3_600 => 60,
3_600..86_400 => 3_600,
_ => 86_400,
};
at + Duration::from_secs((secs / unit + 1) * unit)
}
#[derive(Clone, Debug, PartialEq)]
pub enum Standing {
Following,
Paused,
Ended,
}
pub struct Follow {
id: u64,
path: PathBuf,
shared: Arc<Shared>,
spool: Option<Arc<SpoolHandle>>,
counted: usize,
misfits: usize,
restarted: bool,
shown: usize,
pub(crate) new_below: usize,
pub(crate) standing: Standing,
pub(crate) last_append: Option<Instant>,
pub(crate) settle_at_end: bool,
pub(crate) end_pending: bool,
pub(crate) stale_view: bool,
held: Option<Arc<File>>,
marks: Arc<Marks>,
pipe: bool,
new_fields: Vec<Field>,
pub(crate) described: bool,
}
static NEXT_ID: AtomicU64 = AtomicU64::new(1);
impl Follow {
pub fn start(
mut tail: Tail,
interval: Duration,
events: Sender<AppEvent>,
spool: Option<Arc<SpoolHandle>>,
) -> Follow {
let id = NEXT_ID.fetch_add(1, Ordering::Relaxed);
let path = tail.path.clone();
let shared = Arc::new(Shared::default());
let shown = tail.rows();
if let Some(handle) = &spool {
handle.spool.wake_on_end(shared.clone());
}
let marks = Arc::new(Marks::default());
marks.take_from(&mut tail);
let watcher = Watcher {
marks: marks.clone(),
id,
path: path.clone(),
tail,
shared: shared.clone(),
events,
spool: spool.as_ref().map(|handle| handle.spool.clone()),
interval,
};
let _ = std::thread::Builder::new()
.name("datui-follow".to_string())
.spawn(move || watcher.run());
Follow {
id,
path,
shared,
spool,
counted: shown,
misfits: 0,
restarted: false,
shown,
new_below: 0,
standing: Standing::Following,
last_append: None,
settle_at_end: true,
end_pending: false,
stale_view: false,
held: None,
marks,
pipe: false,
new_fields: Vec::new(),
described: false,
}
}
pub fn as_pipe(mut self) -> Follow {
self.pipe = true;
self.settle_at_end = false;
self
}
pub fn is_pipe(&self) -> bool {
self.pipe
}
pub fn live(&self) -> bool {
self.standing != Standing::Ended
}
pub(crate) fn marks(&self) -> &Arc<Marks> {
&self.marks
}
pub fn id(&self) -> u64 {
self.id
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn shown(&self) -> usize {
self.shown
}
pub fn waiting(&self) -> usize {
self.counted.saturating_sub(self.shown)
}
pub fn misfits(&self) -> usize {
self.misfits
}
pub fn standing(&self) -> &Standing {
&self.standing
}
pub fn new_below(&self) -> usize {
self.new_below
}
pub fn behind(&self) -> bool {
self.standing != Standing::Paused && (self.restarted || self.counted != self.shown)
}
pub fn spool(&self) -> Option<&Arc<Spool>> {
self.spool.as_ref().map(|handle| &handle.spool)
}
pub fn check_now(&self) {
self.shared.wake();
}
pub fn take(&mut self, change: &Change) -> Option<String> {
match change {
Change::Grew { rows, misfits } => {
if *rows > self.counted {
self.last_append = Some(Instant::now());
}
self.counted = *rows;
self.misfits = *misfits;
None
}
Change::Restarted { rows, misfits } => {
self.counted = *rows;
self.misfits = *misfits;
self.restarted = true;
self.last_append = Some(Instant::now());
Some("The file was truncated or replaced: reading it from the start".to_string())
}
Change::NewFields(fields) => {
self.new_fields = fields.clone();
None
}
Change::Gone { handle } => {
self.held = handle.clone();
self.end();
Some("The file was deleted: following stopped, the rows read stay".to_string())
}
Change::Ended(_) if self.spool().is_some_and(|s| s.tee().is_some()) => {
self.end();
None
}
Change::Ended(None) => {
self.end();
Some("Standard input ended".to_string())
}
Change::Ended(Some(reason)) | Change::Failed(reason) => {
self.end();
Some(reason.clone())
}
}
}
pub fn catch_up(&mut self) -> (usize, bool) {
self.shown = self.counted;
(self.shown, std::mem::take(&mut self.restarted))
}
pub(crate) fn take_new_fields(&mut self) -> Vec<Field> {
std::mem::take(&mut self.new_fields)
}
pub fn take_held(&mut self) -> Option<Arc<File>> {
self.held.take()
}
pub fn pause(&mut self) {
if self.standing == Standing::Following {
self.standing = Standing::Paused;
}
}
pub fn resume(&mut self) {
if self.standing == Standing::Paused {
self.standing = Standing::Following;
}
}
pub fn end(&mut self) {
self.standing = Standing::Ended;
self.shared.stop.store(true, Ordering::Relaxed);
self.shared.wake();
if let Some(spool) = self.spool.as_ref().filter(|s| s.spool.tee().is_none()) {
spool.spool.stop();
}
}
}
impl Drop for Follow {
fn drop(&mut self) {
self.shared.stop.store(true, Ordering::Relaxed);
self.shared.wake();
}
}
struct Watcher {
id: u64,
path: PathBuf,
tail: Tail,
shared: Arc<Shared>,
events: Sender<AppEvent>,
spool: Option<Arc<Spool>>,
interval: Duration,
marks: Arc<Marks>,
}
type Identity = (u64, u64);
#[cfg(unix)]
fn identity_of(file: &File) -> Option<Identity> {
use std::os::unix::fs::MetadataExt;
file.metadata().ok().map(|meta| (meta.dev(), meta.ino()))
}
#[cfg(windows)]
fn identity_of(file: &File) -> Option<Identity> {
use std::os::windows::io::AsRawHandle;
use windows_sys::Win32::Storage::FileSystem::{
BY_HANDLE_FILE_INFORMATION, GetFileInformationByHandle,
};
let mut info: BY_HANDLE_FILE_INFORMATION = unsafe { std::mem::zeroed() };
let ok = unsafe { GetFileInformationByHandle(file.as_raw_handle(), &mut info) };
(ok != 0).then(|| {
(
u64::from(info.dwVolumeSerialNumber),
u64::from(info.nFileIndexHigh) << 32 | u64::from(info.nFileIndexLow),
)
})
}
#[cfg(not(any(unix, windows)))]
fn identity_of(_file: &File) -> Option<Identity> {
None
}
#[cfg(unix)]
fn identity_at(_path: &Path, meta: &std::fs::Metadata) -> Option<Identity> {
use std::os::unix::fs::MetadataExt;
Some((meta.dev(), meta.ino()))
}
#[cfg(not(unix))]
fn identity_at(path: &Path, _meta: &std::fs::Metadata) -> Option<Identity> {
File::open(path).ok().as_ref().and_then(identity_of)
}
fn replaced(known: Option<Identity>, now: Option<Identity>) -> bool {
matches!((known, now), (Some(known), Some(now)) if known != now)
}
fn failed_message(path: Option<&Path>, doing: &str, e: &std::io::Error) -> String {
let what = format!(
"{doing}. {}",
crate::error_display::user_message_from_io(e, None)
);
match path {
Some(path) => crate::error_display::file_message(path, &what),
None => crate::error_display::sentence(&format!("standard input: {what}")),
}
}
impl Watcher {
fn failed(&self, doing: &str, e: &std::io::Error) -> String {
failed_message(
self.spool.is_none().then_some(self.path.as_path()),
doing,
e,
)
}
fn run(mut self) {
let mut file = match File::open(&self.path) {
Ok(file) => file,
Err(e) => {
self.send(Change::Failed(self.failed("following it stopped", &e)));
return;
}
};
let mut known = identity_of(&file);
let mut sent = (self.tail.rows(), 0usize);
#[cfg(target_os = "linux")]
let notify = notify::Notify::new(&self.path);
#[cfg(target_os = "linux")]
let mut last = None;
loop {
#[cfg(target_os = "linux")]
let go_on = match ¬ify {
Some(notify) => self
.shared
.wait_for_change(notify, self.interval, &mut last),
None => self.shared.wait(self.interval),
};
#[cfg(not(target_os = "linux"))]
let go_on = self.shared.wait(self.interval);
if !go_on {
return;
}
let spool_ended = self.spool.as_ref().and_then(|spool| spool.ended());
let meta = match std::fs::metadata(&self.path) {
Ok(meta) => meta,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
self.send(Change::Gone {
handle: Some(Arc::new(file)),
});
return;
}
Err(e) => {
self.send(Change::Failed(self.failed("following it stopped", &e)));
return;
}
};
let now = identity_at(&self.path, &meta);
let len = meta.len();
if replaced(known, now) || len < self.tail.complete() {
match File::open(&self.path) {
Ok(reopened) => file = reopened,
Err(e) => {
self.send(Change::Failed(self.failed("following it stopped", &e)));
return;
}
}
known = identity_of(&file);
#[cfg(target_os = "linux")]
if let Some(notify) = ¬ify {
notify.rewatch(&self.path);
}
self.tail.restart();
self.marks.clear();
if let Err(e) = self.tail.read_on(&mut file, len, false) {
self.send(Change::Failed(self.failed("reading it stopped", &e)));
return;
}
self.marks.take_from(&mut self.tail);
sent = (self.tail.rows(), self.tail.misfits());
self.send(Change::Restarted {
rows: sent.0,
misfits: sent.1,
});
continue;
}
let read = if spool_ended.is_some() {
self.tail.read_to_end(&mut file, len, true)
} else {
self.tail.read_on(&mut file, len, true)
};
if let Err(e) = read {
self.send(Change::Failed(self.failed("reading it stopped", &e)));
return;
}
self.marks.take_from(&mut self.tail);
let now_counted = (self.tail.rows(), self.tail.misfits());
if now_counted != sent {
sent = now_counted;
if !self.send(Change::Grew {
rows: sent.0,
misfits: sent.1,
}) {
return;
}
}
if let Some(reason) = spool_ended {
let fields = self.tail.new_fields();
if !fields.is_empty() && !self.send(Change::NewFields(fields)) {
return;
}
self.send(Change::Ended(reason));
return;
}
}
}
fn send(&self, change: Change) -> bool {
self.events
.send(AppEvent::Followed(News {
id: self.id,
change,
}))
.is_ok()
}
}
pub struct Spool {
stop: AtomicBool,
bytes: AtomicU64,
state: Mutex<SpoolState>,
changed: Condvar,
sink: Mutex<Option<File>>,
tee: Option<Tee>,
pass: Mutex<Option<Box<dyn Write + Send>>>,
started: Instant,
}
#[derive(Clone, Debug)]
pub struct Tee {
pub path: PathBuf,
pub raw: bool,
}
impl Tee {
pub fn to_stdout(&self) -> bool {
crate::stdin::is_stdin(&self.path)
}
pub fn name(&self) -> String {
if self.to_stdout() {
return "standard output".to_string();
}
self.path.file_name().map_or_else(
|| self.path.display().to_string(),
|name| name.to_string_lossy().into_owned(),
)
}
}
#[derive(Default)]
struct SpoolState {
lines: usize,
drained: bool,
ended: Option<Option<String>>,
finished: Option<Instant>,
samples: std::collections::VecDeque<(Instant, u64)>,
watcher: Option<Arc<Shared>>,
}
const RATE_WINDOW: Duration = Duration::from_secs(2);
impl Spool {
fn new(sink: File, tee: Option<Tee>, pass: Option<Box<dyn Write + Send>>) -> Spool {
Spool {
stop: AtomicBool::new(false),
bytes: AtomicU64::new(0),
state: Mutex::new(SpoolState::default()),
changed: Condvar::new(),
sink: Mutex::new(Some(sink)),
tee,
pass: Mutex::new(pass),
started: Instant::now(),
}
}
pub fn bytes(&self) -> u64 {
self.bytes.load(Ordering::Relaxed)
}
pub fn rate(&self) -> f64 {
let state = self.lock();
match (state.samples.front(), state.samples.back()) {
(Some((t0, b0)), Some((t1, b1))) if t1 > t0 => {
(b1 - b0) as f64 / t1.duration_since(*t0).as_secs_f64()
}
_ => 0.0,
}
}
pub fn duration(&self) -> Duration {
let finished = self.lock().finished;
finished
.unwrap_or_else(Instant::now)
.duration_since(self.started)
}
pub fn tee(&self) -> Option<&Tee> {
self.tee.as_ref()
}
pub fn stop(&self) {
self.stop.store(true, Ordering::Relaxed);
self.finish(None);
}
pub fn stopped(&self) -> bool {
self.stop.load(Ordering::Relaxed)
}
pub fn ended(&self) -> Option<Option<String>> {
self.lock().ended.clone()
}
pub fn live(&self) -> bool {
self.lock().ended.is_none()
}
pub fn wait(&self) {
let mut state = self.lock();
while state.ended.is_none() {
state = self.changed.wait(state).unwrap_or_else(|e| e.into_inner());
}
}
fn lock(&self) -> std::sync::MutexGuard<'_, SpoolState> {
self.state.lock().unwrap_or_else(|e| e.into_inner())
}
fn write(&self, bytes: &[u8]) -> Result<bool, String> {
let mut sink = self.sink.lock().unwrap_or_else(|e| e.into_inner());
let Some(file) = sink.as_mut() else {
return Ok(false);
};
file.write_all(bytes).map_err(|e| {
format!(
"Could not write {}: {e}",
self.tee
.as_ref()
.filter(|t| !t.to_stdout())
.map_or("what came in".to_string(), |t| t.path.display().to_string())
)
})?;
drop(sink);
if let Some(out) = self.pass.lock().unwrap_or_else(|e| e.into_inner()).as_mut() {
out.write_all(bytes)
.and_then(|()| out.flush())
.map_err(|e| format!("Could not write standard output: {e}"))?;
}
let total =
self.bytes.fetch_add(bytes.len() as u64, Ordering::Relaxed) + bytes.len() as u64;
let now = Instant::now();
let mut state = self.lock();
if state.lines < WANTED_LINES {
state.lines += bytes.iter().filter(|&&b| b == b'\n').count();
}
state.samples.push_back((now, total));
while state
.samples
.front()
.is_some_and(|(at, _)| now.duration_since(*at) > RATE_WINDOW)
&& state.samples.len() > 2
{
state.samples.pop_front();
}
drop(state);
self.changed.notify_all();
Ok(true)
}
fn finish(&self, reason: Option<String>) {
let file = self.sink.lock().unwrap_or_else(|e| e.into_inner()).take();
if let Ok(mut pass) = self.pass.try_lock() {
pass.take();
}
let mut reason = reason;
if let (Some(mut file), Some(tee)) = (file, self.tee.as_ref().filter(|t| !t.to_stdout())) {
let finished = (if tee.raw {
Ok(())
} else {
crate::tee::fix_wav_sizes(&mut file).map(|_| ())
})
.and_then(|()| file.sync_all());
if let Err(e) = finished
&& reason.is_none()
{
reason = Some(format!("Could not finish {}: {e}", tee.path.display()));
}
}
let mut state = self.lock();
if state.ended.is_none() {
state.ended = Some(reason);
state.finished = Some(Instant::now());
}
let watcher = state.watcher.take();
drop(state);
self.changed.notify_all();
if let Some(watcher) = watcher {
watcher.wake();
}
}
fn wake_on_end(&self, shared: Arc<Shared>) {
let mut state = self.lock();
if state.ended.is_some() {
drop(state);
shared.wake();
} else {
state.watcher = Some(shared);
}
}
}
pub struct SpoolHandle {
spool: Arc<Spool>,
}
impl SpoolHandle {
pub fn spool(&self) -> &Arc<Spool> {
&self.spool
}
}
impl Drop for SpoolHandle {
fn drop(&mut self) {
self.spool.stop();
}
}
fn copy_on(mut reader: impl Read + Send + 'static, spool: Arc<Spool>) {
let _ = std::thread::Builder::new()
.name("datui-spool".to_string())
.spawn(move || {
let mut buf = vec![0u8; CHUNK];
let reason = loop {
if spool.stopped() {
break None;
}
let n = match reader.read(&mut buf) {
Ok(0) => break None,
Ok(n) => n,
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => continue,
Err(e) => break Some(format!("Standard input failed: {e}")),
};
match spool.write(&buf[..n]) {
Ok(true) => {}
Ok(false) => break None,
Err(reason) => break Some(reason),
}
spool.lock().drained = n < buf.len();
};
spool.finish(reason);
});
}
const WANTED_LINES: usize = 1000;
pub enum Spooled {
Temp(TempDownload),
Kept(PathBuf),
}
pub(crate) fn spool<R: Read + Send + 'static>(
open: impl FnOnce() -> crate::download::Opened<R>,
options: OpenOptions,
writer: &Writer,
read: &AtomicU64,
stdout: Option<Box<dyn Write + Send>>,
) -> Result<(Spooled, OpenOptions), String> {
let tee = options.tee.clone().map(|path| Tee {
path,
raw: options.tee_raw,
});
let tee = tee.map(|tee| Tee {
raw: tee.raw || tee.to_stdout(),
..tee
});
let pass = match &tee {
Some(tee) if tee.to_stdout() => Some(stdout.ok_or_else(|| {
"--tee - passes the stream on to standard output, which only the datui command has."
.to_string()
})?),
_ => None,
};
let progressive = !options.follow && tee.is_none();
let (reader, _) = open().map_err(|e| format!("Could not read standard input: {e}"))?;
let (spooled, file) = match &tee {
Some(tee) if !tee.to_stdout() => {
let file = crate::tee::create(&tee.path, options.force)?;
(Spooled::Kept(tee.path.clone()), file)
}
_ => {
let dir = crate::stdin::spool_dir(&options);
let Some((named, claim)) = writer
.create(|| TempDownload::create(dir.as_deref(), None))
.map_err(|e| crate::error_display::user_message_from_report(&e, None))?
else {
return Err("Reading standard input was stopped.".to_string());
};
let file = named
.as_file()
.try_clone()
.map_err(|e| format!("Could not write what came in: {e}"))?;
(Spooled::Temp(TempDownload::held(named, Some(claim))), file)
}
};
let spooled_path = match &spooled {
Spooled::Temp(download) => download.path().to_path_buf(),
Spooled::Kept(path) => path.clone(),
};
let spool = Arc::new(Spool::new(file, tee, pass));
let handle = Arc::new(SpoolHandle {
spool: spool.clone(),
});
copy_on(reader, spool.clone());
let wanted = if options.has_header == Some(false) {
1
} else {
2
};
let mut state = spool.lock();
loop {
read.store(spool.bytes(), Ordering::Relaxed);
if writer.stopped() {
drop(state);
spool.stop();
return Err("Reading standard input was stopped.".to_string());
}
let enough = (options.follow || progressive)
&& (state.lines >= WANTED_LINES
|| (state.drained
&& (state.lines >= wanted || stream::begins_with_schema(&spooled_path))));
if enough || state.ended.is_some() {
break;
}
state = spool
.changed
.wait_timeout(state, Duration::from_millis(100))
.unwrap_or_else(|e| e.into_inner())
.0;
}
if let Some(Some(reason)) = &state.ended {
return Err(reason.clone());
}
drop(state);
read.store(spool.bytes(), Ordering::Relaxed);
let path = spooled_path;
let mut head = Vec::new();
File::open(&path)
.and_then(|f| f.take(4096).read_to_end(&mut head))
.map_err(|e| format!("Could not read standard input back: {e}"))?;
if head.is_empty() {
return Err("Nothing came in on standard input.".to_string());
}
let (format, compression, guessed) = crate::stdin::sniff_for(&head, &options);
let asked = options.clone();
let options = OpenOptions {
format_guessed: options.format.is_none() && guessed,
format: options.format.or(Some(format)),
compression: options.compression.or(compression),
..options
};
if progressive {
let read_on = followed_stream(&path, options.format, &options)
|| refusal(options.format, &options).is_none();
if read_on {
return Ok((
spooled,
OpenOptions {
spool: Some(handle),
follow: true,
pipe: true,
..options
},
));
}
while spool.ended().is_none() {
if writer.stopped() {
spool.stop();
return Err("Reading standard input was stopped.".to_string());
}
read.store(spool.bytes(), Ordering::Relaxed);
let state = spool.lock();
if state.ended.is_none() {
let _ = spool
.changed
.wait_timeout(state, Duration::from_millis(100))
.unwrap_or_else(|e| e.into_inner());
}
}
if let Some(Some(reason)) = spool.ended() {
return Err(reason);
}
read.store(spool.bytes(), Ordering::Relaxed);
let options = crate::stdin::described(&path, asked)?;
return Ok((spooled, options));
}
if options.follow
&& options.tee.is_none()
&& !followed_stream(&path, options.format, &options)
&& let Some(refusal) = refusal(options.format, &options)
{
return Err(refusal);
}
Ok((
spooled,
OpenOptions {
spool: Some(handle),
tee: None,
..options
},
))
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn errors_name_the_file() {
let path = Path::new("/data/app.log");
for (doing, e) in [
(
"following it stopped",
std::io::Error::from(std::io::ErrorKind::NotFound),
),
(
"reading it stopped",
std::io::Error::from(std::io::ErrorKind::PermissionDenied),
),
] {
let message = failed_message(Some(path), doing, &e);
eprintln!("{message}");
crate::readers::bad_input::assert_shape(&message, path);
assert!(message.contains(&doing[1..]), "{message}");
}
let e = std::io::Error::from(std::io::ErrorKind::UnexpectedEof);
let message = failed_message(None, "reading it stopped", &e);
assert_eq!(
message,
"Standard input: reading it stopped. Unexpected end of file."
);
}
fn tail_of(text: &[u8], format: FileFormat, options: &OpenOptions) -> Tail {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("t");
std::fs::write(&path, text).unwrap();
let mut tail = Tail::new(format, options, &Schema::default());
let mut file = File::open(&path).unwrap();
tail.read_on(&mut file, text.len() as u64, false).unwrap();
tail
}
fn polars_rows(text: &[u8], options: &OpenOptions) -> usize {
let complete = text.iter().rposition(|&b| b == b'\n').map_or(0, |i| i + 1);
let mut read = CsvReadOptions::default().with_has_header(options.has_header != Some(false));
read = read.map_parse_options(|p| {
p.with_comment_prefix(
options
.comment_char
.as_deref()
.map(polars::io::csv::read::CommentPrefix::new_from_str),
)
});
if let Some(skip) = options.skip_lines {
read.skip_lines = skip;
}
CsvReader::new(std::io::Cursor::new(text[..complete].to_vec()))
.with_options(read)
.finish()
.map(|df| df.height())
.unwrap_or(0)
}
#[test]
fn records_are_counted_as_polars_reads_them() {
let cases: [(&[u8], OpenOptions); 6] = [
(b"a,b\n1,2\n3,4\n", OpenOptions::default()),
(b"a,b\n1,2\n\n3,4\n5,", OpenOptions::default()),
(b"a,b\n1,\"x\ny\"\n3,4\n", OpenOptions::default()),
(b"a,b\r\n1,2\r\n3,4\r\n", OpenOptions::default()),
(
b"#c\na,b\n#x\n1,2\n3,4\n",
OpenOptions {
comment_char: Some("#".to_string()),
..Default::default()
},
),
(
b"junk\na,b\n1,2\n",
OpenOptions::default().with_skip_lines(1),
),
];
for (text, options) in cases {
let tail = tail_of(text, FileFormat::Csv, &options);
assert_eq!(
tail.rows(),
polars_rows(text, &options),
"{}",
String::from_utf8_lossy(text)
);
}
let lines = tail_of(
b"{\"a\":1}\n\n{\"a\":2}\n{\"a\":",
FileFormat::Jsonl,
&Default::default(),
);
assert_eq!(lines.rows(), 2);
assert_eq!(lines.complete(), 17);
}
#[cfg(target_os = "linux")]
#[test]
fn an_append_is_heard_of_without_a_check() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("log.csv");
std::fs::write(&path, "t\n1\n").unwrap();
let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
.finish()
.unwrap();
let (_, tail) =
bound_to_complete(scan, &path, FileFormat::Csv, &OpenOptions::default()).unwrap();
let (tx, rx) = std::sync::mpsc::channel();
let follow = Follow::start(tail, Duration::from_secs(3_600), tx, None);
let guard = Duration::from_secs(30);
let mut file = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
let mut rows = 1;
let deadline = Instant::now() + guard;
let news = loop {
assert!(Instant::now() < deadline, "the watcher never heard");
file.write_all(format!("{}\n", rows + 1).as_bytes())
.unwrap();
rows += 1;
if let Ok(AppEvent::Followed(news)) = rx.recv_timeout(Duration::from_millis(50)) {
break news;
}
};
assert!(matches!(news.change, Change::Grew { rows: 2.., .. }));
drop(follow);
}
#[test]
fn a_tail_reads_on_from_where_it_stopped() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("grow.csv");
let mut out = File::create(&path).unwrap();
out.write_all(b"t,n\n1.5,2\n2.5,").unwrap();
let schema = Schema::from_iter([
Field::new("t".into(), DataType::Float64),
Field::new("n".into(), DataType::Int64),
]);
let mut tail = Tail::new(FileFormat::Csv, &OpenOptions::default(), &schema);
let mut file = File::open(&path).unwrap();
let len = |p: &Path| std::fs::metadata(p).unwrap().len();
tail.read_on(&mut file, len(&path), true).unwrap();
assert_eq!((tail.rows(), tail.complete()), (1, 10));
out.write_all(b"3\nx,4\n4.5,5,6\n").unwrap();
tail.read_on(&mut file, len(&path), true).unwrap();
assert_eq!(tail.rows(), 4);
assert_eq!(
tail.misfits(),
2,
"a word for a number, and a field too many"
);
}
#[test]
fn moving_the_bound_reads_the_new_rows_through_the_view() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("grow.csv");
std::fs::write(&path, "a,b\n1,x\n2,y\n3,").unwrap();
let scan = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
.with_truncate_ragged_lines(true)
.with_ignore_errors(true)
.finish()
.unwrap();
let mut root = scan.clone();
bound(&mut root, &path, 2);
let mut view = root.clone().filter(col("a").gt(lit(1)));
view.collect_schema().unwrap();
assert_eq!(view.clone().collect().unwrap().height(), 1);
let mut out = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
out.write_all(b"z\n4,w\n5").unwrap();
bound(&mut view, &path, 4);
let df = view.clone().collect().unwrap();
assert_eq!(df.height(), 3, "{df}");
bound(&mut root, &path, 4);
assert_eq!(root.collect().unwrap().height(), 4, "the partial row waits");
}
fn marked(
text: &[u8],
format: FileFormat,
options: &OpenOptions,
every: u64,
scan: impl Fn(&Path) -> LazyFrame,
) -> (tempfile::TempDir, PathBuf, LazyFrame, Arc<Marks>, usize) {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("marked");
std::fs::write(&path, text).unwrap();
let mut lf = scan(&path);
let schema = lf.collect_schema().unwrap();
let mut tail = Tail::new(format, options, &schema);
tail.mark_every = (every, u64::MAX);
tail.read_on(&mut File::open(&path).unwrap(), text.len() as u64, false)
.unwrap();
let marks = Arc::new(Marks::default());
marks.take_from(&mut tail);
bound(&mut lf, &path, tail.rows());
(dir, path, lf, marks, tail.rows())
}
fn csv_scan(path: &Path, options: &OpenOptions) -> LazyFrame {
let mut reader = LazyCsvReader::new(PlRefPath::try_from_path(path).unwrap())
.with_ignore_errors(true)
.with_truncate_ragged_lines(true)
.with_has_header(options.has_header != Some(false))
.with_comment_prefix(options.comment_char.as_deref().map(PlSmallStr::from_str));
if let Some(skip) = options.skip_lines {
reader = reader.with_skip_lines(skip);
}
reader.finish().unwrap()
}
#[test]
fn a_window_from_a_mark_reads_what_a_read_from_the_start_does() {
let mut csv = b"skipped\nt,s,n\n".to_vec();
let mut crlf = b"t,s,n\r\n".to_vec();
let mut lines = Vec::new();
for i in 0..120 {
let row = match i % 9 {
0 => format!("{i},\"two\nlines\",{}\n", i * 2),
3 => "\n".to_string(),
5 => "# a comment\n".to_string(),
7 => format!("{i},x,oops\n"),
_ => format!("{i},s{i},{}\n", i * 2),
};
csv.extend(row.as_bytes());
crlf.extend(format!("{i},s{i},{}\r\n", i * 2).as_bytes());
lines.extend(format!("{{\"t\":{i},\"s\":\"s{i}\"}}\n").as_bytes());
if i % 4 == 1 {
lines.extend(b"\n \n");
}
}
csv.extend(b"999,partial");
let commented = OpenOptions {
comment_char: Some("#".to_string()),
..OpenOptions::default().with_skip_lines(1)
};
let (schema, batches) = stream_messages(&arrow_rows(0, 100), 3);
let mut arrows = schema;
batches.iter().for_each(|batch| arrows.extend(batch));
arrows.extend(&batches[0][..20]);
let cases: Vec<(&[u8], FileFormat, OpenOptions)> = vec![
(&csv, FileFormat::Csv, commented),
(&crlf, FileFormat::Csv, OpenOptions::default()),
(&lines, FileFormat::Jsonl, OpenOptions::default()),
(&arrows, FileFormat::Arrow, OpenOptions::default()),
];
for (text, format, options) in cases {
let scan = |path: &Path| match format {
FileFormat::Jsonl => scan_lines(path, &options, false, &mut Vec::new()).unwrap(),
FileFormat::Arrow => stream::scan(path).unwrap(),
_ => csv_scan(path, &options),
};
let (_dir, path, lf, marks, rows) = marked(text, format, &options, 7, scan);
let window = Window {
lf: lf.clone(),
path: path.clone(),
marks: marks.clone(),
known: None,
};
let whole = lf.clone().collect().unwrap();
assert_eq!(whole.height(), rows);
for start in (0..rows + 3).step_by(5) {
for len in [1, 6, 40] {
let read = crate::pushdown::Windowed::window(&window, start, len)
.unwrap()
.collect()
.unwrap();
let expected = whole.slice(start as i64, len);
assert!(
read.equals_missing(&expected),
"{format:?} rows {start}+{len}:\n{read:?}\n{expected:?}"
);
}
if start < rows {
assert!(
from_marks(&lf, &path, &marks, start, Some(start + 1)).is_some(),
"{format:?} row {start} is read from a mark"
);
}
}
}
}
fn arrow_rows(from: i64, n: i64) -> DataFrame {
df!(
"t" => (from..from + n).collect::<Vec<_>>(),
"s" => (from..from + n).map(|i| format!("s{i}")).collect::<Vec<_>>(),
"x" => (from..from + n).map(|i| i as f64 / 2.0).collect::<Vec<_>>(),
)
.unwrap()
}
#[test]
fn an_arrow_stream_is_counted_and_read_by_its_batches() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("live.arrows");
let (schema, batches) = stream_messages(&arrow_rows(0, 50), 4);
let mut head = schema.clone();
head.extend(&batches[0]);
head.extend(&batches[1][..batches[1].len() - 3]);
std::fs::write(&path, &head).unwrap();
let lf = stream::scan(&path).unwrap();
let (mut lf, mut tail) =
bound_to_complete(lf, &path, FileFormat::Arrow, &OpenOptions::default()).unwrap();
assert_eq!(tail.rows(), 4, "the second batch is not all there");
assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 4));
let mut rest = batches[1][batches[1].len() - 3..].to_vec();
batches[2..].iter().for_each(|batch| rest.extend(batch));
let mut file = std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap();
file.write_all(&rest).unwrap();
let size = file.metadata().unwrap().len();
tail.read_on(&mut File::open(&path).unwrap(), size, true)
.unwrap();
assert_eq!((tail.rows(), tail.misfits()), (50, 0));
bound(&mut lf, &path, tail.rows());
assert_eq!(lf.clone().collect().unwrap(), arrow_rows(0, 50));
let kept = lf
.clone()
.filter(col("t").gt_eq(lit(45)))
.select([col("s")])
.collect()
.unwrap();
assert_eq!(kept, arrow_rows(45, 5).select(["s"]).unwrap());
let count = lf.select([len()]).collect().unwrap();
assert_eq!(count.column("len").unwrap().u32().unwrap().get(0), Some(50));
let legacy = dir.path().join("legacy.arrows");
std::fs::write(
&legacy,
crate::ipc_stream::tests::stream(&arrow_rows(0, 10), None, true),
)
.unwrap();
let (lf, tail) = bound_to_complete(
stream::scan(&legacy).unwrap(),
&legacy,
FileFormat::Arrow,
&OpenOptions::default(),
)
.unwrap();
assert_eq!(tail.rows(), 10);
assert_eq!(lf.collect().unwrap(), arrow_rows(0, 10));
let cats = df!("c" => ["a", "b", "a"])
.unwrap()
.lazy()
.with_column(col("c").cast(DataType::from_categories(Categories::global())))
.collect()
.unwrap();
let (schema, batches) = stream_messages(&cats, 3);
let dictionary = dir.path().join("dict.arrows");
std::fs::write(&dictionary, [schema, batches.concat()].concat()).unwrap();
assert!(
stream::scan(&dictionary).is_err_and(|e| e.contains("dictionary")),
"refused"
);
}
#[test]
fn a_filtered_view_reads_and_counts_on_from_what_is_known() {
let mut text = b"t,n\n".to_vec();
for i in 0..300 {
text.extend(format!("{i},{}\n", i % 5).as_bytes());
}
let options = OpenOptions::default();
let (_dir, path, lf, marks, rows) = marked(&text, FileFormat::Csv, &options, 16, |p| {
csv_scan(p, &options)
});
let view = lf.filter(col("n").eq(lit(3)));
let whole = view.clone().collect().unwrap();
let known = (whole.column("t").unwrap().i64().unwrap().to_vec())
.into_iter()
.filter(|t| t.unwrap() < 200)
.count();
let rest = from_marks(&view, &path, &marks, 200, None).unwrap();
let after = rest.collect().unwrap().height();
assert_eq!(known + after, whole.height());
assert_eq!(rows, 300);
let window = Window {
lf: view.clone(),
path,
marks,
known: Some(vec![(0, 0), (known, 200)]),
};
for start in [0, 10, known - 1, known, known + 5, whole.height() - 3] {
let read = crate::pushdown::Windowed::window(&window, start, 4)
.unwrap()
.collect()
.unwrap();
assert!(
read.equals_missing(&whole.slice(start as i64, 4)),
"{start}"
);
}
}
#[test]
fn a_replaced_file_is_told_from_a_grown_one() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("log.csv");
std::fs::write(&path, "t\n1\n").unwrap();
let held = File::open(&path).unwrap();
let known = identity_of(&held);
assert!(cfg!(not(any(unix, windows))) || known.is_some());
let now = |path: &Path| identity_at(path, &std::fs::metadata(path).unwrap());
std::fs::OpenOptions::new()
.append(true)
.open(&path)
.unwrap()
.write_all(b"2\n")
.unwrap();
assert!(!replaced(known, now(&path)), "grown in place");
let other = dir.path().join("next.csv");
std::fs::write(&other, "t\n1\n2\n3\n").unwrap();
std::fs::rename(&other, &path).unwrap();
assert_eq!(
replaced(known, now(&path)),
known.is_some(),
"renamed over it"
);
assert!(!replaced(None, now(&path)), "unknown is no replacement");
drop(held);
}
#[cfg(unix)]
#[test]
fn a_deleted_file_reads_through_its_handle() {
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("gone.csv");
std::fs::write(&path, "a\n1\n2\n").unwrap();
let mut lf = LazyCsvReader::new(PlRefPath::try_from_path(&path).unwrap())
.finish()
.unwrap();
bound(&mut lf, &path, 2);
let mut view = lf.filter(col("a").gt(lit(0)));
view.collect_schema().unwrap();
let handle = File::open(&path).unwrap();
std::fs::remove_file(&path).unwrap();
read_through(&mut view, &path, &handle);
assert_eq!(view.collect().unwrap().height(), 2);
}
}