codex_rollout/
reverse_jsonl_scanner.rs1use 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 Parsed(T),
14 #[allow(dead_code)]
16 Rejected(serde_json::Error),
17}
18
19pub 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 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 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;