1use std::io::BufRead;
13use std::path::{Path, PathBuf};
14
15use polars::prelude::*;
16
17use crate::OpenOptions;
18
19pub const DEFAULT_HEADER_JOIN: &str = " ";
21
22pub use datui_cli::check_comment_char;
23
24const MAX_HEADER_LINE: u64 = 16 << 20;
27
28pub fn names_of(
37 lines: &[Vec<u8>],
38 read: &[usize],
39 rows: &[usize],
40 join: &str,
41 separator: u8,
42 comment: Option<&str>,
43) -> Vec<String> {
44 let mut columns: Vec<Vec<String>> = Vec::new();
45 for &row in rows {
46 let line = read
47 .iter()
48 .position(|&r| r == row)
49 .map_or(&[][..], |i| lines[i].as_slice());
50 for (i, field) in header_fields(line, row, separator, comment)
51 .into_iter()
52 .enumerate()
53 {
54 if columns.len() <= i {
55 columns.resize_with(i + 1, Vec::new);
56 }
57 if !field.is_empty() {
58 columns[i].push(field);
59 }
60 }
61 }
62 columns
63 .into_iter()
64 .map(|pieces| pieces.join(join))
65 .collect()
66}
67
68pub(crate) struct FileHead {
70 pub file: PathBuf,
71 pub names: Option<Vec<String>>,
73 pub units: Vec<(String, String)>,
75 pub window: Vec<Vec<String>>,
77 pub lossy: bool,
79}
80
81const POLARS_INFER_ROWS: usize = 100;
83
84pub(crate) fn head(
89 file: &Path,
90 options: &OpenOptions,
91 compression: Option<crate::CompressionFormat>,
92 with_window: bool,
93) -> color_eyre::Result<FileHead> {
94 if !with_window && options.header_rows().is_none() {
95 return Ok(FileHead::empty(file));
96 }
97 let source = crate::formats::readers::csv::text_source(file, compression)?;
98 head_of(source, file, options, with_window)
99}
100
101impl FileHead {
102 fn empty(file: &Path) -> FileHead {
103 FileHead {
104 file: file.to_path_buf(),
105 names: None,
106 units: Vec::new(),
107 window: Vec::new(),
108 lossy: false,
109 }
110 }
111}
112
113fn wanted_lines(options: &OpenOptions) -> Vec<usize> {
115 let mut wanted: Vec<usize> = options.header_rows().unwrap_or_default().to_vec();
116 if let Some(read) = &options.delimited {
117 wanted.extend(read.delimited().head_lines());
118 }
119 wanted.sort_unstable();
120 wanted.dedup();
121 wanted
122}
123
124pub(crate) fn head_of(
128 mut source: impl BufRead,
129 file: &Path,
130 options: &OpenOptions,
131 with_window: bool,
132) -> color_eyre::Result<FileHead> {
133 let separator = options.separator_or(b',');
134 let comment = options.comment_char.as_deref();
135 let rows = options.header_rows();
136 let wanted = wanted_lines(options);
137 let lines = named_lines(&mut source, &wanted)?;
138 let names = rows.map(|rows| {
139 names_of(
140 &lines,
141 &wanted,
142 rows,
143 &options.header_join,
144 separator,
145 comment,
146 )
147 });
148 let units = options.delimited.as_ref().map_or_else(Vec::new, |read| {
149 read.delimited()
150 .facts_of(&wanted, &lines, separator, &options.header_join)
151 .units
152 });
153 let mut lossy = lines.iter().any(|line| std::str::from_utf8(line).is_err());
154 let mut rows_seen = Vec::new();
155 if with_window {
156 let read = wanted.last().copied().unwrap_or(0);
159 skip_lines(
160 &mut source,
161 options.skip_lines.unwrap_or(0).saturating_sub(read),
162 )?;
163 if rows.is_none() {
164 window(&mut source, 1, separator, comment)?;
165 }
166 if let Some(n) = options.skip_rows {
167 window(&mut source, n, separator, comment)?;
168 }
169 let infer = options.infer_schema_length.unwrap_or(POLARS_INFER_ROWS);
170 let (seen, window_lossy) = window_of(&mut source, infer, separator, comment)?;
171 rows_seen = seen;
172 lossy |= window_lossy;
173 }
174 Ok(FileHead {
175 file: file.to_path_buf(),
176 names,
177 units,
178 window: rows_seen,
179 lossy,
180 })
181}
182
183pub fn named_lines(mut source: impl BufRead, rows: &[usize]) -> color_eyre::Result<Vec<Vec<u8>>> {
187 use std::io::Read;
188 let last = rows.iter().copied().max().unwrap_or(0);
189 let mut lines: Vec<Vec<u8>> = vec![Vec::new(); last];
191 let mut blank = true;
194 for (i, line) in lines.iter_mut().enumerate() {
195 let n = i + 1;
196 let read = if rows.contains(&n) {
197 let read = (&mut source)
198 .take(MAX_HEADER_LINE + 1)
199 .read_until(b'\n', line)?;
200 blank &= line.iter().all(u8::is_ascii_whitespace);
201 read
202 } else {
203 skip_line(&mut source, &mut blank)?
204 };
205 if read == 0 {
206 return Err(NoHeader { line: last, blank }.into());
207 }
208 if line.len() as u64 > MAX_HEADER_LINE {
209 return Err(color_eyre::eyre::eyre!(
210 "header line {n} is longer than {} MiB",
211 MAX_HEADER_LINE >> 20
212 ));
213 }
214 }
215 Ok(rows
216 .iter()
217 .map(|&row| {
218 row.checked_sub(1)
219 .and_then(|i| lines.get(i))
220 .cloned()
221 .unwrap_or_default()
222 })
223 .collect())
224}
225
226#[derive(Debug)]
229pub struct NoHeader {
230 pub line: usize,
231 pub blank: bool,
232}
233
234impl std::fmt::Display for NoHeader {
235 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
236 write!(f, "header line {} is past the end of the file", self.line)
237 }
238}
239
240impl std::error::Error for NoHeader {}
241
242pub fn is_blank_file(e: &color_eyre::Report) -> bool {
244 e.chain()
245 .any(|cause| cause.downcast_ref::<NoHeader>().is_some_and(|h| h.blank))
246}
247
248fn skip_line(source: &mut impl BufRead, blank: &mut bool) -> std::io::Result<usize> {
250 let mut read = 0;
251 loop {
252 let buf = source.fill_buf()?;
253 if buf.is_empty() {
254 return Ok(read);
255 }
256 let (used, done) = match memchr::memchr(b'\n', buf) {
257 Some(at) => (at + 1, true),
258 None => (buf.len(), false),
259 };
260 *blank &= buf[..used].iter().all(u8::is_ascii_whitespace);
261 source.consume(used);
262 read += used;
263 if done {
264 return Ok(read);
265 }
266 }
267}
268
269const MAX_WINDOW_BYTES: u64 = 1 << 20;
271
272pub fn window(
277 source: impl BufRead,
278 rows: usize,
279 separator: u8,
280 comment: Option<&str>,
281) -> std::io::Result<Vec<Vec<String>>> {
282 Ok(window_of(source, rows, separator, comment)?.0)
283}
284
285pub fn window_of(
287 source: impl BufRead,
288 rows: usize,
289 separator: u8,
290 comment: Option<&str>,
291) -> std::io::Result<(Vec<Vec<String>>, bool)> {
292 let mut lossy = false;
293 let mut source = source.take(MAX_WINDOW_BYTES);
294 let comment = comment.filter(|c| !c.is_empty()).map(str::as_bytes);
295 let mut out = Vec::new();
296 let mut line = Vec::new();
297 while out.len() < rows {
298 line.clear();
299 if source.read_until(b'\n', &mut line)? == 0 {
300 break;
301 }
302 if !line.ends_with(b"\n") && source.limit() == 0 {
304 break;
305 }
306 let text = line.strip_suffix(b"\n").unwrap_or(&line);
307 let text = text.strip_suffix(b"\r").unwrap_or(text);
308 if text.iter().all(u8::is_ascii_whitespace) || comment.is_some_and(|c| text.starts_with(c))
309 {
310 continue;
311 }
312 lossy |= std::str::from_utf8(text).is_err();
313 out.push(
314 split_fields(text, separator)
315 .into_iter()
316 .map(|f| f.trim().to_string())
317 .collect(),
318 );
319 }
320 Ok((out, lossy))
321}
322
323pub fn skip_lines(source: &mut impl BufRead, n: usize) -> std::io::Result<()> {
325 let mut blank = true;
326 for _ in 0..n {
327 if skip_line(source, &mut blank)? == 0 {
328 break;
329 }
330 }
331 Ok(())
332}
333
334pub fn header_fields(line: &[u8], row: usize, separator: u8, comment: Option<&str>) -> Vec<String> {
337 let mut line = line;
338 if row == 1 {
339 line = line.strip_prefix(b"\xEF\xBB\xBF").unwrap_or(line);
340 }
341 line = line.strip_suffix(b"\n").unwrap_or(line);
342 line = line.strip_suffix(b"\r").unwrap_or(line);
343 if let Some(prefix) = comment.filter(|c| !c.is_empty()) {
344 line = line.strip_prefix(prefix.as_bytes()).unwrap_or(line);
345 }
346 split_fields(line, separator)
347 .into_iter()
348 .map(|f| f.trim().to_string())
349 .collect()
350}
351
352fn split_fields(line: &[u8], separator: u8) -> Vec<String> {
356 let mut fields = Vec::new();
357 let mut field: Vec<u8> = Vec::new();
358 let mut quoted = false;
359 let mut i = 0;
360 while i < line.len() {
361 let b = line[i];
362 if quoted {
363 if b == b'"' {
364 if line.get(i + 1) == Some(&b'"') {
365 field.push(b'"');
366 i += 1;
367 } else {
368 quoted = false;
369 }
370 } else {
371 field.push(b);
372 }
373 } else if b == separator {
374 fields.push(String::from_utf8_lossy(&field).into_owned());
375 field.clear();
376 } else if b == b'"' && field.iter().all(u8::is_ascii_whitespace) {
377 field.clear();
378 quoted = true;
379 } else {
380 field.push(b);
381 }
382 i += 1;
383 }
384 fields.push(String::from_utf8_lossy(&field).into_owned());
385 fields
386}
387
388pub fn shown_names(raw: &[PlSmallStr], header: Option<&[String]>) -> Vec<String> {
393 let names: Vec<String> = raw
394 .iter()
395 .enumerate()
396 .map(|(i, name)| {
397 let name = match header {
398 Some(header) => header.get(i).map_or("", String::as_str),
399 None => name.as_str(),
400 }
401 .trim();
402 if name.is_empty() {
403 format!("column_{}", i + 1)
404 } else {
405 name.to_string()
406 }
407 })
408 .collect();
409 let mut taken: PlHashSet<String> = PlHashSet::with_capacity(names.len());
410 let mut seen: PlHashMap<String, usize> = PlHashMap::with_capacity(names.len());
411 let mut out = Vec::with_capacity(names.len());
412 for name in names {
413 let count = seen.entry(name.clone()).or_insert(0);
414 let mut candidate = name.clone();
415 while !taken.insert(candidate.clone()) {
416 candidate = format!("{name}_duplicated_{count}");
417 *count += 1;
418 }
419 out.push(candidate);
420 }
421 out
422}
423
424pub fn name_columns(mut lf: LazyFrame, header: Option<&[String]>) -> PolarsResult<LazyFrame> {
431 let schema = match (lf.collect_schema(), header) {
432 (Err(PolarsError::NoData(_)), Some(header)) => return header_only(header),
433 (Ok(schema), Some(header)) if schema.is_empty() => return header_only(header),
434 (schema, _) => schema?,
435 };
436 let raw: Vec<PlSmallStr> = schema.iter_names().cloned().collect();
437 let shown = shown_names(&raw, header);
438 if raw.iter().zip(&shown).all(|(r, s)| r.as_str() == s) {
439 return Ok(lf);
440 }
441 Ok(lf.rename(raw.iter().map(|s| s.as_str()), shown.iter(), true))
442}
443
444fn header_only(header: &[String]) -> PolarsResult<LazyFrame> {
446 let raw: Vec<PlSmallStr> = (1..=header.len().max(1))
447 .map(|i| format!("column_{i}").into())
448 .collect();
449 let columns: Vec<Column> = shown_names(&raw, Some(header))
450 .into_iter()
451 .map(|name| Column::new_empty(name.into(), &DataType::String))
452 .collect();
453 Ok(DataFrame::new(0, columns)?.lazy())
454}
455
456pub fn read_after_header(
459 read: PolarsResult<DataFrame>,
460 header: Option<&[String]>,
461) -> PolarsResult<DataFrame> {
462 match read {
463 Err(PolarsError::NoData(_)) if header.is_some() => Ok(DataFrame::empty()),
464 read => read,
465 }
466}
467
468pub fn skip_initial_space(
473 mut lf: LazyFrame,
474 nulls: impl Fn(&str) -> Vec<String>,
475) -> PolarsResult<LazyFrame> {
476 let schema = lf.collect_schema()?;
477 let exprs: Vec<Expr> = schema
478 .iter()
479 .filter(|(_, dtype)| **dtype == DataType::String)
480 .map(|(name, _)| {
481 let stripped = col(name.clone())
482 .str()
483 .strip_chars_start(lit(PlSmallStr::from_static(" ")));
484 let null = nulls(name.as_str())
485 .into_iter()
486 .fold(stripped.clone().eq(lit("")), |any, value| {
487 any.or(stripped.clone().eq(lit(value)))
488 });
489 when(null)
490 .then(Null {}.lit().cast(DataType::String))
491 .otherwise(stripped)
492 .alias(name.clone())
493 })
494 .collect();
495 if exprs.is_empty() {
496 return Ok(lf);
497 }
498 Ok(lf.with_columns(exprs))
499}
500
501#[cfg(test)]
502mod tests {
503 use super::*;
504
505 fn header_names(
506 source: &[u8],
507 rows: &[usize],
508 join: &str,
509 separator: u8,
510 comment: Option<&str>,
511 ) -> color_eyre::Result<Vec<String>> {
512 let lines = named_lines(source, rows)?;
513 Ok(names_of(&lines, rows, rows, join, separator, comment))
514 }
515
516 fn names(text: &str, rows: &[usize], comment: Option<&str>) -> Vec<String> {
517 header_names(text.as_bytes(), rows, " ", b',', comment).unwrap()
518 }
519
520 fn spec_options(lines: &str) -> OpenOptions {
523 use crate::formats::delimited_spec::DelimitedRead;
524 let text = format!("name = \"a.log\"\nkind = \"delimited\"\n{lines}");
525 let spec = std::sync::Arc::new(crate::formats::Spec::parse(&text, None).unwrap());
526 let read = DelimitedRead::chosen(spec, crate::formats::Chosen::SpecFile, Vec::new());
527 let mut options = OpenOptions {
528 parse_strings: Some(crate::ParseStringsTarget::All),
529 ..OpenOptions::default()
530 };
531 read.delimited().apply(&mut options);
532 options.delimited = Some(std::sync::Arc::new(read));
533 options
534 }
535
536 #[test]
540 fn the_window_starts_below_the_spec_s_unit_and_metadata_lines() {
541 for (spec, text) in [
542 (
543 "header_rows = { name = 3, unit = 2 }\nmetadata_line = 1",
544 "device=\"x\"\n007,m\nid,len\n1,2\n3,4\n",
545 ),
546 (
547 "header_rows = { name = 1, unit = 2 }",
548 "id,len\n007,m\n1,2\n3,4\n",
549 ),
550 ] {
551 let options = spec_options(spec);
552 let head = head_of(text.as_bytes(), Path::new("a.log"), &options, true).unwrap();
553 assert_eq!(head.names.unwrap(), ["id", "len"], "{spec}");
554 assert_eq!(head.window, [["1", "2"], ["3", "4"]], "{spec}");
555 }
556 }
557
558 #[test]
561 fn a_head_without_names_or_window_opens_nothing() {
562 let missing = Path::new("/nonexistent/datui/a.log");
563 let read = head(
564 missing,
565 &spec_options("metadata_line = 1\nskip_lines = 1"),
566 None,
567 false,
568 )
569 .unwrap();
570 assert!(read.names.is_none() && read.window.is_empty());
571 assert!(head(missing, &OpenOptions::default(), None, false).is_ok());
572 let named = OpenOptions {
573 header_rows: vec![1],
574 ..OpenOptions::default()
575 };
576 assert!(
577 head(missing, &named, None, false).is_err(),
578 "names are read"
579 );
580 }
581
582 #[test]
583 fn one_header_line_is_split_and_trimmed() {
584 let text = "#info\n Lcl Date, Lcl Time, Latitude\n1,2,3\n";
585 assert_eq!(
586 names(text, &[2], None),
587 ["Lcl Date", "Lcl Time", "Latitude"]
588 );
589 }
590
591 #[test]
592 fn several_lines_join_in_the_order_given_and_skip_blank_pieces() {
593 let text = "#yyyy-mm-dd, hh:mm:ss, degrees\n Lcl Date, Lcl Time, Latitude\n";
594 assert_eq!(
595 names(text, &[2, 1], Some("#")),
596 [
597 "Lcl Date yyyy-mm-dd",
598 "Lcl Time hh:mm:ss",
599 "Latitude degrees"
600 ]
601 );
602 let text = "a,,c\nx,y\n";
603 assert_eq!(
604 header_names(text.as_bytes(), &[1, 2], "_", b',', None).unwrap(),
605 ["a_x", "y", "c"]
606 );
607 }
608
609 #[test]
610 fn quotes_bom_and_carriage_returns() {
611 let text = "\u{FEFF}id, \"last, first\",\"say \"\"hi\"\"\"\r\n";
612 assert_eq!(names(text, &[1], None), ["id", "last, first", "say \"hi\""]);
613 }
614
615 #[test]
616 fn a_file_that_ends_before_the_header_is_an_error() {
617 let err = header_names("a,b\n".as_bytes(), &[1, 5], " ", b',', None).unwrap_err();
618 assert!(err.to_string().contains("past the end"), "{err}");
619 assert!(header_names("".as_bytes(), &[1], " ", b',', None).is_err());
620 let blank = |text: &str| is_blank_file(&named_lines(text.as_bytes(), &[3]).unwrap_err());
621 assert!(blank(""), "empty");
622 assert!(blank(" \n\t\r\n"), "white space");
623 assert!(!blank("#a\n"), "text, too short");
624 assert!(!blank("\nx\n"), "text on a line passed over");
625 assert_eq!(names("#u\na,b", &[2], None), ["a", "b"]);
627 }
628
629 #[test]
630 fn the_window_is_the_data_lines_after_the_header() {
631 let text = "a,b\n 1, x\n#note\n\n , 2.5\n3,4\n";
632 let mut source = text.as_bytes();
633 skip_lines(&mut source, 1).unwrap();
634 let rows = window(source, 2, b',', Some("#")).unwrap();
635 assert_eq!(rows, [vec!["1", "x"], vec!["", "2.5"]]);
636 }
637
638 #[test]
639 fn a_header_line_is_read_up_to_a_bound() {
640 let wide = "x".repeat(MAX_HEADER_LINE as usize + 1);
641 let err = header_names(wide.as_bytes(), &[1], " ", b',', None).unwrap_err();
642 assert!(err.to_string().contains("header line 1"), "{err}");
643 let text = format!("{wide}\na,b\n");
645 assert_eq!(names(&text, &[2], None), ["a", "b"]);
646 }
647
648 #[test]
649 fn shown_names_trim_fill_and_deduplicate() {
650 let raw: Vec<PlSmallStr> = [" a", "a", " ", "b"].map(PlSmallStr::from).to_vec();
651 assert_eq!(
652 shown_names(&raw, None),
653 ["a", "a_duplicated_0", "column_3", "b"]
654 );
655 let raw: Vec<PlSmallStr> = (1..=4).map(|i| format!("column_{i}").into()).collect();
656 let header = ["x".to_string(), String::new(), "column_1".to_string()];
657 assert_eq!(
658 shown_names(&raw, Some(&header)),
659 ["x", "column_2", "column_1", "column_4"]
660 );
661 }
662
663 #[test]
664 fn padding_is_skipped_and_blank_or_null_values_are_null() {
665 let df = df!(
666 "a" => [" 1.5", " ", " NA", " x y "],
667 "n" => [1i64, 2, 3, 4],
668 )
669 .unwrap();
670 let out = skip_initial_space(df.lazy(), |_| vec!["NA".into()])
671 .unwrap()
672 .collect()
673 .unwrap();
674 let a: Vec<Option<&str>> = out.column("a").unwrap().str().unwrap().iter().collect();
675 assert_eq!(a, [Some("1.5"), None, None, Some("x y ")]);
676 assert_eq!(out.column("n").unwrap().dtype(), &DataType::Int64);
677 }
678}