Skip to main content

rudb_csv/
split.rs

1//! One CSV file read on many threads.
2//!
3//! A Parquet file comes cut into row groups and says where each one starts. A CSV file does not,
4//! and where a record starts is only known by reading every byte before it, since a newline inside
5//! a quoted field is part of a value and not the end of a line. Read that way a large file is read
6//! on one thread however many the query was given.
7//!
8//! So the file is cut into ranges of about [`RANGE`] bytes and each range guesses. A range that
9//! is not the first assumes its first record starts just after the first line ending in it, which
10//! is right unless that line ending is inside a quoted field, and it splits records from there to
11//! learn where its last record ends. That is where the next range really starts if the guess was
12//! right, so the true starts are known a range at a time as fast as the guesses can be checked,
13//! and a range converts nothing until its own start is one of them. Converting is most of the
14//! work, and it is done once, from the right place, on every thread at the same time.
15//!
16//! A wrong guess costs a second split of that range from its true start, which is also what a
17//! range that nobody guessed for gets, so a guess only ever saves time. A range whose owner is
18//! slow to guess is split by whichever thread is waiting on it, which also means no range can be
19//! left waiting for one that is never read.
20//!
21//! Errors are where the ranges have to agree with one thread exactly. A value that does not fit
22//! its column is reported with its line number, which is a count of the records before it, and
23//! when several ranges fail at once the one reported has to be the first a single thread would
24//! have met. A range that fails therefore does not report its own error. The file is read again
25//! on one thread from the end of the last range that converted without one, in the same chunks a
26//! whole read makes, and the first error that read meets is the one every range reports. That is
27//! slow, and it only happens to a query that is failing.
28
29use std::sync::atomic::{AtomicU64, Ordering};
30use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock, PoisonError};
31use std::time::{Duration, Instant};
32
33use rudb_common::{Error, Result};
34use rudb_vector::Chunk;
35
36use crate::Reader;
37
38/// How many bytes of the file a range covers.
39///
40/// Small enough that a few hundred megabytes gives every thread several ranges, so a thread that
41/// is slowed by something else does not leave the rest waiting on its last range. Large enough
42/// that the one record each range finishes past its end, and the guess at where it starts, are
43/// nothing next to the range, and that a range of TPC-H lineitem is about the 131,072 rows a
44/// native stripe holds, so ranges do not leave many short stripes behind them.
45pub const RANGE: u64 = 16 << 20;
46
47/// The range size [`size`] answers, which only tests change.
48static SIZE: AtomicU64 = AtomicU64::new(RANGE);
49
50/// How many bytes of a file a range covers, which is [`RANGE`] unless a test said otherwise.
51#[must_use]
52pub fn size() -> u64 {
53    SIZE.load(Ordering::Relaxed)
54}
55
56/// Makes every later split use ranges of `bytes`, so that a test can cut a small file into many.
57///
58/// This is for the whole process, which is why a test that calls it lives in a test binary of its
59/// own.
60#[doc(hidden)]
61pub fn set_size(bytes: u64) {
62    SIZE.store(bytes.max(1), Ordering::Relaxed);
63}
64
65/// How much of the file is looked at a time for the line ending a guess starts after.
66const LOOK: usize = 64 << 10;
67
68impl Reader {
69    /// How many ranges of about `size` bytes the rest of the file is read in, which is one for a
70    /// file that is not at least two of them long.
71    #[must_use]
72    pub fn ranges(&self, size: u64) -> usize {
73        let Ok(length) = self.file().len() else { return 1 };
74        let rest = length.saturating_sub(self.here());
75        if size == 0 || rest / 2 < size {
76            return 1;
77        }
78        usize::try_from(rest.div_ceil(size)).unwrap_or(1)
79    }
80}
81
82/// A file cut into ranges that are read on separate threads and come back in file order.
83///
84/// Each range is read through a [`Part`], and the rows of part 0, then part 1 and so on are the
85/// rows a single [`Reader`] would have handed back, in the same order.
86#[derive(Debug)]
87pub struct Split {
88    /// A reader positioned at the first record, which every range's reader is copied from.
89    base: Reader,
90    /// Where each range starts as a count of bytes, not as a record.
91    starts: Vec<u64>,
92    /// How long a range is, which is also how far past its end a guess is followed.
93    length: u64,
94    chain: Mutex<Chain>,
95    moved: Condvar,
96    failure: OnceLock<Error>,
97}
98
99/// What is known so far about where the ranges really start.
100#[derive(Debug)]
101struct Chain {
102    /// The true start of each range from the first, as far as it is known.
103    known: Vec<u64>,
104    /// For each range that has guessed, where it guessed it starts and where that guess ends.
105    guessed: Vec<Option<(u64, u64)>>,
106    /// Whether some thread is splitting the range from its true start to learn the next one.
107    taken: Vec<bool>,
108    /// Whether the range has handed back all of its rows without an error.
109    clean: Vec<bool>,
110}
111
112impl Chain {
113    /// Takes every guess that turned out to start where the range before it ended.
114    fn settle(&mut self) {
115        while self.known.len() < self.guessed.len() {
116            let last = self.known.len() - 1;
117            match self.guessed[last] {
118                Some((guess, stop)) if guess == self.known[last] => self.known.push(stop),
119                _ => break,
120            }
121        }
122    }
123}
124
125impl Split {
126    /// Cuts what is left of `reader`'s file into `ranges` ranges of about the same size.
127    ///
128    /// A file whose length cannot be read is one range, which is the whole file read on one thread
129    /// and is right whatever the file holds.
130    #[must_use]
131    pub fn new(reader: Reader, ranges: usize) -> Self {
132        let first = reader.here();
133        let length = reader.file().len().ok();
134        let rest = length.unwrap_or(first).saturating_sub(first);
135        let count = if length.is_some() { ranges.max(1) } else { 1 };
136        let each = rest / count as u64;
137        let starts = (0..count).map(|at| first + each * at as u64).collect();
138        let mut known = Vec::with_capacity(count);
139        known.push(first);
140        let chain = Chain {
141            known,
142            guessed: vec![None; count],
143            taken: vec![false; count],
144            clean: vec![false; count],
145        };
146        Self {
147            base: reader.stretch(first, u64::MAX, reader.line()),
148            starts,
149            length: each.max(1),
150            chain: Mutex::new(chain),
151            moved: Condvar::new(),
152            failure: OnceLock::new(),
153        }
154    }
155
156    /// How many ranges the file was cut into.
157    #[must_use]
158    pub fn ranges(&self) -> usize {
159        self.starts.len()
160    }
161
162    /// The reader of range `index`, which reads nothing until it is first asked for a chunk.
163    #[must_use]
164    pub fn part(self: &Arc<Self>, index: usize) -> Part {
165        Part { split: Arc::clone(self), index, reader: None, finished: false }
166    }
167
168    /// Where range `index` stops taking records, which for the last one is the end of the file.
169    fn end(&self, index: usize) -> u64 {
170        self.starts.get(index + 1).copied().unwrap_or(u64::MAX)
171    }
172
173    fn lock(&self) -> MutexGuard<'_, Chain> {
174        self.chain.lock().unwrap_or_else(PoisonError::into_inner)
175    }
176
177    /// Works out where range `index` really starts and hands back a reader of it.
178    fn begin(&self, index: usize) -> Result<Reader> {
179        let clock = Instant::now();
180        if self.lock().known.len() <= index {
181            self.speculate(index);
182        }
183        // Waiting as long as this thread's own guess took is about how long the range before it
184        // takes to guess, so a thread only splits somebody else's range when that one is late.
185        let from = self.wait(index, clock.elapsed())?;
186        // A guess that was wrong ends in the wrong place too, so whoever reads the next range is
187        // waiting on this one to be split from where it really starts.
188        self.learn(index)?;
189        Ok(self.base.stretch(from, self.end(index), 0))
190    }
191
192    /// Guesses where range `index` starts, follows the guess to where it ends, and writes both
193    /// down.
194    ///
195    /// Nothing that goes wrong here is an error, because a guess that fails was a wrong guess and
196    /// the range is split again from where it really starts, which reports the error if it is real.
197    fn speculate(&self, index: usize) {
198        let guess = if index == 0 { self.starts[0] } else { self.guess(index) };
199        let end = self.end(index);
200        let mut reader = self.base.stretch(guess, end, 0);
201        reader.give_up_at(end.saturating_add(self.length));
202        let Ok(stop) = reader.skim() else { return };
203        let mut chain = self.lock();
204        chain.guessed[index] = Some((guess, stop));
205        chain.settle();
206        self.moved.notify_all();
207    }
208
209    /// Where the first record of range `index` would start if no line ending near it is inside a
210    /// quoted field, which is just after the first line ending that finishes in the range.
211    ///
212    /// A range with no line ending in it guesses that it holds no records, which is right for a
213    /// range inside one enormous record and is checked like any other guess.
214    fn guess(&self, index: usize) -> u64 {
215        let end = self.end(index);
216        let mut at = self.starts[index].saturating_sub(1);
217        let mut bytes = vec![0; LOOK];
218        while at < end {
219            let Ok(read) = self.base.file().read_at(at, &mut bytes) else { return end };
220            if read == 0 {
221                return end;
222            }
223            if let Some(found) = bytes[..read].iter().position(|&b| b == b'\n' || b == b'\r') {
224                let line = at + found as u64;
225                if bytes[found] == b'\r' {
226                    let mut next = [0];
227                    let read = self.base.file().read_at(line + 1, &mut next).unwrap_or(0);
228                    if read == 1 && next[0] == b'\n' {
229                        return line + 2;
230                    }
231                }
232                return line + 1;
233            }
234            at += read as u64;
235        }
236        end
237    }
238
239    /// Waits until the true start of range `index` is known and answers it.
240    ///
241    /// A wait that runs out splits the first range whose end is not known, from its true start,
242    /// unless some thread is doing that already. The one it helps is at or before `index`, so every
243    /// wait that runs out moves the known starts on by a range, and a range that no thread ever
244    /// reads cannot hold up the ones after it.
245    fn wait(&self, index: usize, patience: Duration) -> Result<u64> {
246        let patience = patience.max(Duration::from_millis(1));
247        let mut chain = self.lock();
248        loop {
249            if let Some(error) = self.failure.get() {
250                return Err(error.clone());
251            }
252            if let Some(&from) = chain.known.get(index) {
253                return Ok(from);
254            }
255            let (guard, waited) =
256                self.moved.wait_timeout(chain, patience).unwrap_or_else(PoisonError::into_inner);
257            chain = guard;
258            if waited.timed_out() && chain.known.len() <= index {
259                let frontier = chain.known.len() - 1;
260                if !chain.taken[frontier] {
261                    drop(chain);
262                    self.learn(frontier)?;
263                    chain = self.lock();
264                }
265            }
266        }
267    }
268
269    /// Splits range `index` from its true start to learn where the next one starts, unless that is
270    /// known already or another thread is on it.
271    ///
272    /// This is a range read from where it really starts, so an error here is a real one.
273    fn learn(&self, index: usize) -> Result<()> {
274        let from = {
275            let mut chain = self.lock();
276            if chain.known.len() != index + 1 || index + 1 == self.ranges() || chain.taken[index] {
277                return Ok(());
278            }
279            chain.taken[index] = true;
280            chain.known[index]
281        };
282        match self.base.stretch(from, self.end(index), 0).skim() {
283            Ok(stop) => {
284                let mut chain = self.lock();
285                if chain.known.len() == index + 1 {
286                    chain.known.push(stop);
287                    chain.settle();
288                }
289                self.moved.notify_all();
290                Ok(())
291            }
292            Err(found) => Err(self.fail(found)),
293        }
294    }
295
296    /// Writes down that range `index` handed back all of its rows and stopped at `stop`.
297    fn finish(&self, index: usize, stop: u64) {
298        let mut chain = self.lock();
299        chain.clean[index] = true;
300        if chain.known.len() == index + 1 && index + 1 < self.ranges() {
301            chain.known.push(stop);
302            chain.settle();
303            self.moved.notify_all();
304        }
305    }
306
307    /// The error a single thread reading the whole file would have reported, given that some range
308    /// met `found`.
309    fn fail(&self, found: Error) -> Error {
310        let error = self.failure.get_or_init(|| self.replay(found)).clone();
311        let _chain = self.lock();
312        self.moved.notify_all();
313        error
314    }
315
316    /// Reads the file again on one thread from the end of the ranges that converted cleanly, in the
317    /// chunks a whole read makes, and answers the first error it meets.
318    ///
319    /// The ranges before that point handed back every row without an error, so a whole read gets
320    /// through them too and its first error is at or after it. `found` is what is reported if the
321    /// read again somehow finds nothing, since some range did fail.
322    fn replay(&self, found: Error) -> Error {
323        let target = {
324            let chain = self.lock();
325            let clean = chain.clean.iter().take_while(|&&clean| clean).count();
326            chain.known.get(clean).copied().unwrap_or(u64::MAX)
327        };
328        let mut reader = self.base.stretch(self.starts[0], u64::MAX, self.base.line());
329        if let Err(error) = reader.skip_to(target) {
330            return error;
331        }
332        loop {
333            match reader.next_chunk() {
334                Ok(Some(_)) => {}
335                Ok(None) => return found,
336                Err(error) => return error,
337            }
338        }
339    }
340}
341
342/// One range of a [`Split`], read a chunk at a time like a whole [`Reader`].
343#[derive(Debug)]
344pub struct Part {
345    split: Arc<Split>,
346    index: usize,
347    reader: Option<Reader>,
348    finished: bool,
349}
350
351impl Part {
352    /// The next chunk of this range, or `None` once it has handed back every record that starts in
353    /// it.
354    ///
355    /// The first call waits until the range's true start is known, which usually means until the
356    /// range before it has guessed.
357    ///
358    /// # Errors
359    ///
360    /// The error a single [`Reader`] over the whole file would have reported first, whichever range
361    /// it is in, once any range meets one.
362    pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
363        if let Some(error) = self.split.failure.get() {
364            return Err(error.clone());
365        }
366        if self.finished {
367            return Ok(None);
368        }
369        let reader = match &mut self.reader {
370            Some(reader) => reader,
371            None => self.reader.insert(self.split.begin(self.index)?),
372        };
373        match reader.next_chunk() {
374            Ok(Some(chunk)) => Ok(Some(chunk)),
375            Ok(None) => {
376                self.finished = true;
377                self.split.finish(self.index, reader.here());
378                Ok(None)
379            }
380            Err(found) => Err(self.split.fail(found)),
381        }
382    }
383
384    /// How many bytes of the file this range's reader has read.
385    ///
386    /// The guess and a second split after a wrong one read the range too, and are not counted,
387    /// so that the ranges of a file add up to about the file.
388    #[must_use]
389    pub fn bytes_read(&self) -> u64 {
390        self.reader.as_ref().map_or(0, Reader::bytes_read)
391    }
392}
393
394#[cfg(test)]
395mod tests {
396    use super::*;
397    use rudb_common::LogicalType;
398    use rudb_io::{Filesystem, OpenMode, SimFilesystem};
399    use std::path::Path;
400
401    /// A small generator, since the crate has no dependencies to take one from.
402    struct Rng(u64);
403
404    impl Rng {
405        fn next(&mut self) -> u64 {
406            self.0 ^= self.0 << 13;
407            self.0 ^= self.0 >> 7;
408            self.0 ^= self.0 << 17;
409            self.0
410        }
411
412        fn below(&mut self, n: usize) -> usize {
413            (self.next() % n as u64) as usize
414        }
415    }
416
417    /// A file whose ranges are hard to guess: quoted fields with line endings and delimiters in
418    /// them, all three line endings, empty lines, sometimes a header, and sometimes no line ending
419    /// on the last line. Each column holds one kind of value, answered as its type, and now and
420    /// then a value that kind cannot take or a quote with rubbish after it, so that a read at those
421    /// types fails somewhere in the middle of the file.
422    fn file(rng: &mut Rng) -> (Vec<u8>, Vec<LogicalType>) {
423        const TEXT: [&str; 12] = [
424            "plain text",
425            "",
426            "\"a,b\"",
427            "\"one\ntwo\"",
428            "\"\r\n\"",
429            "\"\n\n\n\"",
430            "\"q\"\"\n\"\"q\"",
431            "\"\"",
432            "\"x\ry\"",
433            "\"a long quoted value that runs, over a line\nand then some more\"",
434            "\"1\n2,3\n\"",
435            "x",
436        ];
437        const ODD: [&str; 5] = ["abc", "1.5.5", "2020-13-01", "\"x\"y", "maybe"];
438        let width = 1 + rng.below(5);
439        let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
440        let most = if rng.below(10) == 0 { 3000 } else { 300 };
441        let rows = rng.below(most);
442        let mut text = String::new();
443        if rng.below(2) == 0 {
444            let names: Vec<String> = (0..width).map(|at| format!("c{at}")).collect();
445            text.push_str(&names.join(","));
446            text.push('\n');
447        }
448        let ending = ["\n", "\r\n", "\r"][rng.below(3)];
449        for _ in 0..rows {
450            if rng.below(40) == 0 {
451                text.push_str(ending);
452                continue;
453            }
454            let fields: Vec<String> = kinds
455                .iter()
456                .map(|&kind| {
457                    let n = rng.next();
458                    if rng.below(2000) == 0 {
459                        return ODD[rng.below(ODD.len())].to_string();
460                    }
461                    match kind {
462                        0 => format!("{}", (n % 2001) as i64 - 1000),
463                        1 => format!("\"{}.{:02}\"", n % 1000, n % 100),
464                        2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
465                        3 => ["true", "false", ""][(n % 3) as usize].to_string(),
466                        _ => TEXT[(n % TEXT.len() as u64) as usize].to_string(),
467                    }
468                })
469                .collect();
470            text.push_str(&fields.join(","));
471            text.push_str(if rng.below(30) == 0 { "\r\n" } else { ending });
472        }
473        if rng.below(4) == 0 {
474            text.pop();
475        }
476        let types = kinds
477            .iter()
478            .map(|&kind| match kind {
479                0 => LogicalType::BigInt,
480                1 => LogicalType::Double,
481                2 => LogicalType::Date,
482                3 => LogicalType::Boolean,
483                _ => LogicalType::Varchar,
484            })
485            .collect();
486        (text.into_bytes(), types)
487    }
488
489    fn open(text: &[u8], block: usize) -> Result<Reader> {
490        let filesystem = SimFilesystem::new();
491        let path = Path::new("/t.csv");
492        let file = filesystem.open(path, OpenMode::Create).expect("creates");
493        file.write_at(0, text).expect("writes");
494        drop(file);
495        let file = filesystem.open(path, OpenMode::Read).expect("opens");
496        Reader::open_sized(file, "/t.csv", crate::Given::default(), block)
497    }
498
499    /// Every row as the debug text of its values, so that a NaN equals itself, and the error the
500    /// read stopped on.
501    type Read = (Vec<String>, Option<String>);
502
503    /// What a read that pulls chunks from `next` until it stops comes to.
504    fn rows(mut next: impl FnMut() -> Result<Option<Chunk>>) -> Read {
505        let mut rows = Vec::new();
506        loop {
507            match next() {
508                Ok(Some(chunk)) => {
509                    for row in 0..chunk.len() {
510                        let values: Vec<_> =
511                            (0..chunk.width()).map(|at| chunk.value_at(row, at)).collect();
512                        rows.push(format!("{values:?}"));
513                    }
514                }
515                Ok(None) => return (rows, None),
516                Err(error) => return (rows, Some(error.to_string())),
517            }
518        }
519    }
520
521    /// Opens the file twice, read at the types its columns were written as or at what the sniffer
522    /// made of them, and sometimes projected, and answers the reader to split and what one reader
523    /// over it gives.
524    fn readers(
525        rng: &mut Rng,
526        text: &[u8],
527        types: &[LogicalType],
528        block: usize,
529    ) -> Option<(Reader, Read)> {
530        let (Ok(mut whole), Ok(mut split)) = (open(text, block), open(text, block)) else {
531            return None;
532        };
533        let width = whole.fields().len();
534        if width == types.len() && rng.below(4) > 0 {
535            let columns: Vec<usize> = if rng.below(2) == 0 {
536                (0..width).collect()
537            } else {
538                (0..1 + rng.below(width)).map(|_| rng.below(width)).collect()
539            };
540            let wanted: Vec<LogicalType> = columns.iter().map(|&at| types[at].clone()).collect();
541            for reader in [&mut whole, &mut split] {
542                reader.project(&columns).expect("projects");
543                reader.retype(&wanted).expect("retypes");
544            }
545        }
546        let expected = rows(|| whole.next_chunk());
547        Some((split, expected))
548    }
549
550    /// Compares what the parts gave, in part order, with what one reader gave: the same rows when
551    /// it read the whole file, and the same error from every part that failed when it did not.
552    fn check(parts: Vec<Read>, expected: &Read) {
553        let failures: Vec<&String> = parts.iter().filter_map(|part| part.1.as_ref()).collect();
554        match &expected.1 {
555            None => {
556                assert!(failures.is_empty(), "{failures:?}");
557                let found: Vec<String> = parts.into_iter().flat_map(|part| part.0).collect();
558                assert_eq!(found.len(), expected.0.len());
559                assert_eq!(&found, &expected.0);
560            }
561            Some(error) => {
562                assert!(!failures.is_empty(), "no part failed, expected {error}");
563                for failure in failures {
564                    assert_eq!(failure, error);
565                }
566            }
567        }
568    }
569
570    /// Split files read on a thread per part come back with the rows and the errors of one reader,
571    /// whatever the ranges cut through.
572    #[test]
573    fn parts_on_many_threads_read_the_same_as_one_reader() {
574        let mut rng = Rng(0x9e37_79b9_7f4a_7c15);
575        for _ in 0..400 {
576            let (text, types) = file(&mut rng);
577            let block = if rng.below(2) == 0 { 1 << 20 } else { 16 + rng.below(300) };
578            let Some((reader, expected)) = readers(&mut rng, &text, &types, block) else {
579                continue;
580            };
581            let split = Arc::new(Split::new(reader, 1 + rng.below(40)));
582            let parts = std::thread::scope(|scope| {
583                let handles: Vec<_> = (0..split.ranges())
584                    .map(|index| {
585                        let mut part = split.part(index);
586                        scope.spawn(move || rows(|| part.next_chunk()))
587                    })
588                    .collect();
589                handles.into_iter().map(|handle| handle.join().expect("joins")).collect()
590            });
591            check(parts, &expected);
592        }
593    }
594
595    /// Parts read one after another in any order, on one thread, which is a pipeline whose
596    /// threads are all busy elsewhere. Each part waits, gives up waiting, and splits the ranges
597    /// before it itself, so it still finishes and still agrees.
598    #[test]
599    fn parts_read_in_any_order_on_one_thread_read_the_same() {
600        let mut rng = Rng(0x2545_f491_4f6c_dd1d);
601        for _ in 0..100 {
602            let (text, types) = file(&mut rng);
603            let block = 16 + rng.below(300);
604            let Some((reader, expected)) = readers(&mut rng, &text, &types, block) else {
605                continue;
606            };
607            let split = Arc::new(Split::new(reader, 1 + rng.below(12)));
608            let mut order: Vec<usize> = (0..split.ranges()).collect();
609            for at in (1..order.len()).rev() {
610                order.swap(at, rng.below(at + 1));
611            }
612            let mut parts = vec![(Vec::new(), None); split.ranges()];
613            for index in order {
614                let mut part = split.part(index);
615                parts[index] = rows(|| part.next_chunk());
616            }
617            check(parts, &expected);
618        }
619    }
620
621    /// A part whose range comes after ranges nobody reads still finishes, which is what happens
622    /// when a query fails between being handed a range and reading it.
623    #[test]
624    fn a_part_finishes_when_the_ranges_before_it_are_never_read() {
625        let mut text = String::from("a,b\n");
626        for row in 0..2000 {
627            text.push_str(&format!("{row},\"x\ny\"\n"));
628        }
629        let reader = open(text.as_bytes(), 1 << 20).expect("opens");
630        let split = Arc::new(Split::new(reader, 10));
631        let mut part = split.part(9);
632        let (rows, error) = rows(|| part.next_chunk());
633        assert_eq!(error, None);
634        assert!(!rows.is_empty());
635        assert!(rows.len() < 2000);
636    }
637
638    /// A value that fails in a late range is reported with the line a single reader gives it, and
639    /// every range that fails reports the same one.
640    #[test]
641    fn an_error_in_a_late_range_names_the_line_one_reader_names() {
642        let mut text = String::from("n\n");
643        for row in 0..50_000 {
644            text.push_str(if row == 41_234 { "oops\n" } else { "12\n" });
645        }
646        let mut whole = open(text.as_bytes(), 4096).expect("opens");
647        whole.retype(&[LogicalType::Integer]).expect("retypes");
648        let expected = rows(|| whole.next_chunk());
649        let error = expected.1.clone().expect("fails");
650        assert!(error.contains("Line: 41236"), "{error}");
651        let mut reader = open(text.as_bytes(), 4096).expect("opens");
652        reader.retype(&[LogicalType::Integer]).expect("retypes");
653        let split = Arc::new(Split::new(reader, 7));
654        let parts: Vec<_> = (0..7)
655            .map(|index| {
656                let mut part = split.part(index);
657                rows(|| part.next_chunk())
658            })
659            .collect();
660        check(parts, &expected);
661    }
662
663    /// A file too short for two ranges is read as one, and a longer one is cut about evenly.
664    #[test]
665    fn a_file_is_cut_into_ranges_only_when_it_is_long_enough() {
666        let text = "a,b\n".to_string() + &"1,2\n".repeat(1000);
667        let reader = open(text.as_bytes(), 1 << 20).expect("opens");
668        assert_eq!(reader.ranges(4000), 1);
669        assert_eq!(reader.ranges(2000), 2);
670        assert_eq!(reader.ranges(1000), 4);
671        assert_eq!(reader.ranges(0), 1);
672    }
673}