Files
codex/codex-rs/rollout/src/reverse_jsonl_scanner.rs
Owen Lin 592467fb96 Load model context from a bounded rollout suffix (#32896)
## Why

Reconstructing the latest model-visible context does not require replaying an
entire paginated rollout when a usable compaction checkpoint and the associated
completed-turn metadata are available.

## What changed

- Add `ThreadStore::load_latest_model_context` and `StoredModelContext` for
  loading replay-ready model context independently of full thread history.
- Reverse-scan plain paginated JSONL rollouts until the newest safe bounded
  suffix is found, while preserving canonical session metadata and chronological
  replay order.
- Fall back to complete history for legacy or compressed rollouts and whenever
  compaction or rollback records make a bounded cutoff unsafe.

## Testing

- Cover checkpoint selection, turn-metadata boundaries, agent messages,
  contextual user fragments, and full-history fallbacks.

GitOrigin-RevId: 3572f4ecc7aa4099a6d9f0d3e72d0ef432cd9497
2026-07-13 23:29:59 +00:00

106 lines
2.9 KiB
Rust

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> {
/// The record was valid JSON and deserialized as the requested type.
Parsed(T),
/// The record was not valid JSON for the requested type.
#[allow(dead_code)]
Rejected(serde_json::Error),
}
/// Read-only scanner for newline-delimited JSON records, starting from the end.
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))?;
Ok(Self {
reader,
next_chunk_end,
chunk_position: 0,
chunk: vec![0; READ_CHUNK_SIZE],
record_reversed: Vec::new(),
})
}
/// Scans the next nonblank record.
///
/// I/O failures are returned as [`Err`]. Invalid JSON records are returned as
/// [`ScanOutcome::Rejected`], and the scanner remains usable.
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;