From 98260013059608992bf0c6056b95d55f5c0e2691 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Wed, 1 Apr 2026 20:18:41 -0700 Subject: [PATCH] Use store-backed run summaries and attach replay --- .../fabro-cli/src/commands/run/attach.rs | 60 ++++++++++++++++++- .../fabro-cli/src/commands/run/output.rs | 46 ++++++++------ 2 files changed, 88 insertions(+), 18 deletions(-) diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index 74a48e3a5..28d27020e 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -16,6 +16,7 @@ use fabro_util::terminal::Styles; use fabro_workflow::outcome::StageStatus; use fabro_workflow::records::{Conclusion, ConclusionExt, RunRecord, RunRecordExt}; use fabro_workflow::run_status::{RunStatus, RunStatusRecord, RunStatusRecordExt}; +use serde_json::{Map, Value}; use tokio::signal::ctrl_c; use tokio::time::{self, sleep}; @@ -148,6 +149,7 @@ async fn attach_run_store( let mut stream = run_store .watch_events_from(if last_seq == 0 { 1 } else { last_seq + 1 }) .await?; + let mut next_seq = if last_seq == 0 { 1 } else { last_seq + 1 }; let mut cached_pid: Option = None; let attach_started = Instant::now(); @@ -189,6 +191,7 @@ async fn attach_run_store( Ok(Some(Ok(event))) => { let line = event_payload_line(&event)?; emit_progress_line(&mut progress_ui, &line, json_output)?; + next_seq = event.seq.saturating_add(1); saw_event = true; } Ok(Some(Err(err))) => return Err(err.into()), @@ -258,10 +261,14 @@ async fn attach_run_store( if let Some(child_alive) = child_alive_via_handle { if !child_alive && !saw_event { + flush_remaining_store_events(run_store, next_seq, &mut progress_ui, json_output) + .await?; break; } } else { if terminal_status.is_some() && !saw_event { + flush_remaining_store_events(run_store, next_seq, &mut progress_ui, json_output) + .await?; break; } @@ -277,6 +284,8 @@ async fn attach_run_store( } }; if !engine_alive { + flush_remaining_store_events(run_store, next_seq, &mut progress_ui, json_output) + .await?; break; } } @@ -287,6 +296,38 @@ async fn attach_run_store( Ok(determine_exit_code_with_store(run_store).await) } +async fn flush_remaining_store_events( + run_store: &dyn RunStore, + mut next_seq: u32, + progress_ui: &mut run_progress::ProgressUI, + json_output: bool, +) -> Result<()> { + let deadline = Instant::now() + ATTACH_FINAL_STATUS_GRACE; + loop { + let mut saw_new_event = false; + let events = run_store + .list_events() + .await? + .into_iter() + .filter(|e| e.seq >= next_seq) + .collect::>(); + for event in events { + let line = event_payload_line(&event)?; + emit_progress_line(progress_ui, &line, json_output)?; + next_seq = event.seq.saturating_add(1); + saw_new_event = true; + } + + if !saw_new_event || Instant::now() >= deadline { + break; + } + + sleep(Duration::from_millis(100)).await; + } + + Ok(()) +} + async fn attach_run_files( run_dir: &Path, verbose: bool, @@ -559,7 +600,24 @@ fn show_progress(progress_ui: &mut run_progress::ProgressUI, json_output: bool) } fn event_payload_line(event: &EventEnvelope) -> Result { - serde_json::to_string(event.payload.as_value()).map_err(Into::into) + serde_json::to_string(&normalize_json_value(event.payload.as_value().clone())) + .map_err(Into::into) +} + +fn normalize_json_value(value: Value) -> Value { + match value { + Value::Object(map) => Value::Object( + map.into_iter() + .map(|(key, value)| (key, normalize_json_value(value))) + .collect::>() + .into_iter() + .collect::>(), + ), + Value::Array(values) => { + Value::Array(values.into_iter().map(normalize_json_value).collect()) + } + other => other, + } } fn read_status_record(path: &Path) -> Option { diff --git a/lib/crates/fabro-cli/src/commands/run/output.rs b/lib/crates/fabro-cli/src/commands/run/output.rs index dff033ab5..5a5798ba6 100644 --- a/lib/crates/fabro-cli/src/commands/run/output.rs +++ b/lib/crates/fabro-cli/src/commands/run/output.rs @@ -10,7 +10,7 @@ use fabro_util::text::strip_goal_decoration; use fabro_workflow::asset_snapshot::collect_asset_paths; use fabro_workflow::outcome::{StageStatus, format_cost}; use fabro_workflow::pipeline::{Persisted, Validated}; -use fabro_workflow::records::{Checkpoint, CheckpointExt, Conclusion, ConclusionExt}; +use fabro_workflow::records::{Checkpoint, CheckpointExt, Conclusion}; use indicatif::HumanDuration; use crate::shared::{format_tokens_human, print_diagnostics, relative_path, tilde_path}; @@ -78,23 +78,26 @@ pub(crate) async fn print_run_summary( styles: &Styles, ) -> Result<()> { let run_id = run_id.to_string(); - let conclusion_path = run_dir.join("conclusion.json"); - let Ok(conclusion) = Conclusion::load(&conclusion_path) else { - return Ok(()); - }; - - let pr_url = match run_id.parse() { + let (run_store, conclusion, pr_url) = match run_id.parse() { Ok(parsed_run_id) => { - if let Some(run_store) = store::open_run_reader(storage_dir, &parsed_run_id).await? { - run_store + let run_store = store::open_run_reader(storage_dir, &parsed_run_id).await?; + let conclusion = match run_store.as_deref() { + Some(run_store) => run_store.get_conclusion().await?, + None => None, + }; + let pr_url = match run_store.as_deref() { + Some(run_store) => run_store .get_pull_request() .await? - .map(|record: PullRequestRecord| record.html_url) - } else { - None - } + .map(|record: PullRequestRecord| record.html_url), + None => None, + }; + (run_store, conclusion, pr_url) } - Err(_) => None, + Err(_) => (None, None, None), + }; + let Some(conclusion) = conclusion else { + return Ok(()); }; print_run_conclusion( @@ -105,7 +108,7 @@ pub(crate) async fn print_run_summary( pr_url.as_deref(), styles, ); - print_final_output(run_dir, styles); + print_final_output(run_store.as_deref(), run_dir, styles).await; print_assets(run_dir, styles); Ok(()) } @@ -199,8 +202,17 @@ pub(crate) fn print_run_conclusion( } } -pub(crate) fn print_final_output(run_dir: &Path, styles: &Styles) { - let Ok(checkpoint) = Checkpoint::load(&run_dir.join("checkpoint.json")) else { +pub(crate) async fn print_final_output( + run_store: Option<&dyn fabro_store::RunStore>, + run_dir: &Path, + styles: &Styles, +) { + let checkpoint = match run_store { + Some(run_store) => run_store.get_checkpoint().await.ok().flatten(), + None => None, + } + .or_else(|| Checkpoint::load(&run_dir.join("checkpoint.json")).ok()); + let Some(checkpoint) = checkpoint else { return; };