Use store-backed run summaries and attach replay

This commit is contained in:
Bryan Helmkamp 2026-04-01 20:18:41 -07:00
parent abffd09f41
commit 9826001305
No known key found for this signature in database
2 changed files with 88 additions and 18 deletions

View file

@ -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<u32> = 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::<Vec<_>>();
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<String> {
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::<std::collections::BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
Value::Array(values) => {
Value::Array(values.into_iter().map(normalize_json_value).collect())
}
other => other,
}
}
fn read_status_record(path: &Path) -> Option<RunStatusRecord> {

View file

@ -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;
};