use std::collections::HashSet;
use std::fs::File;
use std::io::{BufRead, BufReader, Seek, SeekFrom};
use std::path::{Path, PathBuf};
use octl_core::Event;
use crate::error::CliError;
const CORRUPT_EXCERPT_BYTES: usize = 100;
#[derive(Debug, Clone)]
pub struct CorruptLine {
pub byte_offset: u64,
byte_len: u64,
pub line_excerpt: String,
}
pub struct EventTail {
path: PathBuf,
pos: u64,
last_seq: u64,
corrupt: Option<CorruptLine>,
reported_corrupt: HashSet<(u64, u64)>,
}
impl EventTail {
pub fn new(path: impl Into<PathBuf>, since_seq: u64) -> Self {
Self {
path: path.into(),
pos: 0,
last_seq: since_seq,
corrupt: None,
reported_corrupt: HashSet::new(),
}
}
pub fn poll(&mut self) -> Result<Vec<Event>, CliError> {
let mut f = match File::open(&self.path) {
Ok(f) => f,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
Err(e) => {
return Err(CliError::system(
"io_error",
format!("open {}: {}", self.path.display(), e),
));
}
};
let len = f
.metadata()
.map_err(|e| {
CliError::system(
"io_error",
format!("metadata {}: {}", self.path.display(), e),
)
})?
.len();
if len < self.pos {
self.pos = 0;
self.corrupt = None;
}
if self.corrupt.is_some() {
return Ok(Vec::new());
}
if len == self.pos {
return Ok(Vec::new());
}
f.seek(SeekFrom::Start(self.pos)).map_err(|e| {
CliError::system("io_error", format!("seek {}: {}", self.path.display(), e))
})?;
let mut reader = BufReader::new(f);
let mut out = Vec::new();
let mut consumed: u64 = 0;
let mut buf: Vec<u8> = Vec::new();
loop {
buf.clear();
let n = reader.read_until(b'\n', &mut buf).map_err(|e| {
CliError::system("io_error", format!("read {}: {}", self.path.display(), e))
})?;
if n == 0 {
break;
}
if buf.last() != Some(&b'\n') {
break;
}
let trimmed = trim_line_end(&buf);
if trimmed.is_empty() {
consumed += n as u64;
continue;
}
let Ok(ev) = serde_json::from_slice::<Event>(trimmed) else {
self.pos += consumed;
self.corrupt = Some(CorruptLine {
byte_offset: self.pos,
byte_len: n as u64,
line_excerpt: excerpt(trimmed),
});
return Ok(out);
};
consumed += n as u64;
if ev.seq <= self.last_seq {
continue;
}
self.last_seq = ev.seq;
out.push(ev);
}
self.pos += consumed;
Ok(out)
}
pub fn take_new_corrupt(&mut self) -> Option<CorruptLine> {
let c = self.corrupt.take()?;
self.pos += c.byte_len;
if self
.reported_corrupt
.insert((c.byte_offset, excerpt_hash(&c.line_excerpt)))
{
Some(c)
} else {
None
}
}
pub fn restart(&mut self) {
self.pos = 0;
self.corrupt = None;
}
pub fn path(&self) -> &Path {
&self.path
}
#[allow(dead_code)]
pub fn last_seq(&self) -> u64 {
self.last_seq
}
}
fn excerpt_hash(s: &str) -> u64 {
use std::hash::{Hash, Hasher};
let mut h = std::collections::hash_map::DefaultHasher::new();
s.hash(&mut h);
h.finish()
}
fn trim_line_end(buf: &[u8]) -> &[u8] {
let mut end = buf.len();
if end > 0 && buf[end - 1] == b'\n' {
end -= 1;
if end > 0 && buf[end - 1] == b'\r' {
end -= 1;
}
}
&buf[..end]
}
fn excerpt(line: &[u8]) -> String {
let shown = &line[..line.len().min(CORRUPT_EXCERPT_BYTES)];
let mut out: String = String::from_utf8_lossy(shown).escape_debug().to_string();
if line.len() > CORRUPT_EXCERPT_BYTES {
out.push('…');
}
out
}
#[cfg(test)]
mod tests {
use super::*;
use std::io::Write;
use tempfile::TempDir;
fn write_line(path: &Path, seq: u64, kind: &str) {
let line = format!(
"{{\"ts\":\"2026-06-12T00:00:00Z\",\"seq\":{seq},\"kind\":\"{kind}\",\"run_id\":\"01jxsnap000000000000000000\",\"data\":{{}}}}\n"
);
let mut f = std::fs::OpenOptions::new()
.create(true)
.append(true)
.open(path)
.unwrap();
f.write_all(line.as_bytes()).unwrap();
}
#[test]
fn missing_file_yields_empty() {
let dir = TempDir::new().unwrap();
let mut t = EventTail::new(dir.path().join("missing.jsonl"), 0);
assert!(t.poll().unwrap().is_empty());
}
#[test]
fn reads_new_events_only() {
let dir = TempDir::new().unwrap();
let p = dir.path().join("events.jsonl");
write_line(&p, 1, "node.created");
write_line(&p, 2, "node.report");
let mut t = EventTail::new(&p, 0);
let first = t.poll().unwrap();
assert_eq!(first.len(), 2);
assert_eq!(first[0].seq, 1);
assert_eq!(first[1].seq, 2);
assert!(t.poll().unwrap().is_empty());
write_line(&p, 3, "child.spawned");
let next = t.poll().unwrap();
assert_eq!(next.len(), 1);
assert_eq!(next[0].seq, 3);
}
#[test]
fn skips_already_consumed_seqs() {
let dir = TempDir::new().unwrap();
let p = dir.path().join("events.jsonl");
write_line(&p, 1, "a");
write_line(&p, 2, "b");
write_line(&p, 3, "c");
let mut t = EventTail::new(&p, 2);
let evs = t.poll().unwrap();
assert_eq!(evs.len(), 1);
assert_eq!(evs[0].seq, 3);
}
#[test]
fn parks_at_corrupt_line_then_skips_without_looping() {
let dir = TempDir::new().unwrap();
let p = dir.path().join("events.jsonl");
write_line(&p, 1, "a");
{
let mut f = std::fs::OpenOptions::new().append(true).open(&p).unwrap();
f.write_all(b"{garbage not json\n").unwrap();
}
write_line(&p, 2, "b");
let mut t = EventTail::new(&p, 0);
let evs = t.poll().unwrap();
assert_eq!(evs.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![1]);
assert!(t.poll().unwrap().is_empty());
assert!(t.poll().unwrap().is_empty());
let c = t.take_new_corrupt().expect("a corrupt line is parked");
assert!(c.byte_offset > 0);
assert!(
c.line_excerpt.contains("garbage"),
"excerpt: {}",
c.line_excerpt
);
let evs = t.poll().unwrap();
assert_eq!(evs.iter().map(|e| e.seq).collect::<Vec<_>>(), vec![2]);
assert!(t.take_new_corrupt().is_none());
assert!(t.poll().unwrap().is_empty());
}
#[test]
fn ignores_partial_trailing_line() {
let dir = TempDir::new().unwrap();
let p = dir.path().join("events.jsonl");
write_line(&p, 1, "a");
{
let mut f = std::fs::OpenOptions::new().append(true).open(&p).unwrap();
f.write_all(b"{\"ts\":\"2026").unwrap();
}
let mut t = EventTail::new(&p, 0);
let evs = t.poll().unwrap();
assert_eq!(evs.len(), 1);
assert_eq!(evs[0].seq, 1);
{
let mut f = std::fs::OpenOptions::new().append(true).open(&p).unwrap();
f.write_all(
b"-06-12T00:00:00Z\",\"seq\":2,\"kind\":\"b\",\"run_id\":\"01jxsnap000000000000000000\",\"data\":{}}\n",
)
.unwrap();
}
let evs = t.poll().unwrap();
assert_eq!(evs.len(), 1);
assert_eq!(evs[0].seq, 2);
}
}