1use std::path::{Path, PathBuf};
16use std::sync::{Arc, RwLock};
17
18use color_eyre::Result;
19use polars::prelude::*;
20
21use crate::FileFormat;
22use crate::formats::fixed_records::Bytes;
23use crate::formats::indexed::Offsets;
24
25pub const FILE: &str = "file";
27pub const LINE: &str = "line";
28
29pub const FIRST_BYTES: usize = 8 << 20;
32
33const CHUNK: usize = 4 << 20;
36
37pub(crate) const READER: crate::formats::readers::Reader = crate::formats::readers::Reader {
39 scan,
40 preview: Some(crate::formats::readers::Preview::Scan),
41 python: Some(crate::export::python_script::Python {
42 call: "pl.LazyFrame",
43 eager: false,
44 glob_flag: false,
45 arguments: Some(crate::export::python_script::lines_arguments),
46 }),
47 export: Some(crate::export::export_modal::ExportFormat::Csv),
48 ..crate::formats::readers::BASE
49};
50
51pub fn schema(several: bool) -> Schema {
53 let mut fields = Vec::with_capacity(2);
54 if several {
55 fields.push(Field::new(FILE.into(), DataType::String));
56 }
57 fields.push(Field::new(LINE.into(), DataType::String));
58 Schema::from_iter(fields)
59}
60
61#[derive(Debug, Clone, Default)]
64pub struct LineIndex {
65 offsets: Offsets,
66 end: usize,
68 partial: Option<(bool, bool)>,
71 pub invalid: usize,
73 pub past_limit: usize,
75 before: usize,
78}
79
80impl LineIndex {
81 pub fn of(bytes: &[u8]) -> Self {
83 let mut index = Self {
84 offsets: Offsets::for_file(bytes.len()),
85 ..Default::default()
86 };
87 index.extend(bytes);
88 index
89 }
90
91 pub fn lines(&self) -> usize {
93 self.offsets.len()
94 }
95
96 pub fn complete(&self) -> usize {
98 self.lines() - usize::from(self.partial.is_some_and(|(indexed, _)| indexed))
99 }
100
101 pub fn end(&self) -> usize {
103 self.end
104 }
105
106 pub fn whole(&self, bytes: &[u8]) -> bool {
108 self.end >= bytes.len() || self.partial.is_some()
109 }
110
111 pub fn extend(&mut self, bytes: &[u8]) {
114 self.extend_by(bytes, usize::MAX);
115 }
116
117 pub fn extend_by(&mut self, bytes: &[u8], budget: usize) -> bool {
120 if let Some((indexed, invalid)) = self.partial.take() {
121 if indexed {
122 self.offsets.truncate(self.offsets.len() - 1);
123 } else {
124 self.past_limit -= 1;
125 }
126 self.invalid -= usize::from(invalid);
127 }
128 if bytes.len() <= self.end {
129 return true;
130 }
131 self.offsets.widen_for(bytes.len());
132 let stop = self.end.saturating_add(budget).min(bytes.len());
133 while self.end < stop && self.partial.is_none() {
134 let cut = (self.end + CHUNK).min(stop);
136 let to = if cut < bytes.len() {
137 memchr::memchr(b'\n', &bytes[cut - 1..]).map_or(bytes.len(), |i| cut + i)
138 } else {
139 cut
140 };
141 self.index_chunk(bytes, to);
142 }
143 self.whole(bytes)
144 }
145
146 fn index_chunk(&mut self, bytes: &[u8], to: usize) {
148 let all_utf8 = std::str::from_utf8(&bytes[self.end..to]).is_ok();
150 let limit = crate::limits::get().indexed_records;
151 let mut at = self.end;
152 while at < to {
153 let newline = memchr::memchr(b'\n', &bytes[at..to]).map(|i| at + i);
154 let end = newline.unwrap_or(to);
155 let indexed = self.before + self.offsets.len() < limit;
156 if indexed {
157 self.offsets.push(at);
158 } else {
159 self.past_limit += 1;
160 }
161 let invalid = !all_utf8 && std::str::from_utf8(&bytes[at..end]).is_err();
162 self.invalid += usize::from(invalid);
163 match newline {
164 Some(newline) => {
165 at = newline + 1;
166 self.end = at;
167 }
168 None => {
169 self.partial = Some((indexed, invalid));
170 at = to;
171 }
172 }
173 }
174 }
175
176 fn step_after(&self, len: usize) -> Self {
180 Self {
181 offsets: Offsets::for_file(len),
182 end: self.end,
183 before: self.before + self.offsets.len(),
184 ..Default::default()
185 }
186 }
187
188 fn take_step(&mut self, step: LineIndex) {
190 debug_assert_eq!(step.before, self.offsets.len());
191 self.offsets.append(&step.offsets);
192 self.end = step.end;
193 self.partial = step.partial;
194 self.invalid += step.invalid;
195 self.past_limit += step.past_limit;
196 }
197
198 fn line<'a>(&self, bytes: &'a [u8], i: usize) -> &'a [u8] {
200 let start = self.offsets.get(i);
201 let end = if i + 1 < self.offsets.len() {
202 self.offsets.get(i + 1) - 1
203 } else {
204 memchr::memchr(b'\n', &bytes[start..]).map_or(bytes.len(), |n| start + n)
205 };
206 let line = &bytes[start..end.min(bytes.len())];
207 line.strip_suffix(b"\r").unwrap_or(line)
208 }
209}
210
211struct Mapped {
213 bytes: Arc<Bytes>,
214 index: LineIndex,
215}
216
217struct LineFile {
218 name: String,
220 grows: Option<PathBuf>,
222 mapped: RwLock<Mapped>,
223}
224
225impl LineFile {
226 fn rows(&self) -> usize {
227 self.mapped
228 .read()
229 .unwrap_or_else(|e| e.into_inner())
230 .index
231 .lines()
232 }
233
234 fn grow(&self) -> PolarsResult<()> {
237 let Some(path) = &self.grows else {
238 return Ok(());
239 };
240 let mut mapped = self.mapped.write().unwrap_or_else(|e| e.into_inner());
241 let bytes = match Bytes::map(path) {
243 Ok(bytes) => bytes,
244 Err(_) => match mapped.bytes.as_ref() {
245 Bytes::Mapped(_, file) => remap(file)?,
246 Bytes::Owned(_) => return Ok(()),
247 },
248 };
249 if bytes.len() < mapped.index.end() {
250 mapped.index = LineIndex::default();
251 }
252 mapped.index.extend(bytes.as_slice());
253 mapped.bytes = Arc::new(bytes);
254 Ok(())
255 }
256}
257
258fn remap(file: &std::fs::File) -> PolarsResult<Bytes> {
259 let file = file.try_clone()?;
260 if file.metadata()?.len() == 0 {
261 return Ok(Bytes::Owned(Vec::new()));
262 }
263 let map = unsafe { memmap2::Mmap::map(&file)? };
265 Ok(Bytes::Mapped(map, file))
266}
267
268pub struct Lines {
270 files: Vec<LineFile>,
271 schema: SchemaRef,
272 indexing: (std::sync::Mutex<bool>, std::sync::Condvar),
275 stepping: std::sync::Mutex<()>,
278 shrank: std::sync::atomic::AtomicBool,
280}
281
282pub const SHRANK: &str = "the file shrank while it was indexed; open it again to read it";
284
285impl Lines {
286 pub fn from_bytes(parts: Vec<(String, Arc<Bytes>)>) -> Self {
289 Self::from_bytes_first(parts, usize::MAX)
290 }
291
292 pub fn from_bytes_first(parts: Vec<(String, Arc<Bytes>)>, first: usize) -> Self {
296 let several = parts.len() > 1;
297 let mut indexing = false;
298 let files = parts
299 .into_iter()
300 .map(|(name, bytes)| {
301 let mut index = LineIndex {
302 offsets: Offsets::for_file(bytes.len()),
303 ..Default::default()
304 };
305 let budget = if several { usize::MAX } else { first };
306 indexing |= !index.extend_by(bytes.as_slice(), budget);
307 LineFile {
308 name,
309 grows: None,
310 mapped: RwLock::new(Mapped { index, bytes }),
311 }
312 })
313 .collect();
314 Self {
315 files,
316 schema: Arc::new(schema(several)),
317 indexing: (indexing.into(), Default::default()),
318 stepping: Default::default(),
319 shrank: Default::default(),
320 }
321 }
322
323 pub fn whole(&self) -> bool {
325 self.files.iter().all(|f| {
326 let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
327 m.index.whole(m.bytes.as_slice())
328 })
329 }
330
331 pub fn shrank(&self) -> bool {
333 self.shrank.load(std::sync::atomic::Ordering::Relaxed)
334 }
335
336 pub fn resume_indexing(&self) -> bool {
339 let more = !self.whole() && !self.shrank();
340 *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner()) = more;
341 more
342 }
343
344 pub fn line_in_file(&self, place: usize) -> Option<usize> {
347 let mut start = 0;
348 for f in &self.files {
349 let rows = f.rows();
350 if place < start + rows {
351 return Some(place - start);
352 }
353 start += rows;
354 }
355 None
356 }
357
358 pub fn several(&self) -> bool {
360 self.files.len() > 1
361 }
362
363 pub fn indexing(&self) -> bool {
365 *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner())
366 }
367
368 pub fn stop_indexing(&self) {
371 *self.indexing.0.lock().unwrap_or_else(|e| e.into_inner()) = false;
372 self.indexing.1.notify_all();
373 }
374
375 fn wait_indexed(&self) -> PolarsResult<()> {
379 let mut indexing = self.indexing.0.lock().unwrap_or_else(|e| e.into_inner());
380 while *indexing {
381 polars_ensure!(
382 !crate::app::jobs::superseded(),
383 ComputeError: "the read is no longer wanted"
384 );
385 indexing = self
386 .indexing
387 .1
388 .wait_timeout(indexing, std::time::Duration::from_millis(100))
389 .unwrap_or_else(|e| e.into_inner())
390 .0;
391 }
392 drop(indexing);
393 polars_ensure!(!self.shrank(), ComputeError: "{SHRANK}");
394 polars_ensure!(
395 self.whole(),
396 ComputeError: "the file's lines were not all indexed; open it again"
397 );
398 Ok(())
399 }
400
401 pub fn index_more(&self, budget: usize) -> bool {
405 if !self.indexing() {
406 return true;
407 }
408 let _step = self.stepping.lock().unwrap_or_else(|e| e.into_inner());
409 let whole = self.files.iter().all(|f| {
410 let (bytes, mut step) = {
413 let mapped = f.mapped.read().unwrap_or_else(|e| e.into_inner());
414 if mapped.index.whole(mapped.bytes.as_slice()) {
415 return true;
416 }
417 if mapped.bytes.still_whole().is_err() {
420 self.shrank
421 .store(true, std::sync::atomic::Ordering::Relaxed);
422 return true;
423 }
424 let bytes = mapped.bytes.clone();
425 let step = mapped.index.step_after(bytes.len());
426 (bytes, step)
427 };
428 let whole = step.extend_by(bytes.as_slice(), budget);
429 f.mapped
430 .write()
431 .unwrap_or_else(|e| e.into_inner())
432 .index
433 .take_step(step);
434 whole
435 });
436 if whole {
437 self.stop_indexing();
438 }
439 whole
440 }
441
442 pub fn indexed_bytes(&self) -> (u64, u64) {
444 self.files.iter().fold((0, 0), |(done, all), f| {
445 let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
446 let len = m.bytes.as_slice().len();
447 let indexed = if m.index.whole(m.bytes.as_slice()) {
448 len
449 } else {
450 m.index.end()
451 };
452 (done + indexed as u64, all + len as u64)
453 })
454 }
455
456 pub fn open(paths: &[PathBuf], follow: bool) -> Result<Self> {
458 Self::open_first(paths, follow, usize::MAX)
459 }
460
461 pub fn open_first(paths: &[PathBuf], follow: bool, first: usize) -> Result<Self> {
463 let parts = paths
464 .iter()
465 .map(|path| {
466 let bytes =
467 Bytes::map(path).map_err(|e| crate::error_display::in_file(path, e.into()))?;
468 Ok((file_name(path), Arc::new(bytes)))
469 })
470 .collect::<Result<Vec<_>>>()?;
471 let first = if follow { usize::MAX } else { first };
472 let mut lines = Self::from_bytes_first(parts, first);
473 if follow && let ([file], [path]) = (lines.files.as_mut_slice(), paths) {
474 file.grows = Some(path.clone());
475 }
476 Ok(lines)
477 }
478
479 pub fn grows(&self) -> bool {
481 self.files.iter().any(|f| f.grows.is_some())
482 }
483
484 pub fn rows(&self) -> usize {
486 self.files
487 .iter()
488 .map(LineFile::rows)
489 .sum::<usize>()
490 .min(crate::formats::row_index::MAX_ROWS)
491 }
492
493 fn complete_rows(&self) -> usize {
495 self.files
496 .iter()
497 .map(|f| {
498 f.mapped
499 .read()
500 .unwrap_or_else(|e| e.into_inner())
501 .index
502 .complete()
503 })
504 .sum()
505 }
506
507 fn counts(&self) -> (usize, usize) {
510 self.files.iter().fold((0, 0), |(invalid, past), f| {
511 let m = f.mapped.read().unwrap_or_else(|e| e.into_inner());
512 (invalid + m.index.invalid, past + m.index.past_limit)
513 })
514 }
515
516 pub fn lazy(self: &Arc<Self>) -> LazyFrame {
519 if self.indexing() {
520 let lines = self.clone();
524 let height = DataFrame::empty_with_height(0).lazy().map(
525 move |_| {
526 lines.wait_indexed()?;
527 Ok(DataFrame::empty_with_height(lines.rows()))
528 },
529 AllowedOptimizations::empty(),
530 None,
531 Some("every line"),
532 );
533 return crate::formats::row_index::lazy_numbered_over(self, height);
534 }
535 let height = if self.grows() {
536 self.complete_rows()
537 } else {
538 self.rows()
539 };
540 crate::formats::row_index::lazy_numbered(self, height)
541 }
542
543 fn place(&self, rows: &[IdxSize]) -> PolarsResult<Vec<(usize, usize)>> {
545 let mut starts = Vec::with_capacity(self.files.len());
546 let mut total = 0usize;
547 for f in &self.files {
548 starts.push(total);
549 total += f.rows();
550 }
551 rows.iter()
552 .map(|&r| {
553 let r = r as usize;
554 polars_ensure!(r < total, OutOfBounds: "row {r} is past the {total} lines on hand");
555 let file = starts.partition_point(|&s| s <= r) - 1;
556 Ok((file, r - starts[file]))
557 })
558 .collect()
559 }
560
561 fn column(&self, column: usize, rows: &[IdxSize]) -> PolarsResult<Column> {
562 if self.grows()
564 && let Some(max) = rows.iter().max()
565 && *max as usize >= self.complete_rows()
566 {
567 for f in &self.files {
568 f.grow()?;
569 }
570 }
571 let placed = self.place(rows)?;
572 let several = self.files.len() > 1;
573 let name = self
574 .schema
575 .get_at_index(column)
576 .expect("a column")
577 .0
578 .clone();
579 let column = if several { column } else { column + 1 };
580 let series = match column {
581 0 => placed
582 .iter()
583 .map(|&(f, _)| Some(self.files[f].name.as_str()))
584 .collect::<StringChunked>()
585 .into_series(),
586 _ => {
587 let held: Vec<_> = self
588 .files
589 .iter()
590 .map(|f| f.mapped.read().unwrap_or_else(|e| e.into_inner()))
591 .collect();
592 if !self.grows() {
593 for m in &held {
594 m.bytes.still_whole()?;
595 }
596 }
597 let mut builder = StringChunkedBuilder::new(name.clone(), placed.len());
598 for &(f, i) in &placed {
599 let m = &held[f];
600 let bytes = m.bytes.as_slice();
601 match std::str::from_utf8(m.index.line(bytes, i)) {
602 Ok(text) => builder.append_value(text),
603 Err(_) => builder
604 .append_value(String::from_utf8_lossy(m.index.line(bytes, i)).as_ref()),
605 }
606 }
607 builder.finish().into_series()
608 }
609 };
610 Ok(series.with_name(name).into_column())
611 }
612
613 pub fn collect_window(&self, start: usize, len: usize) -> PolarsResult<DataFrame> {
615 let rows = self.rows();
616 let start = start.min(rows);
617 let len = len.min(rows - start);
618 let index: Vec<IdxSize> = (start..start + len).map(|r| r as IdxSize).collect();
619 let mut columns = (0..self.schema.len())
620 .map(|c| self.column(c, &index))
621 .collect::<PolarsResult<Vec<_>>>()?;
622 columns.push(IdxCa::from_vec(crate::formats::row_index::INDEX.into(), index).into_column());
624 DataFrame::new(len, columns)
625 }
626}
627
628impl crate::formats::row_index::RowSource for Lines {
629 fn height(&self) -> usize {
630 self.rows()
631 }
632
633 fn schema(&self) -> SchemaRef {
634 self.schema.clone()
635 }
636
637 fn decode(&self, column: usize, index: &IdxCa) -> PolarsResult<Column> {
638 polars_ensure!(
639 index.null_count() == 0,
640 ComputeError: "a row index has a missing row"
641 );
642 let rows: Vec<IdxSize> = index.into_no_null_iter().collect();
643 self.column(column, &rows)
644 }
645}
646
647impl crate::formats::pushdown::Windowed for Lines {
648 fn window(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
649 Ok(self.collect_window(start, len)?.lazy())
650 }
651}
652
653fn file_name(path: &Path) -> String {
654 path.file_name().map_or_else(
655 || path.display().to_string(),
656 |n| n.to_string_lossy().into_owned(),
657 )
658}
659
660pub(crate) fn bound(plan: &mut polars::lazy::dsl::DslPlan, rows: IdxSize) -> bool {
664 use polars::lazy::dsl::DslPlan;
665 match plan {
666 DslPlan::DataFrameScan { df, .. } if df.width() == 0 => {
667 *df = Arc::new(DataFrame::empty_with_height(rows as usize));
668 true
669 }
670 _ => false,
671 }
672}
673
674fn scan(input: crate::formats::readers::ScanIn<'_>) -> Result<crate::loading::scan::Scan> {
676 let lines = Arc::new(Lines::open_first(
677 input.paths,
678 input.options.follow,
679 FIRST_BYTES,
680 )?);
681 let lf = lines.lazy();
682 input.report.opened = Some(Arc::new(opened(&lines, input.options)));
683 Ok(lf.into())
684}
685
686pub(crate) fn opened(
689 lines: &Arc<Lines>,
690 options: &crate::OpenOptions,
691) -> crate::formats::members::Opened {
692 crate::formats::members::Opened {
693 window: (!lines.grows()).then(|| {
695 (
696 lines.clone() as Arc<dyn crate::formats::pushdown::Windowed>,
697 lines.rows(),
698 )
699 }),
700 notes: notes(lines, options.format_guessed),
701 indexing: lines.indexing().then(|| lines.clone()),
702 numbering: lines.several().then(|| lines.clone()),
703 ..Default::default()
704 }
705}
706
707pub(crate) const GUESSED: &str = "no format detected";
709
710pub(crate) fn notes(lines: &Lines, format_guessed: bool) -> Vec<crate::notes::Note> {
712 let (invalid, past_limit) = lines.counts();
713 let scope = match lines.files.len() {
714 1 => "the file".to_string(),
715 n => format!("the {n} files"),
716 };
717 let middot = crate::glyphs::get().middot;
718 let mut notes = vec![format!(
719 "read as lines {middot} a row per line, blank lines included"
720 )];
721 if format_guessed {
722 notes.push(format!("{GUESSED} {middot} --format csv reads it as CSV"));
723 }
724 if invalid > 0 {
725 notes.push(format!(
726 "{} with invalid UTF-8 {middot} shown as \u{fffd}",
727 crate::formats::text_formats::count(invalid as u64, "line", "lines")
728 ));
729 }
730 if past_limit > 0 {
731 notes.push(crate::limits::left_out(
732 &crate::formats::text_formats::count(past_limit as u64, "line", "lines"),
733 crate::limits::get().indexed_records,
734 "indexed_records",
735 ));
736 }
737 notes
738 .into_iter()
739 .map(|n| crate::formats::text_formats::note(n, scope.clone()))
740 .collect()
741}
742
743const MOST_RECORDS: usize = 50;
747
748pub fn guess(head: &[u8], whole: bool) -> Option<FileFormat> {
754 let text = head.strip_prefix(b"\xef\xbb\xbf").unwrap_or(head);
755 if !is_text(text, whole) {
756 return None;
757 }
758 let start = text
759 .iter()
760 .position(|b| !b.is_ascii_whitespace())
761 .unwrap_or(text.len());
762 let trimmed = &text[start..];
763 match trimmed.first() {
764 Some(b'[') if json_so_far(trimmed, whole) => return Some(FileFormat::Json),
765 Some(b'{') => {
766 let first = trimmed.split(|&b| b == b'\n').next().unwrap_or_default();
767 let line_whole = first.len() < trimmed.len() || whole;
768 if line_whole
769 && matches!(
770 serde_json::from_slice::<serde_json::Value>(first),
771 Ok(serde_json::Value::Object(_))
772 )
773 {
774 return Some(FileFormat::Jsonl);
775 }
776 if json_so_far(trimmed, whole) {
777 return Some(FileFormat::Json);
778 }
779 }
780 _ => {}
781 }
782 let mut best: Option<(FileFormat, usize)> = None;
783 for (separator, format) in [(b',', FileFormat::Csv), (b'\t', FileFormat::Tsv)] {
784 if let Some(fields) = consistent_fields(text, separator, whole)
785 && best.is_none_or(|(_, most)| fields > most)
786 {
787 best = Some((format, fields));
788 }
789 }
790 Some(best.map_or(FileFormat::Text, |(format, _)| format))
791}
792
793pub fn as_asked(format: FileFormat, options: &crate::OpenOptions) -> FileFormat {
797 let delimited = options.delimiter.is_some()
798 || options.comment_char.is_some()
799 || options.header_rows().is_some()
800 || options.skip_initial_space;
801 match format {
802 FileFormat::Text if delimited => FileFormat::Csv,
803 other => other,
804 }
805}
806
807pub fn guess_file(
810 path: &Path,
811 compression: Option<crate::CompressionFormat>,
812) -> Option<FileFormat> {
813 let compression = compression.or_else(|| crate::CompressionFormat::from_extension(path));
814 let reach = crate::formats::readers::HEAD as u64;
815 let head = crate::formats::head_of(path, compression, reach)?;
816 if head.is_empty() {
818 return None;
819 }
820 let whole = match compression {
821 Some(_) => (head.len() as u64) < reach,
822 None => std::fs::metadata(path).is_ok_and(|m| m.len() <= head.len() as u64),
823 };
824 guess(&head, whole).map(|f| match compression {
825 Some(_) if !f.decompressed_once() => FileFormat::Text,
826 _ => f,
827 })
828}
829
830fn is_text(bytes: &[u8], whole: bool) -> bool {
833 if bytes.contains(&0) {
834 return false;
835 }
836 let mut odd = 0usize;
837 let mut chunks = bytes.utf8_chunks().peekable();
838 while let Some(chunk) = chunks.next() {
839 odd += chunk
840 .valid()
841 .chars()
842 .filter(|c| {
843 c.is_control() && !matches!(c, '\t' | '\n' | '\r' | '\x0c' | '\x1b' | '\x08')
844 })
845 .count();
846 let cut_off = !whole && chunks.peek().is_none();
847 if !cut_off {
848 odd += chunk.invalid().len();
849 }
850 }
851 odd * 10 <= bytes.len()
852}
853
854fn json_so_far(bytes: &[u8], whole: bool) -> bool {
856 match serde_json::from_slice::<serde::de::IgnoredAny>(bytes) {
857 Ok(_) => true,
858 Err(e) => !whole && e.is_eof(),
859 }
860}
861
862fn consistent_fields(text: &[u8], separator: u8, whole: bool) -> Option<usize> {
867 let (mut counts, last, spaced) = field_counts(text, separator)?;
868 if whole
871 && let Some(last) = last
872 && counts.first().is_none_or(|&first| first == last)
873 {
874 counts.push(last);
875 }
876 let first = *counts.first()?;
877 let enough = if whole {
878 counts.len() >= 2
879 } else {
880 counts.len() >= 3
881 };
882 let prose = separator == b',' && first == 2 && spaced;
883 (enough && first >= 2 && !prose && counts.iter().all(|&n| n == first)).then_some(first)
884}
885
886fn field_counts(text: &[u8], separator: u8) -> Option<(Vec<usize>, Option<usize>, bool)> {
890 let mut counts = Vec::new();
891 let mut spaced = true;
892 let mut fields = 1usize;
893 let mut blank = true;
894 let mut field_start = true;
895 let mut in_quotes = false;
896 let mut i = 0;
897 while i < text.len() && counts.len() < MOST_RECORDS {
898 let b = text[i];
899 i += 1;
900 if in_quotes {
901 if b == b'"' {
902 if text.get(i) == Some(&b'"') {
903 i += 1;
904 } else {
905 in_quotes = false;
906 match text.get(i) {
908 Some(&next) if next == separator || next == b'\n' || next == b'\r' => {}
909 None => {}
910 Some(_) => return None,
911 }
912 }
913 }
914 continue;
915 }
916 match b {
917 b'"' if field_start => {
918 in_quotes = true;
919 field_start = false;
920 blank = false;
921 }
922 b'\n' => {
923 if !blank {
924 counts.push(fields);
925 }
926 fields = 1;
927 blank = true;
928 field_start = true;
929 }
930 b'\r' if text.get(i) == Some(&b'\n') => {}
931 _ if b == separator => {
932 spaced &= text.get(i) == Some(&b' ');
933 fields += 1;
934 blank = false;
935 field_start = true;
936 }
937 _ => {
938 if !b.is_ascii_whitespace() {
939 blank = false;
940 }
941 field_start = false;
942 }
943 }
944 }
945 let last = (!in_quotes && !blank && counts.len() < MOST_RECORDS).then_some(fields);
947 Some((counts, last, spaced))
948}
949
950#[cfg(test)]
951mod tests {
952 use super::*;
953
954 fn lines_of(bytes: &[u8]) -> Vec<String> {
955 let lines = Arc::new(Lines::from_bytes(vec![(
956 "a.log".to_string(),
957 Arc::new(Bytes::Owned(bytes.to_vec())),
958 )]));
959 let df = lines.lazy().collect().unwrap();
960 let places: Vec<u32> = df
961 .column(crate::formats::row_index::INDEX)
962 .unwrap()
963 .u32()
964 .unwrap()
965 .into_no_null_iter()
966 .collect();
967 assert_eq!(places, (0..df.height() as u32).collect::<Vec<_>>());
968 df.column(LINE)
969 .unwrap()
970 .str()
971 .unwrap()
972 .iter()
973 .map(|v| v.unwrap().to_string())
974 .collect()
975 }
976
977 #[test]
980 fn every_line_as_written() {
981 assert_eq!(lines_of(b"a\n\nb\n\n\nc\n"), ["a", "", "b", "", "", "c"]);
982 assert_eq!(lines_of(b"a\r\nb\r\n\r\nc"), ["a", "b", "", "c"]);
983 assert_eq!(
984 lines_of(b"x\x1fy,\"z\tw\n#c\x00\x1b[1m\r\nq\rr\n"),
985 ["x\x1fy,\"z\tw", "#c\0\x1b[1m", "q\rr"]
986 );
987 assert_eq!(
988 lines_of(b"ok\n\xff\xfe bad\n"),
989 ["ok", "\u{fffd}\u{fffd} bad"]
990 );
991 assert_eq!(lines_of(b"\n\n"), ["", ""]);
992 assert!(lines_of(b"").is_empty());
993 }
994
995 #[test]
998 fn an_index_reads_on_as_bytes_grow() {
999 let all = b"one\ntw\xffo\npartial line\n\nlast";
1000 for cut in 0..all.len() {
1001 let mut index = LineIndex::of(&all[..cut]);
1002 index.extend(all);
1003 let whole = LineIndex::of(all);
1004 assert_eq!(index.lines(), whole.lines(), "{cut}");
1005 assert_eq!(index.invalid, whole.invalid, "{cut}");
1006 assert_eq!(index.complete(), 4);
1007 for i in 0..whole.lines() {
1008 assert_eq!(index.line(all, i), whole.line(all, i));
1009 }
1010 }
1011 }
1012
1013 #[test]
1014 fn several_files_name_theirs() {
1015 let lines = Arc::new(Lines::from_bytes(vec![
1016 ("a.log".into(), Arc::new(Bytes::Owned(b"1\n2\n".to_vec()))),
1017 ("b.log".into(), Arc::new(Bytes::Owned(b"3\n".to_vec()))),
1018 ]));
1019 let df = lines.lazy().collect().unwrap();
1020 assert_eq!(
1021 df.get_column_names(),
1022 [FILE, LINE, crate::formats::row_index::INDEX]
1023 );
1024 let file: Vec<Option<&str>> = df.column(FILE).unwrap().str().unwrap().iter().collect();
1025 assert_eq!(file, [Some("a.log"), Some("a.log"), Some("b.log")]);
1026 let w = lines.collect_window(1, 5).unwrap();
1027 assert_eq!(w.height(), 2);
1028 assert_eq!(w.get_column_names(), df.get_column_names());
1029 let places: Vec<u32> = w
1030 .column(crate::formats::row_index::INDEX)
1031 .unwrap()
1032 .u32()
1033 .unwrap()
1034 .into_no_null_iter()
1035 .collect();
1036 assert_eq!(places, [1, 2]);
1037 }
1038
1039 #[test]
1042 fn an_index_in_steps_is_the_index_whole() {
1043 let mut bytes = Vec::new();
1044 let mut bad = 0;
1045 for i in 0..200_000u32 {
1046 match i % 7 {
1047 0 => bytes.extend_from_slice(b"\r\n"),
1048 3 => {
1049 bytes.extend_from_slice(b"bad \xff\xfe line\n");
1050 bad += 1;
1051 }
1052 _ => bytes.extend_from_slice(
1053 format!("line {i} {}\n", "x".repeat((i % 50) as usize)).as_bytes(),
1054 ),
1055 }
1056 }
1057 bytes.extend_from_slice(b"no newline at the end \xff");
1058 let whole = LineIndex::of(&bytes);
1059 assert_eq!(whole.lines(), 200_001);
1060 assert_eq!(whole.invalid, bad + 1);
1061 for step in [1, 4096, CHUNK - 3, CHUNK * 2 + 17] {
1062 let mut index = LineIndex {
1063 offsets: Offsets::for_file(bytes.len()),
1064 ..Default::default()
1065 };
1066 let mut steps = 0;
1067 while !index.extend_by(&bytes, step) {
1068 steps += 1;
1069 assert!(index.lines() < whole.lines(), "{step}");
1070 }
1071 assert!(steps > 0 || step > bytes.len(), "{step}");
1072 assert_eq!(index.lines(), whole.lines(), "{step}");
1073 assert_eq!(index.invalid, whole.invalid, "{step}");
1074 assert_eq!(index.complete(), whole.complete(), "{step}");
1075 for i in [0, 1, 3, 1000, whole.lines() - 1] {
1076 assert_eq!(index.line(&bytes, i), whole.line(&bytes, i), "{step} {i}");
1077 }
1078 }
1079 }
1080
1081 #[test]
1085 fn a_stopped_index_is_no_count_and_a_paused_one_reads_on() {
1086 let bytes: Vec<u8> = (0..50_000u32)
1087 .flat_map(|i| format!("{i}\n").into_bytes())
1088 .collect();
1089 let lines = Arc::new(Lines::from_bytes_first(
1090 vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1091 1000,
1092 ));
1093 let lf = lines.lazy();
1094 let waiting = {
1095 let lf = lf.clone();
1096 std::thread::spawn(move || lf.collect())
1097 };
1098 assert!(!lines.index_more(1000));
1100 assert!(lines.resume_indexing());
1101 while !lines.index_more(100_000) {}
1102 assert_eq!(waiting.join().unwrap().unwrap().height(), 50_000);
1103
1104 let bytes: Vec<u8> = (0..50_000u32)
1105 .flat_map(|i| format!("{i}\n").into_bytes())
1106 .collect();
1107 let lines = Arc::new(Lines::from_bytes_first(
1108 vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1109 1000,
1110 ));
1111 let lf = lines.lazy();
1112 let waiting = std::thread::spawn(move || lf.collect());
1113 lines.stop_indexing();
1114 let error = waiting.join().unwrap().unwrap_err().to_string();
1115 assert!(error.contains("not all indexed"), "{error}");
1116 assert!(!lines.resume_indexing() || !lines.whole());
1117 }
1118
1119 #[test]
1124 #[cfg(unix)]
1125 fn a_file_that_shrinks_while_indexed_has_no_count() {
1126 let dir = tempfile::tempdir().unwrap();
1127 let path = dir.path().join("rotated.log");
1128 let bytes: Vec<u8> = (0..50_000u32)
1129 .flat_map(|i| format!("{i}\n").into_bytes())
1130 .collect();
1131 std::fs::write(&path, &bytes).unwrap();
1132 let lines = Arc::new(Lines::open_first(std::slice::from_ref(&path), false, 1000).unwrap());
1133 assert!(lines.indexing());
1134 std::fs::OpenOptions::new()
1136 .write(true)
1137 .open(&path)
1138 .unwrap()
1139 .set_len(10)
1140 .unwrap();
1141 assert!(lines.index_more(100_000), "the indexing stops");
1142 assert!(lines.shrank());
1143 let error = lines.lazy().collect().unwrap_err().to_string();
1144 assert!(
1146 error.contains(SHRANK) || error.contains("shorter"),
1147 "{error}"
1148 );
1149 }
1150
1151 #[test]
1153 fn several_files_number_their_own_lines() {
1154 let lines = Lines::from_bytes(vec![
1155 ("a.log".into(), Arc::new(Bytes::Owned(b"1\n2\n".to_vec()))),
1156 (
1157 "b.log".into(),
1158 Arc::new(Bytes::Owned(b"3\n4\n5\n".to_vec())),
1159 ),
1160 ]);
1161 assert!(lines.several());
1162 let numbered: Vec<Option<usize>> = (0..6).map(|p| lines.line_in_file(p)).collect();
1163 assert_eq!(
1164 numbered,
1165 [Some(0), Some(1), Some(0), Some(1), Some(2), None]
1166 );
1167 }
1168
1169 #[test]
1173 fn a_large_file_indexes_behind_its_first_rows() {
1174 let mut bytes = Vec::new();
1175 for i in 0..100_000u32 {
1176 bytes.extend_from_slice(format!("{i}\n").as_bytes());
1177 }
1178 let lines = Arc::new(Lines::from_bytes_first(
1179 vec![("a.log".into(), Arc::new(Bytes::Owned(bytes)))],
1180 1000,
1181 ));
1182 assert!(lines.indexing());
1183 let first = lines.rows();
1184 assert!(first > 0 && first < 100_000, "{first}");
1185 assert_eq!(lines.collect_window(0, 200_000).unwrap().height(), first);
1186 let lf = lines.lazy();
1187 let waiting = {
1188 let lf = lf.clone().filter(col(LINE).eq(lit("99999")));
1189 std::thread::spawn(move || lf.collect().unwrap())
1190 };
1191 while !lines.index_more(50_000) {}
1192 assert!(!lines.indexing());
1193 assert_eq!(lines.rows(), 100_000);
1194 let (done, all) = lines.indexed_bytes();
1195 assert_eq!(done, all);
1196 let df = waiting.join().unwrap();
1197 assert_eq!(df.height(), 1);
1198 assert_eq!(
1199 df.column(crate::formats::row_index::INDEX)
1200 .unwrap()
1201 .u32()
1202 .unwrap()
1203 .get(0),
1204 Some(99_999)
1205 );
1206 for streaming in [false, true] {
1207 let df = crate::analysis::statistics::collect_lazy(lf.clone(), streaming).unwrap();
1208 assert_eq!(df.height(), 100_000, "streaming {streaming}");
1209 }
1210 }
1211
1212 #[test]
1213 fn delimited_only_on_evidence() {
1214 let cases: [(&[u8], bool, Option<FileFormat>); 18] = [
1215 (b"a,b,c\n1,2,3\n4,5,6\n", true, Some(FileFormat::Csv)),
1216 (b"a,b\n1,2\n", true, Some(FileFormat::Csv)),
1217 (b"a\tb\n1\t2\n3\t4\n", true, Some(FileFormat::Tsv)),
1218 (
1219 b"name,note\n\"x\",\"a, b\nc\"\ny,z\n",
1220 true,
1221 Some(FileFormat::Csv),
1222 ),
1223 (b"a,b,c\n1,2,3\n4,5,6\n7,8", false, Some(FileFormat::Csv)),
1225 (b"a,b,c\n1,2,3\n", false, Some(FileFormat::Text)),
1226 (
1228 b"2024-01-01 INFO started, ok\n2024-01-01 WARN slow, very, slow\nINFO done\n",
1229 true,
1230 Some(FileFormat::Text),
1231 ),
1232 (
1233 b"[2024-01-01 12:00] INFO hi\n[2024-01-01 12:01] INFO there\n",
1234 true,
1235 Some(FileFormat::Text),
1236 ),
1237 (b"[INFO] a\n", true, Some(FileFormat::Text)),
1238 (
1239 b"say \"hi\", she said\nok, then\nno, way\n",
1240 true,
1241 Some(FileFormat::Text),
1242 ),
1243 (b"one line, with a comma\n", true, Some(FileFormat::Text)),
1244 (b"[1, 2,", false, Some(FileFormat::Json)),
1245 (b"[1, 2]", true, Some(FileFormat::Json)),
1246 (b"{\"a\": 1}\n{\"a\": 2}\n", true, Some(FileFormat::Jsonl)),
1247 (b"{\n \"a\": 1\n}\n", true, Some(FileFormat::Json)),
1248 (b"{not json}\n", true, Some(FileFormat::Text)),
1249 (b"\x89PNG\r\n\x1a\n\0\0\0\rIHDR", false, None),
1250 (b"\xef\xbb\xbfid,v\n1,2\n", true, Some(FileFormat::Csv)),
1251 ];
1252 for (head, whole, format) in cases {
1253 assert_eq!(
1254 guess(head, whole),
1255 format,
1256 "{}",
1257 String::from_utf8_lossy(head)
1258 );
1259 }
1260 assert_eq!(guess(b"caf\xe9 au lait\n", true), Some(FileFormat::Text));
1262 let noise: Vec<u8> = (0..512u32).map(|i| (i * 97 % 255 + 1) as u8).collect();
1263 assert_eq!(guess(&noise, true), None);
1264 }
1265}