Fix run progress replay durations

This commit is contained in:
Bryan Helmkamp 2026-03-30 14:58:51 -04:00
parent e93c13b9db
commit 62786ccf4b
No known key found for this signature in database
5 changed files with 188 additions and 50 deletions

View file

@ -1,5 +1,6 @@
use std::convert::TryFrom;
use chrono::{DateTime, Utc};
use fabro_workflow::event::RunNoticeLevel;
use fabro_workflow::outcome::{StageUsage, compute_stage_cost};
use serde_json::{Map, Value};
@ -167,11 +168,14 @@ pub(super) enum ProgressEvent {
tool_name: String,
tool_call_id: String,
arguments: Value,
timestamp: Option<DateTime<Utc>>,
},
ToolCallCompleted {
stage_node_id: String,
tool_call_id: String,
is_error: bool,
duration_ms: Option<u64>,
timestamp: Option<DateTime<Utc>>,
},
ContextWindowWarning {
stage_node_id: String,
@ -235,6 +239,7 @@ pub(super) enum ProgressEvent {
},
}
#[allow(clippy::needless_pass_by_value)]
pub(super) fn from_flattened_fields(
event_name: &str,
fields: Map<String, Value>,
@ -381,6 +386,7 @@ pub(super) fn from_flattened_fields(
.get("arguments")
.cloned()
.unwrap_or_else(|| Value::Object(Map::new())),
timestamp: timestamp_field(&fields, "ts"),
}),
"Agent.ToolCallCompleted" => Some(ProgressEvent::ToolCallCompleted {
stage_node_id: string_field(&fields, "node_id")
@ -388,6 +394,8 @@ pub(super) fn from_flattened_fields(
.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"),
duration_ms: optional_u64_field(&fields, "duration_ms"),
timestamp: timestamp_field(&fields, "ts"),
}),
"Agent.Warning" if string_field(&fields, "kind").as_deref() == Some("context_window") => {
let usage_percent = fields
@ -532,6 +540,10 @@ fn u64_field(fields: &Map<String, Value>, key: &str) -> u64 {
fields.get(key).and_then(Value::as_u64).unwrap_or(0)
}
fn optional_u64_field(fields: &Map<String, Value>, key: &str) -> Option<u64> {
fields.get(key).and_then(Value::as_u64)
}
fn i64_field(fields: &Map<String, Value>, key: &str) -> i64 {
fields.get(key).and_then(Value::as_i64).unwrap_or(0)
}
@ -544,6 +556,13 @@ fn bool_field(fields: &Map<String, Value>, key: &str) -> bool {
fields.get(key).and_then(Value::as_bool).unwrap_or(false)
}
fn timestamp_field(fields: &Map<String, Value>, key: &str) -> Option<DateTime<Utc>> {
let value = fields.get(key)?.as_str()?;
DateTime::parse_from_rfc3339(value)
.ok()
.map(|timestamp| timestamp.with_timezone(&Utc))
}
#[cfg(test)]
mod tests {
use fabro_agent::AgentEvent;
@ -631,6 +650,47 @@ mod tests {
));
}
#[test]
fn parse_tool_call_timestamps_from_jsonl_envelope() {
let started_fields = json_map(serde_json::json!({
"ts": "2026-03-30T12:00:00.000Z",
"node_id": "code",
"tool_name": "read_file",
"tool_call_id": "tc1",
"arguments": {"path": "src/main.rs"}
}));
let completed_fields = json_map(serde_json::json!({
"ts": "2026-03-30T12:00:00.500Z",
"node_id": "code",
"tool_call_id": "tc1",
"is_error": false,
"duration_ms": 500
}));
let started = from_flattened_fields("Agent.ToolCallStarted", started_fields).unwrap();
let completed = from_flattened_fields("Agent.ToolCallCompleted", completed_fields).unwrap();
assert!(matches!(
started,
ProgressEvent::ToolCallStarted {
timestamp: Some(timestamp),
..
} if timestamp == DateTime::parse_from_rfc3339("2026-03-30T12:00:00.000Z")
.unwrap()
.with_timezone(&Utc)
));
assert!(matches!(
completed,
ProgressEvent::ToolCallCompleted {
duration_ms: Some(500),
timestamp: Some(timestamp),
..
} if timestamp == DateTime::parse_from_rfc3339("2026-03-30T12:00:00.500Z")
.unwrap()
.with_timezone(&Utc)
));
}
#[test]
fn round_trip_sandbox_ready() {
let event = WorkflowRunEvent::Sandbox {

View file

@ -15,26 +15,20 @@ impl InfoDisplay {
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_worktree(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,
) {
pub(super) fn show_base_info(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);
Self::insert_info_line(renderer, &text);
}
pub(super) fn on_run_notice(
&self,
renderer: &ProgressRenderer,
level: RunNoticeLevel,
code: &str,
@ -51,24 +45,19 @@ impl InfoDisplay {
} else {
format!(" {}", styles.dim.apply_to(format!("[{code}]")))
};
self.insert_info_line(renderer, &format!("{label} {message}{code_suffix}"));
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,
) {
pub(super) fn on_pull_request_created(renderer: &ProgressRenderer, pr_url: &str, draft: bool) {
let label = if draft { "Draft PR:" } else { "PR:" };
self.insert_info_line(
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(
pub(super) fn on_pull_request_failed(renderer: &ProgressRenderer, error: &str) {
Self::insert_info_line(
renderer,
&format!("{} {error}", renderer.styles().red.apply_to("PR failed:")),
);
@ -93,7 +82,7 @@ impl InfoDisplay {
} else {
String::new()
};
self.insert_info_line(
Self::insert_info_line(
renderer,
&format!("\u{2192} {from_node} \u{2192} {to_node}{detail}"),
);
@ -109,7 +98,7 @@ impl InfoDisplay {
return;
}
self.insert_info_line(
Self::insert_info_line(
renderer,
&format!("\u{21ba} {from_node} \u{2192} {to_node} (loop restart)"),
);
@ -127,7 +116,7 @@ impl InfoDisplay {
return;
}
self.insert_info_line(
Self::insert_info_line(
renderer,
&format!(
"\u{21bb} {name}: retrying (attempt {attempt}/{max_attempts}, delay {})",
@ -136,7 +125,7 @@ impl InfoDisplay {
);
}
fn insert_info_line(&self, renderer: &ProgressRenderer, message: &str) {
fn insert_info_line(renderer: &ProgressRenderer, message: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_static_dim());

View file

@ -99,12 +99,10 @@ impl ProgressUI {
base_sha,
} => {
if let Some(worktree_dir) = worktree_dir {
self.info
.show_worktree(renderer, std::path::Path::new(&worktree_dir));
InfoDisplay::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);
InfoDisplay::show_base_info(renderer, base_branch.as_deref(), &base_sha);
}
}
ProgressEvent::WorkingDirectorySet { working_directory } => {
@ -132,7 +130,7 @@ impl ProgressUI {
);
}
ProgressEvent::SshAccessReady { ssh_command } => {
self.setup.on_ssh_access_ready(renderer, &ssh_command);
SetupDisplay::on_ssh_access_ready(renderer, &ssh_command);
}
ProgressEvent::SetupStarted { command_count } => {
self.setup.on_setup_started(renderer, command_count);
@ -178,7 +176,7 @@ impl ProgressUI {
lifecycle_command_count,
workspace_folder,
} => {
self.setup.on_devcontainer_resolved(
SetupDisplay::on_devcontainer_resolved(
renderer,
dockerfile_lines,
environment_count,
@ -291,6 +289,7 @@ impl ProgressUI {
tool_name,
tool_call_id,
arguments,
timestamp,
} => {
self.stage.on_tool_call_started(
renderer,
@ -298,18 +297,23 @@ impl ProgressUI {
&tool_name,
&tool_call_id,
&arguments,
timestamp,
);
}
ProgressEvent::ToolCallCompleted {
stage_node_id,
tool_call_id,
is_error,
duration_ms,
timestamp,
} => {
self.stage.on_tool_call_completed(
renderer,
&stage_node_id,
&tool_call_id,
is_error,
duration_ms,
timestamp,
);
}
ProgressEvent::ContextWindowWarning {
@ -405,13 +409,13 @@ impl ProgressUI {
code,
message,
} => {
self.info.on_run_notice(renderer, level, &code, &message);
InfoDisplay::on_run_notice(renderer, level, &code, &message);
}
ProgressEvent::PullRequestCreated { pr_url, draft } => {
self.info.on_pull_request_created(renderer, &pr_url, draft);
InfoDisplay::on_pull_request_created(renderer, &pr_url, draft);
}
ProgressEvent::PullRequestFailed { error } => {
self.info.on_pull_request_failed(renderer, &error);
InfoDisplay::on_pull_request_failed(renderer, &error);
}
}
}
@ -962,4 +966,46 @@ mod tests {
PR failed: auth token expired
");
}
#[test]
fn tty_parallel_branch_completion_uses_recorded_duration() {
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,
});
ui.handle_event(&WorkflowRunEvent::ParallelBranchCompleted {
branch: "security".into(),
index: 0,
duration_ms: 500,
status: "success".into(),
});
let stage = &ui.stage.active_stages["fork1"];
assert_eq!(stage.tool_calls[0].bar.prefix(), "500ms");
}
#[test]
fn tty_tool_call_completion_uses_jsonl_timestamps() {
let mut ui = ProgressUI::new(true, false);
ui.handle_json_line(
r#"{"ts":"2026-03-30T12:00:00.000Z","event":"StageStarted","node_id":"code","node_label":"Code"}"#,
);
ui.handle_json_line(
r#"{"ts":"2026-03-30T12:00:00.000Z","event":"Agent.ToolCallStarted","node_id":"code","tool_name":"read_file","tool_call_id":"tc1","arguments":{"path":"src/main.rs"}}"#,
);
ui.handle_json_line(
r#"{"ts":"2026-03-30T12:00:00.500Z","event":"Agent.ToolCallCompleted","node_id":"code","tool_call_id":"tc1","is_error":false}"#,
);
let stage = &ui.stage.active_stages["code"];
assert_eq!(stage.tool_calls[0].bar.prefix(), "500ms");
}
}

View file

@ -99,7 +99,7 @@ impl SetupDisplay {
}
}
pub(super) fn on_ssh_access_ready(&self, renderer: &ProgressRenderer, ssh_command: &str) {
pub(super) fn on_ssh_access_ready(renderer: &ProgressRenderer, ssh_command: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_sandbox_detail());
@ -240,7 +240,6 @@ impl SetupDisplay {
}
pub(super) fn on_devcontainer_resolved(
&self,
renderer: &ProgressRenderer,
dockerfile_lines: u64,
environment_count: u64,

View file

@ -2,6 +2,7 @@ use std::collections::{HashMap, VecDeque};
use std::convert::TryFrom;
use std::time::Duration;
use chrono::{DateTime, Utc};
use indicatif::ProgressBar;
use fabro_workflow::outcome::{StageStatus, format_cost};
@ -27,6 +28,7 @@ pub(super) struct ToolCallEntry {
pub(super) status: ToolCallStatus,
pub(super) bar: ProgressBar,
pub(super) is_branch: bool,
pub(super) started_at: Option<DateTime<Utc>>,
}
#[derive(Debug)]
@ -137,10 +139,10 @@ impl StageDisplay {
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 succeeded = status.parse::<StageStatus>().map_or_else(
|_| matches!(status, "success" | "partial_success"),
|status| matches!(status, StageStatus::Success | StageStatus::PartialSuccess),
);
let cost_str = usage
.and_then(ProgressUsage::display_cost)
.map(|cost| format!("{} ", format_cost(cost)))
@ -190,7 +192,7 @@ impl StageDisplay {
"",
);
let summary = styles::last_line_truncated(error, 120);
self.insert_global_info_line(
Self::insert_global_info_line(
renderer,
&format!("{} {summary}", renderer.styles().red.apply_to("Error:")),
);
@ -232,6 +234,7 @@ impl StageDisplay {
status: ToolCallStatus::Running,
bar,
is_branch: true,
started_at: None,
});
}
@ -271,9 +274,7 @@ impl StageDisplay {
if renderer.is_tty() {
entry.bar.set_style(styles::style_branch_done());
entry
.bar
.set_prefix(styles::format_duration_short(entry.bar.elapsed()));
set_duration_prefix(&entry.bar, Some(duration_ms));
entry
.bar
.finish_with_message(format!("{glyph} {}", entry.display_name));
@ -312,6 +313,7 @@ impl StageDisplay {
tool_name: &str,
tool_call_id: &str,
arguments: &serde_json::Value,
timestamp: Option<DateTime<Utc>>,
) {
let display_name = self.tool_display_name(renderer, tool_name, arguments);
let Some(stage) = self.active_stages.get_mut(stage_node_id) else {
@ -341,6 +343,7 @@ impl StageDisplay {
status: ToolCallStatus::Running,
bar,
is_branch: false,
started_at: timestamp,
});
}
@ -350,6 +353,8 @@ impl StageDisplay {
stage_node_id: &str,
tool_call_id: &str,
is_error: bool,
duration_ms: Option<u64>,
timestamp: Option<DateTime<Utc>>,
) {
if let Some(counts) = self.stage_counts.get_mut(stage_node_id) {
counts.1 += 1;
@ -378,9 +383,20 @@ impl StageDisplay {
};
if renderer.is_tty() {
entry.bar.set_style(styles::style_tool_done());
entry
.bar
.set_prefix(styles::format_duration_short(entry.bar.elapsed()));
let computed_duration_ms = duration_ms.or_else(|| {
entry
.started_at
.zip(timestamp)
.and_then(|(started_at, completed_at)| {
u64::try_from(
completed_at
.signed_duration_since(started_at)
.num_milliseconds(),
)
.ok()
})
});
set_duration_prefix(&entry.bar, computed_duration_ms);
entry
.bar
.finish_with_message(format!("{glyph} {}", entry.display_name));
@ -564,7 +580,7 @@ impl StageDisplay {
) {
let Some(stage) = self.active_stages.remove(node_id) else {
if !renderer.is_tty() {
self.print_plain_stage_completion(renderer, name, glyph, prefix);
Self::print_plain_stage_completion(renderer, name, glyph, prefix);
}
return;
};
@ -587,12 +603,11 @@ impl StageDisplay {
.spinner
.finish_with_message(format!("{glyph} {}", stage.display_name));
} else {
self.print_plain_stage_completion(renderer, name, glyph, prefix);
Self::print_plain_stage_completion(renderer, name, glyph, prefix);
}
}
fn print_plain_stage_completion(
&self,
renderer: &ProgressRenderer,
name: &str,
glyph: &str,
@ -605,7 +620,7 @@ impl StageDisplay {
}
}
fn insert_global_info_line(&self, renderer: &ProgressRenderer, message: &str) {
fn insert_global_info_line(renderer: &ProgressRenderer, message: &str) {
if renderer.is_tty() {
let bar = renderer.add_spinner();
bar.set_style(styles::style_static_dim());
@ -697,3 +712,32 @@ impl StageDisplay {
}
}
}
fn set_duration_prefix(bar: &ProgressBar, duration_ms: Option<u64>) {
let prefix = duration_ms.map_or_else(
|| styles::format_duration_short(bar.elapsed()),
|duration_ms| styles::format_duration_short(Duration::from_millis(duration_ms)),
);
bar.set_prefix(prefix);
}
#[cfg(test)]
mod tests {
use super::*;
use crate::commands::run::run_progress::renderer::ProgressRenderer;
#[test]
fn tool_display_name_shortens_paths_relative_to_working_directory() {
let renderer = ProgressRenderer::new_plain(Box::new(std::io::sink()), false);
let mut stage = StageDisplay::new(false);
stage.set_working_directory("/workspace".into());
let display_name = stage.tool_display_name(
&renderer,
"read_file",
&serde_json::json!({"file_path": "/workspace/src/main.rs"}),
);
assert_eq!(display_name, "read_file(src/main.rs)");
}
}