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 <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-09-18 16:00:11 -04:00
parent a5d4e96bbf
commit ae70aa6ccf
No known key found for this signature in database
7 changed files with 302 additions and 37 deletions

View file

@ -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<String> {
)
}
/// The patch as text: the projection holds a large patch as a
/// `blob://sha256/<hex>` reference into the run's blob table.
async fn resolve_patch_text(client: &Client, run_id: &RunId, patch: String) -> Result<String> {
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")

View file

@ -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<Option<String>> {
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(

View file

@ -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<serde_json::Value>,
pub checkpoint: Option<serde_json::Value>,
pub sandbox: Option<serde_json::Value>,
/// The stages by their id, with what each produced. A large output or
/// response is a `blob://sha256/<hex>` reference into the run's blob
/// table, as the projection holds it.
#[serde(skip_serializing_if = "BTreeMap::is_empty")]
pub stages: BTreeMap<String, InspectStage>,
}
#[derive(Debug, Serialize)]
pub(crate) struct InspectStage {
pub handler: Option<StageHandler>,
pub state: StageState,
#[serde(skip_serializing_if = "Option::is_none")]
pub output: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub response: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub diff: Option<String>,
}
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,
}
}

View file

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

View file

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

View file

@ -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)

View file

@ -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<dyn FnMut(BlobHash) -> BoxFuture<'static, Result<Option<Bytes>>> + 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<PathBuf> {
Ok(normalized)
}
fn collect_blob_refs_in_value(value: &serde_json::Value, blob_hashes: &mut Vec<BlobHash>) {
/// 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<serde_json::Value> {
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 })),