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