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