From a31a7295f2f8bc08a1288aa98e94f0321114b984 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 30 Apr 2026 14:19:10 -0400 Subject: [PATCH] fix(retro): hydrate command timing log refs Resolve command stdout and stderr blob refs inside retro script_timing exports and the nested run.json projection, while leaving script invocation metadata unchanged. --- lib/crates/fabro-retro/src/retro_agent.rs | 120 +++++++++++++++++++++- 1 file changed, 116 insertions(+), 4 deletions(-) diff --git a/lib/crates/fabro-retro/src/retro_agent.rs b/lib/crates/fabro-retro/src/retro_agent.rs index ece61c55c..4ec077292 100644 --- a/lib/crates/fabro-retro/src/retro_agent.rs +++ b/lib/crates/fabro-retro/src/retro_agent.rs @@ -14,6 +14,7 @@ use fabro_llm::provider::Provider; use fabro_llm::types::ToolDefinition; use fabro_store::{EventEnvelope, RunProjection, SerializableProjection}; use fabro_types::{RunBlobId, parse_blob_ref}; +use serde_json::Value; use tokio::task::JoinHandle; use crate::retro::{RetroNarrative, SmoothnessRating}; @@ -338,7 +339,9 @@ async fn upload_data_files( ) .await?; - let run_content = serde_json::to_string_pretty(&SerializableProjection(state))?; + let retro_state = hydrate_retro_projection(state, blob_reader).await?; + + let run_content = serde_json::to_string_pretty(&SerializableProjection(&retro_state))?; upload_file( sandbox, target_dir, @@ -361,7 +364,7 @@ async fn upload_data_files( stage_ids.sort(); for stage_id in stage_ids { - let Some(node) = state.node(&stage_id) else { + let Some(node) = retro_state.node(&stage_id) else { continue; }; let base = PathBuf::from("stages").join(stage_id.to_string()); @@ -440,6 +443,60 @@ async fn upload_data_files( Ok(()) } +async fn hydrate_retro_projection( + state: &RunProjection, + blob_reader: Option<&RetroBlobReader>, +) -> anyhow::Result { + let mut hydrated = state.clone(); + let stage_ids: Vec<_> = hydrated + .iter_nodes() + .map(|(stage_id, _)| stage_id.clone()) + .collect(); + + for stage_id in stage_ids { + let Some(mut node) = hydrated.node(&stage_id).cloned() else { + continue; + }; + node.script_timing = + resolve_script_timing_value(node.script_timing.as_ref(), blob_reader).await?; + hydrated.set_node(stage_id, node); + } + + Ok(hydrated) +} + +async fn resolve_script_timing_value( + value: Option<&Value>, + blob_reader: Option<&RetroBlobReader>, +) -> anyhow::Result> { + let Some(value) = value else { + return Ok(None); + }; + let mut value = value.clone(); + let Value::Object(fields) = &mut value else { + return Ok(Some(value)); + }; + + for key in ["stdout", "stderr"] { + let current = fields + .get(key) + .and_then(Value::as_str) + .map(ToOwned::to_owned); + if let Some(current) = current { + fields.insert( + key.to_string(), + Value::String( + resolve_text_file_content(Some(current), blob_reader) + .await? + .unwrap_or_default(), + ), + ); + } + } + + Ok(Some(value)) +} + async fn resolve_text_file_content( content: Option, blob_reader: Option<&RetroBlobReader>, @@ -691,9 +748,21 @@ mod tests { let stage_id = StageId::new("build", 1); let mut state = RunProjection::default(); + let stdout_ref = fabro_types::format_blob_ref(&stdout_id); + let stderr_ref = fabro_types::format_blob_ref(&stderr_id); state.set_node(stage_id, NodeState { - stdout: Some(fabro_types::format_blob_ref(&stdout_id)), - stderr: Some(fabro_types::format_blob_ref(&stderr_id)), + script_invocation: Some(serde_json::json!({ + "command": "cargo test", + "stdout": stdout_ref, + "stderr": stderr_ref, + })), + script_timing: Some(serde_json::json!({ + "exit_code": 0, + "stdout": stdout_ref, + "stderr": stderr_ref, + })), + stdout: Some(stdout_ref), + stderr: Some(stderr_ref), ..NodeState::default() }); @@ -734,5 +803,48 @@ mod tests { .expect("stderr file should exist"), "resolved stderr" ); + + let script_timing: serde_json::Value = serde_json::from_str( + &fs::read_to_string(target_dir.join("stages/build@1/script_timing.json")) + .await + .expect("script timing should exist"), + ) + .expect("script timing should parse"); + assert_eq!(script_timing["stdout"], "resolved stdout"); + assert_eq!(script_timing["stderr"], "resolved stderr"); + + let script_invocation: serde_json::Value = serde_json::from_str( + &fs::read_to_string(target_dir.join("stages/build@1/script_invocation.json")) + .await + .expect("script invocation should exist"), + ) + .expect("script invocation should parse"); + assert_eq!( + script_invocation["stdout"], + fabro_types::format_blob_ref(&stdout_id) + ); + assert_eq!( + script_invocation["stderr"], + fabro_types::format_blob_ref(&stderr_id) + ); + + let run_json: serde_json::Value = serde_json::from_str( + &fs::read_to_string(target_dir.join("run.json")) + .await + .expect("run.json should exist"), + ) + .expect("run.json should parse"); + assert_eq!( + run_json["nodes"]["build@1"]["script_timing"]["stdout"], + "resolved stdout" + ); + assert_eq!( + run_json["nodes"]["build@1"]["script_timing"]["stderr"], + "resolved stderr" + ); + assert_eq!( + run_json["nodes"]["build@1"]["script_invocation"]["stdout"], + fabro_types::format_blob_ref(&stdout_id) + ); } }