use std::io::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 struct JsonSink {
seq: AtomicU64,
out: Mutex<()>,
agent_output: bool,
workspace_root: PathBuf,
}
impl JsonSink {
pub fn new(agent_output: bool, workspace_root: impl AsRef<Path>) -> Self {
Self {
seq: AtomicU64::new(0),
out: Mutex::new(()),
agent_output,
workspace_root: workspace_root.as_ref().to_path_buf(),
}
}
}
impl EventSink for JsonSink {
fn emit(&self, event: RunEvent) {
if !self.agent_output && !event_json::is_structural(&event) {
return;
}
let at = event_json::event_wall_clock(&event).unwrap_or_else(SystemTime::now);
let guard = match self.out.lock() {
Ok(g) => g,
Err(p) => p.into_inner(),
};
let seq =
event_json::is_structural(&event).then(|| self.seq.fetch_add(1, Ordering::SeqCst) + 1);
let mut stdout = std::io::stdout().lock();
let _ =
writeln!(stdout, "{}", event_json::encode(seq, &event, at, Some(&self.workspace_root)));
let _ = stdout.flush();
drop(guard);
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::rhei_tui::event::{AgentStream, MessageLevel};
fn output_event() -> RunEvent {
RunEvent::AgentOutput {
slot: 0,
task: "auth.1".to_string(),
stream: AgentStream::Stdout,
line: "hello".to_string(),
wall_clock: SystemTime::now(),
}
}
#[test]
fn agent_output_is_dropped_unless_it_was_asked_for() {
let quiet = JsonSink::new(false, "/nowhere");
assert!(!quiet.agent_output);
quiet.emit(output_event());
assert_eq!(quiet.seq.load(Ordering::SeqCst), 0, "a dropped event burns no sequence number");
let loud = JsonSink::new(true, "/nowhere");
loud.emit(output_event());
assert_eq!(loud.seq.load(Ordering::SeqCst), 0);
loud.emit(RunEvent::Message { level: MessageLevel::Info, text: "structural".to_string() });
assert_eq!(loud.seq.load(Ordering::SeqCst), 1, "structural records number from 1");
}
#[test]
fn structural_events_are_always_written() {
let sink = JsonSink::new(false, "/nowhere");
sink.emit(RunEvent::Message { level: MessageLevel::Info, text: "hi".to_string() });
assert_eq!(sink.seq.load(Ordering::SeqCst), 1);
}
}