Skip to main content

rudb_csv/
reader.rs

1//! A CSV file as chunks.
2//!
3//! The same shape as `rudb-parquet`'s reader on purpose, because the operator above them is the same
4//! operator with a different constructor: open, ask what the columns are, say which of them you
5//! want, then pull chunks until there are none. A caller that can read one can read the other.
6//!
7//! The file is read in blocks and a record that straddles a block boundary is carried into the next
8//! one, so a file larger than memory reads the same as a small one. The sample the sniffer looks at
9//! is the first block, which is also the first block the reader then goes on to use, so opening a
10//! file reads its front once.
11//!
12//! A reader can also cover a stretch of a file rather than all of it, which is what [`crate::split`]
13//! hands each thread. It starts where it is told a record starts and takes the records that start
14//! before its end, reading past the end only to finish the last of them.
15
16use 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
27/// How much is read at a time, and how much the sniffer gets to look at.
28///
29/// A megabyte holds well over the twenty thousand rows of the sample for any file with ordinary
30/// rows in it, and for a file with enormous rows the sniffer sees fewer of them and says so by
31/// getting a wider type rather than by failing.
32const BLOCK: usize = 1 << 20;
33
34/// How much is read at a time once a reader of a stretch is past its end.
35///
36/// All it is finishing there is the one record that crosses the end, which is usually a line, so a
37/// whole block would be read to be thrown away. Each read after the first is as large as what is
38/// held, so a record that turns out to be long still arrives in a handful of reads.
39const TAIL: usize = 64 << 10;
40
41/// How many rows of a chunk are converted, every projected column of them, before the next.
42///
43/// A row's ranges are eight bytes a field, so for a sixteen column table like `lineitem` a thousand
44/// rows is an eighth of a megabyte of ranges and about as much again of the bytes they point at,
45/// which fits in a core's second level cache with room for the columns being written.
46const BLOCK_ROWS: usize = 1024;
47
48/// A CSV file, positioned at a record boundary.
49#[derive(Debug)]
50pub struct Reader {
51    file: Arc<dyn File>,
52    path: String,
53    given: Given,
54    dialect: Dialect,
55    fields: Vec<Field>,
56    projection: Vec<usize>,
57    buffer: Vec<u8>,
58    at: usize,
59    offset: u64,
60    drained: bool,
61    line: u64,
62    scratch: Vec<String>,
63    records: Records,
64    block: usize,
65    /// Where the reading started, so that what it has read is the offset less this.
66    origin: u64,
67    /// The first byte a record may not start at, which is the end of the file for a whole one.
68    end: u64,
69    /// Where a reader gives up rather than reading on to finish a record, which only a guess
70    /// sets. See [`crate::split`].
71    cap: u64,
72}
73
74impl Reader {
75    /// Opens a file, works out how it is written, and positions it at the first row.
76    ///
77    /// The path is kept because the error a bad value produces names it, the way DuckDB's does.
78    ///
79    /// # Errors
80    ///
81    /// When the file cannot be read, and when the first block of it does not hold one whole record,
82    /// which is a single line longer than a megabyte and is not a CSV file anybody meant to write.
83    pub fn open(file: Box<dyn File>, path: &str) -> Result<Self> {
84        Self::open_with(file, path, Given::default())
85    }
86
87    /// The same, with whatever the caller already knows about how the file is written.
88    ///
89    /// This is where `read_csv('f.csv', delim=';', header=true)` arrives. A given value replaces the
90    /// sniffer's answer rather than seeding it, and it replaces it before the sample is split, so the
91    /// types and the column names come out of the file read the way the caller said it is written.
92    ///
93    /// # Errors
94    ///
95    /// Everything [`Reader::open`] reports.
96    pub fn open_with(file: Box<dyn File>, path: &str, given: Given) -> Result<Self> {
97        Self::open_sized(file, path, given, BLOCK)
98    }
99
100    /// The same, reading `block` bytes at a time, which the tests make small so that a record
101    /// crosses a refill every few lines rather than once a megabyte.
102    pub(crate) fn open_sized(
103        file: Box<dyn File>,
104        path: &str,
105        given: Given,
106        block: usize,
107    ) -> Result<Self> {
108        let mut reader = Self {
109            file: Arc::from(file),
110            path: path.to_string(),
111            given,
112            dialect: Dialect::comma_separated(),
113            fields: Vec::new(),
114            projection: Vec::new(),
115            buffer: Vec::new(),
116            at: 0,
117            offset: 0,
118            drained: false,
119            line: 1,
120            scratch: Vec::new(),
121            records: Records::default(),
122            block,
123            origin: 0,
124            end: u64::MAX,
125            cap: u64::MAX,
126        };
127        reader.fill(0)?;
128        let sample = reader.buffer.clone();
129        let quote = given.quote.or_else(|| dialect::quote(&sample));
130        let delimiter = match given.delimiter {
131            Some(byte) => byte,
132            None => dialect::delimiter(&sample, quote)?,
133        };
134        let escape = given.escape.or(quote);
135        reader.dialect = Dialect { delimiter, quote, escape, header: false };
136        let rows = reader.sample_rows(&sample)?;
137        let (header, fields) = describe(&rows, given.header);
138        reader.dialect.header = header;
139        reader.fields = fields;
140        reader.projection = (0..reader.fields.len()).collect();
141        if header {
142            reader.skip_record()?;
143        }
144        Ok(reader)
145    }
146
147    /// The columns this reader will produce, in order.
148    #[must_use]
149    pub fn fields(&self) -> Vec<Field> {
150        self.projection.iter().map(|&at| self.fields[at].clone()).collect()
151    }
152
153    /// Reads only these columns, by position in the file, in this order.
154    ///
155    /// # Errors
156    ///
157    /// When a position is past the end of the file's columns.
158    pub fn project(&mut self, columns: &[usize]) -> Result<()> {
159        for &column in columns {
160            if column >= self.fields.len() {
161                return Err(Error::io(format!(
162                    "column {column} is past the {} the file has",
163                    self.fields.len()
164                )));
165            }
166        }
167        self.projection = columns.to_vec();
168        Ok(())
169    }
170
171    /// Reads the projected columns as these types rather than as the ones the sample chose.
172    ///
173    /// A read that covers several files produces one stream and a stream has one schema, and no
174    /// single file's sample is that schema. Every file is sniffed on its own and the answers are
175    /// combined by [`crate::across`], so each file is then told what the whole read settled on,
176    /// including the first one. Without it a file whose column happens to hold nothing but whole
177    /// numbers hands up a BIGINT column into a stream that is DOUBLE because some other file in the
178    /// set held a decimal.
179    ///
180    /// This is not a cast of what was read. The type is what the text is converted with, so saying
181    /// it before any row is read converts once rather than converting to the wrong type and again to
182    /// the right one. A value that then does not fit is the conversion error, named and lined the
183    /// way any other one is.
184    ///
185    /// # Errors
186    ///
187    /// When the list is not as long as the projection.
188    pub fn retype(&mut self, types: &[LogicalType]) -> Result<()> {
189        if types.len() != self.projection.len() {
190            return Err(Error::io(format!(
191                "{} types for a projection of {} columns",
192                types.len(),
193                self.projection.len()
194            )));
195        }
196        for (&at, ty) in self.projection.iter().zip(types) {
197            self.fields[at].ty = ty.clone();
198        }
199        Ok(())
200    }
201
202    /// How this file is punctuated, which is what the sniffer decided.
203    #[must_use]
204    pub const fn dialect(&self) -> Dialect {
205        self.dialect
206    }
207
208    /// The next chunk, or `None` at the end of the file.
209    ///
210    /// The records are split first, all of them, and converted afterwards a column at a time, which
211    /// is also the order the errors come out in: a malformed record anywhere in the chunk is
212    /// reported ahead of a value that does not convert, and among values the first projected
213    /// column's first bad row is the one named.
214    ///
215    /// # Errors
216    ///
217    /// A read error, a malformed record, or a value that does not fit the type the sample chose
218    /// for its column.
219    pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
220        let rows = self.next_records()?;
221        if rows == 0 {
222            return Ok(None);
223        }
224        let first = self.line;
225        self.line += rows as u64;
226        let cells = Cells { bytes: &self.buffer, records: &self.records, dialect: self.dialect };
227        let projected: Vec<_> =
228            self.projection.iter().map(|&at| (at, &self.fields[at].ty)).collect();
229        let mut builders = convert::builders(&cells, &projected);
230        let mut start = 0;
231        while start < rows {
232            let end = rows.min(start + BLOCK_ROWS);
233            for (build, &at) in builders.iter_mut().zip(&self.projection) {
234                let field = &self.fields[at];
235                let refuse = |text: &str, row: usize| {
236                    Error::conversion(self.conversion_error(text, field, first + row as u64))
237                };
238                if let Err(error) = build.rows(&cells, at, start..end, &refuse) {
239                    return Err(self.first_bad_value(&cells, first).unwrap_or(error));
240                }
241            }
242            start = end;
243        }
244        let columns = builders.into_iter().map(|build| build.finish()).collect::<Result<_>>()?;
245        Ok(Some(Chunk::with_rows(columns, rows)?))
246    }
247
248    /// The error for the first projected column's first value that does not convert, found by
249    /// converting the chunk a column at a time the way it used to be.
250    ///
251    /// A chunk is converted a block of rows at a time, so the first value it trips on can be in a
252    /// later column than one with a bad value further down. This is only run once one has been
253    /// found, to name the same value the reader always named.
254    fn first_bad_value(&self, cells: &Cells<'_>, first: u64) -> Option<Error> {
255        self.projection.iter().find_map(|&at| {
256            let field = &self.fields[at];
257            let refuse = |text: &str, row: usize| {
258                Error::conversion(self.conversion_error(text, field, first + row as u64))
259            };
260            convert::column(cells, at, &field.ty, &refuse).err()
261        })
262    }
263
264    /// Splits the next chunk's worth of records into [`Self::records`] and answers how many there
265    /// are, which is none at the end.
266    ///
267    /// The end of the file is the end for a whole file. A reader of a stretch stops at the first
268    /// record that starts at or after its end, and it cannot see where a record starts until it has
269    /// split it, so a chunk that went past the end is split again a record at a time. That happens
270    /// once per stretch, on its last chunk.
271    fn next_records(&mut self) -> Result<usize> {
272        self.records.clear();
273        let mut start = self.at;
274        let mut careful = false;
275        loop {
276            if self.here() >= self.end {
277                break;
278            }
279            let limit = if careful { self.records.len() + 1 } else { VECTOR_SIZE };
280            self.at = crate::scan::records(
281                &self.buffer,
282                self.at,
283                self.dialect,
284                self.drained,
285                limit,
286                &mut self.records,
287            )?;
288            if !careful && self.here() > self.end {
289                self.records.clear();
290                self.at = start;
291                careful = true;
292                continue;
293            }
294            if self.records.len() == VECTOR_SIZE {
295                break;
296            }
297            if careful && self.records.len() == limit {
298                continue;
299            }
300            if self.drained {
301                break;
302            }
303            // The records already read point into the buffer, so the refill keeps everything from
304            // the start of the chunk and moves their ranges down by whatever it dropped in front.
305            // A range only reaches so far, so a chunk that would outgrow that ends early, and a
306            // single record that would is refused.
307            if self.buffer.len() - start + self.block > Span::MOST {
308                if self.records.is_empty() {
309                    return Err(Error::io("a record is longer than two gigabytes"));
310                }
311                break;
312            }
313            self.fill(start)?;
314            self.records.shift(start);
315            start = 0;
316        }
317        Ok(self.records.len())
318    }
319
320    /// Where in the file the next record starts.
321    pub(crate) fn here(&self) -> u64 {
322        self.offset - self.buffer.len() as u64 + self.at as u64
323    }
324
325    /// How many bytes of the file this reader has read, sniffing included.
326    #[must_use]
327    pub fn bytes_read(&self) -> u64 {
328        self.offset - self.origin
329    }
330
331    /// A reader over the same file with the same answers, that starts at `from`, which has to be
332    /// where a record starts, and takes the records that start before `end`.
333    ///
334    /// `line` is what the line of the record at `from` is called in an error, which is right only
335    /// for a caller that knows it. [`crate::split`] does not, and so never shows the error such a
336    /// reader makes.
337    pub(crate) fn stretch(&self, from: u64, end: u64, line: u64) -> Self {
338        Self {
339            file: Arc::clone(&self.file),
340            path: self.path.clone(),
341            given: self.given,
342            dialect: self.dialect,
343            fields: self.fields.clone(),
344            projection: self.projection.clone(),
345            buffer: Vec::new(),
346            at: 0,
347            offset: from,
348            drained: false,
349            line,
350            scratch: Vec::new(),
351            records: Records::default(),
352            block: self.block,
353            origin: from,
354            end,
355            cap: u64::MAX,
356        }
357    }
358
359    /// Makes the reader give up once it has read up to `cap`, with an error nobody is shown.
360    pub(crate) fn give_up_at(&mut self, cap: u64) {
361        self.cap = cap;
362    }
363
364    /// Splits records to the end without converting any of them, and answers where the first
365    /// record not taken starts.
366    pub(crate) fn skim(&mut self) -> Result<u64> {
367        while self.next_records()? > 0 {}
368        Ok(self.here())
369    }
370
371    /// Steps over the chunks that end at or before `target` without converting them, counting their
372    /// lines, so that the chunks after are the ones and the lines a whole read would give.
373    ///
374    /// The chunk that runs past `target` is left to be read, since it is the same chunk a whole read
375    /// would convert and a value in it may be the one that fails.
376    pub(crate) fn skip_to(&mut self, target: u64) -> Result<()> {
377        loop {
378            let from = self.here();
379            let rows = self.next_records()?;
380            if rows == 0 {
381                return Ok(());
382            }
383            if self.here() > target {
384                let front = self.offset - self.buffer.len() as u64;
385                self.at = usize::try_from(from - front)
386                    .map_err(|_| Error::internal("a chunk start outside the buffer"))?;
387                return Ok(());
388            }
389            self.line += rows as u64;
390        }
391    }
392
393    /// The line the next record is on, as a conversion error counts them.
394    pub(crate) const fn line(&self) -> u64 {
395        self.line
396    }
397
398    /// The file this reads, for a caller that wants to look at bytes the reader has not.
399    pub(crate) fn file(&self) -> &dyn File {
400        self.file.as_ref()
401    }
402
403    /// DuckDB's message for a value that does not fit the type its column was sniffed as.
404    ///
405    /// Reproduced whole, including the block of settings at the bottom, because that block is the
406    /// answer to the question the message raises. Somebody reading it wants to know what was
407    /// guessed and how to override the guess, and a shorter message would send them to the
408    /// documentation to find out.
409    ///
410    /// A line of that block says where its value came from, and a value the call gave is `(Set By
411    /// User)` rather than `(Auto-Detected)`, measured on `v2.0.0-dev84237` by reading a file with
412    /// `delim=';'` past the sample. Telling somebody that what they wrote down was auto-detected is
413    /// the one thing the block could say that would send them looking in the wrong place.
414    fn conversion_error(&self, text: &str, field: &Field, line: u64) -> String {
415        format!(
416            "CSV Error on Line: {line}\nOriginal Line: {text}\nError when converting column \
417             \"{}\". Could not convert string \"{text}\" to '{}'\n\nColumn {} is being converted \
418             as type {}\nThis type was auto-detected from the CSV file.\nPossible solutions:\n* \
419             Override the type for this column manually by setting the type explicitly, e.g., \
420             types={{'{}': 'VARCHAR'}}\n* Set the sample size to a larger value to enable the \
421             auto-detection to scan more values, e.g., sample_size=-1\n* Use a COPY statement to \
422             automatically derive types from an existing table.\n* Check whether the null string \
423             value is set correctly (e.g., nullstr = 'N/A')\n\n  file = {}\n  delimiter = {}\n  \
424             quote = {}\n  escape = {}\n  header = {} {}\n  sample_size = {}\n",
425            field.name,
426            field.ty,
427            field.name,
428            field.ty,
429            field.name,
430            self.path,
431            Given::shown(self.given.delimiter, Some(self.dialect.delimiter)),
432            Given::shown(self.given.quote, self.dialect.quote),
433            Given::shown(self.given.escape, self.dialect.escape),
434            self.dialect.header,
435            Given::source(self.given.header.is_some()),
436            infer::SAMPLE,
437        )
438    }
439
440    /// The next chunk the way it was read before the records were split a chunk at a time, one
441    /// record into owned strings and one cell into a `Value` at a time.
442    #[cfg(test)]
443    fn next_chunk_by_record(&mut self) -> Result<Option<Chunk>> {
444        let mut rows: Vec<Vec<Option<String>>> = Vec::new();
445        while rows.len() < VECTOR_SIZE {
446            match self.next_record()? {
447                Some(fields) => rows.push(fields),
448                None => break,
449            }
450        }
451        if rows.is_empty() {
452            return Ok(None);
453        }
454        let mut columns = Vec::with_capacity(self.projection.len());
455        for &at in &self.projection {
456            let field = &self.fields[at];
457            let mut values = Vec::with_capacity(rows.len());
458            for (row, held) in rows.iter().enumerate() {
459                let text = held.get(at).and_then(Option::as_deref);
460                values.push(self.convert(
461                    text,
462                    field,
463                    self.line - rows.len() as u64 + row as u64,
464                )?);
465            }
466            columns.push(rudb_vector::Vector::from_values(field.ty.clone(), &values)?);
467        }
468        Ok(Some(Chunk::with_rows(columns, rows.len())?))
469    }
470
471    /// One value, cast from its text to the column's type.
472    #[cfg(test)]
473    fn convert(&self, text: Option<&str>, field: &Field, line: u64) -> Result<rudb_common::Value> {
474        let Some(text) = text else { return Ok(rudb_common::Value::Null) };
475        if field.ty == LogicalType::Varchar {
476            return Ok(rudb_common::Value::Varchar(text.to_string()));
477        }
478        let value = rudb_common::Value::Varchar(text.to_string());
479        match rudb_kernels::cast_value(&value, &field.ty, false) {
480            Ok(converted) => Ok(converted),
481            Err(_) => Err(Error::conversion(self.conversion_error(text, field, line))),
482        }
483    }
484
485    /// The next record, as one entry per field, with an empty field as a null.
486    ///
487    /// This and [`Reader::next_chunk_by_record`] are how chunks were read before
488    /// [`crate::scan::records`], kept for the tests to hold the new path to the old answers.
489    #[cfg(test)]
490    fn next_record(&mut self) -> Result<Option<Vec<Option<String>>>> {
491        let Some(()) = self.advance()? else { return Ok(None) };
492        Ok(Some(
493            self.scratch
494                .iter()
495                .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
496                .collect(),
497        ))
498    }
499
500    /// Reads one record into the scratch, filling the buffer when it has to.
501    fn advance(&mut self) -> Result<Option<()>> {
502        loop {
503            let mut scratch = std::mem::take(&mut self.scratch);
504            let outcome = crate::scan::record(
505                &self.buffer,
506                self.at,
507                self.dialect,
508                self.drained,
509                &mut scratch,
510            );
511            self.scratch = scratch;
512            match outcome? {
513                Some(next) => {
514                    self.at = next;
515                    self.line += 1;
516                    return Ok(Some(()));
517                }
518                None if self.drained => return Ok(None),
519                None => self.fill(self.at)?,
520            }
521        }
522    }
523
524    /// Reads one record and throws it away, which is what a header is.
525    fn skip_record(&mut self) -> Result<()> {
526        self.advance()?;
527        Ok(())
528    }
529
530    /// Drops the bytes before `keep` and reads another block onto the end.
531    ///
532    /// `keep` is where the first record anything still points at starts, which is the record being
533    /// read for one record at a time and the start of the chunk for a chunk.
534    fn fill(&mut self, keep: usize) -> Result<()> {
535        if self.offset >= self.cap {
536            return Err(Error::io("a record runs further than a guessed start is followed"));
537        }
538        self.buffer.drain(..keep);
539        self.at -= keep;
540        let held = self.buffer.len();
541        let want = match self.end.checked_sub(self.offset) {
542            Some(left) if left > 0 => self.block.min(usize::try_from(left).unwrap_or(usize::MAX)),
543            _ => self.block.min(TAIL.max(held)),
544        };
545        self.buffer.resize(held + want, 0);
546        let read = self.file.read_at(self.offset, &mut self.buffer[held..])?;
547        self.buffer.truncate(held + read);
548        self.offset += read as u64;
549        if read == 0 {
550            self.drained = true;
551        }
552        Ok(())
553    }
554
555    /// The records the sniffer gets to look at, which is the sample or the file, whichever is
556    /// shorter.
557    fn sample_rows(&self, sample: &[u8]) -> Result<Vec<Vec<Option<String>>>> {
558        let mut rows = Vec::new();
559        let mut fields = Vec::new();
560        let mut at = 0;
561        while rows.len() <= infer::SAMPLE {
562            // The end of the block is not the end of the file, so a record the block cut in half is
563            // simply not part of the sample.
564            let Some(next) = crate::scan::record(sample, at, self.dialect, false, &mut fields)?
565            else {
566                break;
567            };
568            at = next;
569            rows.push(
570                fields
571                    .iter()
572                    .map(|text| if text.is_empty() { None } else { Some(text.clone()) })
573                    .collect(),
574            );
575        }
576        Ok(rows)
577    }
578}
579
580/// Whether the first row is a header, and what the columns are called and typed.
581///
582/// The rule is DuckDB's and both halves of it were measured. A file whose columns are all `VARCHAR`
583/// once the first row is set aside has a header, because two rows of words is a header and a row.
584/// Otherwise the first row is a header exactly when it does not fit the types the rest of the file
585/// has, which is what makes `1,2` over `3,4` a file of two rows and `a,b` over `1,2` a file of one.
586///
587/// `told` is the caller answering the question instead, which is `header=true` or `header=false` on
588/// the call. It decides the names and the types as well as the row count, since a first row that is
589/// data is a row the types have to fit and a first row that is a header is not.
590fn describe(rows: &[Vec<Option<String>>], told: Option<bool>) -> (bool, Vec<Field>) {
591    let width = rows.iter().map(Vec::len).max().unwrap_or(0);
592    let body = types(&rows[1.min(rows.len())..], width);
593    let all_text = body.iter().all(|ty| *ty == LogicalType::Varchar);
594    let first_fits = rows.first().is_some_and(|first| {
595        first.iter().zip(&body).all(|(text, ty)| match text {
596            None => true,
597            Some(text) => infer::fits(text, ty),
598        })
599    });
600    // An empty file has no row to take names from, so it has no header whatever it was told.
601    let header = !rows.is_empty() && told.unwrap_or(rows.len() > 1 && (all_text || !first_fits));
602    if !header {
603        let types = types(rows, width);
604        let fields = types
605            .into_iter()
606            .enumerate()
607            .map(|(at, ty)| Field::new(format!("column{at}"), ty))
608            .collect();
609        return (false, fields);
610    }
611    let names = unique(&rows[0], width);
612    let fields = body.into_iter().zip(names).map(|(ty, name)| Field::new(name, ty)).collect();
613    (true, fields)
614}
615
616/// The column names a header row gives, with the collisions resolved the way DuckDB resolves them.
617///
618/// A header is text somebody typed and nothing stops it naming two columns the same thing, so the
619/// second one gets `_1`, and the count goes up until the name is free. It has to count rather than
620/// stop at one, because the suffix can collide too: a file whose header is `a,a,a_1` comes back as
621/// `a`, `a_1`, `a_1_1` from the binary, and it is the second column that took the name the third one
622/// was written with.
623///
624/// The comparison ignores case and the written case is kept, which was measured: `a,a,A` comes back
625/// as `a`, `a_1`, `A_2`, so `A` collided with `a` and then `A_1` collided with `a_1`. An empty
626/// header cell is a column with no name, and it falls back to the generated one rather than to an
627/// empty string that no query could write.
628fn unique(header: &[Option<String>], width: usize) -> Vec<String> {
629    let mut taken: Vec<String> = Vec::with_capacity(width);
630    for at in 0..width {
631        let base = match header.get(at).and_then(Option::as_deref) {
632            Some(written) => written.to_string(),
633            None => format!("column{at}"),
634        };
635        let mut name = base.clone();
636        let mut next = 1;
637        while taken.iter().any(|held| held.eq_ignore_ascii_case(&name)) {
638            name = format!("{base}_{next}");
639            next += 1;
640        }
641        taken.push(name);
642    }
643    taken
644}
645
646/// The type of each of `width` columns, over these rows.
647fn types(rows: &[Vec<Option<String>>], width: usize) -> Vec<LogicalType> {
648    (0..width)
649        .map(|at| {
650            let values: Vec<Option<&str>> =
651                rows.iter().map(|row| row.get(at).and_then(Option::as_deref)).collect();
652            infer::column(&values)
653        })
654        .collect()
655}
656
657#[cfg(test)]
658mod tests {
659    use super::*;
660    use rudb_common::Value;
661    use rudb_io::{Filesystem, OpenMode, SimFilesystem};
662    use std::path::Path;
663
664    fn read(text: &str) -> Reader {
665        let filesystem = SimFilesystem::new();
666        let path = Path::new("/t.csv");
667        let file = filesystem.open(path, OpenMode::Create).expect("creates");
668        file.write_at(0, text.as_bytes()).expect("writes");
669        drop(file);
670        let file = filesystem.open(path, OpenMode::Read).expect("opens");
671        Reader::open(file, "/t.csv").expect("sniffs")
672    }
673
674    /// The same file opened with something already known about how it is written.
675    fn read_with(text: &str, given: Given) -> Reader {
676        let filesystem = SimFilesystem::new();
677        let path = Path::new("/t.csv");
678        let file = filesystem.open(path, OpenMode::Create).expect("creates");
679        file.write_at(0, text.as_bytes()).expect("writes");
680        drop(file);
681        let file = filesystem.open(path, OpenMode::Read).expect("opens");
682        Reader::open_with(file, "/t.csv", given).expect("reads")
683    }
684
685    fn names_and_types(reader: &Reader) -> Vec<(String, String)> {
686        reader.fields().iter().map(|f| (f.name.clone(), f.ty.to_string())).collect()
687    }
688
689    fn all(reader: &mut Reader) -> Vec<Vec<Value>> {
690        let mut rows = Vec::new();
691        while let Some(chunk) = reader.next_chunk().expect("reads") {
692            for row in 0..chunk.len() {
693                rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
694            }
695        }
696        rows
697    }
698
699    #[test]
700    fn a_header_that_names_two_columns_the_same_thing_counts_the_second_one_up() {
701        let names: Vec<String> =
702            read("a,a,A,a_1\n1,2,3,4\nx,y,z,w\n").fields().into_iter().map(|f| f.name).collect();
703        // Measured against the binary, all four of them. The last one is the interesting one: the
704        // second column took `a_1`, which is the name the fourth column was written with, so the
705        // fourth has to keep counting from its own name rather than from `a`.
706        assert_eq!(names, ["a", "a_1", "A_2", "a_1_1"]);
707    }
708
709    #[test]
710    fn a_header_and_three_types_are_what_duckdb_sniffs_for_the_same_bytes() {
711        let reader = read("a,b,c\n1,x,2.5\n2,y,3.5\n");
712        assert_eq!(
713            names_and_types(&reader),
714            [
715                ("a".to_string(), "BIGINT".to_string()),
716                ("b".to_string(), "VARCHAR".to_string()),
717                ("c".to_string(), "DOUBLE".to_string()),
718            ]
719        );
720    }
721
722    #[test]
723    fn a_file_with_no_header_gets_the_names_duckdb_gives_it() {
724        let reader = read("1,x\n2,y\n");
725        assert_eq!(
726            names_and_types(&reader),
727            [
728                ("column0".to_string(), "BIGINT".to_string()),
729                ("column1".to_string(), "VARCHAR".to_string()),
730            ]
731        );
732    }
733
734    #[test]
735    fn two_rows_of_words_are_a_header_and_a_row() {
736        let reader = read("a,b\nc,d\n");
737        assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "b"]);
738    }
739
740    #[test]
741    fn one_column_of_words_under_a_row_of_numbers_is_still_a_header() {
742        // `a,2` over `3,4`. One column disagreeing is enough, and the second column is then named
743        // `2`, which is the text that was in it.
744        let reader = read("a,2\n3,4\n");
745        assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["a", "2"]);
746    }
747
748    #[test]
749    fn the_rows_are_the_rows_of_the_file() {
750        let mut reader = read("a,b\n1,x\n2,y\n");
751        assert_eq!(
752            all(&mut reader),
753            [
754                vec![Value::BigInt(1), Value::Varchar("x".into())],
755                vec![Value::BigInt(2), Value::Varchar("y".into())],
756            ]
757        );
758    }
759
760    #[test]
761    fn an_empty_field_is_a_null_whether_it_was_quoted_or_not() {
762        // Measured. `allow_quoted_nulls` is on by default, so `""` is a null and not the empty
763        // string, which is the one place a quoted field and a bare one agree about being nothing.
764        let mut reader = read("a,b\n1,\n\"\",y\n");
765        assert_eq!(
766            all(&mut reader),
767            [vec![Value::BigInt(1), Value::Null], vec![Value::Null, Value::Varchar("y".into())],]
768        );
769    }
770
771    #[test]
772    fn a_projection_picks_columns_out_by_position_and_can_reorder_them() {
773        let mut reader = read("a,b,c\n1,x,2.5\n");
774        reader.project(&[2, 0]).expect("projects");
775        assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["c", "a"]);
776        assert_eq!(all(&mut reader), [vec![Value::Double(2.5), Value::BigInt(1)]]);
777    }
778
779    #[test]
780    fn a_projection_of_nothing_still_counts_the_rows() {
781        let mut reader = read("a,b\n1,x\n2,y\n3,z\n");
782        reader.project(&[]).expect("projects");
783        let chunk = reader.next_chunk().expect("reads").expect("a chunk");
784        assert_eq!(chunk.len(), 3);
785        assert_eq!(chunk.width(), 0);
786    }
787
788    #[test]
789    fn a_pipe_separated_file_reads_as_one() {
790        let mut reader = read("a|b\n1|x\n");
791        assert_eq!(reader.dialect().delimiter, b'|');
792        assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x".into())]]);
793    }
794
795    #[test]
796    fn a_quoted_field_with_a_delimiter_in_it_is_one_value() {
797        let mut reader = read("a,b\n1,\"x,y\"\n");
798        assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
799    }
800
801    #[test]
802    fn more_rows_than_fit_one_chunk_arrive_as_more_than_one_chunk() {
803        let mut text = String::from("a\n");
804        for row in 0..VECTOR_SIZE + 5 {
805            text.push_str(&format!("{row}\n"));
806        }
807        let mut reader = read(&text);
808        let first = reader.next_chunk().expect("reads").expect("a chunk");
809        assert_eq!(first.len(), VECTOR_SIZE);
810        let second = reader.next_chunk().expect("reads").expect("a second chunk");
811        assert_eq!(second.len(), 5);
812        assert!(reader.next_chunk().expect("reads").is_none());
813    }
814
815    #[test]
816    fn a_value_the_sniffer_never_saw_is_an_error_rather_than_a_wider_column() {
817        // The value has to be past the sample, because a value inside it would have widened the
818        // column to VARCHAR and there would be nothing to fail. Widening after the fact is not an
819        // option: the chunks before this one have already gone out with the narrow type on them.
820        let mut text = String::from("c\n");
821        for row in 0..infer::SAMPLE {
822            text.push_str(&format!("{row}\n"));
823        }
824        text.push_str("oops\n");
825        let mut reader = read(&text);
826        assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
827        let error = all_or_error(&mut reader).unwrap_err();
828        let line = infer::SAMPLE + 2;
829        assert!(error.message().starts_with(&format!("CSV Error on Line: {line}")), "{error}");
830        assert!(
831            error.message().contains("Could not convert string \"oops\" to 'BIGINT'"),
832            "{error}"
833        );
834        assert!(error.message().contains("sample_size = 20480"), "{error}");
835    }
836
837    /// A chunk is converted a block of rows at a time, which meets the second column's bad value
838    /// in the first block before the first column's in a later one. The error still names the
839    /// first column's, as it did when a chunk was converted a column at a time.
840    #[test]
841    fn the_error_names_the_first_columns_bad_value_even_when_a_later_column_is_bad_sooner() {
842        let mut text = String::from("a,b\n");
843        for row in 0..infer::SAMPLE {
844            text.push_str(&format!("{row},{row}\n"));
845        }
846        text.push_str("1,late\n");
847        for row in 0..BLOCK_ROWS * 2 {
848            text.push_str(&format!("{row},{row}\n"));
849        }
850        text.push_str("early,1\n");
851        let mut reader = read(&text);
852        let error = all_or_error(&mut reader).unwrap_err();
853        assert!(error.message().contains("Could not convert string \"early\""), "{error}");
854    }
855
856    /// Measured. The header row becomes a row, so the names are the generated ones and the first
857    /// column holds `a`, `1` and `2`, which is text rather than the BIGINT the sniffer would say.
858    #[test]
859    fn a_file_told_it_has_no_header_reads_its_first_line_as_a_row() {
860        let mut reader =
861            read_with("a,b\n1,x\n2,y\n", Given { header: Some(false), ..Given::default() });
862        assert_eq!(
863            names_and_types(&reader),
864            [
865                ("column0".to_string(), "VARCHAR".to_string()),
866                ("column1".to_string(), "VARCHAR".to_string()),
867            ]
868        );
869        assert_eq!(all(&mut reader).len(), 3);
870    }
871
872    /// The other way round, on a file the sniffer would call two rows of data.
873    #[test]
874    fn a_file_told_it_has_a_header_takes_its_first_line_as_the_names() {
875        let reader = read_with("1,2\n3,4\n", Given { header: Some(true), ..Given::default() });
876        assert_eq!(reader.fields().iter().map(|f| f.name.clone()).collect::<Vec<_>>(), ["1", "2"]);
877    }
878
879    /// A given delimiter is used rather than tried, so a file that is really commas is one column.
880    #[test]
881    fn a_given_delimiter_is_the_delimiter_whatever_the_file_looks_like() {
882        let reader = read_with("a,b\n1,x\n", Given { delimiter: Some(b';'), ..Given::default() });
883        assert_eq!(reader.dialect().delimiter, b';');
884        assert_eq!(reader.fields().len(), 1);
885    }
886
887    /// A quote the sniffer would never find, since it only ever looks for the double quote.
888    #[test]
889    fn a_given_quote_makes_a_field_that_holds_the_delimiter_one_value() {
890        let mut reader =
891            read_with("a,b\n1,'x,y'\n", Given { quote: Some(b'\''), ..Given::default() });
892        assert_eq!(all(&mut reader), [vec![Value::BigInt(1), Value::Varchar("x,y".into())]]);
893    }
894
895    /// The block under a conversion error says which of its lines the caller wrote down.
896    #[test]
897    fn the_block_says_set_by_user_for_what_the_call_gave_it() {
898        let mut text = String::from("c;d\n");
899        for row in 0..infer::SAMPLE {
900            text.push_str(&format!("{row};x\n"));
901        }
902        text.push_str("oops;x\n");
903        let given = Given { delimiter: Some(b';'), ..Given::default() };
904        let mut reader = read_with(&text, given);
905        let error = all_or_error(&mut reader).unwrap_err();
906        assert!(error.message().contains("delimiter = ; (Set By User)"), "{error}");
907        assert!(error.message().contains("header = true (Auto-Detected)"), "{error}");
908    }
909
910    #[test]
911    fn a_file_told_a_wider_type_than_it_sniffed_reads_its_whole_numbers_as_that_type() {
912        // What a glob does to every file it names. This file on its own is BIGINT and the set it
913        // belongs to is DOUBLE because some other file in it holds a decimal, so the column comes
914        // out DOUBLE and the rows come with it rather than the reader being overruled afterwards.
915        let mut reader = read("a\n1\n2\n");
916        assert_eq!(reader.fields()[0].ty, LogicalType::BigInt);
917        reader.retype(&[LogicalType::Double]).expect("one type for one column");
918        assert_eq!(reader.fields()[0].ty, LogicalType::Double);
919        assert_eq!(all(&mut reader), [[Value::Double(1.0)], [Value::Double(2.0)]]);
920    }
921
922    #[test]
923    fn a_type_list_that_is_not_as_long_as_the_projection_is_refused() {
924        let mut reader = read("a,b\n1,two\n");
925        let error = reader.retype(&[LogicalType::Double]).unwrap_err();
926        assert!(error.message().contains("1 types for a projection of 2 columns"), "{error}");
927    }
928
929    /// Every chunk a reader hands back, as its length and its values written out, or the error
930    /// it stopped on. The values are compared as their debug text so that a NaN equals itself.
931    fn drained(reader: &mut Reader, old: bool) -> (Vec<(usize, Vec<String>)>, Option<String>) {
932        let mut chunks = Vec::new();
933        loop {
934            let next = if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
935            match next {
936                Ok(Some(chunk)) => {
937                    let mut values = Vec::new();
938                    for row in 0..chunk.len() {
939                        for at in 0..chunk.width() {
940                            values.push(format!("{:?}", chunk.value_at(row, at)));
941                        }
942                    }
943                    chunks.push((chunk.len(), values));
944                }
945                Ok(None) => return (chunks, None),
946                Err(error) => return (chunks, Some(error.to_string())),
947            }
948        }
949    }
950
951    /// A small generator, since the crate has no dependencies to take one from.
952    struct Rng(u64);
953
954    impl Rng {
955        fn next(&mut self) -> u64 {
956            self.0 ^= self.0 << 13;
957            self.0 ^= self.0 >> 7;
958            self.0 ^= self.0 << 17;
959            self.0
960        }
961
962        fn below(&mut self, n: usize) -> usize {
963            (self.next() % n as u64) as usize
964        }
965    }
966
967    /// A typed CSV file with a header, whose values are mostly what their column says and now and
968    /// then something only the cast knows what to do with, or something nothing can convert.
969    fn typed_file(rng: &mut Rng, rows: usize) -> String {
970        const ODD: [&str; 27] = [
971            "",
972            " 1",
973            "1 ",
974            "1e3",
975            "0x10",
976            "inf",
977            "-nan",
978            "abc",
979            "\"12\"",
980            "\"a\"\"b\"",
981            "\"x,y\"",
982            "\"x\ny\"",
983            "h\u{e9}llo",
984            "99999999999999999999",
985            "9999999999999999999",
986            "-",
987            "+5",
988            "007",
989            "2020-02-30",
990            "2020-02-29",
991            "0000-01-01",
992            "TRUE",
993            "no",
994            "1_000",
995            "1.5e-3",
996            "-0",
997            "\"\"",
998        ];
999        let width = 1 + rng.below(6);
1000        let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
1001        let mut text: String = (0..width).map(|at| format!("c{at}")).collect::<Vec<_>>().join(",");
1002        text.push('\n');
1003        for _ in 0..rows {
1004            let mut fields = Vec::with_capacity(width);
1005            for &kind in &kinds {
1006                let odd = rng.below(60) == 0;
1007                fields.push(if odd {
1008                    ODD[rng.below(ODD.len())].to_string()
1009                } else {
1010                    let n = rng.next();
1011                    match kind {
1012                        0 => format!("{}", (n % 2_000_001) as i64 - 1_000_000),
1013                        1 => format!("{}.{:02}", n % 100_000, n % 100),
1014                        2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
1015                        3 => ["true", "false", "t", "F"][(n % 4) as usize].to_string(),
1016                        _ => ["x", "hello world", "a longer piece of text", "\"q,\"\"q\""]
1017                            [(n % 4) as usize]
1018                            .to_string(),
1019                    }
1020                });
1021            }
1022            if rng.below(200) == 0 {
1023                fields.pop();
1024            }
1025            if rng.below(200) == 0 {
1026                fields.push("extra".to_string());
1027            }
1028            text.push_str(&fields.join(","));
1029            text.push_str(["\n", "\n", "\n", "\r\n", "\r"][rng.below(5)]);
1030        }
1031        if rng.below(4) == 0 {
1032            text.pop();
1033        }
1034        text
1035    }
1036
1037    fn open_sized(text: &[u8], block: usize) -> Result<Reader> {
1038        let filesystem = SimFilesystem::new();
1039        let path = Path::new("/t.csv");
1040        let file = filesystem.open(path, OpenMode::Create).expect("creates");
1041        file.write_at(0, text).expect("writes");
1042        drop(file);
1043        let file = filesystem.open(path, OpenMode::Read).expect("opens");
1044        Reader::open_sized(file, "/t.csv", Given::default(), block)
1045    }
1046
1047    /// The chunk at a time reader and the record at a time one, over generated files, read with
1048    /// blocks small enough that records straddle the refills, projected and retyped at random so
1049    /// that every type with a parser of its own gets values it takes and values it has to pass on.
1050    #[test]
1051    fn generated_files_read_the_same_a_chunk_at_a_time_as_a_record_at_a_time() {
1052        let types = [
1053            LogicalType::BigInt,
1054            LogicalType::Integer,
1055            LogicalType::SmallInt,
1056            LogicalType::TinyInt,
1057            LogicalType::UBigInt,
1058            LogicalType::UInteger,
1059            LogicalType::USmallInt,
1060            LogicalType::UTinyInt,
1061            LogicalType::Double,
1062            LogicalType::Float,
1063            LogicalType::Date,
1064            LogicalType::Boolean,
1065            LogicalType::Varchar,
1066            LogicalType::Timestamp,
1067            LogicalType::Decimal { width: 18, scale: 3 },
1068        ];
1069        let mut rng = Rng(0x2545_f491_4f6c_dd1d);
1070        for case in 0..200 {
1071            let rows = if case % 100 == 0 { 8192 + rng.below(1000) } else { rng.below(200) };
1072            let mut text = typed_file(&mut rng, rows).into_bytes();
1073            if rng.below(20) == 0 {
1074                // A quoted field with rubbish after it somewhere, which is a malformed record.
1075                let at = rng.below(text.len() + 1);
1076                text.splice(at..at, *b",\"x\"y,");
1077            }
1078            for block in [1 << 20, 32 + rng.below(400)] {
1079                let (Ok(mut new), Ok(mut old)) =
1080                    (open_sized(&text, block), open_sized(&text, block))
1081                else {
1082                    continue;
1083                };
1084                assert_eq!(new.fields(), old.fields());
1085                let width = new.fields().len();
1086                if width > 0 && rng.below(2) == 0 {
1087                    let columns: Vec<usize> =
1088                        (0..rng.below(width + 2)).map(|_| rng.below(width)).collect();
1089                    new.project(&columns).expect("projects");
1090                    old.project(&columns).expect("projects");
1091                    let wanted: Vec<LogicalType> =
1092                        columns.iter().map(|_| types[rng.below(types.len())].clone()).collect();
1093                    new.retype(&wanted).expect("retypes");
1094                    old.retype(&wanted).expect("retypes");
1095                }
1096                let expected = drained(&mut old, true);
1097                let found = drained(&mut new, false);
1098                assert_eq!(found.1, expected.1, "case {case}, block {block}");
1099                assert_eq!(found.0, expected.0, "case {case}, block {block}");
1100            }
1101        }
1102    }
1103
1104    /// How much faster the chunk at a time reader is, on a file shaped like TPC-H's `lineitem`.
1105    ///
1106    /// Not run by default, because it is a measurement rather than a check. Run it with
1107    /// `cargo test --release -p rudb-csv -- --ignored --nocapture reads_lineitem`.
1108    #[test]
1109    #[ignore = "a measurement, run by hand"]
1110    fn reads_lineitem_faster_a_chunk_at_a_time() {
1111        use rudb_io::RealFilesystem;
1112        use std::time::Instant;
1113
1114        const ROWS: usize = 200_000;
1115        let mut rng = Rng(0x1234_5678_9abc_def1);
1116        let mut text = String::from(
1117            "l_orderkey,l_partkey,l_suppkey,l_linenumber,l_quantity,l_extendedprice,l_discount,\
1118             l_tax,l_returnflag,l_linestatus,l_shipdate,l_commitdate,l_receiptdate,\
1119             l_shipinstruct,l_shipmode,l_comment\n",
1120        );
1121        let words = ["carefully", "final", "deposits", "furiously", "regular", "ideas", "sleep"];
1122        for row in 0..ROWS {
1123            let n = rng.next();
1124            let date = |shift: u64| {
1125                format!(
1126                    "{}-{:02}-{:02}",
1127                    1992 + (n >> shift) % 7,
1128                    1 + (n >> shift) % 12,
1129                    1 + (n >> shift) % 28
1130                )
1131            };
1132            let comment: Vec<&str> =
1133                (0..3 + n % 4).map(|k| words[((n >> (k * 3)) % 7) as usize]).collect();
1134            text.push_str(&format!(
1135                "{},{},{},{},{}.00,{}.{:02},0.0{},0.0{},{},{},{},{},{},{},{},{}\n",
1136                row / 4 + 1,
1137                n % 200_000,
1138                n % 10_000,
1139                row % 4 + 1,
1140                1 + n % 50,
1141                900 + n % 100_000,
1142                n % 100,
1143                n % 10,
1144                (n >> 7) % 9,
1145                ["A", "N", "R"][(n % 3) as usize],
1146                ["O", "F"][(n % 2) as usize],
1147                date(3),
1148                date(11),
1149                date(19),
1150                ["DELIVER IN PERSON", "NONE", "TAKE BACK RETURN"][(n % 3) as usize],
1151                ["TRUCK", "MAIL", "AIR", "SHIP"][(n % 4) as usize],
1152                comment.join(" "),
1153            ));
1154        }
1155        let path = std::env::temp_dir().join(format!("rudb-lineitem-{}.csv", std::process::id()));
1156        std::fs::write(&path, &text).expect("writes");
1157        let megabytes = text.len() as f64 / 1e6;
1158        let filesystem = RealFilesystem::new();
1159        let mut best = [f64::MAX; 2];
1160        for _ in 0..5 {
1161            for (slot, old) in [(0, true), (1, false)] {
1162                let file = filesystem.open(&path, OpenMode::Read).expect("opens");
1163                let mut reader = Reader::open(file, "lineitem.csv").expect("sniffs");
1164                let started = Instant::now();
1165                let mut rows = 0;
1166                loop {
1167                    let next =
1168                        if old { reader.next_chunk_by_record() } else { reader.next_chunk() };
1169                    let Some(chunk) = next.expect("reads") else { break };
1170                    rows += chunk.len();
1171                }
1172                assert_eq!(rows, ROWS);
1173                best[slot] = best[slot].min(started.elapsed().as_secs_f64());
1174            }
1175        }
1176        std::fs::remove_file(&path).expect("removes");
1177        let (old, new) = (megabytes / best[0], megabytes / best[1]);
1178        println!(
1179            "{ROWS} rows, {megabytes:.1} MB: a record at a time {old:.1} MB/s, a chunk at a time"
1180        );
1181        println!("{new:.1} MB/s, {:.2}x", best[0] / best[1]);
1182    }
1183
1184    fn all_or_error(reader: &mut Reader) -> Result<Vec<Vec<Value>>> {
1185        let mut rows = Vec::new();
1186        while let Some(chunk) = reader.next_chunk()? {
1187            for row in 0..chunk.len() {
1188                rows.push((0..chunk.width()).map(|at| chunk.value_at(row, at)).collect());
1189            }
1190        }
1191        Ok(rows)
1192    }
1193}