1use 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
38pub const RANGE: u64 = 16 << 20;
46
47static SIZE: AtomicU64 = AtomicU64::new(RANGE);
49
50#[must_use]
52pub fn size() -> u64 {
53 SIZE.load(Ordering::Relaxed)
54}
55
56#[doc(hidden)]
61pub fn set_size(bytes: u64) {
62 SIZE.store(bytes.max(1), Ordering::Relaxed);
63}
64
65const LOOK: usize = 64 << 10;
67
68impl Reader {
69 #[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#[derive(Debug)]
87pub struct Split {
88 base: Reader,
90 starts: Vec<u64>,
92 length: u64,
94 chain: Mutex<Chain>,
95 moved: Condvar,
96 failure: OnceLock<Error>,
97}
98
99#[derive(Debug)]
101struct Chain {
102 known: Vec<u64>,
104 guessed: Vec<Option<(u64, u64)>>,
106 taken: Vec<bool>,
108 clean: Vec<bool>,
110}
111
112impl Chain {
113 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 #[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 #[must_use]
158 pub fn ranges(&self) -> usize {
159 self.starts.len()
160 }
161
162 #[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 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 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 let from = self.wait(index, clock.elapsed())?;
186 self.learn(index)?;
189 Ok(self.base.stretch(from, self.end(index), 0))
190 }
191
192 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 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 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 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 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 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 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#[derive(Debug)]
344pub struct Part {
345 split: Arc<Split>,
346 index: usize,
347 reader: Option<Reader>,
348 finished: bool,
349}
350
351impl Part {
352 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 #[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 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 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 type Read = (Vec<String>, Option<String>);
502
503 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 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 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 #[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 #[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 #[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 #[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 #[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}