From ae70aa6ccf7922523b6e83cffe7516f03ef50f7d Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Fri, 18 Sep 2026 16:00:11 -0400 Subject: [PATCH] Resolve blob references in inspect, diff, output and dump `fabro inspect` lists the stages with their output, response and diff as the projection holds them; `fabro diff` resolves a patch reference through the run's blob endpoint; the final output and `dump` decode a plain reference as text and a `#json` reference as a value. The diff tests run over a git-backed Petri run again, and the large-output dump tests print many lines rather than one, since Petri caps a single line at 64 KiB. Co-Authored-By: Claude Fable 5.1 --- lib/apps/fabro-cli/src/commands/run/diff.rs | 17 ++- lib/apps/fabro-cli/src/commands/run/output.rs | 22 ++-- .../fabro-cli/src/commands/runs/inspect.rs | 33 ++++++ lib/apps/fabro-cli/tests/it/cmd/diff.rs | 101 +++++++++++++++++ lib/apps/fabro-cli/tests/it/cmd/dump.rs | 4 +- lib/apps/fabro-cli/tests/it/cmd/support.rs | 60 +++++++++++ lib/components/fabro-dump/src/lib.rs | 102 +++++++++++++----- 7 files changed, 302 insertions(+), 37 deletions(-) diff --git a/lib/apps/fabro-cli/src/commands/run/diff.rs b/lib/apps/fabro-cli/src/commands/run/diff.rs index 41e86954a..15f8fdca6 100644 --- a/lib/apps/fabro-cli/src/commands/run/diff.rs +++ b/lib/apps/fabro-cli/src/commands/run/diff.rs @@ -10,11 +10,12 @@ use std::io::{self, IsTerminal, Write}; use anyhow::{Context, Result, bail}; +use fabro_types::{RunId, parse_blob_ref}; use tracing::{debug, info}; use crate::args::DiffArgs; use crate::command_context::CommandContext; -use crate::server_client::RunProjection; +use crate::server_client::{Client, RunProjection}; use crate::shared::print_json_pretty; pub(crate) async fn run(args: DiffArgs, base_ctx: &CommandContext) -> Result<()> { @@ -25,6 +26,7 @@ pub(crate) async fn run(args: DiffArgs, base_ctx: &CommandContext) -> Result<()> let state = client.get_run_state(&run_id).await?; let patch = resolve_diff(&state, &args)?; + let patch = resolve_patch_text(client.as_ref(), &run_id, patch).await?; if ctx.json_output() { let value = serde_json::json!({ @@ -92,6 +94,19 @@ fn resolve_diff(state: &RunProjection, args: &DiffArgs) -> Result { ) } +/// The patch as text: the projection holds a large patch as a +/// `blob://sha256/` reference into the run's blob table. +async fn resolve_patch_text(client: &Client, run_id: &RunId, patch: String) -> Result { + let Some(blob_hash) = parse_blob_ref(patch.trim()) else { + return Ok(patch); + }; + let bytes = client + .read_run_blob(run_id, &blob_hash) + .await? + .with_context(|| format!("the diff's patch blob {blob_hash} is missing from the store"))?; + Ok(String::from_utf8_lossy(&bytes).into_owned()) +} + fn colorize_diff_line(line: &str) -> String { if line.starts_with("+++") || line.starts_with("---") { format!("\x1b[1m{line}\x1b[0m") diff --git a/lib/apps/fabro-cli/src/commands/run/output.rs b/lib/apps/fabro-cli/src/commands/run/output.rs index 50fa18de1..c1a935aa8 100644 --- a/lib/apps/fabro-cli/src/commands/run/output.rs +++ b/lib/apps/fabro-cli/src/commands/run/output.rs @@ -6,7 +6,7 @@ use cli_table::format::{Border, Justify, Separator}; use cli_table::{Cell, CellStruct, Style, Table}; use fabro_api::types; use fabro_types::diagnostic::{Diagnostic, RelatedDiagnostic, Severity}; -use fabro_types::{PullRequestLink, RunId, StageId, parse_blob_ref}; +use fabro_types::{BlobRefEncoding, PullRequestLink, RunId, StageId, parse_blob_ref_encoded}; use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus}; use fabro_util::error::render_with_causes; use fabro_util::printer::Printer; @@ -324,20 +324,24 @@ async fn resolve_response_string( run_id: &RunId, response: &str, ) -> Result> { - let Some(blob_hash) = parse_blob_ref(response) else { + let Some((blob_hash, encoding)) = parse_blob_ref_encoded(response) else { return Ok(Some(response.to_string())); }; let Some(bytes) = client.read_run_blob(run_id, &blob_hash).await? else { return Ok(None); }; - let value: serde_json::Value = - serde_json::from_slice(&bytes).context("blob-backed final output should be valid JSON")?; - - Ok(Some(match value { - serde_json::Value::String(text) => text, - other => other.to_string(), - })) + match encoding { + BlobRefEncoding::Text => Ok(Some(String::from_utf8_lossy(&bytes).into_owned())), + BlobRefEncoding::Json => { + let value: serde_json::Value = serde_json::from_slice(&bytes) + .context("blob-backed final output should be valid JSON")?; + Ok(Some(match value { + serde_json::Value::String(text) => text, + other => other.to_string(), + })) + } + } } async fn list_artifact_display_entries_with_client( diff --git a/lib/apps/fabro-cli/src/commands/runs/inspect.rs b/lib/apps/fabro-cli/src/commands/runs/inspect.rs index 71f55cab2..de06e8684 100644 --- a/lib/apps/fabro-cli/src/commands/runs/inspect.rs +++ b/lib/apps/fabro-cli/src/commands/runs/inspect.rs @@ -1,4 +1,7 @@ +use std::collections::BTreeMap; + use anyhow::Result; +use fabro_types::{StageHandler, StageState}; use fabro_workflow::run_status::RunStatus; use serde::Serialize; @@ -17,6 +20,23 @@ pub(crate) struct InspectOutput { pub conclusion: Option, pub checkpoint: Option, pub sandbox: Option, + /// The stages by their id, with what each produced. A large output or + /// response is a `blob://sha256/` reference into the run's blob + /// table, as the projection holds it. + #[serde(skip_serializing_if = "BTreeMap::is_empty")] + pub stages: BTreeMap, +} + +#[derive(Debug, Serialize)] +pub(crate) struct InspectStage { + pub handler: Option, + pub state: StageState, + #[serde(skip_serializing_if = "Option::is_none")] + pub output: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub response: Option, + #[serde(skip_serializing_if = "Option::is_none")] + pub diff: Option, } pub(crate) async fn run(args: &InspectArgs, base_ctx: &CommandContext) -> Result<()> { @@ -36,6 +56,18 @@ fn inspect_run_state(run: &ServerRunInfo, state: RunProjection) -> InspectOutput let checkpoint = state .current_checkpoint() .and_then(|record| serde_json::to_value(record).ok()); + let stages = state + .iter_stages() + .map(|(stage_id, stage)| { + (stage_id.to_string(), InspectStage { + handler: stage.handler, + state: stage.state, + output: stage.output.clone(), + response: stage.response.clone(), + diff: stage.diff.clone(), + }) + }) + .collect(); InspectOutput { run_id: run.run_id().to_string(), parent_id: state.parent_id.map(|parent_id| parent_id.to_string()), @@ -51,5 +83,6 @@ fn inspect_run_state(run: &ServerRunInfo, state: RunProjection) -> InspectOutput sandbox: state .sandbox .and_then(|record| serde_json::to_value(record).ok()), + stages, } } diff --git a/lib/apps/fabro-cli/tests/it/cmd/diff.rs b/lib/apps/fabro-cli/tests/it/cmd/diff.rs index 0a8370413..c70d87918 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/diff.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/diff.rs @@ -1,5 +1,7 @@ use fabro_test::{fabro_snapshot, test_context}; +use super::support::{git_filters, setup_git_backed_changed_run, setup_git_backed_noop_run}; + #[test] fn help() { let context = test_context!(); @@ -28,3 +30,102 @@ fn help() { ----- stderr ----- "); } + +#[test] +fn diff_completed_run_without_changes_reports_no_patch() { + let context = test_context!(); + let setup = setup_git_backed_noop_run(&context); + let mut cmd = context.command(); + cmd.args(["diff", &setup.run.run_id]); + + fabro_snapshot!(git_filters(&context), cmd, @" + success: false + exit_code: 1 + ----- stdout ----- + ----- stderr ----- + × Run completed but no stored diff exists — the run may not have produced any changes + "); +} + +#[test] +fn diff_missing_node_diff_reports_helpful_error() { + let context = test_context!(); + let setup = setup_git_backed_changed_run(&context); + let mut cmd = context.command(); + cmd.args(["diff", &setup.run.run_id, "--node", "missing"]); + + fabro_snapshot!(git_filters(&context), cmd, @" + success: false + exit_code: 1 + ----- stdout ----- + ----- stderr ----- + × No diff found for node 'missing' — check the node ID and try again + "); +} + +#[test] +fn diff_completed_run_with_changes_prints_patch() { + let context = test_context!(); + let setup = setup_git_backed_changed_run(&context); + let mut cmd = context.command(); + cmd.args(["diff", &setup.run.run_id]); + + fabro_snapshot!(git_filters(&context), cmd, @" + success: true + exit_code: 0 + ----- stdout ----- + diff --git a/story.txt b/story.txt + index [SHA]..[SHA] 100644 + --- a/story.txt + +++ b/story.txt + @@ -1 +1,3 @@ + line 1 + +line 2 + +line 3 + ----- stderr ----- + "); +} + +#[test] +fn diff_node_outputs_specific_patch() { + let context = test_context!(); + let setup = setup_git_backed_changed_run(&context); + let mut cmd = context.command(); + cmd.args(["diff", &setup.run.run_id, "--node", "step_one"]); + + fabro_snapshot!(git_filters(&context), cmd, @" + success: true + exit_code: 0 + ----- stdout ----- + diff --git a/story.txt b/story.txt + index [SHA]..[SHA] 100644 + --- a/story.txt + +++ b/story.txt + @@ -1 +1,2 @@ + line 1 + +line 2 + ----- stderr ----- + "); +} + +#[test] +fn diff_json_carries_the_resolved_patch() { + let context = test_context!(); + let setup = setup_git_backed_changed_run(&context); + let output = context + .command() + .args(["diff", "--json", &setup.run.run_id]) + .output() + .expect("diff should execute"); + assert!( + output.status.success(), + "diff failed\nstdout:\n{}\nstderr:\n{}", + String::from_utf8_lossy(&output.stdout), + String::from_utf8_lossy(&output.stderr) + ); + let value: serde_json::Value = + serde_json::from_slice(&output.stdout).expect("diff JSON should parse"); + let patch = value["diff"].as_str().expect("the patch is text"); + assert!(patch.contains("+line 2\n+line 3"), "{patch}"); + assert!(!patch.contains("blob://"), "{patch}"); +} diff --git a/lib/apps/fabro-cli/tests/it/cmd/dump.rs b/lib/apps/fabro-cli/tests/it/cmd/dump.rs index e23a1712c..9beafb095 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/dump.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/dump.rs @@ -91,7 +91,7 @@ fn dump_exports_large_command_output_backed_by_blob_refs() { start [shape=Mdiamond, label="Start"] exit [shape=Msquare, label="Exit"] - big [shape=parallelogram, label="Big", script="printf '%*s' 120000 '' | tr ' ' x"] + big [shape=parallelogram, label="Big", script="yes xxxxxxxxxxxxxxxx | head -n 8000"] start -> big -> exit } @@ -160,7 +160,7 @@ fn dump_exports_blob_refs_and_artifacts_together() { start [shape=Mdiamond, label="Start"] exit [shape=Msquare, label="Exit"] - big [shape=parallelogram, label="Big", script="mkdir -p assets/shared && printf exported > assets/shared/report.txt && printf '%*s' 120000 '' | tr ' ' x"] + big [shape=parallelogram, label="Big", script="mkdir -p assets/shared && printf exported > assets/shared/report.txt && yes xxxxxxxxxxxxxxxx | head -n 8000"] start -> big -> exit } diff --git a/lib/apps/fabro-cli/tests/it/cmd/support.rs b/lib/apps/fabro-cli/tests/it/cmd/support.rs index fb4d73002..f6cc024fe 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/support.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/support.rs @@ -504,6 +504,66 @@ impl WorkflowGate { } } +/// A git-backed workspace whose run appends two lines to `story.txt`, one +/// per stage: `step_one` adds `line 2`, `step_two` adds `line 3`. +pub(crate) fn setup_git_backed_changed_run(context: &TestContext) -> WorkspaceRunSetup { + git_backed_run( + context, + "changed", + "step_one [shape=parallelogram, script=\"printf 'line 2\\n' >> story.txt\"]\n \ + step_two [shape=parallelogram, script=\"printf 'line 3\\n' >> story.txt\"]", + "start -> step_one -> step_two -> exit", + ) +} + +/// A git-backed workspace whose run changes nothing. +pub(crate) fn setup_git_backed_noop_run(context: &TestContext) -> WorkspaceRunSetup { + git_backed_run( + context, + "noop", + "step_one [shape=parallelogram, script=\"cat story.txt\"]", + "start -> step_one -> exit", + ) +} + +/// A run on the local provider from a workspace with one commit, so the +/// run branch starts from a base the run's diff is measured against. +fn git_backed_run( + context: &TestContext, + name: &str, + stages: &str, + edges: &str, +) -> WorkspaceRunSetup { + let workspace_dir = context.temp_dir.join(format!("git-{name}")); + std::fs::create_dir_all(&workspace_dir) + .unwrap_or_else(|err| panic!("failed to create {}: {err}", workspace_dir.display())); + write_text_file(&workspace_dir.join("story.txt"), "line 1\n"); + write_text_file( + &workspace_dir.join("story.fabro"), + &format!( + "digraph Story {{\n graph [goal=\"Change the story\", default_max_retries=0]\n \ + start [shape=Mdiamond]\n exit [shape=Msquare]\n {stages}\n {edges}\n}}\n" + ), + ); + write_text_file( + &workspace_dir.join("workflow.toml"), + "_version = 1\n\n[workflow]\ngraph = \"story.fabro\"\n\n[run]\ngoal = \"Change the \ + story\"\n\n[run.environment]\nid = \"local\"\n\n[environments.local]\nprovider = \ + \"local\"\n", + ); + init_remote_fixture(&workspace_dir, "main"); + let run = run_local_workflow(context, &workspace_dir, "workflow.toml"); + WorkspaceRunSetup { run, workspace_dir } +} + +/// The run output filters plus one for commit shas, which a patch names in +/// its index lines. +pub(crate) fn git_filters(context: &TestContext) -> Vec<(String, String)> { + let mut filters = context.filters(); + filters.push((r"\b[0-9a-f]{7,40}\b".to_string(), "[SHA]".to_string())); + filters +} + pub(crate) fn setup_local_sandbox_run(context: &TestContext) -> WorkspaceRunSetup { let workspace_dir = context.temp_dir.join("local-sandbox"); std::fs::create_dir_all(&workspace_dir) diff --git a/lib/components/fabro-dump/src/lib.rs b/lib/components/fabro-dump/src/lib.rs index 6062d23f0..4834e65dd 100644 --- a/lib/components/fabro-dump/src/lib.rs +++ b/lib/components/fabro-dump/src/lib.rs @@ -15,7 +15,7 @@ use std::path::{Component, Path, PathBuf}; use anyhow::{Context, Result, bail}; use bytes::Bytes; use fabro_store::{RunProjection, SerializableProjection, StageId, retry_storage_segment}; -use fabro_types::{BlobHash, RunStreamItem, parse_blob_ref}; +use fabro_types::{BlobHash, BlobRefEncoding, RunStreamItem, parse_blob_ref_encoded}; use futures::future::BoxFuture; pub type BlobReader = Box BoxFuture<'static, Result>> + Send>; @@ -215,23 +215,21 @@ impl RunDump { for entry in &mut self.entries { match &mut entry.contents { RunDumpContents::Json(value) => { - let mut blob_hashes = Vec::new(); - collect_blob_refs_in_value(value, &mut blob_hashes); - for blob_hash in blob_hashes { + let mut blob_refs = Vec::new(); + collect_blob_refs_in_value(value, &mut blob_refs); + for (blob_hash, encoding) in blob_refs { if cache.contains_key(&blob_hash) { continue; } let blob = read_blob(blob_hash).await?.with_context(|| { format!("blob {blob_hash:?} is missing from the store") })?; - let hydrated: serde_json::Value = serde_json::from_slice(&blob) - .with_context(|| format!("blob {blob_hash:?} is not valid JSON"))?; - cache.insert(blob_hash, hydrated); + cache.insert(blob_hash, decode_blob(blob_hash, encoding, &blob)?); } replace_blob_refs_in_value(value, &cache)?; } RunDumpContents::Text(text) => { - let Some(blob_hash) = parse_blob_ref(text) else { + let Some((blob_hash, encoding)) = parse_blob_ref_encoded(text.trim()) else { continue; }; let hydrated = match cache.entry(blob_hash) { @@ -240,17 +238,13 @@ impl RunDump { let blob = read_blob(blob_hash).await?.with_context(|| { format!("blob {blob_hash:?} is missing from the store") })?; - let hydrated: serde_json::Value = serde_json::from_slice(&blob) - .with_context(|| format!("blob {blob_hash:?} is not valid JSON"))?; - entry.insert(hydrated) + entry.insert(decode_blob(blob_hash, encoding, &blob)?) } }; - *text = hydrated - .as_str() - .with_context(|| { - format!("blob {blob_hash:?} is not a JSON string text log") - })? - .to_string(); + *text = match hydrated { + serde_json::Value::String(text) => text.clone(), + other => serde_json::to_string_pretty(other)?, + }; } RunDumpContents::Bytes(_) => {} } @@ -398,21 +392,40 @@ fn validate_relative_path(kind: &str, value: &str) -> Result { Ok(normalized) } -fn collect_blob_refs_in_value(value: &serde_json::Value, blob_hashes: &mut Vec) { +/// The value a blob's bytes stand for: the text itself for a text blob, the +/// parsed document for a JSON one. +fn decode_blob( + blob_hash: BlobHash, + encoding: BlobRefEncoding, + bytes: &[u8], +) -> Result { + match encoding { + BlobRefEncoding::Text => Ok(serde_json::Value::String( + String::from_utf8_lossy(bytes).into_owned(), + )), + BlobRefEncoding::Json => serde_json::from_slice(bytes) + .with_context(|| format!("blob {blob_hash:?} is not valid JSON")), + } +} + +fn collect_blob_refs_in_value( + value: &serde_json::Value, + blob_refs: &mut Vec<(BlobHash, BlobRefEncoding)>, +) { match value { serde_json::Value::String(current) => { - if let Some(blob_hash) = parse_blob_ref(current) { - blob_hashes.push(blob_hash); + if let Some(blob_ref) = parse_blob_ref_encoded(current) { + blob_refs.push(blob_ref); } } serde_json::Value::Array(items) => { for item in items { - collect_blob_refs_in_value(item, blob_hashes); + collect_blob_refs_in_value(item, blob_refs); } } serde_json::Value::Object(map) => { for item in map.values() { - collect_blob_refs_in_value(item, blob_hashes); + collect_blob_refs_in_value(item, blob_refs); } } serde_json::Value::Null | serde_json::Value::Bool(_) | serde_json::Value::Number(_) => {} @@ -425,7 +438,7 @@ fn replace_blob_refs_in_value( ) -> Result<()> { match value { serde_json::Value::String(current) => { - let Some(blob_hash) = parse_blob_ref(current) else { + let Some((blob_hash, _)) = parse_blob_ref_encoded(current) else { return Ok(()); }; let hydrated = cache.get(&blob_hash).cloned().with_context(|| { @@ -750,10 +763,49 @@ mod tests { } #[test] - fn hydrate_referenced_blobs_fetches_shared_blobs_once() { - let blob = serde_json::to_vec("offloaded response text").unwrap(); + fn hydrate_referenced_blobs_takes_a_plain_reference_as_text() { + // A large string leaves the run context as its own bytes, under a + // plain reference: the bytes are the text, not JSON. + let blob = b"x".repeat(12).to_vec(); let blob_hash = fabro_types::BlobHash::new(&blob); let blob_ref = fabro_types::format_blob_ref(&blob_hash); + let mut dump = RunDump { + entries: vec![ + RunDumpEntry::json("run.json", serde_json::json!({ "output": blob_ref })), + RunDumpEntry::text("stages/001-big@1/output.log", blob_ref.clone()), + ], + stage_ranks: HashMap::new(), + dump_log_index: None, + }; + + executor::block_on(async { + dump.hydrate_referenced_blobs_with_reader(|read_blob_hash| { + let blob = blob.clone(); + Box::pin(async move { + assert_eq!(read_blob_hash, blob_hash); + Ok(Some(bytes::Bytes::from(blob))) + }) + }) + .await + }) + .unwrap(); + + let RunDumpContents::Json(value) = &dump.entries[0].contents else { + panic!("entry should be JSON"); + }; + assert_eq!(value["output"], "xxxxxxxxxxxx"); + let RunDumpContents::Text(text) = &dump.entries[1].contents else { + panic!("entry should be text"); + }; + assert_eq!(text, "xxxxxxxxxxxx"); + } + + #[test] + fn hydrate_referenced_blobs_fetches_shared_blobs_once() { + // A structured value's blob is JSON, and its reference says so. + let blob = serde_json::to_vec("offloaded response text").unwrap(); + let blob_hash = fabro_types::BlobHash::new(&blob); + let blob_ref = format!("{}#json", fabro_types::format_blob_ref(&blob_hash)); let mut dump = RunDump { entries: vec![ RunDumpEntry::json("run.json", serde_json::json!({ "response": blob_ref })),