use std::io::Write;
use std::sync::{Arc, Mutex};
use std::time::Instant;
#[derive(Debug, Clone)]
pub enum ProgressEvent {
StepStarted {
step: usize,
model: String,
tools_available: usize,
},
LlmRequestSent { tokens: usize },
LlmResponseReceived {
finish_reason: String,
completion_tokens: u32,
},
ToolCallStarted {
tool: String,
args_short: String,
},
ToolCallCompleted {
tool: String,
ok: bool,
elapsed_ms: u64,
},
GuardFired {
kind: String,
count: usize,
},
StepCompleted {
step: usize,
mutating_tools_so_far: usize,
},
SubprocessStarted { name: String },
SubprocessCompleted {
name: String,
exit: i32,
elapsed_ms: u64,
},
TaskCompleted { outcome: String },
TaskFailed { reason: String },
TurnDecision { decision: String, detail: String },
}
pub fn short_args_for(tool_name: &str, args: &serde_json::Value) -> String {
fn shorten(s: &str, max: usize) -> String {
if s.chars().count() <= max {
s.replace('\n', " ")
} else {
let mut out: String = s.chars().take(max).collect();
out.push('…');
out.replace('\n', " ")
}
}
let pull = |keys: &[&str]| -> Option<String> {
for k in keys {
if let Some(v) = args.get(*k).and_then(|v| v.as_str()) {
return Some(format!("{}={}", k, shorten(v, 60)));
}
}
None
};
if let Some(s) = pull(&["path", "file_path", "file"]) {
return s;
}
if let Some(s) = pull(&["command", "cmd"]) {
return s;
}
if let Some(s) = pull(&["pattern", "query", "search"]) {
return s;
}
if let Some(s) = pull(&["url"]) {
return s;
}
let _ = tool_name;
String::new()
}
impl ProgressEvent {
pub fn kind(&self) -> &'static str {
match self {
ProgressEvent::StepStarted { .. } => "step_started",
ProgressEvent::LlmRequestSent { .. } => "llm_request_sent",
ProgressEvent::LlmResponseReceived { .. } => "llm_response_received",
ProgressEvent::ToolCallStarted { .. } => "tool_call_started",
ProgressEvent::ToolCallCompleted { .. } => "tool_call_completed",
ProgressEvent::GuardFired { .. } => "guard_fired",
ProgressEvent::StepCompleted { .. } => "step_completed",
ProgressEvent::SubprocessStarted { .. } => "subprocess_started",
ProgressEvent::SubprocessCompleted { .. } => "subprocess_completed",
ProgressEvent::TaskCompleted { .. } => "task_completed",
ProgressEvent::TaskFailed { .. } => "task_failed",
ProgressEvent::TurnDecision { .. } => "turn_decision",
}
}
}
pub trait ProgressEmitter: Send + Sync {
fn emit(&self, event: ProgressEvent);
}
pub struct NoopProgressEmitter;
impl ProgressEmitter for NoopProgressEmitter {
fn emit(&self, _event: ProgressEvent) {}
}
pub struct MultiProgressEmitter {
inner: Vec<Arc<dyn ProgressEmitter>>,
}
impl MultiProgressEmitter {
pub fn new(inner: Vec<Arc<dyn ProgressEmitter>>) -> Self {
Self { inner }
}
}
impl ProgressEmitter for MultiProgressEmitter {
fn emit(&self, event: ProgressEvent) {
for em in &self.inner {
em.emit(event.clone());
}
}
}
pub struct StderrProgressEmitter {
start: Instant,
use_system_clock: bool,
lock: Mutex<()>,
}
impl Default for StderrProgressEmitter {
fn default() -> Self {
Self::new()
}
}
impl StderrProgressEmitter {
pub fn new() -> Self {
Self {
start: Instant::now(),
use_system_clock: true,
lock: Mutex::new(()),
}
}
fn timestamp(&self) -> String {
if self.use_system_clock {
chrono::Local::now().format("%H:%M:%S").to_string()
} else {
let elapsed = self.start.elapsed();
format!(
"{:02}:{:02}:{:02}",
elapsed.as_secs() / 3600,
(elapsed.as_secs() / 60) % 60,
elapsed.as_secs() % 60
)
}
}
}
pub fn render_event_kv(event: &ProgressEvent) -> String {
match event {
ProgressEvent::StepStarted {
step,
model,
tools_available,
} => format!(
"kind=step_started step={} model={} tools_available={}",
step, model, tools_available
),
ProgressEvent::LlmRequestSent { tokens } => {
format!("kind=llm_request_sent prompt_tokens={}", tokens)
}
ProgressEvent::LlmResponseReceived {
finish_reason,
completion_tokens,
} => format!(
"kind=llm_response_received finish_reason={} completion_tokens={}",
finish_reason, completion_tokens
),
ProgressEvent::ToolCallStarted { tool, args_short } => {
if args_short.is_empty() {
format!("kind=tool_call_started tool={}", tool)
} else {
format!("kind=tool_call_started tool={} args={}", tool, args_short)
}
}
ProgressEvent::ToolCallCompleted {
tool,
ok,
elapsed_ms,
} => format!(
"kind=tool_call_completed tool={} ok={} {}ms",
tool, ok, elapsed_ms
),
ProgressEvent::GuardFired { kind, count } => {
format!("kind=guard_fired guard={} count={}", kind, count)
}
ProgressEvent::StepCompleted {
step,
mutating_tools_so_far,
} => format!(
"kind=step_completed step={} mutating_tools_so_far={}",
step, mutating_tools_so_far
),
ProgressEvent::SubprocessStarted { name } => {
format!("kind=subprocess_started name={}", name)
}
ProgressEvent::SubprocessCompleted {
name,
exit,
elapsed_ms,
} => format!(
"kind=subprocess_completed name={} exit={} {}ms",
name, exit, elapsed_ms
),
ProgressEvent::TaskCompleted { outcome } => {
format!("kind=task_completed outcome={}", outcome)
}
ProgressEvent::TaskFailed { reason } => {
format!("kind=task_failed reason={:?}", reason)
}
ProgressEvent::TurnDecision { decision, detail } => {
if detail.is_empty() {
format!("kind=turn_decision decision={}", decision)
} else {
format!(
"kind=turn_decision decision={} detail=\"{}\"",
decision, detail
)
}
}
}
}
impl ProgressEmitter for StderrProgressEmitter {
fn emit(&self, event: ProgressEvent) {
if crate::output::is_tui_active() {
return;
}
let line = format!("[{}] {}", self.timestamp(), render_event_kv(&event));
let _global = crate::output::OUTPUT_LOCK
.lock()
.unwrap_or_else(|e| e.into_inner());
let _local = self.lock.lock().unwrap_or_else(|e| e.into_inner());
let mut stderr = std::io::stderr().lock();
let _ = writeln!(stderr, "{}", line);
let _ = stderr.flush();
}
}
#[cfg(test)]
#[derive(Default, Clone)]
pub struct RecordingProgressEmitter {
pub events: Arc<Mutex<Vec<ProgressEvent>>>,
}
#[cfg(test)]
impl RecordingProgressEmitter {
pub fn new() -> Self {
Self {
events: Arc::new(Mutex::new(Vec::new())),
}
}
pub fn snapshot(&self) -> Vec<ProgressEvent> {
self.events.lock().unwrap().clone()
}
pub fn kinds(&self) -> Vec<&'static str> {
self.snapshot().iter().map(|e| e.kind()).collect()
}
}
#[cfg(test)]
impl ProgressEmitter for RecordingProgressEmitter {
fn emit(&self, event: ProgressEvent) {
self.events.lock().unwrap().push(event);
}
}
#[cfg(test)]
#[path = "../../tests/unit/agent/progress/progress_test.rs"]
mod tests;