use std::path::Path;
use crate::transcript::{self, Entry};
const MAX_PARTIAL: usize = 8 * 1024 * 1024;
#[derive(Debug, Default)]
pub struct TailState {
pub offset: u64,
pub partial: Vec<u8>,
overflowed: bool,
identity: Option<(u64, u64)>,
}
pub(crate) enum ReadResult {
Missing,
NoChange,
Reset,
Entries(Vec<Entry>),
}
pub(crate) fn read_appended(path: &Path, state: &mut TailState) -> ReadResult {
use std::io::{Read, Seek, SeekFrom};
let metadata = match std::fs::metadata(path) {
Ok(m) => m,
Err(_) => return ReadResult::Missing,
};
let len = metadata.len();
let identity = file_identity(&metadata);
let replaced = matches!((identity, state.identity), (Some(new), Some(old)) if new != old);
if len < state.offset || replaced {
state.offset = 0;
state.partial.clear();
state.overflowed = false;
state.identity = identity;
return ReadResult::Reset;
}
if len == state.offset {
return ReadResult::NoChange;
}
state.identity = identity;
let mut file = match std::fs::File::open(path) {
Ok(f) => f,
Err(_) => return ReadResult::Missing,
};
if file.seek(SeekFrom::Start(state.offset)).is_err() {
return ReadResult::Missing;
}
let to_read = (len - state.offset) as usize;
let mut buf = vec![0u8; to_read];
let n = match file.read(&mut buf) {
Ok(n) => n,
Err(_) => return ReadResult::Missing,
};
buf.truncate(n);
state.offset += n as u64;
ReadResult::Entries(consume_bytes(state, &buf))
}
pub(crate) fn consume_bytes(state: &mut TailState, appended: &[u8]) -> Vec<Entry> {
let mut entries = Vec::new();
let mut start = 0;
for (i, &byte) in appended.iter().enumerate() {
if byte == b'\n' {
if state.overflowed {
state.overflowed = false;
state.partial.clear();
start = i + 1;
continue;
}
let line_bytes = &appended[start..i];
if state.partial.is_empty() {
if let Some(entry) = parse_bytes(line_bytes) {
entries.push(entry);
}
} else {
state.partial.extend_from_slice(line_bytes);
if let Some(entry) = parse_bytes(&state.partial) {
entries.push(entry);
}
state.partial.clear();
}
start = i + 1;
}
}
if !state.overflowed && start < appended.len() {
state.partial.extend_from_slice(&appended[start..]);
if state.partial.len() > MAX_PARTIAL {
state.partial.clear();
state.overflowed = true;
}
}
entries
}
#[cfg(unix)]
fn file_identity(metadata: &std::fs::Metadata) -> Option<(u64, u64)> {
use std::os::unix::fs::MetadataExt;
Some((metadata.dev(), metadata.ino()))
}
#[cfg(not(unix))]
fn file_identity(_metadata: &std::fs::Metadata) -> Option<(u64, u64)> {
None
}
fn parse_bytes(bytes: &[u8]) -> Option<Entry> {
let bytes = match bytes.last() {
Some(b'\r') => &bytes[..bytes.len() - 1],
_ => bytes,
};
let line = std::str::from_utf8(bytes).ok()?;
transcript::parse_line(line)
}
#[cfg(test)]
mod tests {
use super::*;
use crate::transcript::Entry;
fn is_user(entry: &Entry) -> bool {
matches!(entry, Entry::User(_))
}
#[test]
fn consume_complete_lines() {
let mut state = TailState::default();
let data = b"{\"type\":\"user\"}\n{\"type\":\"assistant\"}\n";
let entries = consume_bytes(&mut state, data);
assert_eq!(entries.len(), 2);
assert!(state.partial.is_empty());
}
#[test]
fn consume_partial_then_complete() {
let mut state = TailState::default();
let first = consume_bytes(&mut state, b"{\"type\":\"us");
assert!(first.is_empty());
assert!(!state.partial.is_empty());
let second = consume_bytes(&mut state, b"er\"}\n");
assert_eq!(second.len(), 1);
assert!(is_user(&second[0]));
assert!(state.partial.is_empty());
}
#[test]
fn consume_partial_carried_across_three_chunks() {
let mut state = TailState::default();
assert!(consume_bytes(&mut state, b"{\"ty").is_empty());
assert!(consume_bytes(&mut state, b"pe\":\"system").is_empty());
let out = consume_bytes(&mut state, b"\"}\n");
assert_eq!(out.len(), 1);
assert!(matches!(out[0], Entry::System(_)));
}
#[test]
fn consume_multiple_with_trailing_partial() {
let mut state = TailState::default();
let data = b"{\"type\":\"user\"}\n{\"type\":\"assistant\"}\n{\"type\":\"sys";
let out = consume_bytes(&mut state, data);
assert_eq!(out.len(), 2);
assert_eq!(state.partial, b"{\"type\":\"sys");
}
#[test]
fn consume_caps_a_runaway_partial_line() {
let mut state = TailState::default();
let huge = vec![b'x'; super::MAX_PARTIAL + 1];
let out = consume_bytes(&mut state, &huge);
assert!(out.is_empty());
assert!(state.partial.is_empty(), "oversized partial is dropped");
consume_bytes(&mut state, b"more-garbage-no-newline");
assert!(state.partial.is_empty());
let out = consume_bytes(&mut state, b"tail-of-garbage\n{\"type\":\"user\"}\n");
assert_eq!(out.len(), 1, "resynced after the runaway line ended");
assert!(is_user(&out[0]));
}
#[test]
fn consume_skips_malformed_lines() {
let mut state = TailState::default();
let out = consume_bytes(&mut state, b"not json at all\n{\"type\":\"user\"}\n");
assert_eq!(out.len(), 1);
assert!(is_user(&out[0]));
}
#[test]
fn consume_tolerates_crlf() {
let mut state = TailState::default();
let out = consume_bytes(&mut state, b"{\"type\":\"user\"}\r\n");
assert_eq!(out.len(), 1);
assert!(is_user(&out[0]));
}
#[test]
fn read_appended_partial_then_complete_over_tempfile() {
use std::io::Write;
let mut tmp = std::env::temp_dir();
tmp.push(format!("zoetrope_tail_test_{}.jsonl", std::process::id()));
let _ = std::fs::remove_file(&tmp);
let mut state = TailState::default();
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"us").unwrap();
f.flush().unwrap();
}
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.is_empty()));
assert!(!state.partial.is_empty());
{
let mut f = std::fs::OpenOptions::new().append(true).open(&tmp).unwrap();
f.write_all(b"er\"}\n").unwrap();
f.flush().unwrap();
}
let r = read_appended(&tmp, &mut state);
match r {
ReadResult::Entries(v) => {
assert_eq!(v.len(), 1);
assert!(is_user(&v[0]));
}
_ => panic!("expected entries"),
}
assert!(state.partial.is_empty());
assert!(matches!(
read_appended(&tmp, &mut state),
ReadResult::NoChange
));
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn read_appended_truncation_resets() {
use std::io::Write;
let mut tmp = std::env::temp_dir();
tmp.push(format!("zoetrope_trunc_test_{}.jsonl", std::process::id()));
let _ = std::fs::remove_file(&tmp);
let mut state = TailState::default();
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"user\"}\n{\"type\":\"assistant\"}\n")
.unwrap();
}
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.len() == 2));
assert!(state.offset > 0);
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"user\"}\n").unwrap();
}
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Reset));
assert_eq!(state.offset, 0);
assert!(state.partial.is_empty());
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.len() == 1));
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn read_appended_detects_rotation_to_longer_file() {
use std::io::Write;
let mut tmp = std::env::temp_dir();
tmp.push(format!("zoetrope_rotate_test_{}.jsonl", std::process::id()));
let _ = std::fs::remove_file(&tmp);
let mut state = TailState::default();
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"user\"}\n").unwrap();
}
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.len() == 1));
std::fs::remove_file(&tmp).unwrap();
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"user\"}\n{\"type\":\"assistant\"}\n")
.unwrap();
}
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Reset));
assert_eq!(state.offset, 0);
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.len() == 2));
let _ = std::fs::remove_file(&tmp);
}
#[test]
fn truncation_reset_clears_runaway_line_skip() {
use std::io::Write;
let mut tmp = std::env::temp_dir();
tmp.push(format!(
"zoetrope_overflow_test_{}.jsonl",
std::process::id()
));
let _ = std::fs::remove_file(&tmp);
let mut state = TailState::default();
consume_bytes(&mut state, &vec![b'x'; MAX_PARTIAL + 1]);
state.offset = (MAX_PARTIAL + 1) as u64;
{
let mut f = std::fs::File::create(&tmp).unwrap();
f.write_all(b"{\"type\":\"user\"}\n").unwrap();
}
assert!(matches!(read_appended(&tmp, &mut state), ReadResult::Reset));
let r = read_appended(&tmp, &mut state);
assert!(matches!(r, ReadResult::Entries(ref v) if v.len() == 1));
let _ = std::fs::remove_file(&tmp);
}
}