use std::io;
use std::io::Read;
use std::io::Seek;
use std::io::SeekFrom;
use serde::de::DeserializeOwned;
const READ_CHUNK_SIZE: usize = 8 * 1024;
#[derive(Debug)]
pub enum ScanOutcome<T> {
Parsed(T),
#[allow(dead_code)]
Rejected(serde_json::Error),
}
pub struct ReverseJsonlScanner<R> {
reader: R,
next_chunk_end: u64,
chunk_position: usize,
chunk: Vec<u8>,
record_reversed: Vec<u8>,
}
impl<R> ReverseJsonlScanner<R>
where
R: Read + Seek,
{
pub fn new(mut reader: R) -> io::Result<Self> {
let next_chunk_end = reader.seek(SeekFrom::End(0))?;
Self::new_at(reader, next_chunk_end)
}
pub fn new_at(mut reader: R, end_byte_offset: u64) -> io::Result<Self> {
let file_len = reader.seek(SeekFrom::End(0))?;
if end_byte_offset > file_len {
return Err(io::Error::new(
io::ErrorKind::InvalidInput,
"reverse JSONL scan end is past the file",
));
}
Ok(Self {
reader,
next_chunk_end: end_byte_offset,
chunk_position: 0,
chunk: vec![0; READ_CHUNK_SIZE],
record_reversed: Vec::new(),
})
}
pub fn scan_next<T>(&mut self) -> io::Result<Option<ScanOutcome<T>>>
where
T: DeserializeOwned,
{
loop {
let Some(byte) = self.read_previous_byte()? else {
return Ok(self.finish_record());
};
if byte != b'\n' {
self.record_reversed.push(byte);
continue;
}
if let Some(outcome) = self.finish_record() {
return Ok(Some(outcome));
}
}
}
fn read_previous_byte(&mut self) -> io::Result<Option<u8>> {
if self.chunk_position == 0 {
if self.next_chunk_end == 0 {
return Ok(None);
}
let read_size = usize::try_from(self.next_chunk_end.min(READ_CHUNK_SIZE as u64))
.map_err(io::Error::other)?;
self.next_chunk_end -= read_size as u64;
self.reader.seek(SeekFrom::Start(self.next_chunk_end))?;
self.reader.read_exact(&mut self.chunk[..read_size])?;
self.chunk_position = read_size;
}
self.chunk_position -= 1;
Ok(Some(self.chunk[self.chunk_position]))
}
fn finish_record<T>(&mut self) -> Option<ScanOutcome<T>>
where
T: DeserializeOwned,
{
self.record_reversed.reverse();
let outcome = if self.record_reversed.iter().all(u8::is_ascii_whitespace) {
None
} else {
Some(match serde_json::from_slice::<T>(&self.record_reversed) {
Ok(value) => ScanOutcome::Parsed(value),
Err(error) => ScanOutcome::Rejected(error),
})
};
self.record_reversed.clear();
outcome
}
}
#[cfg(test)]
#[path = "reverse_jsonl_scanner_tests.rs"]
mod tests;