use std::path::Path;
use crate::{Error, Result};
const WINDOW_BYTES: usize = 16 << 20;
pub(super) fn for_each_window(
path: &Path,
threads: usize,
mut on_window: impl FnMut(usize, &str) -> Result<usize>,
) -> Result<()> {
let io_err = |e: std::io::Error| Error::Io(format!("{}: {e}", path.display()));
let file = std::fs::File::open(path).map_err(io_err)?;
let mut windows = Windows::new(file, WINDOW_BYTES);
let mut first_line = 1;
let mut consume = |buf: &[u8]| -> Result<()> {
let text = std::str::from_utf8(buf).map_err(|_| {
Error::Io(format!(
"{}: stream did not contain valid UTF-8",
path.display()
))
})?;
first_line += on_window(first_line, text)?;
Ok(())
};
if threads <= 1 {
let mut buf = Vec::new();
while windows.next_into(&mut buf).map_err(io_err)? {
consume(&buf)?;
}
return Ok(());
}
let (full_tx, full_rx) = std::sync::mpsc::sync_channel::<std::io::Result<Vec<u8>>>(1);
let (empty_tx, empty_rx) = std::sync::mpsc::sync_channel::<Vec<u8>>(2);
for _ in 0..2 {
empty_tx.send(Vec::new()).expect("the receiver is alive");
}
std::thread::scope(|scope| {
scope.spawn(move || {
while let Ok(mut buf) = empty_rx.recv() {
match windows.next_into(&mut buf) {
Ok(true) => {
if full_tx.send(Ok(buf)).is_err() {
return;
}
}
Ok(false) => return,
Err(e) => {
let _ = full_tx.send(Err(e));
return;
}
}
}
});
let empty_tx = empty_tx;
for buf in full_rx {
let buf = buf.map_err(io_err)?;
consume(&buf)?;
let _ = empty_tx.send(buf);
}
Ok(())
})
}
fn read_up_to(file: &mut std::fs::File, buf: &mut Vec<u8>, size: usize) -> std::io::Result<usize> {
use std::io::Read;
let start = buf.len();
buf.resize(start + size, 0);
let mut filled = start;
while filled < start + size {
match file.read(&mut buf[filled..start + size]) {
Ok(0) => break,
Ok(n) => filled += n,
Err(e) if e.kind() == std::io::ErrorKind::Interrupted => {}
Err(e) => {
buf.truncate(start);
return Err(e);
}
}
}
buf.truncate(filled);
Ok(filled - start)
}
pub(super) fn count_newlines(bytes: &[u8]) -> usize {
let mut chunks = bytes.chunks_exact(8);
let mut count = 0usize;
for chunk in &mut chunks {
let word = u64::from_le_bytes(chunk.try_into().expect("eight bytes"));
let x = word ^ 0x0a0a_0a0a_0a0a_0a0a;
let nonzero = ((x & 0x7f7f_7f7f_7f7f_7f7f) + 0x7f7f_7f7f_7f7f_7f7f) | x;
count += (!nonzero & 0x8080_8080_8080_8080).count_ones() as usize;
}
count + chunks.remainder().iter().filter(|&&b| b == b'\n').count()
}
pub(super) struct Windows {
file: std::fs::File,
size: usize,
remaining: u64,
carry: Vec<u8>,
done: bool,
}
impl Windows {
pub(super) fn new(file: std::fs::File, size: usize) -> Self {
let remaining = file.metadata().map_or(size as u64, |m| m.len());
Self {
file,
size,
remaining,
carry: Vec::new(),
done: false,
}
}
pub(super) fn next_into(&mut self, buf: &mut Vec<u8>) -> std::io::Result<bool> {
buf.clear();
if self.done {
return Ok(false);
}
buf.append(&mut self.carry);
loop {
let want = usize::try_from(self.remaining.saturating_add(1))
.unwrap_or(self.size)
.clamp(1, self.size);
let got = read_up_to(&mut self.file, buf, want)?;
self.remaining = self.remaining.saturating_sub(got as u64);
if got < want {
self.done = true;
return Ok(!buf.is_empty());
}
if let Some(pos) = buf.iter().rposition(|&b| b == b'\n') {
self.carry.extend_from_slice(&buf[pos + 1..]);
buf.truncate(pos + 1);
return Ok(true);
}
}
}
}