1use rudb_common::{Error, Field, LogicalType, Result, Value};
13use rudb_io::File;
14use rudb_kernels::cast_value;
15use rudb_vector::{Chunk, VECTOR_SIZE, Vector};
16
17use crate::dialect::{self, Dialect};
18use crate::infer;
19
20const BLOCK: usize = 1 << 20;
26
27#[derive(Debug)]
29pub struct Reader {
30 file: Box<dyn File>,
31 path: String,
32 dialect: Dialect,
33 fields: Vec<Field>,
34 projection: Vec<usize>,
35 buffer: Vec<u8>,
36 at: usize,
37 offset: u64,
38 drained: bool,
39 line: u64,
40 scratch: Vec<String>,
41}
42
43impl Reader {
44 pub fn open(file: Box<dyn File>, path: &str) -> Result<Self> {
53 let mut reader = Self {
54 file,
55 path: path.to_string(),
56 dialect: Dialect::comma_separated(),
57 fields: Vec::new(),
58 projection: Vec::new(),
59 buffer: Vec::new(),
60 at: 0,
61 offset: 0,
62 drained: false,
63 line: 1,
64 scratch: Vec::new(),
65 };
66 reader.fill()?;
67 let sample = reader.buffer.clone();
68 let quote = dialect::quote(&sample);
69 let delimiter = dialect::delimiter(&sample, quote)?;
70 reader.dialect = Dialect { delimiter, quote, escape: quote, header: false };
71 let rows = reader.sample_rows(&sample)?;
72 let (header, fields) = describe(&rows);
73 reader.dialect.header = header;
74 reader.fields = fields;
75 reader.projection = (0..reader.fields.len()).collect();
76 if header {
77 reader.skip_record()?;
78 }
79 Ok(reader)
80 }
81
82 #[must_use]
84 pub fn fields(&self) -> Vec<Field> {
85 self.projection.iter().map(|&at| self.fields[at].clone()).collect()
86 }
87
88 pub fn project(&mut self, columns: &[usize]) -> Result<()> {
94 for &column in columns {
95 if column >= self.fields.len() {
96 return Err(Error::io(format!(
97 "column {column} is past the {} the file has",
98 self.fields.len()
99 )));
100 }
101 }
102 self.projection = columns.to_vec();
103 Ok(())
104 }
105
106 pub fn retype(&mut self, types: &[LogicalType]) -> Result<()> {
124 if types.len() != self.projection.len() {
125 return Err(Error::io(format!(
126 "{} types for a projection of {} columns",
127 types.len(),
128 self.projection.len()
129 )));
130 }
131 for (&at, ty) in self.projection.iter().zip(types) {
132 self.fields[at].ty = ty.clone();
133 }
134 Ok(())
135 }
136
137 #[must_use]
139 pub const fn dialect(&self) -> Dialect {
140 self.dialect
141 }
142
143 pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
150 let mut rows: Vec<Vec<Option<String>>> = Vec::new();
151 while rows.len() < VECTOR_SIZE {
152 match self.next_record()? {
153 Some(fields) => rows.push(fields),
154 None => break,
155 }
156 }
157 if rows.is_empty() {
158 return Ok(None);
159 }
160 let mut columns = Vec::with_capacity(self.projection.len());
161 for &at in &self.projection {
162 let field = &self.fields[at];
163 let mut values = Vec::with_capacity(rows.len());
164 for (row, held) in rows.iter().enumerate() {
165 let text = held.get(at).and_then(Option::as_deref);
166 values.push(self.convert(
167 text,
168 field,
169 self.line - rows.len() as u64 + row as u64,
170 )?);
171 }
172 columns.push(Vector::from_values(field.ty.clone(), &values)?);
173 }
174 Ok(Some(Chunk::with_rows(columns, rows.len())?))
175 }
176
177 fn convert(&self, text: Option<&str>, field: &Field, line: u64) -> Result<Value> {
179 let Some(text) = text else { return Ok(Value::Null) };
180 if field.ty == LogicalType::Varchar {
181 return Ok(Value::Varchar(text.to_string()));
182 }
183 let value = Value::Varchar(text.to_string());
184 match cast_value(&value, &field.ty, false) {
185 Ok(converted) => Ok(converted),
186 Err(_) => Err(Error::conversion(self.conversion_error(text, field, line))),
187 }
188 }
189
190 fn conversion_error(&self, text: &str, field: &Field, line: u64) -> String {
197 format!(
198 "CSV Error on Line: {line}\nOriginal Line: {text}\nError when converting column \
199 \"{}\". Could not convert string \"{text}\" to '{}'\n\nColumn {} is being converted \
200 as type {}\nThis type was auto-detected from the CSV file.\nPossible solutions:\n* \
201 Override the type for this column manually by setting the type explicitly, e.g., \
202 types={{'{}': 'VARCHAR'}}\n* Set the sample size to a larger value to enable the \
203 auto-detection to scan more values, e.g., sample_size=-1\n* Use a COPY statement to \
204 automatically derive types from an existing table.\n* Check whether the null string \
205 value is set correctly (e.g., nullstr = 'N/A')\n\n file = {}\n delimiter = {} \
206 (Auto-Detected)\n quote = {} (Auto-Detected)\n escape = {} (Auto-Detected)\n \
207 header = {} (Auto-Detected)\n sample_size = {}\n",
208 field.name,
209 field.ty,
210 field.name,
211 field.ty,
212 field.name,
213 self.path,
214 Dialect::shown(Some(self.dialect.delimiter)),
215 Dialect::shown(self.dialect.quote),
216 Dialect::shown(self.dialect.escape),
217 self.dialect.header,
218 infer::SAMPLE,
219 )
220 }
221
222 fn next_record(&mut self) -> Result<Option<Vec<Option<String>>>> {
224 let Some(()) = self.advance()? else { return Ok(None) };
225 Ok(Some(
226 self.scratch
227 .iter()
228 .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
229 .collect(),
230 ))
231 }
232
233 fn advance(&mut self) -> Result<Option<()>> {
235 loop {
236 let mut scratch = std::mem::take(&mut self.scratch);
237 let outcome = crate::scan::record(
238 &self.buffer,
239 self.at,
240 self.dialect,
241 self.drained,
242 &mut scratch,
243 );
244 self.scratch = scratch;
245 match outcome? {
246 Some(next) => {
247 self.at = next;
248 self.line += 1;
249 return Ok(Some(()));
250 }
251 None if self.drained => return Ok(None),
252 None => self.fill()?,
253 }
254 }
255 }
256
257 fn skip_record(&mut self) -> Result<()> {
259 self.advance()?;
260 Ok(())
261 }
262
263 fn fill(&mut self) -> Result<()> {
265 self.buffer.drain(..self.at);
266 self.at = 0;
267 let held = self.buffer.len();
268 self.buffer.resize(held + BLOCK, 0);
269 let read = self.file.read_at(self.offset, &mut self.buffer[held..])?;
270 self.buffer.truncate(held + read);
271 self.offset += read as u64;
272 if read == 0 {
273 self.drained = true;
274 }
275 Ok(())
276 }
277
278 fn sample_rows(&self, sample: &[u8]) -> Result<Vec<Vec<Option<String>>>> {
281 let mut rows = Vec::new();
282 let mut fields = Vec::new();
283 let mut at = 0;
284 while rows.len() <= infer::SAMPLE {
285 let Some(next) = crate::scan::record(sample, at, self.dialect, false, &mut fields)?
288 else {
289 break;
290 };
291 at = next;
292 rows.push(
293 fields
294 .iter()
295 .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
296 .collect(),
297 );
298 }
299 Ok(rows)
300 }
301}
302
303fn describe(rows: &[Vec<Option<String>>]) -> (bool, Vec<Field>) {
310 let width = rows.iter().map(Vec::len).max().unwrap_or(0);
311 let body = types(&rows[1.min(rows.len())..], width);
312 let all_text = body.iter().all(|ty| *ty == LogicalType::Varchar);
313 let first_fits = rows.first().is_some_and(|first| {
314 first.iter().zip(&body).all(|(text, ty)| match text {
315 None => true,
316 Some(text) => infer::fits(text, ty),
317 })
318 });
319 let header = rows.len() > 1 && (all_text || !first_fits);
320 if !header {
321 let types = types(rows, width);
322 let fields = types
323 .into_iter()
324 .enumerate()
325 .map(|(at, ty)| Field::new(format!("column{at}"), ty))
326 .collect();
327 return (false, fields);
328 }
329 let names = unique(&rows[0], width);
330 let fields = body.into_iter().zip(names).map(|(ty, name)| Field::new(name, ty)).collect();
331 (true, fields)
332}
333
334fn unique(header: &[Option<String>], width: usize) -> Vec<String> {
347 let mut taken: Vec<String> = Vec::with_capacity(width);
348 for at in 0..width {
349 let base = match header.get(at).and_then(Option::as_deref) {
350 Some(written) => written.to_string(),
351 None => format!("column{at}"),
352 };
353 let mut name = base.clone();
354 let mut next = 1;
355 while taken.iter().any(|held| held.eq_ignore_ascii_case(&name)) {
356 name = format!("{base}_{next}");
357 next += 1;
358 }
359 taken.push(name);
360 }
361 taken
362}
363
364fn types(rows: &[Vec<Option<String>>], width: usize) -> Vec<LogicalType> {
366 (0..width)
367 .map(|at| {
368 let values: Vec<Option<&str>> =
369 rows.iter().map(|row| row.get(at).and_then(Option::as_deref)).collect();
370 infer::column(&values)
371 })
372 .collect()
373}
374
375#[cfg(test)]
376mod tests {
377 use super::*;
378 use rudb_io::{Filesystem, OpenMode, SimFilesystem};
379 use std::path::Path;
380
381 fn read(text: &str) -> Reader {
382 let filesystem = SimFilesystem::new();
383 let path = Path::new("/t.csv");
384 let file = filesystem.open(path, OpenMode::Create).expect("creates");
385 file.write_at(0, text.as_bytes()).expect("writes");
386 drop(file);
387 let file = filesystem.open(path, OpenMode::Read).expect("opens");
388 Reader::open(file, "/t.csv").expect("sniffs")
389 }
390
391 fn names_and_types(reader: &Reader) -> Vec<(String, String)> {
392 reader.fields().iter().map(|f| (f.name.clone(), f.ty.to_string())).collect()
393 }
394
395 fn all(reader: &mut Reader) -> Vec<Vec<Value>> {
396 let mut rows = Vec::new();
397 while let Some(chunk) = reader.next_chunk().expect("reads") {
398 for row in 0..chunk.len() {
399 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
400 }
401 }
402 rows
403 }
404
405 #[test]
406 fn a_header_that_names_two_columns_the_same_thing_counts_the_second_one_up() {
407 let names: Vec<String> =
408 read("a,a,A,a_1\n1,2,3,4\nx,y,z,w\n").fields().into_iter().map(|f| f.name).collect();
409 assert_eq!(names, ["a", "a_1", "A_2", "a_1_1"]);
413 }
414
415 #[test]
416 fn a_header_and_three_types_are_what_duckdb_sniffs_for_the_same_bytes() {
417 let reader = read("a,b,c\n1,x,2.5\n2,y,3.5\n");
418 assert_eq!(
419 names_and_types(&reader),
420 [
421 ("a".to_string(), "BIGINT".to_string()),
422 ("b".to_string(), "VARCHAR".to_string()),
423 ("c".to_string(), "DOUBLE".to_string()),
424 ]
425 );
426 }
427
428 #[test]
429 fn a_file_with_no_header_gets_the_names_duckdb_gives_it() {
430 let reader = read("1,x\n2,y\n");
431 assert_eq!(
432 names_and_types(&reader),
433 [
434 ("column0".to_string(), "BIGINT".to_string()),
435 ("column1".to_string(), "VARCHAR".to_string()),
436 ]
437 );
438 }
439
440 #[test]
441 fn two_rows_of_words_are_a_header_and_a_row() {
442 let reader = read("a,b\nc,d\n");
443 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "b"]);
444 }
445
446 #[test]
447 fn one_column_of_words_under_a_row_of_numbers_is_still_a_header() {
448 let reader = read("a,2\n3,4\n");
451 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "2"]);
452 }
453
454 #[test]
455 fn the_rows_are_the_rows_of_the_file() {
456 let mut reader = read("a,b\n1,x\n2,y\n");
457 assert_eq!(
458 all(&mut reader),
459 [
460 vec![Value::BigInt(1), Value::Varchar("x".into())],
461 vec![Value::BigInt(2), Value::Varchar("y".into())],
462 ]
463 );
464 }
465
466 #[test]
467 fn an_empty_field_is_a_null_whether_it_was_quoted_or_not() {
468 let mut reader = read("a,b\n1,\n\"\",y\n");
471 assert_eq!(
472 all(&mut reader),
473 [vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::Varchar("y".into())],]
474 );
475 }
476
477 #[test]
478 fn a_projection_picks_columns_out_by_position_and_can_reorder_them() {
479 let mut reader = read("a,b,c\n1,x,2.5\n");
480 reader.project(&[2, 0]).expect("projects");
481 assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["c", "a"]);
482 assert_eq!(all(&mut reader), [vec![Value::Double(2.5), Value::BigInt(1)]]);
483 }
484
485 #[test]
486 fn a_projection_of_nothing_still_counts_the_rows() {
487 let mut reader = read("a,b\n1,x\n2,y\n3,z\n");
488 reader.project(&[]).expect("projects");
489 let chunk = reader.next_chunk().expect("reads").expect("a chunk");
490 assert_eq!(chunk.len(), 3);
491 assert_eq!(chunk.width(), 0);
492 }
493
494 #[test]
495 fn a_pipe_separated_file_reads_as_one() {
496 let mut reader = read("a|b\n1|x\n");
497 assert_eq!(reader.dialect().delimiter, b'|');
498 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x".into())]]);
499 }
500
501 #[test]
502 fn a_quoted_field_with_a_delimiter_in_it_is_one_value() {
503 let mut reader = read("a,b\n1,\"x,y\"\n");
504 assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
505 }
506
507 #[test]
508 fn more_rows_than_fit_one_chunk_arrive_as_more_than_one_chunk() {
509 let mut text = String::from("a\n");
510 for row in 0..VECTOR_SIZE + 5 {
511 text.push_str(&format!("{row}\n"));
512 }
513 let mut reader = read(&text);
514 let first = reader.next_chunk().expect("reads").expect("a chunk");
515 assert_eq!(first.len(), VECTOR_SIZE);
516 let second = reader.next_chunk().expect("reads").expect("a second chunk");
517 assert_eq!(second.len(), 5);
518 assert!(reader.next_chunk().expect("reads").is_none());
519 }
520
521 #[test]
522 fn a_value_the_sniffer_never_saw_is_an_error_rather_than_a_wider_column() {
523 let mut text = String::from("c\n");
527 for row in 0..infer::SAMPLE {
528 text.push_str(&format!("{row}\n"));
529 }
530 text.push_str("oops\n");
531 let mut reader = read(&text);
532 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
533 let error = all_or_error(&mut reader).unwrap_err();
534 let line = infer::SAMPLE + 2;
535 assert!(error.message().starts_with(&format!("CSV Error on Line: {line}")), "{error}");
536 assert!(
537 error.message().contains("Could not convert string \"oops\" to 'BIGINT'"),
538 "{error}"
539 );
540 assert!(error.message().contains("sample_size = 20480"), "{error}");
541 }
542
543 #[test]
544 fn a_file_told_a_wider_type_than_it_sniffed_reads_its_whole_numbers_as_that_type() {
545 let mut reader = read("a\n1\n2\n");
549 assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
550 reader.retype(&[LogicalType::Double]).expect("one type for one column");
551 assert_eq!(reader.fields()[0].ty, LogicalType::Double);
552 assert_eq!(all(&mut reader), [[Value::Double(1.0)], [Value::Double(2.0)]]);
553 }
554
555 #[test]
556 fn a_type_list_that_is_not_as_long_as_the_projection_is_refused() {
557 let mut reader = read("a,b\n1,two\n");
558 let error = reader.retype(&[LogicalType::Double]).unwrap_err();
559 assert!(error.message().contains("1 types for a projection of 2 columns"), "{error}");
560 }
561
562 fn all_or_error(reader: &mut Reader) -> Result<Vec<Vec<Value>>> {
563 let mut rows = Vec::new();
564 while let Some(chunk) = reader.next_chunk()? {
565 for row in 0..chunk.len() {
566 rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
567 }
568 }
569 Ok(rows)
570 }
571}