Refactor run progress rendering

This commit is contained in:
Bryan Helmkamp 2026-03-30 14:40:36 -04:00
parent 63a258829b
commit 727c8ba1d3
No known key found for this signature in database
8 changed files with 3080 additions and 2155 deletions

File diff suppressed because it is too large Load diff

View file

@ -0,0 +1,679 @@
use std::convert::TryFrom;
use fabro_workflow::event::RunNoticeLevel;
use fabro_workflow::outcome::{StageUsage, compute_stage_cost};
use serde_json::{Map, Value};
#[derive(Debug, Clone)]
pub(super) struct ProgressUsage {
pub(super) model: Option<String>,
pub(super) input_tokens: u64,
pub(super) output_tokens: u64,
pub(super) speed: Option<String>,
pub(super) cost: Option<f64>,
}
impl ProgressUsage {
pub(super) fn from_value(value: &Value) -> Option<Self> {
let Value::Object(fields) = value else {
return None;
};
Some(Self {
model: string_field(fields, "model"),
input_tokens: u64_field(fields, "input_tokens"),
output_tokens: u64_field(fields, "output_tokens"),
speed: string_field(fields, "speed"),
cost: f64_field(fields, "cost"),
})
}
pub(super) fn total_tokens(&self) -> u64 {
self.input_tokens.saturating_add(self.output_tokens)
}
pub(super) fn display_cost(&self) -> Option<f64> {
self.cost.or_else(|| {
let model = self.model.clone()?;
let input_tokens = i64::try_from(self.input_tokens).ok()?;
let output_tokens = i64::try_from(self.output_tokens).ok()?;
let usage = StageUsage {
model,
input_tokens,
output_tokens,
cache_read_tokens: None,
cache_write_tokens: None,
reasoning_tokens: None,
speed: self.speed.clone(),
cost: None,
};
compute_stage_cost(&usage)
})
}
}
#[derive(Debug, Clone)]
pub(super) enum ProgressEvent {
WorkflowStarted {
worktree_dir: Option<String>,
base_branch: Option<String>,
base_sha: Option<String>,
},
WorkingDirectorySet {
working_directory: String,
},
SandboxInitializing {
provider: String,
},
SandboxReady {
provider: String,
duration_ms: u64,
name: Option<String>,
cpu: Option<f64>,
memory: Option<f64>,
url: Option<String>,
},
SshAccessReady {
ssh_command: String,
},
SetupStarted {
command_count: u64,
},
SetupCompleted {
duration_ms: u64,
},
SetupCommandCompleted {
command: String,
command_index: u64,
exit_code: i64,
duration_ms: u64,
},
CliEnsureStarted {
cli_name: String,
},
CliEnsureCompleted {
cli_name: String,
already_installed: bool,
duration_ms: u64,
},
CliEnsureFailed {
cli_name: String,
},
DevcontainerResolved {
dockerfile_lines: u64,
environment_count: u64,
lifecycle_command_count: u64,
workspace_folder: String,
},
DevcontainerLifecycleStarted {
phase: String,
command_count: u64,
},
DevcontainerLifecycleCompleted {
phase: String,
duration_ms: u64,
},
DevcontainerLifecycleFailed {
phase: String,
command: String,
exit_code: i64,
stderr: String,
},
DevcontainerLifecycleCommandCompleted {
command: String,
command_index: u64,
exit_code: i64,
duration_ms: u64,
},
StageStarted {
node_id: String,
name: String,
script: Option<String>,
},
StageCompleted {
node_id: String,
name: String,
duration_ms: u64,
status: String,
usage: Option<ProgressUsage>,
},
StageFailed {
node_id: String,
name: String,
error: String,
},
StageRetrying {
name: String,
attempt: u64,
max_attempts: u64,
delay_ms: u64,
},
ParallelStarted,
ParallelBranchStarted {
branch: String,
},
ParallelBranchCompleted {
branch: String,
duration_ms: u64,
status: String,
},
ParallelCompleted,
AssistantMessage {
stage_node_id: String,
model: String,
},
ToolCallStarted {
stage_node_id: String,
tool_name: String,
tool_call_id: String,
arguments: Value,
},
ToolCallCompleted {
stage_node_id: String,
tool_call_id: String,
is_error: bool,
},
ContextWindowWarning {
stage_node_id: String,
usage_percent: u64,
},
CompactionStarted {
stage_node_id: String,
},
CompactionCompleted {
stage_node_id: String,
original_turn_count: u64,
preserved_turn_count: u64,
tracked_file_count: u64,
},
LlmRetry {
stage_node_id: String,
model: String,
attempt: u64,
delay_ms: u64,
error: String,
},
SubagentSpawned {
stage_node_id: String,
agent_id: String,
task: String,
},
SubagentCompleted {
stage_node_id: String,
agent_id: String,
success: bool,
turns_used: u64,
},
EdgeSelected {
from_node: String,
to_node: String,
label: Option<String>,
condition: Option<String>,
},
LoopRestart {
from_node: String,
to_node: String,
},
RetroStarted,
RetroCompleted {
duration_ms: u64,
},
RetroFailed {
duration_ms: u64,
},
RunNotice {
level: RunNoticeLevel,
code: String,
message: String,
},
PullRequestCreated {
pr_url: String,
draft: bool,
},
PullRequestFailed {
error: String,
},
}
pub(super) fn from_flattened_fields(
event_name: &str,
fields: Map<String, Value>,
) -> Option<ProgressEvent> {
match event_name {
"WorkflowRunStarted" => Some(ProgressEvent::WorkflowStarted {
worktree_dir: string_field(&fields, "worktree_dir"),
base_branch: string_field(&fields, "base_branch"),
base_sha: string_field(&fields, "base_sha"),
}),
"SandboxInitialized" => Some(ProgressEvent::WorkingDirectorySet {
working_directory: string_field(&fields, "working_directory")?,
}),
"Sandbox.Initializing" => Some(ProgressEvent::SandboxInitializing {
provider: string_field(&fields, "sandbox_provider")
.or_else(|| string_field(&fields, "provider"))
.unwrap_or_else(|| "unknown".to_string()),
}),
"Sandbox.Ready" => Some(ProgressEvent::SandboxReady {
provider: string_field(&fields, "sandbox_provider")
.or_else(|| string_field(&fields, "provider"))
.unwrap_or_else(|| "unknown".to_string()),
duration_ms: u64_field(&fields, "duration_ms"),
name: string_field(&fields, "name"),
cpu: f64_field(&fields, "cpu"),
memory: f64_field(&fields, "memory"),
url: string_field(&fields, "url"),
}),
"SshAccessReady" => Some(ProgressEvent::SshAccessReady {
ssh_command: string_field(&fields, "ssh_command")?,
}),
"SetupStarted" => Some(ProgressEvent::SetupStarted {
command_count: u64_field(&fields, "command_count"),
}),
"SetupCompleted" => Some(ProgressEvent::SetupCompleted {
duration_ms: u64_field(&fields, "duration_ms"),
}),
"SetupCommandCompleted" => Some(ProgressEvent::SetupCommandCompleted {
command: string_field(&fields, "command").unwrap_or_else(|| "?".to_string()),
command_index: u64_field(&fields, "command_index").max(u64_field(&fields, "index")),
exit_code: i64_field(&fields, "exit_code"),
duration_ms: u64_field(&fields, "duration_ms"),
}),
"CliEnsureStarted" => Some(ProgressEvent::CliEnsureStarted {
cli_name: string_field(&fields, "cli_name").unwrap_or_else(|| "?".to_string()),
}),
"CliEnsureCompleted" => Some(ProgressEvent::CliEnsureCompleted {
cli_name: string_field(&fields, "cli_name").unwrap_or_else(|| "?".to_string()),
already_installed: bool_field(&fields, "already_installed"),
duration_ms: u64_field(&fields, "duration_ms"),
}),
"CliEnsureFailed" => Some(ProgressEvent::CliEnsureFailed {
cli_name: string_field(&fields, "cli_name").unwrap_or_else(|| "?".to_string()),
}),
"DevcontainerResolved" => Some(ProgressEvent::DevcontainerResolved {
dockerfile_lines: u64_field(&fields, "dockerfile_lines"),
environment_count: u64_field(&fields, "environment_count"),
lifecycle_command_count: u64_field(&fields, "lifecycle_command_count"),
workspace_folder: string_field(&fields, "workspace_folder")
.unwrap_or_else(|| "?".to_string()),
}),
"DevcontainerLifecycleStarted" => Some(ProgressEvent::DevcontainerLifecycleStarted {
phase: string_field(&fields, "phase").unwrap_or_else(|| "?".to_string()),
command_count: u64_field(&fields, "command_count"),
}),
"DevcontainerLifecycleCompleted" => Some(ProgressEvent::DevcontainerLifecycleCompleted {
phase: string_field(&fields, "phase").unwrap_or_else(|| "?".to_string()),
duration_ms: u64_field(&fields, "duration_ms"),
}),
"DevcontainerLifecycleFailed" => Some(ProgressEvent::DevcontainerLifecycleFailed {
phase: string_field(&fields, "phase").unwrap_or_else(|| "?".to_string()),
command: string_field(&fields, "command").unwrap_or_else(|| "?".to_string()),
exit_code: i64_field(&fields, "exit_code"),
stderr: display_field(&fields, "stderr").unwrap_or_default(),
}),
"DevcontainerLifecycleCommandCompleted" => {
Some(ProgressEvent::DevcontainerLifecycleCommandCompleted {
command: string_field(&fields, "command").unwrap_or_else(|| "?".to_string()),
command_index: u64_field(&fields, "command_index").max(u64_field(&fields, "index")),
exit_code: i64_field(&fields, "exit_code"),
duration_ms: u64_field(&fields, "duration_ms"),
})
}
"StageStarted" => Some(ProgressEvent::StageStarted {
node_id: string_field(&fields, "node_id").unwrap_or_else(|| "?".to_string()),
name: string_field(&fields, "node_label")
.or_else(|| string_field(&fields, "name"))
.unwrap_or_else(|| "?".to_string()),
script: string_field(&fields, "script"),
}),
"StageCompleted" => Some(ProgressEvent::StageCompleted {
node_id: string_field(&fields, "node_id").unwrap_or_else(|| "?".to_string()),
name: string_field(&fields, "node_label")
.or_else(|| string_field(&fields, "name"))
.unwrap_or_else(|| "?".to_string()),
duration_ms: u64_field(&fields, "duration_ms"),
status: string_field(&fields, "status").unwrap_or_else(|| "success".to_string()),
usage: fields.get("usage").and_then(ProgressUsage::from_value),
}),
"StageFailed" => Some(ProgressEvent::StageFailed {
node_id: string_field(&fields, "node_id").unwrap_or_else(|| "?".to_string()),
name: string_field(&fields, "node_label")
.or_else(|| string_field(&fields, "name"))
.unwrap_or_else(|| "?".to_string()),
error: display_field(&fields, "error")
.or_else(|| display_field(&fields, "failure_reason"))
.unwrap_or_else(|| "unknown error".to_string()),
}),
"StageRetrying" => Some(ProgressEvent::StageRetrying {
name: string_field(&fields, "node_label")
.or_else(|| string_field(&fields, "name"))
.unwrap_or_else(|| "?".to_string()),
attempt: u64_field(&fields, "attempt"),
max_attempts: u64_field(&fields, "max_attempts"),
delay_ms: u64_field(&fields, "delay_ms"),
}),
"ParallelStarted" => Some(ProgressEvent::ParallelStarted),
"ParallelBranchStarted" => Some(ProgressEvent::ParallelBranchStarted {
branch: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "branch"))
.unwrap_or_else(|| "?".to_string()),
}),
"ParallelBranchCompleted" => Some(ProgressEvent::ParallelBranchCompleted {
branch: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "branch"))
.unwrap_or_else(|| "?".to_string()),
duration_ms: u64_field(&fields, "duration_ms"),
status: string_field(&fields, "status").unwrap_or_else(|| "success".to_string()),
}),
"ParallelCompleted" => Some(ProgressEvent::ParallelCompleted),
"Agent.AssistantMessage" => Some(ProgressEvent::AssistantMessage {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
model: string_field(&fields, "model").unwrap_or_else(|| "?".to_string()),
}),
"Agent.ToolCallStarted" => Some(ProgressEvent::ToolCallStarted {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
tool_name: string_field(&fields, "tool_name").unwrap_or_else(|| "?".to_string()),
tool_call_id: string_field(&fields, "tool_call_id").unwrap_or_else(|| "?".to_string()),
arguments: fields
.get("arguments")
.cloned()
.unwrap_or_else(|| Value::Object(Map::new())),
}),
"Agent.ToolCallCompleted" => Some(ProgressEvent::ToolCallCompleted {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
tool_call_id: string_field(&fields, "tool_call_id").unwrap_or_else(|| "?".to_string()),
is_error: bool_field(&fields, "is_error"),
}),
"Agent.Warning" if string_field(&fields, "kind").as_deref() == Some("context_window") => {
let usage_percent = fields
.get("details")
.and_then(Value::as_object)
.and_then(|details| details.get("usage_percent"))
.and_then(Value::as_u64)
.unwrap_or(0);
Some(ProgressEvent::ContextWindowWarning {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
usage_percent,
})
}
"Agent.CompactionStarted" => Some(ProgressEvent::CompactionStarted {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
}),
"Agent.CompactionCompleted" => Some(ProgressEvent::CompactionCompleted {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
original_turn_count: u64_field(&fields, "original_turn_count"),
preserved_turn_count: u64_field(&fields, "preserved_turn_count"),
tracked_file_count: u64_field(&fields, "tracked_file_count"),
}),
"Agent.LlmRetry" => {
let delay_secs = f64_field(&fields, "delay_secs").unwrap_or(0.0);
#[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)]
let delay_ms = (delay_secs * 1000.0) as u64;
Some(ProgressEvent::LlmRetry {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
model: string_field(&fields, "model").unwrap_or_else(|| "?".to_string()),
attempt: u64_field(&fields, "attempt"),
delay_ms,
error: display_field(&fields, "error")
.unwrap_or_else(|| "unknown error".to_string()),
})
}
"Agent.SubAgentSpawned" => Some(ProgressEvent::SubagentSpawned {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
agent_id: string_field(&fields, "agent_id").unwrap_or_else(|| "?".to_string()),
task: string_field(&fields, "task").unwrap_or_default(),
}),
"Agent.SubAgentCompleted" => Some(ProgressEvent::SubagentCompleted {
stage_node_id: string_field(&fields, "node_id")
.or_else(|| string_field(&fields, "stage"))
.unwrap_or_else(|| "?".to_string()),
agent_id: string_field(&fields, "agent_id").unwrap_or_else(|| "?".to_string()),
success: bool_field(&fields, "success"),
turns_used: u64_field(&fields, "turns_used"),
}),
"EdgeSelected" => Some(ProgressEvent::EdgeSelected {
from_node: string_field(&fields, "from_node_id")
.or_else(|| string_field(&fields, "from_node"))
.unwrap_or_else(|| "?".to_string()),
to_node: string_field(&fields, "to_node_id")
.or_else(|| string_field(&fields, "to_node"))
.unwrap_or_else(|| "?".to_string()),
label: string_field(&fields, "label"),
condition: string_field(&fields, "condition"),
}),
"LoopRestart" => Some(ProgressEvent::LoopRestart {
from_node: string_field(&fields, "from_node_id")
.or_else(|| string_field(&fields, "from_node"))
.unwrap_or_else(|| "?".to_string()),
to_node: string_field(&fields, "to_node_id")
.or_else(|| string_field(&fields, "to_node"))
.unwrap_or_else(|| "?".to_string()),
}),
"RetroStarted" => Some(ProgressEvent::RetroStarted),
"RetroCompleted" => Some(ProgressEvent::RetroCompleted {
duration_ms: u64_field(&fields, "duration_ms"),
}),
"RetroFailed" => Some(ProgressEvent::RetroFailed {
duration_ms: u64_field(&fields, "duration_ms"),
}),
"RunNotice" => Some(ProgressEvent::RunNotice {
level: parse_run_notice_level(string_field(&fields, "level").as_deref()),
code: string_field(&fields, "code").unwrap_or_default(),
message: string_field(&fields, "message").unwrap_or_default(),
}),
"PullRequestCreated" => Some(ProgressEvent::PullRequestCreated {
pr_url: string_field(&fields, "pr_url").unwrap_or_else(|| "?".to_string()),
draft: bool_field(&fields, "draft"),
}),
"PullRequestFailed" => Some(ProgressEvent::PullRequestFailed {
error: display_field(&fields, "error").unwrap_or_else(|| "unknown error".to_string()),
}),
_ => None,
}
}
fn parse_run_notice_level(level: Option<&str>) -> RunNoticeLevel {
match level.unwrap_or("info") {
"warn" => RunNoticeLevel::Warn,
"error" => RunNoticeLevel::Error,
_ => RunNoticeLevel::Info,
}
}
fn string_field(fields: &Map<String, Value>, key: &str) -> Option<String> {
fields.get(key).and_then(Value::as_str).map(str::to_owned)
}
fn display_field(fields: &Map<String, Value>, key: &str) -> Option<String> {
let value = fields.get(key)?;
match value {
Value::Null => None,
Value::String(value) => Some(value.clone()),
Value::Object(map) => map
.get("message")
.and_then(Value::as_str)
.map(str::to_owned)
.or_else(|| {
map.get("detail")
.and_then(Value::as_object)
.and_then(|detail| detail.get("message"))
.and_then(Value::as_str)
.map(str::to_owned)
})
.or_else(|| {
map.get("data")
.and_then(Value::as_object)
.and_then(|detail| detail.get("message"))
.and_then(Value::as_str)
.map(str::to_owned)
})
.or_else(|| map.get("data").and_then(Value::as_str).map(str::to_owned))
.or_else(|| Some(value.to_string())),
_ => Some(value.to_string()),
}
}
fn u64_field(fields: &Map<String, Value>, key: &str) -> u64 {
fields.get(key).and_then(Value::as_u64).unwrap_or(0)
}
fn i64_field(fields: &Map<String, Value>, key: &str) -> i64 {
fields.get(key).and_then(Value::as_i64).unwrap_or(0)
}
fn f64_field(fields: &Map<String, Value>, key: &str) -> Option<f64> {
fields.get(key).and_then(Value::as_f64)
}
fn bool_field(fields: &Map<String, Value>, key: &str) -> bool {
fields.get(key).and_then(Value::as_bool).unwrap_or(false)
}
#[cfg(test)]
mod tests {
use fabro_agent::AgentEvent;
use fabro_workflow::event::WorkflowRunEvent;
use fabro_workflow::event::flatten_event;
use super::*;
fn json_map(value: Value) -> Map<String, Value> {
value.as_object().cloned().expect("json object")
}
#[test]
fn parse_edge_selected() {
let fields = json_map(serde_json::json!({
"from_node_id": "a",
"to_node_id": "b",
"label": "yes"
}));
let event = from_flattened_fields("EdgeSelected", fields).unwrap();
assert!(matches!(
event,
ProgressEvent::EdgeSelected {
from_node,
to_node,
label,
..
} if from_node == "a" && to_node == "b" && label.as_deref() == Some("yes")
));
}
#[test]
fn round_trip_stage_completed() {
let event = WorkflowRunEvent::StageCompleted {
node_id: "plan".into(),
name: "Plan".into(),
index: 0,
duration_ms: 5000,
status: "success".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
usage: None,
failure: None,
notes: None,
files_touched: Vec::new(),
attempt: 1,
max_attempts: 1,
};
let (name, fields) = flatten_event(&event);
let parsed = from_flattened_fields(&name, fields).unwrap();
assert!(matches!(
parsed,
ProgressEvent::StageCompleted {
node_id,
name,
duration_ms,
..
} if node_id == "plan" && name == "Plan" && duration_ms == 5000
));
}
#[test]
fn round_trip_agent_tool_call() {
let event = WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::ToolCallStarted {
tool_name: "read_file".into(),
tool_call_id: "tc1".into(),
arguments: serde_json::json!({"path": "src/main.rs"}),
},
};
let (name, fields) = flatten_event(&event);
let parsed = from_flattened_fields(&name, fields).unwrap();
assert!(matches!(
parsed,
ProgressEvent::ToolCallStarted {
stage_node_id,
tool_name,
tool_call_id,
..
} if stage_node_id == "code" && tool_name == "read_file" && tool_call_id == "tc1"
));
}
#[test]
fn round_trip_sandbox_ready() {
let event = WorkflowRunEvent::Sandbox {
event: fabro_agent::SandboxEvent::Ready {
provider: "daytona".into(),
duration_ms: 2500,
name: Some("sandbox-1".into()),
cpu: Some(4.0),
memory: Some(8.0),
url: Some("https://example.test".into()),
},
};
let (name, fields) = flatten_event(&event);
let parsed = from_flattened_fields(&name, fields).unwrap();
assert!(matches!(
parsed,
ProgressEvent::SandboxReady {
provider,
duration_ms,
name,
..
} if provider == "daytona" && duration_ms == 2500 && name.as_deref() == Some("sandbox-1")
));
}
#[test]
fn round_trip_run_notice() {
let event = WorkflowRunEvent::RunNotice {
level: RunNoticeLevel::Warn,
code: "sandbox_cleanup_failed".into(),
message: "sandbox cleanup failed".into(),
};
let (name, fields) = flatten_event(&event);
let parsed = from_flattened_fields(&name, fields).unwrap();
assert!(matches!(
parsed,
ProgressEvent::RunNotice {
level: RunNoticeLevel::Warn,
code,
message,
} if code == "sandbox_cleanup_failed" && message == "sandbox cleanup failed"
));
}
}

View file

@ -0,0 +1,148 @@
use std::path::Path;
use fabro_workflow::event::RunNoticeLevel;
use super::renderer::ProgressRenderer;
use super::styles;
use crate::shared::{format_duration_ms, tilde_path};
pub(super) struct InfoDisplay {
verbose: bool,
}
impl InfoDisplay {
pub(super) fn new(verbose: bool) -> Self {
Self { verbose }
}
pub(super) fn show_worktree(&self, renderer: &ProgressRenderer, path: &Path) {
self.insert_info_line(renderer, &format!("Worktree: {}", tilde_path(path)));
}
pub(super) fn show_base_info(
&self,
renderer: &ProgressRenderer,
branch: Option<&str>,
sha: &str,
) {
let short_sha = &sha[..sha.len().min(12)];
let text = match branch {
Some(branch) => format!("Base: {branch} ({short_sha})"),
None => format!("Base: {short_sha}"),
};
self.insert_info_line(renderer, &text);
}
pub(super) fn on_run_notice(
&self,
renderer: &ProgressRenderer,
level: RunNoticeLevel,
code: &str,
message: &str,
) {
let styles = renderer.styles();
let label = match level {
RunNoticeLevel::Info => styles.bold.apply_to("Info:").to_string(),
RunNoticeLevel::Warn => styles.yellow.apply_to("Warning:").to_string(),
RunNoticeLevel::Error => styles.red.apply_to("Error:").to_string(),
};
let code_suffix = if code.is_empty() {
String::new()
} else {
format!(" {}", styles.dim.apply_to(format!("[{code}]")))
};
self.insert_info_line(renderer, &format!("{label} {message}{code_suffix}"));
}
pub(super) fn on_pull_request_created(
&self,
renderer: &ProgressRenderer,
pr_url: &str,
draft: bool,
) {
let label = if draft { "Draft PR:" } else { "PR:" };
self.insert_info_line(
renderer,
&format!("{} {pr_url}", renderer.styles().bold.apply_to(label)),
);
}
pub(super) fn on_pull_request_failed(&self, renderer: &ProgressRenderer, error: &str) {
self.insert_info_line(
renderer,
&format!("{} {error}", renderer.styles().red.apply_to("PR failed:")),
);
}
pub(super) fn on_edge_selected(
&self,
renderer: &ProgressRenderer,
from_node: &str,
to_node: &str,
label: Option<&str>,
condition: Option<&str>,
) {
if !self.verbose {
return;
}
let detail = if let Some(condition) = condition {
format!(" [{condition}]")
} else if let Some(label) = label {
format!(" \"{label}\"")
} else {
String::new()
};
self.insert_info_line(
renderer,
&format!("\u{2192} {from_node} \u{2192} {to_node}{detail}"),
);
}
pub(super) fn on_loop_restart(
&self,
renderer: &ProgressRenderer,
from_node: &str,
to_node: &str,
) {
if !self.verbose {
return;
}
self.insert_info_line(
renderer,
&format!("\u{21ba} {from_node} \u{2192} {to_node} (loop restart)"),
);
}
pub(super) fn on_stage_retrying(
&self,
renderer: &ProgressRenderer,
name: &str,
attempt: u64,
max_attempts: u64,
delay_ms: u64,
) {
if !self.verbose {
return;
}
self.insert_info_line(
renderer,
&format!(
"\u{21bb} {name}: retrying (attempt {attempt}/{max_attempts}, delay {})",
format_duration_ms(delay_ms)
),
);
}
fn insert_info_line(&self, renderer: &ProgressRenderer, message: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_static_dim());
bar.finish_with_message(message.to_string());
} else {
renderer.print_line(4, message);
}
}
}

View file

@ -0,0 +1,965 @@
use serde_json::Value;
use fabro_workflow::event::{WorkflowRunEvent, flatten_event};
mod event;
mod info_display;
mod renderer;
mod setup_display;
mod stage_display;
mod styles;
use event::{ProgressEvent, from_flattened_fields};
use info_display::InfoDisplay;
use renderer::ProgressRenderer;
use setup_display::SetupDisplay;
use stage_display::StageDisplay;
pub(crate) struct ProgressUI {
renderer: ProgressRenderer,
stage: StageDisplay,
setup: SetupDisplay,
info: InfoDisplay,
}
impl ProgressUI {
pub(crate) fn new(is_tty: bool, verbose: bool) -> Self {
let renderer = if is_tty {
ProgressRenderer::new_tty()
} else {
ProgressRenderer::new_plain(
Box::new(std::io::stderr()),
console::colors_enabled_stderr(),
)
};
Self::with_renderer(renderer, verbose)
}
fn with_renderer(renderer: ProgressRenderer, verbose: bool) -> Self {
Self {
renderer,
stage: StageDisplay::new(verbose),
setup: SetupDisplay::new(verbose),
info: InfoDisplay::new(verbose),
}
}
#[cfg(test)]
fn new_plain_test(out: Box<dyn std::io::Write + Send>, verbose: bool, colors: bool) -> Self {
Self::with_renderer(ProgressRenderer::new_plain(out, colors), verbose)
}
pub(crate) fn set_working_directory(&mut self, dir: String) {
self.stage.set_working_directory(dir);
}
pub(crate) fn hide_bars(&self) {
self.renderer.hide();
}
pub(crate) fn show_bars(&self) {
self.renderer.show();
}
pub(crate) fn finish(&mut self) {
self.stage.finish();
self.setup.finish();
self.renderer.finish();
}
#[cfg_attr(not(test), allow(dead_code))]
pub(crate) fn handle_event(&mut self, event: &WorkflowRunEvent) {
let (event_name, fields) = flatten_event(event);
if let Some(progress_event) = from_flattened_fields(&event_name, fields) {
self.dispatch(progress_event);
}
}
pub(crate) fn handle_json_line(&mut self, line: &str) {
let Ok(Value::Object(mut envelope)) = serde_json::from_str(line) else {
return;
};
let Some(event_name) = envelope
.remove("event")
.and_then(|value| value.as_str().map(str::to_owned))
else {
return;
};
if let Some(progress_event) = from_flattened_fields(&event_name, envelope) {
self.dispatch(progress_event);
}
}
fn dispatch(&mut self, event: ProgressEvent) {
let renderer = &self.renderer;
match event {
ProgressEvent::WorkflowStarted {
worktree_dir,
base_branch,
base_sha,
} => {
if let Some(worktree_dir) = worktree_dir {
self.info
.show_worktree(renderer, std::path::Path::new(&worktree_dir));
}
if let Some(base_sha) = base_sha {
self.info
.show_base_info(renderer, base_branch.as_deref(), &base_sha);
}
}
ProgressEvent::WorkingDirectorySet { working_directory } => {
self.set_working_directory(working_directory);
}
ProgressEvent::SandboxInitializing { provider } => {
self.setup.on_sandbox_initializing(renderer, &provider);
}
ProgressEvent::SandboxReady {
provider,
duration_ms,
name,
cpu,
memory,
url,
} => {
self.setup.on_sandbox_ready(
renderer,
&provider,
duration_ms,
name.as_deref(),
cpu,
memory,
url.as_deref(),
);
}
ProgressEvent::SshAccessReady { ssh_command } => {
self.setup.on_ssh_access_ready(renderer, &ssh_command);
}
ProgressEvent::SetupStarted { command_count } => {
self.setup.on_setup_started(renderer, command_count);
}
ProgressEvent::SetupCompleted { duration_ms } => {
self.setup.on_setup_completed(renderer, duration_ms);
}
ProgressEvent::SetupCommandCompleted {
command,
command_index,
exit_code,
duration_ms,
} => {
self.setup.on_setup_command_completed(
renderer,
&command,
command_index,
exit_code,
duration_ms,
);
}
ProgressEvent::CliEnsureStarted { cli_name } => {
self.setup.on_cli_ensure_started(renderer, &cli_name);
}
ProgressEvent::CliEnsureCompleted {
cli_name,
already_installed,
duration_ms,
} => {
self.setup.on_cli_ensure_completed(
renderer,
&cli_name,
already_installed,
duration_ms,
);
}
ProgressEvent::CliEnsureFailed { cli_name } => {
self.setup.on_cli_ensure_failed(renderer, &cli_name);
}
ProgressEvent::DevcontainerResolved {
dockerfile_lines,
environment_count,
lifecycle_command_count,
workspace_folder,
} => {
self.setup.on_devcontainer_resolved(
renderer,
dockerfile_lines,
environment_count,
lifecycle_command_count,
&workspace_folder,
);
}
ProgressEvent::DevcontainerLifecycleStarted {
phase,
command_count,
} => {
self.setup
.on_devcontainer_lifecycle_started(renderer, &phase, command_count);
}
ProgressEvent::DevcontainerLifecycleCompleted { phase, duration_ms } => {
self.setup
.on_devcontainer_lifecycle_completed(renderer, &phase, duration_ms);
}
ProgressEvent::DevcontainerLifecycleFailed {
phase,
command,
exit_code,
stderr,
} => {
self.setup.on_devcontainer_lifecycle_failed(
renderer, &phase, &command, exit_code, &stderr,
);
}
ProgressEvent::DevcontainerLifecycleCommandCompleted {
command,
command_index,
exit_code,
duration_ms,
} => {
self.setup.on_devcontainer_lifecycle_command_completed(
renderer,
&command,
command_index,
exit_code,
duration_ms,
);
}
ProgressEvent::StageStarted {
node_id,
name,
script,
} => {
self.stage
.on_stage_started(renderer, &node_id, &name, script.as_deref());
}
ProgressEvent::StageCompleted {
node_id,
name,
duration_ms,
status,
usage,
} => {
self.stage.on_stage_completed(
renderer,
&node_id,
&name,
duration_ms,
&status,
usage.as_ref(),
);
}
ProgressEvent::StageFailed {
node_id,
name,
error,
} => {
self.stage
.on_stage_failed(renderer, &node_id, &name, &error);
}
ProgressEvent::StageRetrying {
name,
attempt,
max_attempts,
delay_ms,
} => {
self.info
.on_stage_retrying(renderer, &name, attempt, max_attempts, delay_ms);
}
ProgressEvent::ParallelStarted => {
self.stage.on_parallel_started();
}
ProgressEvent::ParallelBranchStarted { branch } => {
self.stage.on_parallel_branch_started(renderer, &branch);
}
ProgressEvent::ParallelBranchCompleted {
branch,
duration_ms,
status,
} => {
self.stage
.on_parallel_branch_completed(renderer, &branch, duration_ms, &status);
}
ProgressEvent::ParallelCompleted => {
self.stage.on_parallel_completed();
}
ProgressEvent::AssistantMessage {
stage_node_id,
model,
} => {
self.stage
.on_assistant_message(renderer, &stage_node_id, &model);
}
ProgressEvent::ToolCallStarted {
stage_node_id,
tool_name,
tool_call_id,
arguments,
} => {
self.stage.on_tool_call_started(
renderer,
&stage_node_id,
&tool_name,
&tool_call_id,
&arguments,
);
}
ProgressEvent::ToolCallCompleted {
stage_node_id,
tool_call_id,
is_error,
} => {
self.stage.on_tool_call_completed(
renderer,
&stage_node_id,
&tool_call_id,
is_error,
);
}
ProgressEvent::ContextWindowWarning {
stage_node_id,
usage_percent,
} => {
self.stage
.on_context_window_warning(renderer, &stage_node_id, usage_percent);
}
ProgressEvent::CompactionStarted { stage_node_id } => {
self.stage.on_compaction_started(renderer, &stage_node_id);
}
ProgressEvent::CompactionCompleted {
stage_node_id,
original_turn_count,
preserved_turn_count,
tracked_file_count,
} => {
self.stage.on_compaction_completed(
renderer,
&stage_node_id,
original_turn_count,
preserved_turn_count,
tracked_file_count,
);
}
ProgressEvent::LlmRetry {
stage_node_id,
model,
attempt,
delay_ms,
error,
} => {
self.stage.on_llm_retry(
renderer,
&stage_node_id,
&model,
attempt,
delay_ms,
&error,
);
}
ProgressEvent::SubagentSpawned {
stage_node_id,
agent_id,
task,
} => {
self.stage
.on_subagent_spawned(renderer, &stage_node_id, &agent_id, &task);
}
ProgressEvent::SubagentCompleted {
stage_node_id,
agent_id,
success,
turns_used,
} => {
self.stage.on_subagent_completed(
renderer,
&stage_node_id,
&agent_id,
success,
turns_used,
);
}
ProgressEvent::EdgeSelected {
from_node,
to_node,
label,
condition,
} => {
self.info.on_edge_selected(
renderer,
&from_node,
&to_node,
label.as_deref(),
condition.as_deref(),
);
}
ProgressEvent::LoopRestart { from_node, to_node } => {
self.info.on_loop_restart(renderer, &from_node, &to_node);
}
ProgressEvent::RetroStarted => {
self.stage.on_retro_started(renderer);
}
ProgressEvent::RetroCompleted { duration_ms } => {
self.stage.on_retro_completed(renderer, duration_ms);
}
ProgressEvent::RetroFailed { duration_ms } => {
self.stage.on_retro_failed(renderer, duration_ms);
}
ProgressEvent::RunNotice {
level,
code,
message,
} => {
self.info.on_run_notice(renderer, level, &code, &message);
}
ProgressEvent::PullRequestCreated { pr_url, draft } => {
self.info.on_pull_request_created(renderer, &pr_url, draft);
}
ProgressEvent::PullRequestFailed { error } => {
self.info.on_pull_request_failed(renderer, &error);
}
}
}
}
#[cfg(test)]
mod tests {
use std::io::{self, Write};
use std::sync::{Arc, Mutex};
use fabro_agent::{AgentEvent, SandboxEvent};
use fabro_llm::types::Usage;
use fabro_workflow::event::{RunNoticeLevel, flatten_event};
use fabro_workflow::outcome::StageUsage;
use super::*;
use crate::commands::run::run_progress::stage_display::ToolCallStatus;
struct SharedBuffer {
inner: Arc<Mutex<Vec<u8>>>,
}
impl Write for SharedBuffer {
fn write(&mut self, buf: &[u8]) -> io::Result<usize> {
self.inner
.lock()
.expect("buffer lock poisoned")
.extend_from_slice(buf);
Ok(buf.len())
}
fn flush(&mut self) -> io::Result<()> {
Ok(())
}
}
fn capture_ui(verbose: bool) -> (ProgressUI, Arc<Mutex<Vec<u8>>>) {
let buffer = Arc::new(Mutex::new(Vec::new()));
let ui = ProgressUI::new_plain_test(
Box::new(SharedBuffer {
inner: Arc::clone(&buffer),
}),
verbose,
false,
);
(ui, buffer)
}
fn rendered(buffer: &Arc<Mutex<Vec<u8>>>) -> String {
String::from_utf8(buffer.lock().expect("buffer lock poisoned").clone())
.expect("valid utf-8")
}
fn stage_started(node_id: &str, name: &str) -> WorkflowRunEvent {
WorkflowRunEvent::StageStarted {
node_id: node_id.into(),
name: name.into(),
index: 0,
handler_type: None,
script: None,
attempt: 1,
max_attempts: 1,
}
}
fn assistant_message(stage: &str, model: &str) -> WorkflowRunEvent {
WorkflowRunEvent::Agent {
stage: stage.into(),
event: AgentEvent::AssistantMessage {
text: "done".into(),
model: model.into(),
usage: Usage::default(),
tool_call_count: 0,
},
}
}
fn stage_completed(node_id: &str, name: &str) -> WorkflowRunEvent {
WorkflowRunEvent::StageCompleted {
node_id: node_id.into(),
name: name.into(),
index: 0,
duration_ms: 5000,
status: "success".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
usage: Some(StageUsage {
model: "gpt-5-mini".into(),
input_tokens: 1200,
output_tokens: 300,
cache_read_tokens: None,
cache_write_tokens: None,
reasoning_tokens: None,
speed: None,
cost: Some(0.12),
}),
failure: None,
notes: None,
files_touched: Vec::new(),
attempt: 1,
max_attempts: 1,
}
}
#[test]
fn parallel_branches_tracked_as_tool_calls() {
let mut ui = ProgressUI::new(true, false);
ui.handle_event(&stage_started("fork1", "Fork Analysis"));
assert!(ui.stage.active_stages.contains_key("fork1"));
assert!(ui.stage.parallel_parent.is_none());
ui.handle_event(&WorkflowRunEvent::ParallelStarted {
branch_count: 2,
join_policy: "wait_all".into(),
});
assert_eq!(ui.stage.parallel_parent.as_deref(), Some("fork1"));
ui.handle_event(&WorkflowRunEvent::ParallelBranchStarted {
branch: "security".into(),
index: 0,
});
let stage = &ui.stage.active_stages["fork1"];
assert_eq!(stage.tool_calls.len(), 1);
assert_eq!(stage.tool_calls[0].tool_call_id, "security");
assert!(matches!(
stage.tool_calls[0].status,
ToolCallStatus::Running
));
ui.handle_event(&WorkflowRunEvent::ParallelBranchCompleted {
branch: "security".into(),
index: 0,
duration_ms: 2000,
status: "success".into(),
});
let stage = &ui.stage.active_stages["fork1"];
assert!(matches!(
stage.tool_calls[0].status,
ToolCallStatus::Succeeded
));
}
#[test]
fn parallel_branch_running_shows_triangle_glyph() {
let mut ui = ProgressUI::new(true, false);
ui.handle_event(&stage_started("fork1", "Fork"));
ui.handle_event(&WorkflowRunEvent::ParallelStarted {
branch_count: 1,
join_policy: "wait_all".into(),
});
ui.handle_event(&WorkflowRunEvent::ParallelBranchStarted {
branch: "security".into(),
index: 0,
});
let stage = &ui.stage.active_stages["fork1"];
let message = stage.tool_calls[0].bar.message();
assert!(
message.contains('\u{25b8}'),
"expected branch message to contain ▸, got: {message:?}"
);
}
#[test]
fn compaction_sets_and_clears_bar() {
let mut ui = ProgressUI::new(true, false);
ui.handle_event(&stage_started("s1", "Build"));
assert!(ui.stage.active_stages["s1"].compaction_bar.is_none());
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "s1".into(),
event: AgentEvent::CompactionStarted {
estimated_tokens: 5000,
context_window_size: 8000,
},
});
assert!(ui.stage.active_stages["s1"].compaction_bar.is_some());
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "s1".into(),
event: AgentEvent::CompactionCompleted {
original_turn_count: 20,
preserved_turn_count: 6,
summary_token_estimate: 500,
tracked_file_count: 3,
},
});
assert!(ui.stage.active_stages["s1"].compaction_bar.is_none());
}
#[test]
fn handle_json_line_ignores_invalid_json() {
let (mut ui, buffer) = capture_ui(false);
ui.handle_json_line("not valid json");
ui.handle_json_line("");
ui.handle_json_line("{}");
assert!(rendered(&buffer).is_empty());
}
#[test]
fn handle_json_line_matches_handle_event_for_verbose_events() {
let events = vec![
stage_started("code", "Code"),
WorkflowRunEvent::SandboxInitialized {
working_directory: "/home/daytona/workspace".into(),
},
WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::ToolCallStarted {
tool_name: "read_file".into(),
tool_call_id: "tc1".into(),
arguments: serde_json::json!({
"file_path": "/home/daytona/workspace/src/main.rs"
}),
},
},
assistant_message("code", "gpt-5-mini"),
WorkflowRunEvent::EdgeSelected {
from_node: "code".into(),
to_node: "review".into(),
label: Some("ship".into()),
condition: None,
reason: "condition".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
stage_status: "success".into(),
is_jump: false,
},
WorkflowRunEvent::StageRetrying {
node_id: "code".into(),
name: "Code".into(),
index: 0,
attempt: 2,
max_attempts: 3,
delay_ms: 1500,
},
WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::Warning {
kind: "context_window".into(),
message: "high usage".into(),
details: serde_json::json!({"usage_percent": 92}),
},
},
WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::LlmRetry {
provider: "openai".into(),
model: "gpt-5-mini".into(),
attempt: 2,
delay_secs: 1.5,
error: fabro_llm::error::SdkError::Configuration {
message: "busy".into(),
source: None,
},
},
},
WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::SubAgentSpawned {
agent_id: "a1".into(),
depth: 1,
task: "review recent changes".into(),
},
},
WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::SubAgentCompleted {
agent_id: "a1".into(),
depth: 1,
success: true,
turns_used: 3,
},
},
WorkflowRunEvent::SetupStarted { command_count: 1 },
WorkflowRunEvent::SetupCommandCompleted {
command: "bun install".into(),
index: 0,
exit_code: 0,
duration_ms: 2200,
},
WorkflowRunEvent::SetupCompleted { duration_ms: 2200 },
WorkflowRunEvent::DevcontainerLifecycleStarted {
phase: "postCreate".into(),
command_count: 1,
},
WorkflowRunEvent::DevcontainerLifecycleCommandCompleted {
phase: "postCreate".into(),
command: "npm run setup".into(),
index: 0,
exit_code: 0,
duration_ms: 1400,
},
WorkflowRunEvent::DevcontainerLifecycleCompleted {
phase: "postCreate".into(),
duration_ms: 1400,
},
];
let (mut event_ui, event_buffer) = capture_ui(true);
for event in &events {
event_ui.handle_event(event);
}
let (mut json_ui, json_buffer) = capture_ui(true);
for event in &events {
let (event_name, fields) = flatten_event(event);
let mut envelope = serde_json::Map::new();
envelope.insert("event".into(), event_name.into());
envelope.extend(fields);
let line = serde_json::to_string(&serde_json::Value::Object(envelope)).unwrap();
json_ui.handle_json_line(&line);
}
assert_eq!(rendered(&event_buffer), rendered(&json_buffer));
}
#[test]
fn plain_default_stage_snapshot() {
let (mut ui, buffer) = capture_ui(false);
ui.handle_event(&stage_started("plan", "Plan"));
ui.handle_event(&assistant_message("plan", "gpt-5-mini"));
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "plan".into(),
event: AgentEvent::ToolCallStarted {
tool_name: "read_file".into(),
tool_call_id: "tc1".into(),
arguments: serde_json::json!({"path": "src/main.rs"}),
},
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "plan".into(),
event: AgentEvent::ToolCallCompleted {
tool_name: "read_file".into(),
tool_call_id: "tc1".into(),
output: serde_json::json!({"ok": true}),
is_error: false,
},
});
ui.handle_event(&stage_completed("plan", "Plan"));
insta::assert_snapshot!(rendered(&buffer), @r"
Plan $0.12 5s
");
}
#[test]
fn plain_default_setup_snapshot() {
let (mut ui, buffer) = capture_ui(false);
ui.handle_event(&WorkflowRunEvent::Sandbox {
event: SandboxEvent::Initializing {
provider: "daytona".into(),
},
});
ui.handle_event(&WorkflowRunEvent::Sandbox {
event: SandboxEvent::Ready {
provider: "daytona".into(),
duration_ms: 2500,
name: Some("sandbox-1".into()),
cpu: Some(4.0),
memory: Some(8.0),
url: None,
},
});
ui.handle_event(&WorkflowRunEvent::SshAccessReady {
ssh_command: "ssh daytona@example".into(),
});
ui.handle_event(&WorkflowRunEvent::SetupStarted { command_count: 2 });
ui.handle_event(&WorkflowRunEvent::SetupCompleted { duration_ms: 8200 });
ui.handle_event(&WorkflowRunEvent::CliEnsureCompleted {
cli_name: "gh".into(),
provider: "github".into(),
already_installed: false,
node_installed: false,
duration_ms: 600,
});
ui.handle_event(&WorkflowRunEvent::DevcontainerResolved {
dockerfile_lines: 24,
environment_count: 3,
lifecycle_command_count: 2,
workspace_folder: "/workspace".into(),
});
ui.handle_event(&WorkflowRunEvent::DevcontainerLifecycleStarted {
phase: "postCreate".into(),
command_count: 2,
});
ui.handle_event(&WorkflowRunEvent::DevcontainerLifecycleCompleted {
phase: "postCreate".into(),
duration_ms: 1800,
});
insta::assert_snapshot!(rendered(&buffer), @r"
Sandbox: daytona (ready in 2s)
sandbox-1 (4 cpu, 8 GB)
ssh daytona@example
Setup: 2 commands (8s)
CLI: gh (installed, 600ms)
Devcontainer: resolved
24 Dockerfile lines, 3 env vars, 2 lifecycle cmds, /workspace
Running devcontainer postCreate (2 commands)...
Devcontainer: postCreate (1s)
");
}
#[test]
fn plain_verbose_snapshot() {
let (mut ui, buffer) = capture_ui(true);
ui.handle_event(&stage_started("code", "Code"));
ui.handle_event(&WorkflowRunEvent::SandboxInitialized {
working_directory: "/home/daytona/workspace".into(),
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::ToolCallStarted {
tool_name: "read_file".into(),
tool_call_id: "tc1".into(),
arguments: serde_json::json!({
"file_path": "/home/daytona/workspace/src/main.rs"
}),
},
});
ui.handle_event(&assistant_message("code", "gpt-5-mini"));
ui.handle_event(&WorkflowRunEvent::EdgeSelected {
from_node: "code".into(),
to_node: "review".into(),
label: Some("ship".into()),
condition: None,
reason: "condition".into(),
preferred_label: None,
suggested_next_ids: Vec::new(),
stage_status: "success".into(),
is_jump: false,
});
ui.handle_event(&WorkflowRunEvent::StageRetrying {
node_id: "code".into(),
name: "Code".into(),
index: 0,
attempt: 2,
max_attempts: 3,
delay_ms: 1500,
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::Warning {
kind: "context_window".into(),
message: "high usage".into(),
details: serde_json::json!({"usage_percent": 92}),
},
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::LlmRetry {
provider: "openai".into(),
model: "gpt-5-mini".into(),
attempt: 2,
delay_secs: 1.5,
error: fabro_llm::error::SdkError::Configuration {
message: "busy".into(),
source: None,
},
},
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::SubAgentSpawned {
agent_id: "a1".into(),
depth: 1,
task: "review recent changes".into(),
},
});
ui.handle_event(&WorkflowRunEvent::Agent {
stage: "code".into(),
event: AgentEvent::SubAgentCompleted {
agent_id: "a1".into(),
depth: 1,
success: true,
turns_used: 3,
},
});
ui.handle_event(&WorkflowRunEvent::SetupStarted { command_count: 1 });
ui.handle_event(&WorkflowRunEvent::SetupCommandCompleted {
command: "bun install".into(),
index: 0,
exit_code: 0,
duration_ms: 2200,
});
ui.handle_event(&WorkflowRunEvent::SetupCompleted { duration_ms: 2200 });
ui.handle_event(&WorkflowRunEvent::DevcontainerLifecycleStarted {
phase: "postCreate".into(),
command_count: 1,
});
ui.handle_event(&WorkflowRunEvent::DevcontainerLifecycleCommandCompleted {
phase: "postCreate".into(),
command: "npm run setup".into(),
index: 0,
exit_code: 0,
duration_ms: 1400,
});
ui.handle_event(&WorkflowRunEvent::DevcontainerLifecycleCompleted {
phase: "postCreate".into(),
duration_ms: 1400,
});
ui.handle_event(&stage_completed("code", "Code"));
insta::assert_snapshot!(rendered(&buffer), @r#"
code review "ship"
Code: retrying (attempt 2/3, delay 1s)
context window: 92% used
retry: gpt-5-mini attempt 2 (busy, delay 1s)
subagent[a1] "review recent changes"
subagent[a1] (3 turns)
[1/1] bun install 2s
Setup: 1 command (2s)
Running devcontainer postCreate (1 commands)...
[1/1] npm run setup 1s
Devcontainer: postCreate (1s)
Code $0.12 5s (1 turns, 0 tools, 1.5k toks)
"#);
}
#[test]
fn plain_notice_snapshot() {
let (mut ui, buffer) = capture_ui(false);
ui.handle_event(&WorkflowRunEvent::RunNotice {
level: RunNoticeLevel::Warn,
code: "sandbox_cleanup_failed".into(),
message: "sandbox cleanup failed".into(),
});
ui.handle_event(&WorkflowRunEvent::PullRequestCreated {
pr_url: "https://github.com/fabro-sh/fabro/pull/42".into(),
pr_number: 42,
draft: true,
});
ui.handle_event(&WorkflowRunEvent::PullRequestFailed {
error: "auth token expired".into(),
});
insta::assert_snapshot!(rendered(&buffer), @r"
Warning: sandbox cleanup failed [sandbox_cleanup_failed]
Draft PR: https://github.com/fabro-sh/fabro/pull/42
PR failed: auth token expired
");
}
}

View file

@ -0,0 +1,94 @@
use std::io::Write;
use std::sync::Mutex;
use fabro_util::terminal::Styles;
use indicatif::{MultiProgress, ProgressBar, ProgressDrawTarget};
use super::styles;
enum RendererInner {
Tty { multi: MultiProgress },
Plain { out: Mutex<Box<dyn Write + Send>> },
}
pub(super) struct ProgressRenderer {
inner: RendererInner,
styles: Styles,
}
impl ProgressRenderer {
pub(super) fn new_tty() -> Self {
Self {
inner: RendererInner::Tty {
multi: MultiProgress::new(),
},
styles: Styles::new(console::colors_enabled_stderr()),
}
}
pub(super) fn new_plain(out: Box<dyn Write + Send>, colors: bool) -> Self {
Self {
inner: RendererInner::Plain {
out: Mutex::new(out),
},
styles: Styles::new(colors),
}
}
pub(super) fn add_spinner(&self) -> ProgressBar {
match &self.inner {
RendererInner::Tty { multi } => multi.add(ProgressBar::new_spinner()),
RendererInner::Plain { .. } => ProgressBar::hidden(),
}
}
pub(super) fn insert_after(&self, after: &ProgressBar) -> ProgressBar {
match &self.inner {
RendererInner::Tty { multi } => multi.insert_after(after, ProgressBar::new_spinner()),
RendererInner::Plain { .. } => ProgressBar::hidden(),
}
}
pub(super) fn insert_before(&self, before: &ProgressBar) -> ProgressBar {
match &self.inner {
RendererInner::Tty { multi } => multi.insert_before(before, ProgressBar::new_spinner()),
RendererInner::Plain { .. } => ProgressBar::hidden(),
}
}
pub(super) fn print_line(&self, indent: usize, message: &str) {
if let RendererInner::Plain { out } = &self.inner {
let mut out = out.lock().expect("plain renderer lock poisoned");
let _ = writeln!(out, "{}{message}", " ".repeat(indent));
}
}
pub(super) fn is_tty(&self) -> bool {
matches!(self.inner, RendererInner::Tty { .. })
}
pub(super) fn styles(&self) -> &Styles {
&self.styles
}
pub(super) fn hide(&self) {
if let RendererInner::Tty { multi } = &self.inner {
multi.set_draw_target(ProgressDrawTarget::hidden());
}
}
pub(super) fn show(&self) {
if let RendererInner::Tty { multi } = &self.inner {
multi.set_draw_target(ProgressDrawTarget::stderr());
}
}
pub(super) fn finish(&self) {
if let RendererInner::Tty { multi } = &self.inner {
let sep = multi.add(ProgressBar::new_spinner());
sep.set_style(styles::style_empty());
sep.finish();
multi.set_draw_target(ProgressDrawTarget::hidden());
}
}
}

View file

@ -0,0 +1,379 @@
use std::time::Duration;
use indicatif::ProgressBar;
use super::renderer::ProgressRenderer;
use super::styles;
use crate::shared::format_duration_ms;
pub(super) struct SetupDisplay {
verbose: bool,
pub(super) sandbox_bar: Option<ProgressBar>,
pub(super) setup_bar: Option<ProgressBar>,
pub(super) setup_command_count: u64,
pub(super) devcontainer_bar: Option<ProgressBar>,
pub(super) devcontainer_command_count: u64,
pub(super) cli_ensure_bar: Option<ProgressBar>,
}
impl SetupDisplay {
pub(super) fn new(verbose: bool) -> Self {
Self {
verbose,
sandbox_bar: None,
setup_bar: None,
setup_command_count: 0,
devcontainer_bar: None,
devcontainer_command_count: 0,
cli_ensure_bar: None,
}
}
pub(super) fn finish(&mut self) {
if let Some(bar) = self.sandbox_bar.take() {
bar.finish_and_clear();
}
if let Some(bar) = self.setup_bar.take() {
bar.finish_and_clear();
}
if let Some(bar) = self.devcontainer_bar.take() {
bar.finish_and_clear();
}
if let Some(bar) = self.cli_ensure_bar.take() {
bar.finish_and_clear();
}
}
pub(super) fn on_sandbox_initializing(&mut self, renderer: &ProgressRenderer, provider: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_header_running());
bar.set_message(format!("Initializing {provider} sandbox..."));
bar.enable_steady_tick(Duration::from_millis(100));
self.sandbox_bar = Some(bar);
}
}
pub(super) fn on_sandbox_ready(
&mut self,
renderer: &ProgressRenderer,
provider: &str,
duration_ms: u64,
name: Option<&str>,
cpu: Option<f64>,
memory: Option<f64>,
url: Option<&str>,
) {
let dur = format_duration_ms(duration_ms);
let detail = match (name, cpu, memory) {
(Some(name), Some(cpu), Some(memory)) => Some(format!(
"{name} ({} cpu, {} GB)",
styles::format_number(cpu),
styles::format_number(memory)
)),
(Some(name), _, _) => Some(name.to_string()),
_ => None,
};
if renderer.is_tty() {
let display_provider = match url {
Some(url) => styles::terminal_hyperlink(url, provider),
None => provider.to_string(),
};
if let Some(bar) = self.sandbox_bar.take() {
bar.set_style(styles::style_header_done());
bar.set_prefix(dur);
bar.finish_with_message(format!("Sandbox: {display_provider}"));
if let Some(detail) = detail {
let detail_bar = renderer.insert_after(&bar);
detail_bar.set_style(styles::style_sandbox_detail());
detail_bar.finish_with_message(detail);
}
}
} else {
renderer.print_line(4, &format!("Sandbox: {provider} (ready in {dur})"));
if let Some(detail) = detail {
renderer.print_line(13, &detail);
}
}
}
pub(super) fn on_ssh_access_ready(&self, renderer: &ProgressRenderer, ssh_command: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_sandbox_detail());
bar.finish_with_message(ssh_command.to_string());
} else {
renderer.print_line(13, ssh_command);
}
}
pub(super) fn on_setup_started(&mut self, renderer: &ProgressRenderer, command_count: u64) {
self.setup_command_count = command_count;
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_header_running());
bar.set_message(format!(
"Setup: {command_count} command{}...",
if command_count == 1 { "" } else { "s" }
));
bar.enable_steady_tick(Duration::from_millis(100));
self.setup_bar = Some(bar);
}
}
pub(super) fn on_setup_completed(&mut self, renderer: &ProgressRenderer, duration_ms: u64) {
let dur = format_duration_ms(duration_ms);
let suffix = if self.setup_command_count == 1 {
""
} else {
"s"
};
if renderer.is_tty() {
if let Some(bar) = self.setup_bar.take() {
bar.set_style(styles::style_header_done());
bar.set_prefix(dur);
bar.finish_with_message(format!(
"Setup: {} command{suffix}",
self.setup_command_count
));
}
} else {
renderer.print_line(
4,
&format!(
"Setup: {} command{suffix} ({dur})",
self.setup_command_count
),
);
}
}
pub(super) fn on_setup_command_completed(
&self,
renderer: &ProgressRenderer,
command: &str,
command_index: u64,
exit_code: i64,
duration_ms: u64,
) {
if !self.verbose {
return;
}
let glyph = if exit_code == 0 {
styles::green_check(renderer.styles())
} else {
styles::red_cross(renderer.styles())
};
let msg = format!(
"{glyph} [{}/{}] {}",
command_index + 1,
self.setup_command_count,
styles::truncate(command, 60)
);
let dur = format_duration_ms(duration_ms);
if renderer.is_tty() {
let bar = match &self.setup_bar {
Some(setup_bar) => renderer.insert_before(setup_bar),
None => renderer.add_spinner(),
};
bar.set_style(styles::style_tool_done());
bar.set_prefix(dur);
bar.finish_with_message(msg);
} else {
renderer.print_line(6, &format!("{msg} {dur}"));
}
}
pub(super) fn on_cli_ensure_started(&mut self, renderer: &ProgressRenderer, cli_name: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_header_running());
bar.set_message(format!("CLI: ensuring {cli_name}..."));
bar.enable_steady_tick(Duration::from_millis(100));
self.cli_ensure_bar = Some(bar);
}
}
pub(super) fn on_cli_ensure_completed(
&mut self,
renderer: &ProgressRenderer,
cli_name: &str,
already_installed: bool,
duration_ms: u64,
) {
let status = if already_installed {
"found"
} else {
"installed"
};
let dur = format_duration_ms(duration_ms);
if renderer.is_tty() {
if let Some(bar) = self.cli_ensure_bar.take() {
bar.set_style(styles::style_header_done());
bar.set_prefix(dur);
bar.finish_with_message(format!("CLI: {cli_name} ({status})"));
}
} else {
renderer.print_line(4, &format!("CLI: {cli_name} ({status}, {dur})"));
}
}
pub(super) fn on_cli_ensure_failed(&mut self, renderer: &ProgressRenderer, cli_name: &str) {
let message = format!(
"{} CLI: {cli_name} install failed",
styles::red_cross(renderer.styles())
);
if renderer.is_tty() {
if let Some(bar) = self.cli_ensure_bar.take() {
bar.set_style(styles::style_header_done());
bar.finish_with_message(message);
}
} else {
renderer.print_line(4, &message);
}
}
pub(super) fn on_devcontainer_resolved(
&self,
renderer: &ProgressRenderer,
dockerfile_lines: u64,
environment_count: u64,
lifecycle_command_count: u64,
workspace_folder: &str,
) {
let detail = format!(
"{dockerfile_lines} Dockerfile lines, {environment_count} env vars, \
{lifecycle_command_count} lifecycle cmds, {workspace_folder}"
);
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_header_done());
bar.finish_with_message("Devcontainer: resolved".to_string());
let detail_bar = renderer.insert_after(&bar);
detail_bar.set_style(styles::style_sandbox_detail());
detail_bar.finish_with_message(detail);
} else {
renderer.print_line(4, "Devcontainer: resolved");
renderer.print_line(13, &detail);
}
}
pub(super) fn on_devcontainer_lifecycle_started(
&mut self,
renderer: &ProgressRenderer,
phase: &str,
command_count: u64,
) {
self.devcontainer_command_count = command_count;
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_header_running());
bar.set_message(format!(
"Running devcontainer {phase} ({command_count} commands)..."
));
bar.enable_steady_tick(Duration::from_millis(100));
self.devcontainer_bar = Some(bar);
} else {
renderer.print_line(
4,
&format!("Running devcontainer {phase} ({command_count} commands)..."),
);
}
}
pub(super) fn on_devcontainer_lifecycle_completed(
&mut self,
renderer: &ProgressRenderer,
phase: &str,
duration_ms: u64,
) {
let dur = format_duration_ms(duration_ms);
if renderer.is_tty() {
if let Some(bar) = self.devcontainer_bar.take() {
bar.set_style(styles::style_header_done());
bar.set_prefix(dur);
bar.finish_with_message(format!("Devcontainer: {phase}"));
}
} else {
renderer.print_line(4, &format!("Devcontainer: {phase} ({dur})"));
}
}
pub(super) fn on_devcontainer_lifecycle_failed(
&mut self,
renderer: &ProgressRenderer,
phase: &str,
command: &str,
exit_code: i64,
stderr: &str,
) {
if let Some(bar) = self.devcontainer_bar.take() {
bar.abandon();
}
let summary = if stderr.len() > 120 {
&stderr[..120]
} else {
stderr
};
let message = format!(
"{} Devcontainer {phase} command failed (exit {exit_code}): {command}\n {summary}",
renderer.styles().red.apply_to("Error:")
);
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_static_dim());
bar.finish_with_message(message);
} else {
renderer.print_line(4, &message);
}
}
pub(super) fn on_devcontainer_lifecycle_command_completed(
&self,
renderer: &ProgressRenderer,
command: &str,
command_index: u64,
exit_code: i64,
duration_ms: u64,
) {
if !self.verbose {
return;
}
let glyph = if exit_code == 0 {
styles::green_check(renderer.styles())
} else {
styles::red_cross(renderer.styles())
};
let msg = format!(
"{glyph} [{}/{}] {}",
command_index + 1,
self.devcontainer_command_count,
styles::truncate(command, 60)
);
let dur = format_duration_ms(duration_ms);
if renderer.is_tty() {
let bar = match &self.devcontainer_bar {
Some(devcontainer_bar) => renderer.insert_before(devcontainer_bar),
None => renderer.add_spinner(),
};
bar.set_style(styles::style_tool_done());
bar.set_prefix(dur);
bar.finish_with_message(msg);
} else {
renderer.print_line(6, &format!("{msg} {dur}"));
}
}
}

View file

@ -0,0 +1,699 @@
use std::collections::{HashMap, VecDeque};
use std::convert::TryFrom;
use std::time::Duration;
use indicatif::ProgressBar;
use fabro_workflow::outcome::{StageStatus, format_cost};
use super::event::ProgressUsage;
use super::renderer::ProgressRenderer;
use super::styles;
use crate::shared::{format_duration_ms, format_tokens_human};
const MAX_TOOL_CALLS: usize = 5;
#[derive(Debug)]
pub(super) enum ToolCallStatus {
Running,
Succeeded,
Failed,
}
#[derive(Debug)]
pub(super) struct ToolCallEntry {
pub(super) display_name: String,
pub(super) tool_call_id: String,
pub(super) status: ToolCallStatus,
pub(super) bar: ProgressBar,
pub(super) is_branch: bool,
}
#[derive(Debug)]
pub(super) struct ActiveStage {
pub(super) display_name: String,
pub(super) has_model: bool,
pub(super) spinner: ProgressBar,
pub(super) tool_calls: VecDeque<ToolCallEntry>,
pub(super) compaction_bar: Option<ProgressBar>,
}
impl ActiveStage {
fn last_bar(&self) -> &ProgressBar {
self.tool_calls
.back()
.map_or(&self.spinner, |entry| &entry.bar)
}
}
pub(super) struct StageDisplay {
verbose: bool,
pub(super) active_stages: HashMap<String, ActiveStage>,
pub(super) stage_counts: HashMap<String, (u64, u64)>,
pub(super) parallel_parent: Option<String>,
any_stage_started: bool,
working_directory: Option<String>,
}
impl StageDisplay {
pub(super) fn new(verbose: bool) -> Self {
Self {
verbose,
active_stages: HashMap::new(),
stage_counts: HashMap::new(),
parallel_parent: None,
any_stage_started: false,
working_directory: None,
}
}
pub(super) fn set_working_directory(&mut self, dir: String) {
self.working_directory = Some(dir);
}
pub(super) fn finish(&mut self) {
for (_node_id, stage) in self.active_stages.drain() {
if let Some(bar) = stage.compaction_bar {
bar.finish_and_clear();
}
for entry in &stage.tool_calls {
if entry.is_branch || self.verbose {
entry.bar.abandon();
} else {
entry.bar.finish_and_clear();
}
}
stage.spinner.finish_and_clear();
}
}
pub(super) fn on_stage_started(
&mut self,
renderer: &ProgressRenderer,
node_id: &str,
name: &str,
script: Option<&str>,
) {
self.stage_counts.insert(node_id.to_string(), (0, 0));
let display_name = match script {
Some(script) => format!(
"{name} {}",
renderer.styles().dim.apply_to(styles::truncate(script, 60))
),
None => name.to_string(),
};
if renderer.is_tty() && !self.any_stage_started {
self.any_stage_started = true;
let sep = renderer.add_spinner();
sep.set_style(styles::style_empty());
sep.finish();
}
let bar = renderer.add_spinner();
bar.set_style(styles::style_stage_running());
bar.set_message(display_name.clone());
if renderer.is_tty() {
bar.enable_steady_tick(Duration::from_millis(100));
}
self.active_stages.insert(
node_id.to_string(),
ActiveStage {
display_name,
has_model: false,
spinner: bar,
tool_calls: VecDeque::new(),
compaction_bar: None,
},
);
}
pub(super) fn on_stage_completed(
&mut self,
renderer: &ProgressRenderer,
node_id: &str,
name: &str,
duration_ms: u64,
status: &str,
usage: Option<&ProgressUsage>,
) {
let succeeded = status
.parse::<StageStatus>()
.map(|status| matches!(status, StageStatus::Success | StageStatus::PartialSuccess))
.unwrap_or_else(|_| matches!(status, "success" | "partial_success"));
let cost_str = usage
.and_then(ProgressUsage::display_cost)
.map(|cost| format!("{} ", format_cost(cost)))
.unwrap_or_default();
let stats_str = if self.verbose {
let (turn_count, tool_call_count) =
self.stage_counts.get(node_id).copied().unwrap_or((0, 0));
let total_tokens = usage.map_or(0, ProgressUsage::total_tokens);
if turn_count > 0 || tool_call_count > 0 || total_tokens > 0 {
let total_tokens = i64::try_from(total_tokens).unwrap_or(i64::MAX);
format!(
" {}",
renderer.styles().dim.apply_to(format!(
"({} turns, {} tools, {} toks)",
turn_count,
tool_call_count,
format_tokens_human(total_tokens),
))
)
} else {
String::new()
}
} else {
String::new()
};
let prefix = format!("{cost_str}{}{stats_str}", format_duration_ms(duration_ms));
let glyph = if succeeded {
styles::green_check(renderer.styles())
} else {
styles::red_cross(renderer.styles())
};
self.finish_stage(renderer, node_id, name, &glyph, &prefix);
}
pub(super) fn on_stage_failed(
&mut self,
renderer: &ProgressRenderer,
node_id: &str,
name: &str,
error: &str,
) {
self.finish_stage(
renderer,
node_id,
name,
&styles::red_cross(renderer.styles()),
"",
);
let summary = styles::last_line_truncated(error, 120);
self.insert_global_info_line(
renderer,
&format!("{} {summary}", renderer.styles().red.apply_to("Error:")),
);
}
pub(super) fn on_parallel_started(&mut self) {
self.parallel_parent = self
.active_stages
.keys()
.next()
.cloned()
.or_else(|| Some(String::new()));
}
pub(super) fn on_parallel_completed(&mut self) {
self.parallel_parent = None;
}
pub(super) fn on_parallel_branch_started(&mut self, renderer: &ProgressRenderer, branch: &str) {
let Some(parent_id) = self.parallel_parent.clone() else {
return;
};
let Some(stage) = self.active_stages.get_mut(&parent_id) else {
return;
};
let bar = renderer.insert_after(stage.last_bar());
bar.set_style(styles::style_subagent_info());
bar.set_message(
renderer
.styles()
.dim
.apply_to(format!("\u{25b8} {branch}"))
.to_string(),
);
stage.tool_calls.push_back(ToolCallEntry {
display_name: branch.to_string(),
tool_call_id: branch.to_string(),
status: ToolCallStatus::Running,
bar,
is_branch: true,
});
}
pub(super) fn on_parallel_branch_completed(
&mut self,
renderer: &ProgressRenderer,
branch: &str,
duration_ms: u64,
status: &str,
) {
let Some(parent_id) = self.parallel_parent.clone() else {
return;
};
let Some(stage) = self.active_stages.get_mut(&parent_id) else {
return;
};
let Some(entry) = stage
.tool_calls
.iter_mut()
.find(|entry| entry.tool_call_id == branch)
else {
return;
};
let succeeded = matches!(status, "success" | "partial_success");
entry.status = if succeeded {
ToolCallStatus::Succeeded
} else {
ToolCallStatus::Failed
};
let glyph = if succeeded {
styles::green_check(renderer.styles())
} else {
styles::red_cross(renderer.styles())
};
if renderer.is_tty() {
entry.bar.set_style(styles::style_branch_done());
entry
.bar
.set_prefix(styles::format_duration_short(entry.bar.elapsed()));
entry
.bar
.finish_with_message(format!("{glyph} {}", entry.display_name));
} else {
renderer.print_line(
8,
&format!("{glyph} {branch} {}", format_duration_ms(duration_ms)),
);
}
}
pub(super) fn on_assistant_message(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
model: &str,
) {
if let Some(counts) = self.stage_counts.get_mut(stage_node_id) {
counts.0 += 1;
}
if let Some(stage) = self.active_stages.get_mut(stage_node_id) {
if !stage.has_model {
stage.has_model = true;
let suffix = format!(" {}", renderer.styles().dim.apply_to(format!("[{model}]")));
stage.display_name.push_str(&suffix);
stage.spinner.set_message(stage.display_name.clone());
}
}
}
pub(super) fn on_tool_call_started(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
tool_name: &str,
tool_call_id: &str,
arguments: &serde_json::Value,
) {
let display_name = self.tool_display_name(renderer, tool_name, arguments);
let Some(stage) = self.active_stages.get_mut(stage_node_id) else {
return;
};
if !self.verbose && stage.tool_calls.len() >= MAX_TOOL_CALLS {
let evict_idx = stage
.tool_calls
.iter()
.position(|entry| !matches!(entry.status, ToolCallStatus::Running))
.unwrap_or(0);
if let Some(evicted) = stage.tool_calls.remove(evict_idx) {
evicted.bar.finish_and_clear();
}
}
let bar = renderer.insert_after(stage.last_bar());
bar.set_style(styles::style_tool_running());
bar.set_message(display_name.clone());
if renderer.is_tty() {
bar.enable_steady_tick(Duration::from_millis(100));
}
stage.tool_calls.push_back(ToolCallEntry {
display_name,
tool_call_id: tool_call_id.to_string(),
status: ToolCallStatus::Running,
bar,
is_branch: false,
});
}
pub(super) fn on_tool_call_completed(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
tool_call_id: &str,
is_error: bool,
) {
if let Some(counts) = self.stage_counts.get_mut(stage_node_id) {
counts.1 += 1;
}
let Some(stage) = self.active_stages.get_mut(stage_node_id) else {
return;
};
let Some(entry) = stage
.tool_calls
.iter_mut()
.find(|entry| entry.tool_call_id == tool_call_id)
else {
return;
};
let glyph = if is_error {
styles::red_cross(renderer.styles())
} else {
styles::green_check(renderer.styles())
};
entry.status = if is_error {
ToolCallStatus::Failed
} else {
ToolCallStatus::Succeeded
};
if renderer.is_tty() {
entry.bar.set_style(styles::style_tool_done());
entry
.bar
.set_prefix(styles::format_duration_short(entry.bar.elapsed()));
entry
.bar
.finish_with_message(format!("{glyph} {}", entry.display_name));
}
}
pub(super) fn on_context_window_warning(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
usage_percent: u64,
) {
if !self.verbose {
return;
}
self.insert_info_line_for_stage(
renderer,
stage_node_id,
&format!(
"{} context window: {usage_percent}% used",
styles::warning_glyph(renderer.styles())
),
);
}
pub(super) fn on_compaction_started(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
) {
if !renderer.is_tty() {
return;
}
let Some(stage) = self.active_stages.get_mut(stage_node_id) else {
return;
};
if let Some(old) = stage.compaction_bar.take() {
old.finish_and_clear();
}
let bar = renderer.insert_after(stage.last_bar());
bar.set_style(styles::style_tool_running());
bar.set_message("\u{27f3} compacting context\u{2026}");
bar.enable_steady_tick(Duration::from_millis(100));
stage.compaction_bar = Some(bar);
}
pub(super) fn on_compaction_completed(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
original_turn_count: u64,
preserved_turn_count: u64,
tracked_file_count: u64,
) {
let message = format!(
"\u{27f3} compaction: {original_turn_count} \u{2192} {preserved_turn_count} turns, {tracked_file_count} files"
);
if renderer.is_tty() {
if let Some(bar) = self
.active_stages
.get_mut(stage_node_id)
.and_then(|stage| stage.compaction_bar.take())
{
bar.set_style(styles::style_tool_done());
bar.finish_with_message(message);
} else {
self.insert_info_line_for_stage(renderer, stage_node_id, &message);
}
} else {
renderer.print_line(6, &message);
}
}
pub(super) fn on_llm_retry(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
model: &str,
attempt: u64,
delay_ms: u64,
error: &str,
) {
if !self.verbose {
return;
}
self.insert_info_line_for_stage(
renderer,
stage_node_id,
&format!(
"{} retry: {model} attempt {attempt} ({error}, delay {})",
styles::warning_glyph(renderer.styles()),
format_duration_ms(delay_ms)
),
);
}
pub(super) fn on_subagent_spawned(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
agent_id: &str,
task: &str,
) {
if !self.verbose {
return;
}
self.insert_subagent_line_for_stage(
renderer,
stage_node_id,
&renderer
.styles()
.dim
.apply_to(format!(
"\u{25b8} subagent[{agent_id}] \"{}\"",
styles::truncate(task, 50)
))
.to_string(),
);
}
pub(super) fn on_subagent_completed(
&mut self,
renderer: &ProgressRenderer,
stage_node_id: &str,
agent_id: &str,
success: bool,
turns_used: u64,
) {
if !self.verbose {
return;
}
let glyph = if success {
styles::green_check(renderer.styles())
} else {
styles::red_cross(renderer.styles())
};
self.insert_subagent_line_for_stage(
renderer,
stage_node_id,
&format!("{glyph} subagent[{agent_id}] ({turns_used} turns)"),
);
}
pub(super) fn on_retro_started(&mut self, renderer: &ProgressRenderer) {
self.on_stage_started(renderer, "retro", "Retro", None);
}
pub(super) fn on_retro_completed(&mut self, renderer: &ProgressRenderer, duration_ms: u64) {
self.finish_stage(
renderer,
"retro",
"Retro",
&styles::green_check(renderer.styles()),
&format_duration_ms(duration_ms),
);
}
pub(super) fn on_retro_failed(&mut self, renderer: &ProgressRenderer, duration_ms: u64) {
self.finish_stage(
renderer,
"retro",
"Retro",
&styles::red_cross(renderer.styles()),
&format_duration_ms(duration_ms),
);
}
fn finish_stage(
&mut self,
renderer: &ProgressRenderer,
node_id: &str,
name: &str,
glyph: &str,
prefix: &str,
) {
let Some(stage) = self.active_stages.remove(node_id) else {
if !renderer.is_tty() {
self.print_plain_stage_completion(renderer, name, glyph, prefix);
}
return;
};
if let Some(bar) = stage.compaction_bar {
bar.finish_and_clear();
}
for entry in &stage.tool_calls {
if entry.is_branch || self.verbose {
entry.bar.abandon();
} else {
entry.bar.finish_and_clear();
}
}
if renderer.is_tty() {
stage.spinner.set_style(styles::style_stage_done());
stage.spinner.set_prefix(prefix.to_string());
stage
.spinner
.finish_with_message(format!("{glyph} {}", stage.display_name));
} else {
self.print_plain_stage_completion(renderer, name, glyph, prefix);
}
}
fn print_plain_stage_completion(
&self,
renderer: &ProgressRenderer,
name: &str,
glyph: &str,
prefix: &str,
) {
if prefix.is_empty() {
renderer.print_line(4, &format!("{glyph} {name}"));
} else {
renderer.print_line(4, &format!("{glyph} {name} {prefix}"));
}
}
fn insert_global_info_line(&self, renderer: &ProgressRenderer, message: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_static_dim());
bar.finish_with_message(message.to_string());
} else {
renderer.print_line(4, message);
}
}
fn insert_info_line_for_stage(
&self,
renderer: &ProgressRenderer,
stage_node_id: &str,
message: &str,
) {
if renderer.is_tty() {
let bar = if let Some(stage) = self.active_stages.get(stage_node_id) {
renderer.insert_after(stage.last_bar())
} else {
renderer.add_spinner()
};
bar.set_style(styles::style_tool_done());
bar.finish_with_message(message.to_string());
} else {
renderer.print_line(6, message);
}
}
fn insert_subagent_line_for_stage(
&self,
renderer: &ProgressRenderer,
stage_node_id: &str,
message: &str,
) {
if renderer.is_tty() {
let bar = if let Some(stage) = self.active_stages.get(stage_node_id) {
renderer.insert_after(stage.last_bar())
} else {
renderer.add_spinner()
};
bar.set_style(styles::style_subagent_info());
bar.finish_with_message(message.to_string());
} else {
renderer.print_line(8, message);
}
}
fn tool_display_name(
&self,
renderer: &ProgressRenderer,
tool_name: &str,
arguments: &serde_json::Value,
) -> String {
let arg = |key: &str| arguments.get(key).and_then(serde_json::Value::as_str);
let working_directory = self.working_directory.as_deref();
let path_arg = || {
arg("path")
.or_else(|| arg("file_path"))
.map(|path| styles::truncate(&styles::shorten_path(path, working_directory), 60))
};
let detail = match tool_name {
"bash" | "shell" | "execute_command" => {
arg("command").map(|command| styles::truncate(command, 60))
}
"glob" => arg("pattern").map(String::from),
"grep" | "ripgrep" => arg("pattern").map(|pattern| styles::truncate(pattern, 40)),
"read_file" | "read" | "write_file" | "write" | "create_file" | "edit_file"
| "edit" | "list_dir" => path_arg(),
"web_search" => arg("query").map(|query| styles::truncate(query, 60)),
"web_fetch" => arg("url").map(|url| styles::truncate(url, 60)),
"spawn_agent" => arg("task").map(|task| styles::truncate(task, 60)),
"wait" | "send_input" | "close_agent" => arg("agent_id").map(String::from),
"use_skill" => arg("skill_name").map(String::from),
"apply_patch" => Some("...".to_string()),
"read_many_files" => arguments
.get("paths")
.and_then(serde_json::Value::as_array)
.map(|paths| format!("{} files", paths.len())),
_ => None,
};
match detail {
Some(detail) => format!(
"{tool_name}{}",
renderer.styles().dim.apply_to(format!("({detail})"))
),
None => tool_name.to_string(),
}
}
}

View file

@ -0,0 +1,116 @@
use std::path::Path;
use std::sync::OnceLock;
use std::time::Duration;
use fabro_util::terminal::Styles;
use indicatif::ProgressStyle;
macro_rules! cached_style {
($name:ident, $template:expr) => {
pub(super) fn $name() -> ProgressStyle {
static STYLE: OnceLock<ProgressStyle> = OnceLock::new();
STYLE
.get_or_init(|| ProgressStyle::with_template($template).expect("valid template"))
.clone()
}
};
}
cached_style!(
style_header_running,
" {spinner:.dim} {wide_msg} {elapsed:.dim}"
);
cached_style!(style_header_done, " {wide_msg:.dim} {prefix:.dim}");
cached_style!(
style_stage_running,
" {spinner:.cyan} {wide_msg} {elapsed:.dim}"
);
cached_style!(style_stage_done, " {wide_msg} {prefix:.dim}");
cached_style!(
style_tool_running,
" {spinner:.dim} {wide_msg} {elapsed:.dim}"
);
cached_style!(style_tool_done, " {wide_msg} {prefix:.dim}");
cached_style!(style_subagent_info, " {wide_msg}");
cached_style!(style_branch_done, " {wide_msg} {prefix:.dim}");
cached_style!(style_static_dim, " {wide_msg:.dim}");
cached_style!(style_sandbox_detail, " {wide_msg:.dim}");
cached_style!(style_empty, " ");
pub(super) fn green_check(styles: &Styles) -> String {
styles.green.apply_to("\u{2713}").to_string()
}
pub(super) fn red_cross(styles: &Styles) -> String {
styles.red.apply_to("\u{2717}").to_string()
}
pub(super) fn warning_glyph(styles: &Styles) -> String {
styles.yellow.apply_to("\u{26a0}").to_string()
}
pub(crate) fn format_duration_short(d: Duration) -> String {
let secs = d.as_secs();
if secs >= 60 {
format!("{}m{:02}s", secs / 60, secs % 60)
} else if d.as_millis() >= 1000 {
format!("{secs}s")
} else {
format!("{}ms", d.as_millis())
}
}
pub(super) fn terminal_hyperlink(url: &str, text: &str) -> String {
format!("\x1b]8;;{url}\x1b\\{text}\x1b]8;;\x1b\\")
}
pub(super) fn format_number(n: f64) -> String {
if (n - n.round()).abs() < f64::EPSILON {
#[allow(clippy::cast_possible_truncation)]
let i = n as i64;
format!("{i}")
} else {
format!("{n:.1}")
}
}
pub(super) fn truncate(s: &str, max: usize) -> String {
let single_line = s.split_whitespace().collect::<Vec<_>>().join(" ");
if single_line.len() > max {
let mut truncated: String = single_line.chars().take(max - 3).collect();
truncated.push_str("...");
truncated
} else {
single_line
}
}
pub(super) fn last_line_truncated(s: &str, max: usize) -> String {
let line = s
.trim()
.lines()
.rfind(|line| !line.trim().is_empty())
.unwrap_or("")
.trim();
if line.len() > max {
let mut truncated: String = line.chars().take(max - 3).collect();
truncated.push_str("...");
truncated
} else {
line.to_string()
}
}
pub(super) fn shorten_path(path: &str, working_directory: Option<&str>) -> String {
if let Some(wd) = working_directory {
if let Ok(rel) = Path::new(path).strip_prefix(wd) {
return rel.display().to_string();
}
}
if let Ok(cwd) = std::env::current_dir() {
if let Ok(rel) = Path::new(path).strip_prefix(&cwd) {
return rel.display().to_string();
}
}
path.to_string()
}