use std::fs::File;
use std::io::{BufRead, BufReader, Read, Seek, SeekFrom};
use std::path::Path;
use anyhow::{Context, Result};
use super::loader;
pub const STRIDE: usize = 65_536;
#[derive(Debug)]
pub struct RowIndex {
checkpoints: Vec<u64>,
rows: usize,
header: Vec<u8>,
len: u64,
lines: bool,
}
impl RowIndex {
pub fn rows(&self) -> usize {
self.rows
}
pub fn header(&self) -> &[u8] {
&self.header
}
pub fn seek(&self, row: usize) -> (usize, u64) {
let k = (row / STRIDE).min(self.checkpoints.len().saturating_sub(1));
match self.checkpoints.get(k) {
Some(&at) => (k * STRIDE, at),
None => (0, 0),
}
}
pub fn end_of(&self, row: usize) -> u64 {
let k = row.div_ceil(STRIDE);
self.checkpoints.get(k).copied().unwrap_or(self.len)
}
pub fn chunk_end(&self, from: usize, bytes: u64) -> usize {
let start = self.seek(from).1;
let mut end = (from / STRIDE + 1) * STRIDE;
while end < self.rows {
let next = (end / STRIDE + 1) * STRIDE;
if self.end_of(next).saturating_sub(start) > bytes {
break;
}
end = next;
}
end.min(self.rows)
}
pub fn build(path: &Path, separator: u8) -> Result<Self> {
let preamble = loader::preamble(path, separator)?.bytes;
let mut file =
File::open(path).with_context(|| format!("cannot read {}", path.display()))?;
let len = file.metadata()?.len();
file.seek(SeekFrom::Start(preamble))?;
let mut reader = BufReader::with_capacity(1 << 20, file);
let mut scan = Records::new(separator);
let mut position: u64 = preamble;
let mut header = Vec::new();
let mut rows = 0usize;
let mut checkpoints = Vec::new();
let mut past_header = false;
let mut started = false;
let mut data_start = len;
loop {
let chunk = reader.fill_buf()?;
if chunk.is_empty() {
break;
}
let read = chunk.len();
let base = position;
if !past_header {
for (i, &byte) in chunk.iter().enumerate() {
header.push(byte);
started = true;
if scan.push(byte) {
past_header = true;
while header.last().is_some_and(|&b| b == b'\n' || b == b'\r') {
header.pop();
}
started = false;
position = base + i as u64 + 1;
data_start = position;
break;
}
}
if !past_header {
position = base + read as u64;
reader.consume(read);
continue;
}
}
let from = (position - base) as usize;
for (i, &byte) in chunk[from..].iter().enumerate() {
started = true;
if scan.push(byte) {
rows += 1;
started = false;
if rows % STRIDE == 0 {
checkpoints.push(base + (from + i) as u64 + 1);
}
}
}
position = base + read as u64;
reader.consume(read);
}
if started && scan.in_record() && past_header {
rows += 1;
}
Ok(Self {
checkpoints: with_first(checkpoints, data_start),
rows,
header,
len,
lines: false,
})
}
pub fn build_lines(path: &Path, observe: &mut impl FnMut(u8)) -> Result<Self> {
let file = File::open(path).with_context(|| format!("cannot read {}", path.display()))?;
let len = file.metadata()?.len();
let mut reader = BufReader::with_capacity(1 << 20, file);
let mut position: u64 = 0;
let mut rows = 0usize;
let mut started = false;
let mut checkpoints = vec![0u64];
loop {
let chunk = reader.fill_buf()?;
if chunk.is_empty() {
break;
}
let read = chunk.len();
for (i, &byte) in chunk.iter().enumerate() {
observe(byte);
if byte == b'\n' {
rows += 1;
started = false;
if rows % STRIDE == 0 {
checkpoints.push(position + i as u64 + 1);
}
} else {
started = true;
}
}
position += read as u64;
reader.consume(read);
}
if started {
rows += 1;
}
Ok(Self {
checkpoints,
rows,
header: Vec::new(),
len,
lines: true,
})
}
}
#[derive(Default)]
pub struct Building {
checkpoints: Vec<u64>,
rows: usize,
start: Option<u64>,
}
impl Building {
pub fn new() -> Self {
Self::default()
}
pub fn begin(&mut self, at: u64) {
self.start.get_or_insert(at);
}
pub fn record(&mut self, end: u64) {
self.rows += 1;
if self.rows % STRIDE == 0 {
self.checkpoints.push(end);
}
}
pub fn finish(self, header: Vec<u8>, len: u64) -> RowIndex {
RowIndex {
checkpoints: with_first(self.checkpoints, self.start.unwrap_or(len)),
rows: self.rows,
header,
len,
lines: false,
}
}
}
fn with_first(mut checkpoints: Vec<u64>, data_start: u64) -> Vec<u64> {
checkpoints.insert(0, data_start);
checkpoints
}
struct Records {
separator: u8,
in_quotes: bool,
pending_quote: bool,
at_field_start: bool,
in_record: bool,
}
impl Records {
fn new(separator: u8) -> Self {
Self {
separator,
in_quotes: false,
pending_quote: false,
at_field_start: true,
in_record: false,
}
}
fn in_record(&self) -> bool {
self.in_record
}
#[inline(always)]
fn push(&mut self, byte: u8) -> bool {
if self.pending_quote {
self.pending_quote = false;
if byte == b'"' {
return false; }
self.in_quotes = false;
} else if self.in_quotes {
if byte == b'"' {
self.pending_quote = true;
}
return false;
}
if byte == b'\n' {
self.at_field_start = true;
self.in_record = false;
return true;
}
self.in_record = true;
if byte == b'"' && self.at_field_start {
self.in_quotes = true;
self.at_field_start = false;
} else {
self.at_field_start = byte == self.separator;
}
false
}
}
pub fn read_span(path: &Path, index: &RowIndex, from: u64, to: u64) -> Result<Vec<u8>> {
let mut file = File::open(path)?;
file.seek(SeekFrom::Start(from))?;
let mut buffer = Vec::with_capacity((to.saturating_sub(from) as usize).saturating_add(64));
if !index.lines {
buffer.extend_from_slice(index.header());
buffer.push(b'\n');
}
file.take(to.saturating_sub(from))
.read_to_end(&mut buffer)?;
Ok(buffer)
}
#[cfg(test)]
mod tests {
use super::*;
fn write(name: &str, contents: &str) -> std::path::PathBuf {
let dir = std::env::temp_dir().join("plv-index-tests");
std::fs::create_dir_all(&dir).unwrap();
let path = dir.join(name);
std::fs::write(&path, contents).unwrap();
path
}
fn rows_of(contents: &str) -> usize {
let name = format!(
"count-{:x}.csv",
contents
.bytes()
.fold(0u64, |h, b| h.wrapping_mul(31).wrapping_add(b as u64))
);
RowIndex::build(&write(&name, contents), b',')
.unwrap()
.rows()
}
#[test]
fn records_are_counted_the_way_the_writer_counts_them() {
assert_eq!(rows_of("a,b\n1,2\n3,4\n"), 2);
assert_eq!(rows_of("a,b\n1,2"), 1, "no trailing newline");
assert_eq!(rows_of("a,b\n"), 0, "header only");
assert_eq!(rows_of("a,b\n1,2\n\n3,4\n"), 3, "a blank line is a row");
assert_eq!(rows_of("a,b\r\n1,2\r\n"), 1, "crlf");
}
#[test]
fn a_newline_inside_quotes_is_not_a_record_boundary() {
assert_eq!(rows_of("a,b\n\"one\ntwo\",x\ny,z\n"), 2);
assert_eq!(rows_of("a,b\n\"say \"\"hi\"\"\",x\n"), 1, "escaped quotes");
}
#[test]
fn a_quote_that_is_not_at_a_field_start_is_an_ordinary_character() {
assert_eq!(rows_of("a,b\nab\"cd,x\ny,z\n"), 2);
}
#[test]
fn the_index_counts_records_the_way_polars_and_the_writer_do() {
use crate::data::edit::Overlay;
use crate::data::{loader, writer};
let cases: &[(&str, &str)] = &[
("plain", "a,b\n1,2\n3,4\n"),
("no-trailing-newline", "a,b\n1,2\n3,4"),
("crlf", "a,b\r\n1,2\r\n3,4\r\n"),
("blank-line-between", "a,b\n1,2\n\n3,4\n"),
("trailing-blank-line", "a,b\n1,2\n\n"),
("quoted-newline", "a,b\n\"one\ntwo\",x\ny,z\n"),
("escaped-quotes", "a,b\n\"say \"\"hi\"\"\",x\ny,z\n"),
("quoted-separator", "a,b\n\"x,y\",z\np,q\n"),
("bare-quote-mid-field", "a,b\nab\"cd,x\ny,z\n"),
("empty-fields", "a,b,c\n,,\n1,,3\n"),
("comment-preamble", "# about\n# this, file\na,b\n1,2\n3,4\n"),
("comment-preamble-crlf", "# about\r\na,b\r\n1,2\r\n"),
("hash-header", "#a,b\n1,2\n"),
("hash-data-row", "a,b\n#1,2\n3,4\n"),
];
for (label, contents) in cases {
let path = write(&format!("agree-{label}.csv"), contents);
let polars = match loader::load(&path).and_then(|lf| Ok(lf.collect()?)) {
Ok(df) => df.height(),
Err(e) => {
println!(" {label}: polars will not read this — {e}");
continue;
}
};
let index = RowIndex::build(&path, b',').unwrap().rows();
assert_eq!(
index, polars,
"{label}: index says {index}, polars {polars}"
);
writer::splice(
contents.as_bytes(),
&mut Vec::new(),
b',',
loader::preamble(&path, b',').unwrap().bytes,
true,
&Overlay::new(),
index,
)
.unwrap_or_else(|e| panic!("{label}: the writer disagrees — {e}"));
}
}
#[test]
fn the_header_is_kept_without_its_terminator() {
let index = RowIndex::build(&write("hdr.csv", "a,b\r\n1,2\n"), b',').unwrap();
assert_eq!(index.header(), b"a,b");
}
#[test]
fn a_page_can_be_read_back_from_its_offset() {
let path = write("span.csv", "a,b\n1,2\n3,4\n5,6\n");
let index = RowIndex::build(&path, b',').unwrap();
let (row, at) = index.seek(0);
assert_eq!((row, at), (0, 4), "data starts after `a,b\\n`");
let bytes = read_span(&path, &index, at, index.end_of(3)).unwrap();
assert_eq!(String::from_utf8(bytes).unwrap(), "a,b\n1,2\n3,4\n5,6\n");
}
#[test]
fn a_page_after_a_preamble_starts_at_the_header() {
let path = write("preamble.csv", "# a note\r\na,b\r\n1,2\r\n3,4\r\n");
let index = RowIndex::build(&path, b',').unwrap();
assert_eq!(index.header(), b"a,b");
assert_eq!(index.rows(), 2);
assert_eq!(index.seek(0), (0, 15), "after the comment and `a,b\\r\\n`");
let bytes = read_span(&path, &index, 15, index.end_of(2)).unwrap();
assert_eq!(String::from_utf8(bytes).unwrap(), "a,b\n1,2\r\n3,4\r\n");
}
#[test]
fn a_chunk_is_bounded_by_bytes_rather_than_by_rows() {
let mut csv = String::from("id\n");
for i in 0..STRIDE * 8 {
csv.push_str(&format!("{i:08}\n"));
}
let path = write("chunked.csv", &csv);
let index = RowIndex::build(&path, b',').unwrap();
assert_eq!(index.rows(), STRIDE * 8);
let wide = index.chunk_end(0, 10 << 20);
assert_eq!(wide, STRIDE * 8, "the whole file fits in 10MB");
let tight = index.chunk_end(0, 1);
assert_eq!(tight, STRIDE, "never stalls, never overshoots");
assert!(index.chunk_end(STRIDE, 1) > STRIDE);
let mut at = 0usize;
let mut steps = 0usize;
while at < index.rows() {
let next = index.chunk_end(at, 200_000);
assert!(next > at, "a chunk must move forward");
at = next;
steps += 1;
}
assert_eq!(at, index.rows());
assert!(steps > 1, "the budget was meant to force several chunks");
}
#[test]
fn checkpoints_land_on_row_boundaries() {
let path = write("stride.csv", "a\n1\n2\n3\n");
let index = RowIndex::build(&path, b',').unwrap();
assert_eq!(index.rows(), 3);
assert_eq!(index.seek(0), (0, 2));
assert_eq!(index.seek(2), (0, 2));
assert_eq!(index.end_of(3), index.len, "past the end is the file's end");
}
}