use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::{Arc, Condvar, Mutex, MutexGuard, OnceLock, PoisonError};
use std::time::{Duration, Instant};
use rudb_common::{Error, Result};
use rudb_vector::Chunk;
use crate::Reader;
pub const RANGE: u64 = 16 << 20;
static SIZE: AtomicU64 = AtomicU64::new(RANGE);
#[must_use]
pub fn size() -> u64 {
SIZE.load(Ordering::Relaxed)
}
#[doc(hidden)]
pub fn set_size(bytes: u64) {
SIZE.store(bytes.max(1), Ordering::Relaxed);
}
const LOOK: usize = 64 << 10;
impl Reader {
#[must_use]
pub fn ranges(&self, size: u64) -> usize {
let Ok(length) = self.file().len() else { return 1 };
let rest = length.saturating_sub(self.here());
if size == 0 || rest / 2 < size {
return 1;
}
usize::try_from(rest.div_ceil(size)).unwrap_or(1)
}
}
#[derive(Debug)]
pub struct Split {
base: Reader,
starts: Vec<u64>,
length: u64,
chain: Mutex<Chain>,
moved: Condvar,
failure: OnceLock<Error>,
}
#[derive(Debug)]
struct Chain {
known: Vec<u64>,
guessed: Vec<Option<(u64, u64)>>,
taken: Vec<bool>,
clean: Vec<bool>,
}
impl Chain {
fn settle(&mut self) {
while self.known.len() < self.guessed.len() {
let last = self.known.len() - 1;
match self.guessed[last] {
Some((guess, stop)) if guess == self.known[last] => self.known.push(stop),
_ => break,
}
}
}
}
impl Split {
#[must_use]
pub fn new(reader: Reader, ranges: usize) -> Self {
let first = reader.here();
let length = reader.file().len().ok();
let rest = length.unwrap_or(first).saturating_sub(first);
let count = if length.is_some() { ranges.max(1) } else { 1 };
let each = rest / count as u64;
let starts = (0..count).map(|at| first + each * at as u64).collect();
let mut known = Vec::with_capacity(count);
known.push(first);
let chain = Chain {
known,
guessed: vec![None; count],
taken: vec![false; count],
clean: vec![false; count],
};
Self {
base: reader.stretch(first, u64::MAX, reader.line()),
starts,
length: each.max(1),
chain: Mutex::new(chain),
moved: Condvar::new(),
failure: OnceLock::new(),
}
}
#[must_use]
pub fn ranges(&self) -> usize {
self.starts.len()
}
#[must_use]
pub fn part(self: &Arc<Self>, index: usize) -> Part {
Part { split: Arc::clone(self), index, reader: None, finished: false }
}
fn end(&self, index: usize) -> u64 {
self.starts.get(index + 1).copied().unwrap_or(u64::MAX)
}
fn lock(&self) -> MutexGuard<'_, Chain> {
self.chain.lock().unwrap_or_else(PoisonError::into_inner)
}
fn begin(&self, index: usize) -> Result<Reader> {
let clock = Instant::now();
if self.lock().known.len() <= index {
self.speculate(index);
}
let from = self.wait(index, clock.elapsed())?;
self.learn(index)?;
Ok(self.base.stretch(from, self.end(index), 0))
}
fn speculate(&self, index: usize) {
let guess = if index == 0 { self.starts[0] } else { self.guess(index) };
let end = self.end(index);
let mut reader = self.base.stretch(guess, end, 0);
reader.give_up_at(end.saturating_add(self.length));
let Ok(stop) = reader.skim() else { return };
let mut chain = self.lock();
chain.guessed[index] = Some((guess, stop));
chain.settle();
self.moved.notify_all();
}
fn guess(&self, index: usize) -> u64 {
let end = self.end(index);
let mut at = self.starts[index].saturating_sub(1);
let mut bytes = vec![0; LOOK];
while at < end {
let Ok(read) = self.base.file().read_at(at, &mut bytes) else { return end };
if read == 0 {
return end;
}
if let Some(found) = bytes[..read].iter().position(|&b| b == b'\n' || b == b'\r') {
let line = at + found as u64;
if bytes[found] == b'\r' {
let mut next = [0];
let read = self.base.file().read_at(line + 1, &mut next).unwrap_or(0);
if read == 1 && next[0] == b'\n' {
return line + 2;
}
}
return line + 1;
}
at += read as u64;
}
end
}
fn wait(&self, index: usize, patience: Duration) -> Result<u64> {
let patience = patience.max(Duration::from_millis(1));
let mut chain = self.lock();
loop {
if let Some(error) = self.failure.get() {
return Err(error.clone());
}
if let Some(&from) = chain.known.get(index) {
return Ok(from);
}
let (guard, waited) =
self.moved.wait_timeout(chain, patience).unwrap_or_else(PoisonError::into_inner);
chain = guard;
if waited.timed_out() && chain.known.len() <= index {
let frontier = chain.known.len() - 1;
if !chain.taken[frontier] {
drop(chain);
self.learn(frontier)?;
chain = self.lock();
}
}
}
}
fn learn(&self, index: usize) -> Result<()> {
let from = {
let mut chain = self.lock();
if chain.known.len() != index + 1 || index + 1 == self.ranges() || chain.taken[index] {
return Ok(());
}
chain.taken[index] = true;
chain.known[index]
};
match self.base.stretch(from, self.end(index), 0).skim() {
Ok(stop) => {
let mut chain = self.lock();
if chain.known.len() == index + 1 {
chain.known.push(stop);
chain.settle();
}
self.moved.notify_all();
Ok(())
}
Err(found) => Err(self.fail(found)),
}
}
fn finish(&self, index: usize, stop: u64) {
let mut chain = self.lock();
chain.clean[index] = true;
if chain.known.len() == index + 1 && index + 1 < self.ranges() {
chain.known.push(stop);
chain.settle();
self.moved.notify_all();
}
}
fn fail(&self, found: Error) -> Error {
let error = self.failure.get_or_init(|| self.replay(found)).clone();
let _chain = self.lock();
self.moved.notify_all();
error
}
fn replay(&self, found: Error) -> Error {
let target = {
let chain = self.lock();
let clean = chain.clean.iter().take_while(|&&clean| clean).count();
chain.known.get(clean).copied().unwrap_or(u64::MAX)
};
let mut reader = self.base.stretch(self.starts[0], u64::MAX, self.base.line());
if let Err(error) = reader.skip_to(target) {
return error;
}
loop {
match reader.next_chunk() {
Ok(Some(_)) => {}
Ok(None) => return found,
Err(error) => return error,
}
}
}
}
#[derive(Debug)]
pub struct Part {
split: Arc<Split>,
index: usize,
reader: Option<Reader>,
finished: bool,
}
impl Part {
pub fn next_chunk(&mut self) -> Result<Option<Chunk>> {
if let Some(error) = self.split.failure.get() {
return Err(error.clone());
}
if self.finished {
return Ok(None);
}
let reader = match &mut self.reader {
Some(reader) => reader,
None => self.reader.insert(self.split.begin(self.index)?),
};
match reader.next_chunk() {
Ok(Some(chunk)) => Ok(Some(chunk)),
Ok(None) => {
self.finished = true;
self.split.finish(self.index, reader.here());
Ok(None)
}
Err(found) => Err(self.split.fail(found)),
}
}
#[must_use]
pub fn bytes_read(&self) -> u64 {
self.reader.as_ref().map_or(0, Reader::bytes_read)
}
}
#[cfg(test)]
mod tests {
use super::*;
use rudb_common::LogicalType;
use rudb_io::{Filesystem, OpenMode, SimFilesystem};
use std::path::Path;
struct Rng(u64);
impl Rng {
fn next(&mut self) -> u64 {
self.0 ^= self.0 << 13;
self.0 ^= self.0 >> 7;
self.0 ^= self.0 << 17;
self.0
}
fn below(&mut self, n: usize) -> usize {
(self.next() % n as u64) as usize
}
}
fn file(rng: &mut Rng) -> (Vec<u8>, Vec<LogicalType>) {
const TEXT: [&str; 12] = [
"plain text",
"",
"\"a,b\"",
"\"one\ntwo\"",
"\"\r\n\"",
"\"\n\n\n\"",
"\"q\"\"\n\"\"q\"",
"\"\"",
"\"x\ry\"",
"\"a long quoted value that runs, over a line\nand then some more\"",
"\"1\n2,3\n\"",
"x",
];
const ODD: [&str; 5] = ["abc", "1.5.5", "2020-13-01", "\"x\"y", "maybe"];
let width = 1 + rng.below(5);
let kinds: Vec<usize> = (0..width).map(|_| rng.below(5)).collect();
let most = if rng.below(10) == 0 { 3000 } else { 300 };
let rows = rng.below(most);
let mut text = String::new();
if rng.below(2) == 0 {
let names: Vec<String> = (0..width).map(|at| format!("c{at}")).collect();
text.push_str(&names.join(","));
text.push('\n');
}
let ending = ["\n", "\r\n", "\r"][rng.below(3)];
for _ in 0..rows {
if rng.below(40) == 0 {
text.push_str(ending);
continue;
}
let fields: Vec<String> = kinds
.iter()
.map(|&kind| {
let n = rng.next();
if rng.below(2000) == 0 {
return ODD[rng.below(ODD.len())].to_string();
}
match kind {
0 => format!("{}", (n % 2001) as i64 - 1000),
1 => format!("\"{}.{:02}\"", n % 1000, n % 100),
2 => format!("{}-{:02}-{:02}", 1990 + n % 20, 1 + n % 12, 1 + n % 28),
3 => ["true", "false", ""][(n % 3) as usize].to_string(),
_ => TEXT[(n % TEXT.len() as u64) as usize].to_string(),
}
})
.collect();
text.push_str(&fields.join(","));
text.push_str(if rng.below(30) == 0 { "\r\n" } else { ending });
}
if rng.below(4) == 0 {
text.pop();
}
let types = kinds
.iter()
.map(|&kind| match kind {
0 => LogicalType::BigInt,
1 => LogicalType::Double,
2 => LogicalType::Date,
3 => LogicalType::Boolean,
_ => LogicalType::Varchar,
})
.collect();
(text.into_bytes(), types)
}
fn open(text: &[u8], block: usize) -> Result<Reader> {
let filesystem = SimFilesystem::new();
let path = Path::new("/t.csv");
let file = filesystem.open(path, OpenMode::Create).expect("creates");
file.write_at(0, text).expect("writes");
drop(file);
let file = filesystem.open(path, OpenMode::Read).expect("opens");
Reader::open_sized(file, "/t.csv", crate::Given::default(), block)
}
type Read = (Vec<String>, Option<String>);
fn rows(mut next: impl FnMut() -> Result<Option<Chunk>>) -> Read {
let mut rows = Vec::new();
loop {
match next() {
Ok(Some(chunk)) => {
for row in 0..chunk.len() {
let values: Vec<_> =
(0..chunk.width()).map(|at| chunk.value_at(row, at)).collect();
rows.push(format!("{values:?}"));
}
}
Ok(None) => return (rows, None),
Err(error) => return (rows, Some(error.to_string())),
}
}
}
fn readers(
rng: &mut Rng,
text: &[u8],
types: &[LogicalType],
block: usize,
) -> Option<(Reader, Read)> {
let (Ok(mut whole), Ok(mut split)) = (open(text, block), open(text, block)) else {
return None;
};
let width = whole.fields().len();
if width == types.len() && rng.below(4) > 0 {
let columns: Vec<usize> = if rng.below(2) == 0 {
(0..width).collect()
} else {
(0..1 + rng.below(width)).map(|_| rng.below(width)).collect()
};
let wanted: Vec<LogicalType> = columns.iter().map(|&at| types[at].clone()).collect();
for reader in [&mut whole, &mut split] {
reader.project(&columns).expect("projects");
reader.retype(&wanted).expect("retypes");
}
}
let expected = rows(|| whole.next_chunk());
Some((split, expected))
}
fn check(parts: Vec<Read>, expected: &Read) {
let failures: Vec<&String> = parts.iter().filter_map(|part| part.1.as_ref()).collect();
match &expected.1 {
None => {
assert!(failures.is_empty(), "{failures:?}");
let found: Vec<String> = parts.into_iter().flat_map(|part| part.0).collect();
assert_eq!(found.len(), expected.0.len());
assert_eq!(&found, &expected.0);
}
Some(error) => {
assert!(!failures.is_empty(), "no part failed, expected {error}");
for failure in failures {
assert_eq!(failure, error);
}
}
}
}
#[test]
fn parts_on_many_threads_read_the_same_as_one_reader() {
let mut rng = Rng(0x9e37_79b9_7f4a_7c15);
for _ in 0..400 {
let (text, types) = file(&mut rng);
let block = if rng.below(2) == 0 { 1 << 20 } else { 16 + rng.below(300) };
let Some((reader, expected)) = readers(&mut rng, &text, &types, block) else {
continue;
};
let split = Arc::new(Split::new(reader, 1 + rng.below(40)));
let parts = std::thread::scope(|scope| {
let handles: Vec<_> = (0..split.ranges())
.map(|index| {
let mut part = split.part(index);
scope.spawn(move || rows(|| part.next_chunk()))
})
.collect();
handles.into_iter().map(|handle| handle.join().expect("joins")).collect()
});
check(parts, &expected);
}
}
#[test]
fn parts_read_in_any_order_on_one_thread_read_the_same() {
let mut rng = Rng(0x2545_f491_4f6c_dd1d);
for _ in 0..100 {
let (text, types) = file(&mut rng);
let block = 16 + rng.below(300);
let Some((reader, expected)) = readers(&mut rng, &text, &types, block) else {
continue;
};
let split = Arc::new(Split::new(reader, 1 + rng.below(12)));
let mut order: Vec<usize> = (0..split.ranges()).collect();
for at in (1..order.len()).rev() {
order.swap(at, rng.below(at + 1));
}
let mut parts = vec![(Vec::new(), None); split.ranges()];
for index in order {
let mut part = split.part(index);
parts[index] = rows(|| part.next_chunk());
}
check(parts, &expected);
}
}
#[test]
fn a_part_finishes_when_the_ranges_before_it_are_never_read() {
let mut text = String::from("a,b\n");
for row in 0..2000 {
text.push_str(&format!("{row},\"x\ny\"\n"));
}
let reader = open(text.as_bytes(), 1 << 20).expect("opens");
let split = Arc::new(Split::new(reader, 10));
let mut part = split.part(9);
let (rows, error) = rows(|| part.next_chunk());
assert_eq!(error, None);
assert!(!rows.is_empty());
assert!(rows.len() < 2000);
}
#[test]
fn an_error_in_a_late_range_names_the_line_one_reader_names() {
let mut text = String::from("n\n");
for row in 0..50_000 {
text.push_str(if row == 41_234 { "oops\n" } else { "12\n" });
}
let mut whole = open(text.as_bytes(), 4096).expect("opens");
whole.retype(&[LogicalType::Integer]).expect("retypes");
let expected = rows(|| whole.next_chunk());
let error = expected.1.clone().expect("fails");
assert!(error.contains("Line: 41236"), "{error}");
let mut reader = open(text.as_bytes(), 4096).expect("opens");
reader.retype(&[LogicalType::Integer]).expect("retypes");
let split = Arc::new(Split::new(reader, 7));
let parts: Vec<_> = (0..7)
.map(|index| {
let mut part = split.part(index);
rows(|| part.next_chunk())
})
.collect();
check(parts, &expected);
}
#[test]
fn a_file_is_cut_into_ranges_only_when_it_is_long_enough() {
let text = "a,b\n".to_string() + &"1,2\n".repeat(1000);
let reader = open(text.as_bytes(), 1 << 20).expect("opens");
assert_eq!(reader.ranges(4000), 1);
assert_eq!(reader.ranges(2000), 2);
assert_eq!(reader.ranges(1000), 4);
assert_eq!(reader.ranges(0), 1);
}
}