Skip to main content

codex_rollout/
reverse_jsonl_scanner.rs

1use std::io;
2use std::io::Read;
3use std::io::Seek;
4use std::io::SeekFrom;
5
6use serde::de::DeserializeOwned;
7
8const READ_CHUNK_SIZE: usize = 8 * 1024;
9
10#[derive(Debug)]
11pub enum ScanOutcome<T> {
12    /// The record was valid JSON and deserialized as the requested type.
13    Parsed(T),
14    /// The record was not valid JSON for the requested type.
15    #[allow(dead_code)]
16    Rejected(serde_json::Error),
17}
18
19/// Read-only scanner for newline-delimited JSON records, starting from the end.
20pub struct ReverseJsonlScanner<R> {
21    reader: R,
22    next_chunk_end: u64,
23    chunk_position: usize,
24    chunk: Vec<u8>,
25    record_reversed: Vec<u8>,
26}
27
28impl<R> ReverseJsonlScanner<R>
29where
30    R: Read + Seek,
31{
32    pub fn new(mut reader: R) -> io::Result<Self> {
33        let next_chunk_end = reader.seek(SeekFrom::End(0))?;
34        Self::new_at(reader, next_chunk_end)
35    }
36
37    /// Creates a reverse scanner whose logical end is the given byte offset.
38    ///
39    /// This lets callers scan a frozen JSONL prefix without reading records appended after that
40    /// prefix was captured.
41    pub fn new_at(mut reader: R, end_byte_offset: u64) -> io::Result<Self> {
42        let file_len = reader.seek(SeekFrom::End(0))?;
43        if end_byte_offset > file_len {
44            return Err(io::Error::new(
45                io::ErrorKind::InvalidInput,
46                "reverse JSONL scan end is past the file",
47            ));
48        }
49        Ok(Self {
50            reader,
51            next_chunk_end: end_byte_offset,
52            chunk_position: 0,
53            chunk: vec![0; READ_CHUNK_SIZE],
54            record_reversed: Vec::new(),
55        })
56    }
57
58    /// Scans the next nonblank record.
59    ///
60    /// I/O failures are returned as [`Err`]. Invalid JSON records are returned as
61    /// [`ScanOutcome::Rejected`], and the scanner remains usable.
62    pub fn scan_next<T>(&mut self) -> io::Result<Option<ScanOutcome<T>>>
63    where
64        T: DeserializeOwned,
65    {
66        loop {
67            let Some(byte) = self.read_previous_byte()? else {
68                return Ok(self.finish_record());
69            };
70
71            if byte != b'\n' {
72                self.record_reversed.push(byte);
73                continue;
74            }
75
76            if let Some(outcome) = self.finish_record() {
77                return Ok(Some(outcome));
78            }
79        }
80    }
81
82    fn read_previous_byte(&mut self) -> io::Result<Option<u8>> {
83        if self.chunk_position == 0 {
84            if self.next_chunk_end == 0 {
85                return Ok(None);
86            }
87
88            let read_size = usize::try_from(self.next_chunk_end.min(READ_CHUNK_SIZE as u64))
89                .map_err(io::Error::other)?;
90            self.next_chunk_end -= read_size as u64;
91            self.reader.seek(SeekFrom::Start(self.next_chunk_end))?;
92            self.reader.read_exact(&mut self.chunk[..read_size])?;
93            self.chunk_position = read_size;
94        }
95
96        self.chunk_position -= 1;
97        Ok(Some(self.chunk[self.chunk_position]))
98    }
99
100    fn finish_record<T>(&mut self) -> Option<ScanOutcome<T>>
101    where
102        T: DeserializeOwned,
103    {
104        self.record_reversed.reverse();
105        let outcome = if self.record_reversed.iter().all(u8::is_ascii_whitespace) {
106            None
107        } else {
108            Some(match serde_json::from_slice::<T>(&self.record_reversed) {
109                Ok(value) => ScanOutcome::Parsed(value),
110                Err(error) => ScanOutcome::Rejected(error),
111            })
112        };
113        self.record_reversed.clear();
114        outcome
115    }
116}
117
118#[cfg(test)]
119#[path = "reverse_jsonl_scanner_tests.rs"]
120mod tests;