use std::fs::{self, File, OpenOptions};
use std::io::{BufRead, BufReader, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Mutex;
use std::time::SystemTime;
use crate::rhei_tui::event::{EventSink, RunEvent};
use crate::rhei_tui::event_json;
pub const EVENT_LOG_NAME: &str = "events.jsonl";
pub fn event_log_path(workspace_root: impl AsRef<Path>) -> PathBuf {
workspace_root.as_ref().join("runtime").join(EVENT_LOG_NAME)
}
pub struct EventLogSink {
inner: Mutex<Option<File>>,
seq: AtomicU64,
path: PathBuf,
workspace_root: PathBuf,
}
impl EventLogSink {
pub fn create(workspace_root: impl AsRef<Path>) -> std::io::Result<Self> {
let workspace_root = workspace_root.as_ref().to_path_buf();
let path = event_log_path(&workspace_root);
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let file = OpenOptions::new().create(true).write(true).truncate(true).open(&path)?;
Ok(Self { inner: Mutex::new(Some(file)), seq: AtomicU64::new(0), path, workspace_root })
}
pub fn path(&self) -> &Path {
&self.path
}
pub fn last_seq(&self) -> u64 {
self.seq.load(Ordering::SeqCst)
}
}
impl EventSink for EventLogSink {
fn emit(&self, event: RunEvent) {
if !event_json::is_structural(&event) {
return;
}
let at = event_json::event_wall_clock(&event).unwrap_or_else(SystemTime::now);
let mut guard = match self.inner.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
let Some(file) = guard.as_mut() else {
return;
};
let seq = self.seq.fetch_add(1, Ordering::SeqCst) + 1;
let line =
format!("{}\n", event_json::encode(Some(seq), &event, at, Some(&self.workspace_root)));
if let Err(err) = file.write_all(line.as_bytes()).and_then(|()| file.flush()) {
eprintln!("warning: event log write failed ({}): {err}", self.path.display());
*guard = None;
}
}
}
pub struct EventLogReader {
path: PathBuf,
offset: u64,
last_seq: u64,
}
impl EventLogReader {
pub fn open(path: impl Into<PathBuf>) -> Self {
Self { path: path.into(), offset: 0, last_seq: 0 }
}
pub fn last_seq(&self) -> u64 {
self.last_seq
}
pub fn poll(&mut self) -> Vec<event_json::DecodedRecord> {
let Ok(file) = File::open(&self.path) else {
return Vec::new();
};
let len = file.metadata().map(|m| m.len()).unwrap_or(0);
if len < self.offset {
self.offset = 0;
self.last_seq = 0;
}
let mut reader = BufReader::new(file);
if reader.seek(SeekFrom::Start(self.offset)).is_err() {
return Vec::new();
}
let mut events = Vec::new();
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line) {
Ok(0) => break,
Ok(read) => {
if !line.ends_with('\n') {
break;
}
self.offset += read as u64;
if let Some(record) = event_json::decode(&line) {
if let Some(seq) = record.seq {
self.last_seq = seq;
}
events.push(record);
}
}
Err(_) => break,
}
}
events
}
}
#[cfg(test)]
mod tests;