use std::fs::{File, OpenOptions};
use std::io::{BufRead, BufReader, Seek, Write};
use std::path::{Path, PathBuf};
use std::process;
use std::sync::atomic::{AtomicBool, Ordering};
use std::sync::Arc;
use std::time::Duration;
use serde::Serialize;
use serde_json::Value;
use tracing::info;
use octl_core::{read_manifest_opt, Event, SCHEMA_VERSION};
use crate::error::CliError;
use crate::event::{resolve_format, FormatArg};
use crate::output::OutputSpec;
use crate::run::{from_core, run_paths};
const POLL_INTERVAL: Duration = Duration::from_millis(500);
const CANCELLED_EXIT_FALLBACK: i32 = 130;
pub struct Args<'a> {
pub run_id: String,
pub from_seq: u64,
pub follow: bool,
pub to_file: Option<PathBuf>,
pub spec: &'a OutputSpec,
pub warnings: &'a [String],
}
pub fn run(args: Args<'_>) -> Result<(), CliError> {
let run_id = &args.run_id;
let format = resolve_format(args.spec.format)?;
let root = crate::home::root_dir()?;
let paths = run_paths(&root, run_id)?;
if read_manifest_opt(&paths).map_err(from_core)?.is_none() {
return Err(
CliError::user("run_not_found", format!("no run with id {run_id}"))
.with_invalid_value(run_id.as_str()),
);
}
let events_path = paths.events();
if let Some(out) = &args.to_file {
reject_output_alias(&events_path, out)?;
}
crate::output::emit_text_warnings(args.warnings);
let cancel = Arc::new(AtomicBool::new(false));
{
let flag = cancel.clone();
let _ = ctrlc::try_set_handler(move || {
flag.store(true, Ordering::Release);
});
}
let mut writer: Box<dyn Write> = match &args.to_file {
None => Box::new(std::io::stdout().lock()),
Some(p) => Box::new(open_output(p, args.follow)?),
};
let mut reader = try_open_events_reader(&events_path)?;
let mut last_seen_seq: Option<u64> = None;
let mut buf = String::new();
if let Some(r) = reader.as_mut() {
drain(
r,
&mut buf,
&mut *writer,
format,
args.from_seq,
&mut last_seen_seq,
&cancel,
)?;
}
if cancel.load(Ordering::Acquire) {
emit_terminal(&mut *writer, format, TerminalKind::Cancelled, last_seen_seq)?;
flush_and_exit(writer, CANCELLED_EXIT_FALLBACK);
}
if !args.follow {
emit_terminal(&mut *writer, format, TerminalKind::Result, last_seen_seq)?;
writer
.flush()
.map_err(|e| CliError::system("io_error", format!("flush: {e}")))?;
return Ok(());
}
loop {
if cancel.load(Ordering::Acquire) {
emit_terminal(&mut *writer, format, TerminalKind::Cancelled, last_seen_seq)?;
flush_and_exit(writer, CANCELLED_EXIT_FALLBACK);
}
std::thread::sleep(POLL_INTERVAL);
if cancel.load(Ordering::Acquire) {
continue;
}
if reader.is_none() {
reader = try_open_events_reader(&events_path)?;
if reader.is_none() {
continue;
}
}
let r = reader.as_mut().unwrap();
let pos = r
.stream_position()
.map_err(|e| CliError::system("io_error", format!("stream_position: {e}")))?;
let len_now = r
.get_ref()
.metadata()
.map_err(|e| CliError::system("io_error", format!("fstat: {e}")))?
.len();
if len_now < pos {
return Err(CliError::system(
"events_log_truncated",
format!(
"{} shrank from {} to {} bytes — append-only contract violated",
events_path.display(),
pos,
len_now
),
));
}
drain(
r,
&mut buf,
&mut *writer,
format,
args.from_seq,
&mut last_seen_seq,
&cancel,
)?;
}
}
fn try_open_events_reader(events_path: &Path) -> Result<Option<BufReader<File>>, CliError> {
match File::open(events_path) {
Ok(f) => Ok(Some(BufReader::new(f))),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => Ok(None),
Err(e) => Err(CliError::system(
"io_error",
format!("open {}: {}", events_path.display(), e),
)),
}
}
fn open_output(p: &Path, follow: bool) -> Result<File, CliError> {
let mut opts = OpenOptions::new();
opts.create(true).write(true);
if follow {
opts.append(true);
} else {
opts.truncate(true);
}
opts.open(p)
.map_err(|e| CliError::system("io_error", format!("open {}: {}", p.display(), e)))
}
fn reject_output_alias(events_path: &Path, output: &Path) -> Result<(), CliError> {
if let (Ok(ev), Ok(out)) = (
std::fs::canonicalize(events_path),
std::fs::canonicalize(output),
) {
if ev == out {
return Err(CliError::user(
"invalid_output",
"--output must not point at the run's events.jsonl",
)
.with_invalid_value(output.display().to_string()));
}
}
Ok(())
}
fn drain(
reader: &mut BufReader<File>,
buf: &mut String,
writer: &mut dyn Write,
format: FormatArg,
from_seq: u64,
last_seen_seq: &mut Option<u64>,
cancel: &AtomicBool,
) -> Result<(), CliError> {
let mut emitted_any = false;
loop {
if cancel.load(Ordering::Acquire) {
break;
}
let n = reader
.read_line(buf)
.map_err(|e| CliError::system("io_error", format!("read_line: {e}")))?;
if n == 0 {
break;
}
if !buf.ends_with('\n') {
break;
}
let line = buf.trim_end_matches(['\n', '\r']);
if !line.is_empty() {
let ev: Event = serde_json::from_str(line).map_err(|e| {
CliError::system(
"events_log_corrupt",
format!("parse event line: {e}: {line}"),
)
})?;
let dedup_ok = match *last_seen_seq {
Some(last) => ev.seq > last,
None => true,
};
if ev.seq >= from_seq && dedup_ok {
emit_event(writer, format, &ev)?;
emitted_any = true;
}
if dedup_ok {
*last_seen_seq = Some(ev.seq);
}
}
buf.clear();
}
if emitted_any {
writer
.flush()
.map_err(|e| CliError::system("io_error", format!("flush: {e}")))?;
}
Ok(())
}
fn emit_event(writer: &mut dyn Write, format: FormatArg, ev: &Event) -> Result<(), CliError> {
match format {
FormatArg::Jsonl => {
let line = serde_json::to_string(ev)
.map_err(|e| CliError::system("internal_serialize", e.to_string()))?;
writeln!(writer, "{line}")
}
FormatArg::Text => writeln!(writer, "{}", text_summary(ev)),
}
.map_err(|e| CliError::system("io_error", format!("write event: {e}")))
}
fn text_summary(ev: &Event) -> String {
let node = ev
.node_id
.as_ref()
.map(|n| format!(" node={n}"))
.unwrap_or_default();
let detail = text_data_detail(&ev.data);
let detail = if detail.is_empty() {
String::new()
} else {
format!(" {detail}")
};
format!(
"[{seq}] {ts} {kind}{node}{detail}",
seq = ev.seq,
ts = ev.ts.to_rfc3339(),
kind = ev.kind,
)
}
fn text_data_detail(data: &Value) -> String {
let obj = match data.as_object() {
Some(o) => o,
None => return String::new(),
};
for key in ["status", "title", "reason", "message"] {
if let Some(v) = obj.get(key).and_then(Value::as_str) {
return format!("{key}={}", crate::output::escape_one_line(v));
}
}
String::new()
}
#[derive(Debug, Clone, Copy)]
enum TerminalKind {
Result,
Cancelled,
}
#[derive(Serialize)]
struct TerminalEnvelope<'a> {
schema_version: u32,
event: &'a str,
#[serde(skip_serializing_if = "Option::is_none")]
status: Option<&'a str>,
#[serde(skip_serializing_if = "Option::is_none")]
last_seq: Option<u64>,
}
fn emit_terminal(
writer: &mut dyn Write,
format: FormatArg,
kind: TerminalKind,
last_seen_seq: Option<u64>,
) -> Result<(), CliError> {
if matches!(format, FormatArg::Text) {
return Ok(());
}
let envelope = match kind {
TerminalKind::Result => TerminalEnvelope {
schema_version: SCHEMA_VERSION,
event: "result",
status: Some("ok"),
last_seq: last_seen_seq,
},
TerminalKind::Cancelled => TerminalEnvelope {
schema_version: SCHEMA_VERSION,
event: "cancelled",
status: None,
last_seq: last_seen_seq,
},
};
let line = serde_json::to_string(&envelope)
.map_err(|e| CliError::system("internal_serialize", e.to_string()))?;
writeln!(writer, "{line}")
.map_err(|e| CliError::system("io_error", format!("write terminal: {e}")))
}
fn flush_and_exit(mut writer: Box<dyn Write>, code: i32) -> ! {
let _ = writer.flush();
drop(writer);
let _ = std::io::stdout().lock().flush();
info!(target: "orchestratectl::event::tail", code, "event tail exiting via signal");
crate::cli::flush_logs();
process::exit(code);
}