1use std::sync::Arc;
17
18use rudb_common::{Error, Field, LogicalType, Result};
19use rudb_io::File;
20use rudb_vector::{Chunk, VECTOR_SIZE};
21
22use crate::convert::{self, Cells};
23use crate::dialect::{self, Dialect, Given};
24use crate::infer;
25use crate::scan::{Records, Span};
26
27const BLOCK: usize = 1 << 20;
33
34const TAIL: usize = 64 << 10;
40
41const BLOCK_ROWS: usize = 1024;
47
48#[derive(Debug)]
50pub struct Reader {
51 file: Arc<dyn File>,
52 path: String,
53 given: Given,
54 dialect: Dialect,
55 nulls: Option<Vec<Vec<u8>>>,
57 fields: Vec<Field>,
58 projection: Vec<usize>,
59 buffer: Vec<u8>,
60 at: usize,
61 offset: u64,
62 drained: bool,
63 line: u64,
64 scratch: Vec<String>,
65 records: Records,
66 block: usize,
67 origin: u64,
69 end: u64,
71 cap: u64,
74}
75
76impl Reader {
77 pub fn open(file: Box<dyn File>, path: &str) -> Result<Self> {
86 Self::open_with(file, path, Given::default())
87 }
88
89 pub fn open_with(file: Box<dyn File>, path: &str, given: Given) -> Result<Self> {
99 Self::open_sized(file, path, given, BLOCK)
100 }
101
102 pub(crate) fn open_sized(
105 file: Box<dyn File>,
106 path: &str,
107 given: Given,
108 block: usize,
109 ) -> Result<Self> {
110 let nulls = given.null_strings();
111 let mut reader = Self {
112 file: Arc::from(file),
113 path: path.to_string(),
114 given,
115 dialect: Dialect::comma_separated(),
116 nulls,
117 fields: Vec::new(),
118 projection: Vec::new(),
119 buffer: Vec::new(),
120 at: 0,
121 offset: 0,
122 drained: false,
123 line: 1,
124 scratch: Vec::new(),
125 records: Records::default(),
126 block,
127 origin: 0,
128 end: u64::MAX,
129 cap: u64::MAX,
130 };
131 reader.fill(0)?;
132 let sample = reader.buffer.clone();
133 let given = &reader.given;
134 let quote = given.quote.or_else(|| dialect::quote(&sample));
135 let delimiter = match given.delimiter {
136 Some(byte) => byte,
137 None => dialect::delimiter(&sample, quote)?,
138 };
139 let escape = given.escape.or(quote);
140 let told = given.header;
141 reader.dialect = Dialect { delimiter, quote, escape, header: false };
142 let rows = reader.sample_rows(&sample)?;
143 let (header, mut fields) = describe(&rows, told);
144 if let Some(names) = &reader.given.names {
145 if names.len() > fields.len() {
146 return Err(Error::invalid_input(format!(
147 "Error when sniffing file \"{path}\".\nIt was not possible to detect the CSV \
148 Header, due to the header having less columns than expected\nNumber of \
149 expected columns: {}. Actual number of columns {}",
150 names.len(),
151 fields.len()
152 )));
153 }
154 for (field, name) in fields.iter_mut().zip(names) {
155 field.name.clone_from(name);
156 }
157 }
158 reader.dialect.header = header;
159 reader.fields = fields;
160 reader.projection = (0..reader.fields.len()).collect();
161 if header {
162 reader.skip_record()?;
163 }
164 Ok(reader)
165 }
166
167 #[must_use]
169 pub fn fields(&self) -> Vec<Field> {
170 self.projection.iter().map(|&at| self.fields[at].clone()).collect()
171 }
172
173 pub fn project(&mut self, columns: &[usize]) -> Result<()> {
179 for &column in columns {
180 if column >= self.fields.len() {
181 return Err(Error::io(format!(
182 "column {column} is past the {} the file has",
183 self.fields.len()
184 )));
185 }
186 }
187 self.projection = columns.to_vec();
188 Ok(())
189 }
190
191 pub fn retype(&mut self, types: &[LogicalType]) -> Result<()> {
209 if types.len() != self.projection.len() {
210 return Err(Error::io(format!(
211 "{} types for a projection of {} columns",
212 types.len(),
213 self.projection.len()
214 )));
215 }
216 for (&at, ty) in self.projection.iter().zip(types) {
217 self.fields[at].ty = ty.clone();
218 }
219 Ok(())
220 }
221
222 #[must_use]
224 pub const fn dialect(&self) -> Dialect {
225 self.dialect
226 }
227
228 pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
240 let rows = self.next_records()?;
241 if rows == 0 {
242 return Ok(None);
243 }
244 let first = self.line;
245 self.line += rows as u64;
246 let cells = Cells {
247 bytes: &self.buffer,
248 records: &self.records,
249 dialect: self.dialect,
250 nulls: self.nulls.as_deref(),
251 };
252 let projected: Vec<_> =
253 self.projection.iter().map(|&at| (at, &self.fields[at].ty)).collect();
254 let mut builders = convert::builders(&cells, &projected);
255 let mut start = 0;
256 while start < rows {
257 let end = rows.min(start + BLOCK_ROWS);
258 for (build, &at) in builders.iter_mut().zip(&self.projection) {
259 let field = &self.fields[at];
260 let refuse = |text: &str, row: usize| {
261 Error::conversion(self.conversion_error(text, field, first + row as u64))
262 };
263 if let Err(error) = build.rows(&cells, at, start..end, &refuse) {
264 return Err(self.first_bad_value(&cells, first).unwrap_or(error));
265 }
266 }
267 start = end;
268 }
269 let columns = builders.into_iter().map(|build| build.finish()).collect::<Result<_>>()?;
270 Ok(Some(Chunk::with_rows(columns, rows)?))
271 }
272
273 fn first_bad_value(&self, cells: &Cells<'_>, first: u64) -> Option<Error> {
280 self.projection.iter().find_map(|&at| {
281 let field = &self.fields[at];
282 let refuse = |text: &str, row: usize| {
283 Error::conversion(self.conversion_error(text, field, first + row as u64))
284 };
285 convert::column(cells, at, &field.ty, &refuse).err()
286 })
287 }
288
289 fn next_records(&mut self) -> Result<usize> {
297 self.records.clear();
298 let mut start = self.at;
299 let mut careful = false;
300 loop {
301 if self.here() >= self.end {
302 break;
303 }
304 let limit = if careful { self.records.len() + 1 } else { VECTOR_SIZE };
305 self.at = crate::scan::records(
306 &self.buffer,
307 self.at,
308 self.dialect,
309 self.drained,
310 limit,
311 &mut self.records,
312 )?;
313 if !careful && self.here() > self.end {
314 self.records.clear();
315 self.at = start;
316 careful = true;
317 continue;
318 }
319 if self.records.len() == VECTOR_SIZE {
320 break;
321 }
322 if careful && self.records.len() == limit {
323 continue;
324 }
325 if self.drained {
326 break;
327 }
328 if self.buffer.len() - start + self.block > Span::MOST {
333 if self.records.is_empty() {
334 return Err(Error::io("a record is longer than two gigabytes"));
335 }
336 break;
337 }
338 self.fill(start)?;
339 self.records.shift(start);
340 start = 0;
341 }
342 Ok(self.records.len())
343 }
344
345 pub(crate) fn here(&self) -> u64 {
347 self.offset - self.buffer.len() as u64 + self.at as u64
348 }
349
350 #[must_use]
352 pub fn bytes_read(&self) -> u64 {
353 self.offset - self.origin
354 }
355
356 pub(crate) fn stretch(&self, from: u64, end: u64, line: u64) -> Self {
363 Self {
364 file: Arc::clone(&self.file),
365 path: self.path.clone(),
366 given: self.given.clone(),
367 dialect: self.dialect,
368 nulls: self.nulls.clone(),
369 fields: self.fields.clone(),
370 projection: self.projection.clone(),
371 buffer: Vec::new(),
372 at: 0,
373 offset: from,
374 drained: false,
375 line,
376 scratch: Vec::new(),
377 records: Records::default(),
378 block: self.block,
379 origin: from,
380 end,
381 cap: u64::MAX,
382 }
383 }
384
385 pub(crate) fn give_up_at(&mut self, cap: u64) {
387 self.cap = cap;
388 }
389
390 pub(crate) fn skim(&mut self) -> Result<u64> {
393 while self.next_records()? > 0 {}
394 Ok(self.here())
395 }
396
397 pub(crate) fn skip_to(&mut self, target: u64) -> Result<()> {
403 loop {
404 let from = self.here();
405 let rows = self.next_records()?;
406 if rows == 0 {
407 return Ok(());
408 }
409 if self.here() > target {
410 let front = self.offset - self.buffer.len() as u64;
411 self.at = usize::try_from(from - front)
412 .map_err(|_| Error::internal("a chunk start outside the buffer"))?;
413 return Ok(());
414 }
415 self.line += rows as u64;
416 }
417 }
418
419 pub(crate) const fn line(&self) -> u64 {
421 self.line
422 }
423
424 pub(crate) fn file(&self) -> &dyn File {
426 self.file.as_ref()
427 }
428
429 fn conversion_error(&self, text: &str, field: &Field, line: u64) -> String {
441 let advice = if self.given.typed {
445 "This type was either manually set or derived from an existing table. Select a \
446 different type to correctly parse this column."
447 .to_string()
448 } else {
449 format!(
450 "This type was auto-detected from the CSV file.\nPossible solutions:\n* Override \
451 the type for this column manually by setting the type explicitly, e.g., \
452 types={{'{}': 'VARCHAR'}}\n* Set the sample size to a larger value to enable the \
453 auto-detection to scan more values, e.g., sample_size=-1\n* Use a COPY statement \
454 to automatically derive types from an existing table.",
455 field.name
456 )
457 };
458 format!(
459 "CSV Error on Line: {line}\nOriginal Line: {text}\nError when converting column \
460 \"{}\". Could not convert string \"{text}\" to '{}'\n\nColumn {} is being converted \
461 as type {}\n{advice}\n* Check whether the null string value is set correctly (e.g., \
462 nullstr = 'N/A')\n\n file = {}\n delimiter = {}\n quote = {}\n escape = {}\n \
463 header = {} {}\n sample_size = {}\n",
464 field.name,
465 field.ty,
466 field.name,
467 field.ty,
468 self.path,
469 Given::shown(self.given.delimiter, Some(self.dialect.delimiter)),
470 Given::shown(self.given.quote, self.dialect.quote),
471 Given::shown(self.given.escape, self.dialect.escape),
472 self.dialect.header,
473 Given::source(self.given.header.is_some()),
474 infer::SAMPLE,
475 )
476 }
477
478 #[cfg(test)]
481 fn next_chunk_by_record(&mut self) -> Result<Option<Chunk>> {
482 let mut rows: Vec<Vec<Option<String>>> = Vec::new();
483 while rows.len() < VECTOR_SIZE {
484 match self.next_record()? {
485 Some(fields) => rows.push(fields),
486 None => break,
487 }
488 }
489 if rows.is_empty() {
490 return Ok(None);
491 }
492 let mut columns = Vec::with_capacity(self.projection.len());
493 for &at in &self.projection {
494 let field = &self.fields[at];
495 let mut values = Vec::with_capacity(rows.len());
496 for (row, held) in rows.iter().enumerate() {
497 let text = held.get(at).and_then(Option::as_deref);
498 values.push(self.convert(
499 text,
500 field,
501 self.line - rows.len() as u64 + row as u64,
502 )?);
503 }
504 columns.push(rudb_vector::Vector::from_values(field.ty.clone(), &values)?);
505 }
506 Ok(Some(Chunk::with_rows(columns, rows.len())?))
507 }
508
509 #[cfg(test)]
511 fn convert(&self, text: Option<&str>, field: &Field, line: u64) -> Result<rudb_common::Value> {
512 let Some(text) = text else { return Ok(rudb_common::Value::Null) };
513 if field.ty == LogicalType::Varchar {
514 return Ok(rudb_common::Value::Varchar(text.to_string()));
515 }
516 let value = rudb_common::Value::Varchar(text.to_string());
517 match rudb_kernels::cast_value(&value, &field.ty, false) {
518 Ok(converted) => Ok(converted),
519 Err(_) => Err(Error::conversion(self.conversion_error(text, field, line))),
520 }
521 }
522
523 #[cfg(test)]
528 fn next_record(&mut self) -> Result<Option<Vec<Option<String>>>> {
529 let Some(()) = self.advance()? else { return Ok(None) };
530 Ok(Some(
531 self.scratch
532 .iter()
533 .map(|text| if self.is_null(text) { None } else { Some(text.clone()) })
534 .collect(),
535 ))
536 }
537
538 fn advance(&mut self) -> Result<Option<()>> {
540 loop {
541 let mut scratch = std::mem::take(&mut self.scratch);
542 let outcome = crate::scan::record(
543 &self.buffer,
544 self.at,
545 self.dialect,
546 self.drained,
547 &mut scratch,
548 );
549 self.scratch = scratch;
550 match outcome? {
551 Some(next) => {
552 self.at = next;
553 self.line += 1;
554 return Ok(Some(()));
555 }
556 None if self.drained => return Ok(None),
557 None => self.fill(self.at)?,
558 }
559 }
560 }
561
562 fn skip_record(&mut self) -> Result<()> {
564 self.advance()?;
565 Ok(())
566 }
567
568 fn fill(&mut self, keep: usize) -> Result<()> {
573 if self.offset >= self.cap {
574 return Err(Error::io("a record runs further than a guessed start is followed"));
575 }
576 self.buffer.drain(..keep);
577 self.at -= keep;
578 let held = self.buffer.len();
579 let want = match self.end.checked_sub(self.offset) {
580 Some(left) if left > 0 => self.block.min(usize::try_from(left).unwrap_or(usize::MAX)),
581 _ => self.block.min(TAIL.max(held)),
582 };
583 self.buffer.resize(held + want, 0);
584 let read = self.file.read_at(self.offset, &mut self.buffer[held..])?;
585 self.buffer.truncate(held + read);
586 self.offset += read as u64;
587 if read == 0 {
588 self.drained = true;
589 }
590 Ok(())
591 }
592
593 fn is_null(&self, text: &str) -> bool {
596 match &self.nulls {
597 None => text.is_empty(),
598 Some(nulls) => nulls.iter().any(|null| null.as_slice() == text.as_bytes()),
599 }
600 }
601
602 fn sample_rows(&self, sample: &[u8]) -> Result<Vec<Vec<Option<String>>>> {
605 let mut rows = Vec::new();
606 let mut fields = Vec::new();
607 let mut at = 0;
608 while rows.len() <= infer::SAMPLE {
609 let Some(next) = crate::scan::record(sample, at, self.dialect, false, &mut fields)?
612 else {
613 break;
614 };
615 at = next;
616 rows.push(
617 fields
618 .iter()
619 .map(|text| if self.is_null(text) { None } else { Some(text.clone()) })
620 .collect(),
621 );
622 }
623 Ok(rows)
624 }
625}
626
627fn describe(rows: &[Vec<Option<String>>], told: Option<bool>) -> (bool, Vec<Field>) {
638 let width = rows.iter().map(Vec::len).max().unwrap_or(0);
639 let body = types(&rows[1.min(rows.len())..], width);
640 let all_text = body.iter().all(|ty| *ty == LogicalType::Varchar);
641 let first_fits = rows.first().is_some_and(|first| {
642 first.iter().zip(&body).all(|(text, ty)| match text {
643 None => true,
644 Some(text) => infer::fits(text, ty),
645 })
646 });
647 let header = !rows.is_empty() && told.unwrap_or(rows.len() > 1 && (all_text || !first_fits));
649 if !header {
650 let types = types(rows, width);
651 let fields = types
652 .into_iter()
653 .enumerate()
654 .map(|(at, ty)| Field::new(format!("column{at}"), ty))
655 .collect();
656 return (false, fields);
657 }
658 let names = unique(&rows[0], width);
659 let fields = body.into_iter().zip(names).map(|(ty, name)| Field::new(name, ty)).collect();
660 (true, fields)
661}
662
663fn unique(header: &[Option<String>], width: usize) -> Vec<String> {
676 let mut taken: Vec<String> = Vec::with_capacity(width);
677 for at in 0..width {
678 let base = match header.get(at).and_then(Option::as_deref) {
679 Some(written) => written.to_string(),
680 None => format!("column{at}"),
681 };
682 let mut name = base.clone();
683 let mut next = 1;
684 while taken.iter().any(|held| held.eq_ignore_ascii_case(&name)) {
685 name = format!("{base}_{next}");
686 next += 1;
687 }
688 taken.push(name);
689 }
690 taken
691}
692
693fn types(rows: &[Vec<Option<String>>], width: usize) -> Vec<LogicalType> {
695 (0..width)
696 .map(|at| {
697 let values: Vec<Option<&str>> =
698 rows.iter().map(|row| row.get(at).and_then(Option::as_deref)).collect();
699 infer::column(&values)
700 })
701 .collect()
702}
703
704#[cfg(test)]
705mod tests {
706 use super::*;
707 use rudb_common::Value;
708 use rudb_io::{Filesystem, OpenMode, SimFilesystem};
709 use std::path::Path;
710
711 fn read(text: &str) -> Reader {
712 let filesystem = SimFilesystem::new();
713 let path = Path::new("/t.csv");
714 let file = filesystem.open(path, OpenMode::Create).expect("creates");
715 file.write_at(0, text.as_bytes()).expect("writes");
716 drop(file);
717 let file = filesystem.open(path, OpenMode::Read).expect("opens");
718 Reader::open(file, "/t.csv").expect("sniffs")
719 }
720
721 fn read_with(text: &str, given: Given) -> Reader {
723 let filesystem = SimFilesystem::new();
724 let path = Path::new("/t.csv");
725 let file = filesystem.open(path, OpenMode::Create).expect("creates");
726 file.write_at(0, text.as_bytes()).expect("writes");
727 drop(file);
728 let file = filesystem.open(path, OpenMode::Read).expect("opens");
729 Reader::open_with(file, "/t.csv", given).expect("reads")
730 }
731
732 fn names_and_types(reader: &Reader) -> Vec<(String, String)> {
733 reader.fields().iter().map(|f| (f.name.clone(), f.ty.to_string())).collect()
734 }
735
736 fn all(reader: &mut Reader) -> Vec<Vec<Value>> {
737 let mut rows = Vec::new();
738 while let Some(chunk) = reader.next_chunk().expect("reads") {
739 for row in 0..chunk.len() {
740 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
741 }
742 }
743 rows
744 }
745
746 #[test]
747 fn a_header_that_names_two_columns_the_same_thing_counts_the_second_one_up() {
748 let names: Vec<String> =
749 read("a,a,A,a_1\n1,2,3,4\nx,y,z,w\n").fields().into_iter().map(|f| f.name).collect();
750 assert_eq!(names, ["a", "a_1", "A_2", "a_1_1"]);
754 }
755
756 #[test]
757 fn a_header_and_three_types_are_what_duckdb_sniffs_for_the_same_bytes() {
758 let reader = read("a,b,c\n1,x,2.5\n2,y,3.5\n");
759 assert_eq!(
760 names_and_types(&reader),
761 [
762 ("a".to_string(), "BIGINT".to_string()),
763 ("b".to_string(), "VARCHAR".to_string()),
764 ("c".to_string(), "DOUBLE".to_string()),
765 ]
766 );
767 }
768
769 #[test]
770 fn a_file_with_no_header_gets_the_names_duckdb_gives_it() {
771 let reader = read("1,x\n2,y\n");
772 assert_eq!(
773 names_and_types(&reader),
774 [
775 ("column0".to_string(), "BIGINT".to_string()),
776 ("column1".to_string(), "VARCHAR".to_string()),
777 ]
778 );
779 }
780
781 #[test]
782 fn two_rows_of_words_are_a_header_and_a_row() {
783 let reader = read("a,b\nc,d\n");
784 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "b"]);
785 }
786
787 #[test]
788 fn one_column_of_words_under_a_row_of_numbers_is_still_a_header() {
789 let reader = read("a,2\n3,4\n");
792 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "2"]);
793 }
794
795 #[test]
796 fn the_rows_are_the_rows_of_the_file() {
797 let mut reader = read("a,b\n1,x\n2,y\n");
798 assert_eq!(
799 all(&mut reader),
800 [
801 vec![Value::BigInt(1), Value::Varchar("x".into())],
802 vec![Value::BigInt(2), Value::Varchar("y".into())],
803 ]
804 );
805 }
806
807 #[test]
808 fn an_empty_field_is_a_null_whether_it_was_quoted_or_not() {
809 let mut reader = read("a,b\n1,\n\"\",y\n");
812 assert_eq!(
813 all(&mut reader),
814 [vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::Varchar("y".into())],]
815 );
816 }
817
818 #[test]
819 fn a_projection_picks_columns_out_by_position_and_can_reorder_them() {
820 let mut reader = read("a,b,c\n1,x,2.5\n");
821 reader.project(&[2, 0]).expect("projects");
822 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["c", "a"]);
823 assert_eq!(all(&mut reader), [vec![Value::Double(2.5), Value::BigInt(1)]]);
824 }
825
826 #[test]
827 fn a_projection_of_nothing_still_counts_the_rows() {
828 let mut reader = read("a,b\n1,x\n2,y\n3,z\n");
829 reader.project(&[]).expect("projects");
830 let chunk = reader.next_chunk().expect("reads").expect("a chunk");
831 assert_eq!(chunk.len(), 3);
832 assert_eq!(chunk.width(), 0);
833 }
834
835 #[test]
836 fn a_pipe_separated_file_reads_as_one() {
837 let mut reader = read("a|b\n1|x\n");
838 assert_eq!(reader.dialect().delimiter, b'|');
839 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x".into())]]);
840 }
841
842 #[test]
843 fn a_quoted_field_with_a_delimiter_in_it_is_one_value() {
844 let mut reader = read("a,b\n1,\"x,y\"\n");
845 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
846 }
847
848 #[test]
849 fn more_rows_than_fit_one_chunk_arrive_as_more_than_one_chunk() {
850 let mut text = String::from("a\n");
851 for row in 0..VECTOR_SIZE + 5 {
852 text.push_str(&format!("{row}\n"));
853 }
854 let mut reader = read(&text);
855 let first = reader.next_chunk().expect("reads").expect("a chunk");
856 assert_eq!(first.len(), VECTOR_SIZE);
857 let second = reader.next_chunk().expect("reads").expect("a second chunk");
858 assert_eq!(second.len(), 5);
859 assert!(reader.next_chunk().expect("reads").is_none());
860 }
861
862 #[test]
863 fn a_value_the_sniffer_never_saw_is_an_error_rather_than_a_wider_column() {
864 let mut text = String::from("c\n");
868 for row in 0..infer::SAMPLE {
869 text.push_str(&format!("{row}\n"));
870 }
871 text.push_str("oops\n");
872 let mut reader = read(&text);
873 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
874 let error = all_or_error(&mut reader).unwrap_err();
875 let line = infer::SAMPLE + 2;
876 assert!(error.message().starts_with(&format!("CSV Error on Line: {line}")), "{error}");
877 assert!(
878 error.message().contains("Could not convert string \"oops\" to 'BIGINT'"),
879 "{error}"
880 );
881 assert!(error.message().contains("sample_size = 20480"), "{error}");
882 }
883
884 #[test]
888 fn the_error_names_the_first_columns_bad_value_even_when_a_later_column_is_bad_sooner() {
889 let mut text = String::from("a,b\n");
890 for row in 0..infer::SAMPLE {
891 text.push_str(&format!("{row},{row}\n"));
892 }
893 text.push_str("1,late\n");
894 for row in 0..BLOCK_ROWS * 2 {
895 text.push_str(&format!("{row},{row}\n"));
896 }
897 text.push_str("early,1\n");
898 let mut reader = read(&text);
899 let error = all_or_error(&mut reader).unwrap_err();
900 assert!(error.message().contains("Could not convert string \"early\""), "{error}");
901 }
902
903 #[test]
906 fn a_file_told_it_has_no_header_reads_its_first_line_as_a_row() {
907 let mut reader =
908 read_with("a,b\n1,x\n2,y\n", Given { header: Some(false), ..Given::default() });
909 assert_eq!(
910 names_and_types(&reader),
911 [
912 ("column0".to_string(), "VARCHAR".to_string()),
913 ("column1".to_string(), "VARCHAR".to_string()),
914 ]
915 );
916 assert_eq!(all(&mut reader).len(), 3);
917 }
918
919 #[test]
921 fn a_file_told_it_has_a_header_takes_its_first_line_as_the_names() {
922 let reader = read_with("1,2\n3,4\n", Given { header: Some(true), ..Given::default() });
923 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["1", "2"]);
924 }
925
926 #[test]
928 fn a_given_delimiter_is_the_delimiter_whatever_the_file_looks_like() {
929 let reader = read_with("a,b\n1,x\n", Given { delimiter: Some(b';'), ..Given::default() });
930 assert_eq!(reader.dialect().delimiter, b';');
931 assert_eq!(reader.fields().len(), 1);
932 }
933
934 #[test]
936 fn a_given_quote_makes_a_field_that_holds_the_delimiter_one_value() {
937 let mut reader =
938 read_with("a,b\n1,'x,y'\n", Given { quote: Some(b'\''), ..Given::default() });
939 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
940 }
941
942 #[test]
944 fn the_block_says_set_by_user_for_what_the_call_gave_it() {
945 let mut text = String::from("c;d\n");
946 for row in 0..infer::SAMPLE {
947 text.push_str(&format!("{row};x\n"));
948 }
949 text.push_str("oops;x\n");
950 let given = Given { delimiter: Some(b';'), ..Given::default() };
951 let mut reader = read_with(&text, given);
952 let error = all_or_error(&mut reader).unwrap_err();
953 assert!(error.message().contains("delimiter = ; (Set By User)"), "{error}");
954 assert!(error.message().contains("header = true (Auto-Detected)"), "{error}");
955 }
956
957 #[test]
958 fn a_file_told_a_wider_type_than_it_sniffed_reads_its_whole_numbers_as_that_type() {
959 let mut reader = read("a\n1\n2\n");
963 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
964 reader.retype(&[LogicalType::Double]).expect("one type for one column");
965 assert_eq!(reader.fields()[0].ty, LogicalType::Double);
966 assert_eq!(all(&mut reader), [[Value::Double(1.0)], [Value::Double(2.0)]]);
967 }
968
969 #[test]
973 fn a_null_string_other_than_the_empty_one_leaves_an_empty_field_as_text() {
974 let given =
975 Given { header: Some(false), nulls: Some(vec!["NA".to_string()]), ..Given::default() };
976 let mut reader = read_with("x,NA\n,\"NA\"\n\"\",y\n", given);
977 let text = |value: &str| Value::Varchar(value.into());
978 assert_eq!(
979 all(&mut reader),
980 [vec![text("x"), Value::Null], vec![text(""), Value::Null], vec![text(""), text("y")],]
981 );
982 }
983
984 #[test]
985 fn a_list_of_null_strings_makes_each_of_them_a_null() {
986 let given = Given {
987 header: Some(false),
988 nulls: Some(vec!["NA".to_string(), "-".to_string()]),
989 ..Given::default()
990 };
991 let mut reader = read_with("1,NA\n-,2\n", given);
992 assert_eq!(
993 all(&mut reader),
994 [vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::BigInt(2)]]
995 );
996 }
997
998 #[test]
999 fn given_names_rename_the_first_columns_and_too_many_of_them_are_refused() {
1000 let names = |list: &[&str]| Some(list.iter().map(ToString::to_string).collect());
1001 let reader = read_with("1,2,3\n", Given { names: names(&["p"]), ..Given::default() });
1002 let found: Vec<String> = reader.fields().iter().map(|f| f.name.clone()).collect();
1003 assert_eq!(found, ["p", "column1", "column2"]);
1004
1005 let filesystem = SimFilesystem::new();
1006 let path = Path::new("/t.csv");
1007 let file = filesystem.open(path, OpenMode::Create).expect("creates");
1008 file.write_at(0, b"1,2\n").expect("writes");
1009 drop(file);
1010 let file = filesystem.open(path, OpenMode::Read).expect("opens");
1011 let given = Given { names: names(&["a", "b", "c"]), ..Given::default() };
1012 let Err(error) = Reader::open_with(file, "/t.csv", given) else {
1013 panic!("three names for two columns are refused")
1014 };
1015 assert!(
1016 error.message().ends_with("Number of expected columns: 3. Actual number of columns 2"),
1017 "{error}"
1018 );
1019 }
1020
1021 #[test]
1024 fn a_bad_value_in_a_column_whose_type_was_set_gets_the_advice_for_a_set_type() {
1025 let given = Given { header: Some(false), typed: true, ..Given::default() };
1026 let mut reader = read_with("1,x\nfoo,y\n", given);
1027 reader.retype(&[LogicalType::Integer, LogicalType::Varchar]).expect("two for two");
1028 let error = all_or_error(&mut reader).unwrap_err();
1029 let message = error.message();
1030 assert!(message.starts_with("CSV Error on Line: 2"), "{error}");
1031 assert!(message.contains("Could not convert string \"foo\" to 'INTEGER'"), "{error}");
1032 assert!(
1033 message.contains(
1034 "Column column0 is being converted as type INTEGER\nThis type was either manually \
1035 set or derived from an existing table."
1036 ),
1037 "{error}"
1038 );
1039 assert!(!message.contains("Possible solutions"), "{error}");
1040 }
1041
1042 #[test]
1043 fn a_type_list_that_is_not_as_long_as_the_projection_is_refused() {
1044 let mut reader = read("a,b\n1,two\n");
1045 let error = reader.retype(&[LogicalType::Double]).unwrap_err();
1046 assert!(error.message().contains("1 types for a projection of 2 columns"), "{error}");
1047 }
1048
1049 fn drained(reader: &mut Reader, old: bool) -> (Vec<(usize, Vec<String>)>, Option<String>) {
1052 let mut chunks = Vec::new();
1053 loop {
1054 let next = if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
1055 match next {
1056 Ok(Some(chunk)) => {
1057 let mut values = Vec::new();
1058 for row in 0..chunk.len() {
1059 for at in 0..chunk.width() {
1060 values.push(format!("{:?}", chunk.value_at(row, at)));
1061 }
1062 }
1063 chunks.push((chunk.len(), values));
1064 }
1065 Ok(None) => return (chunks, None),
1066 Err(error) => return (chunks, Some(error.to_string())),
1067 }
1068 }
1069 }
1070
1071 struct Rng(u64);
1073
1074 impl Rng {
1075 fn next(&mut self) -> u64 {
1076 self.0 ^= self.0 << 13;
1077 self.0 ^= self.0 >> 7;
1078 self.0 ^= self.0 << 17;
1079 self.0
1080 }
1081
1082 fn below(&mut self, n: usize) -> usize {
1083 (self.next() % n as u64) as usize
1084 }
1085 }
1086
1087 fn typed_file(rng: &mut Rng, rows: usize) -> String {
1090 const ODD: [&str; 27] = [
1091 "",
1092 " 1",
1093 "1 ",
1094 "1e3",
1095 "0x10",
1096 "inf",
1097 "-nan",
1098 "abc",
1099 "\"12\"",
1100 "\"a\"\"b\"",
1101 "\"x,y\"",
1102 "\"x\ny\"",
1103 "h\u{e9}llo",
1104 "99999999999999999999",
1105 "9999999999999999999",
1106 "-",
1107 "+5",
1108 "007",
1109 "2020-02-30",
1110 "2020-02-29",
1111 "0000-01-01",
1112 "TRUE",
1113 "no",
1114 "1_000",
1115 "1.5e-3",
1116 "-0",
1117 "\"\"",
1118 ];
1119 let width = 1 + rng.below(6);
1120 let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
1121 let mut text: String = (0..width).map(|at| format!("c{at}")).collect::<Vec<_>>().join(",");
1122 text.push('\n');
1123 for _ in 0..rows {
1124 let mut fields = Vec::with_capacity(width);
1125 for &kind in &kinds {
1126 let odd = rng.below(60) == 0;
1127 fields.push(if odd {
1128 ODD[rng.below(ODD.len())].to_string()
1129 } else {
1130 let n = rng.next();
1131 match kind {
1132 0 => format!("{}", (n % 2_000_001) as i64 - 1_000_000),
1133 1 => format!("{}.{:02}", n % 100_000, n % 100),
1134 2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
1135 3 => ["true", "false", "t", "F"][(n % 4) as usize].to_string(),
1136 _ => ["x", "hello world", "a longer piece of text", "\"q,\"\"q\""]
1137 [(n % 4) as usize]
1138 .to_string(),
1139 }
1140 });
1141 }
1142 if rng.below(200) == 0 {
1143 fields.pop();
1144 }
1145 if rng.below(200) == 0 {
1146 fields.push("extra".to_string());
1147 }
1148 text.push_str(&fields.join(","));
1149 text.push_str(["\n", "\n", "\n", "\r\n", "\r"][rng.below(5)]);
1150 }
1151 if rng.below(4) == 0 {
1152 text.pop();
1153 }
1154 text
1155 }
1156
1157 fn open_sized(text: &[u8], block: usize) -> Result<Reader> {
1158 let filesystem = SimFilesystem::new();
1159 let path = Path::new("/t.csv");
1160 let file = filesystem.open(path, OpenMode::Create).expect("creates");
1161 file.write_at(0, text).expect("writes");
1162 drop(file);
1163 let file = filesystem.open(path, OpenMode::Read).expect("opens");
1164 Reader::open_sized(file, "/t.csv", Given::default(), block)
1165 }
1166
1167 #[test]
1171 fn generated_files_read_the_same_a_chunk_at_a_time_as_a_record_at_a_time() {
1172 let types = [
1173 LogicalType::BigInt,
1174 LogicalType::Integer,
1175 LogicalType::SmallInt,
1176 LogicalType::TinyInt,
1177 LogicalType::UBigInt,
1178 LogicalType::UInteger,
1179 LogicalType::USmallInt,
1180 LogicalType::UTinyInt,
1181 LogicalType::Double,
1182 LogicalType::Float,
1183 LogicalType::Date,
1184 LogicalType::Boolean,
1185 LogicalType::Varchar,
1186 LogicalType::Timestamp,
1187 LogicalType::Decimal { width: 18, scale: 3 },
1188 ];
1189 let mut rng = Rng(0x2545_f491_4f6c_dd1d);
1190 for case in 0..200 {
1191 let rows = if case % 100 == 0 { 8192 + rng.below(1000) } else { rng.below(200) };
1192 let mut text = typed_file(&mut rng, rows).into_bytes();
1193 if rng.below(20) == 0 {
1194 let at = rng.below(text.len() + 1);
1196 text.splice(at..at, *b",\"x\"y,");
1197 }
1198 for block in [1 << 20, 32 + rng.below(400)] {
1199 let (Ok(mut new), Ok(mut old)) =
1200 (open_sized(&text, block), open_sized(&text, block))
1201 else {
1202 continue;
1203 };
1204 assert_eq!(new.fields(), old.fields());
1205 let width = new.fields().len();
1206 if width > 0 && rng.below(2) == 0 {
1207 let columns: Vec<usize> =
1208 (0..rng.below(width + 2)).map(|_| rng.below(width)).collect();
1209 new.project(&columns).expect("projects");
1210 old.project(&columns).expect("projects");
1211 let wanted: Vec<LogicalType> =
1212 columns.iter().map(|_| types[rng.below(types.len())].clone()).collect();
1213 new.retype(&wanted).expect("retypes");
1214 old.retype(&wanted).expect("retypes");
1215 }
1216 let expected = drained(&mut old, true);
1217 let found = drained(&mut new, false);
1218 assert_eq!(found.1, expected.1, "case {case}, block {block}");
1219 assert_eq!(found.0, expected.0, "case {case}, block {block}");
1220 }
1221 }
1222 }
1223
1224 #[test]
1229 #[ignore = "a measurement, run by hand"]
1230 fn reads_lineitem_faster_a_chunk_at_a_time() {
1231 use rudb_io::RealFilesystem;
1232 use std::time::Instant;
1233
1234 const ROWS: usize = 200_000;
1235 let mut rng = Rng(0x1234_5678_9abc_def1);
1236 let mut text = String::from(
1237 "l_orderkey,l_partkey,l_suppkey,l_linenumber,l_quantity,l_extendedprice,l_discount,\
1238 l_tax,l_returnflag,l_linestatus,l_shipdate,l_commitdate,l_receiptdate,\
1239 l_shipinstruct,l_shipmode,l_comment\n",
1240 );
1241 let words = ["carefully", "final", "deposits", "furiously", "regular", "ideas", "sleep"];
1242 for row in 0..ROWS {
1243 let n = rng.next();
1244 let date = |shift: u64| {
1245 format!(
1246 "{}-{:02}-{:02}",
1247 1992 + (n >> shift) % 7,
1248 1 + (n >> shift) % 12,
1249 1 + (n >> shift) % 28
1250 )
1251 };
1252 let comment: Vec<&str> =
1253 (0..3 + n % 4).map(|k| words[((n >> (k * 3)) % 7) as usize]).collect();
1254 text.push_str(&format!(
1255 "{},{},{},{},{}.00,{}.{:02},0.0{},0.0{},{},{},{},{},{},{},{},{}\n",
1256 row / 4 + 1,
1257 n % 200_000,
1258 n % 10_000,
1259 row % 4 + 1,
1260 1 + n % 50,
1261 900 + n % 100_000,
1262 n % 100,
1263 n % 10,
1264 (n >> 7) % 9,
1265 ["A", "N", "R"][(n % 3) as usize],
1266 ["O", "F"][(n % 2) as usize],
1267 date(3),
1268 date(11),
1269 date(19),
1270 ["DELIVER IN PERSON", "NONE", "TAKE BACK RETURN"][(n % 3) as usize],
1271 ["TRUCK", "MAIL", "AIR", "SHIP"][(n % 4) as usize],
1272 comment.join(" "),
1273 ));
1274 }
1275 let path = std::env::temp_dir().join(format!("rudb-lineitem-{}.csv", std::process::id()));
1276 std::fs::write(&path, &text).expect("writes");
1277 let megabytes = text.len() as f64 / 1e6;
1278 let filesystem = RealFilesystem::new();
1279 let mut best = [f64::MAX; 2];
1280 for _ in 0..5 {
1281 for (slot, old) in [(0, true), (1, false)] {
1282 let file = filesystem.open(&path, OpenMode::Read).expect("opens");
1283 let mut reader = Reader::open(file, "lineitem.csv").expect("sniffs");
1284 let started = Instant::now();
1285 let mut rows = 0;
1286 loop {
1287 let next =
1288 if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
1289 let Some(chunk) = next.expect("reads") else { break };
1290 rows += chunk.len();
1291 }
1292 assert_eq!(rows, ROWS);
1293 best[slot] = best[slot].min(started.elapsed().as_secs_f64());
1294 }
1295 }
1296 std::fs::remove_file(&path).expect("removes");
1297 let (old, new) = (megabytes / best[0], megabytes / best[1]);
1298 println!(
1299 "{ROWS} rows, {megabytes:.1} MB: a record at a time {old:.1} MB/s, a chunk at a time"
1300 );
1301 println!("{new:.1} MB/s, {:.2}x", best[0] / best[1]);
1302 }
1303
1304 fn all_or_error(reader: &mut Reader) -> Result<Vec<Vec<Value>>> {
1305 let mut rows = Vec::new();
1306 while let Some(chunk) = reader.next_chunk()? {
1307 for row in 0..chunk.len() {
1308 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
1309 }
1310 }
1311 Ok(rows)
1312 }
1313}