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