From af522d1aae03676f74df9e1865b3b092f135e8ae Mon Sep 17 00:00:00 2001 From: Scott Werner Date: Sun, 16 Aug 2026 10:37:32 -0400 Subject: [PATCH] Share the blob cache across dump Json and Text hydration hydrate_referenced_blobs_with_reader kept a per-call blob cache for the Json entries but the Text branch bypassed it, so offloaded stage responses (referenced by both checkpoint values and response.md) were fetched twice per dump. Both branches now hydrate through the shared cache, and a test pins the single-fetch behavior. Co-Authored-By: Claude Fable 5 --- lib/components/fabro-dump/src/lib.rs | 63 +++++++++++++++++++++++++--- 1 file changed, 57 insertions(+), 6 deletions(-) diff --git a/lib/components/fabro-dump/src/lib.rs b/lib/components/fabro-dump/src/lib.rs index 9210a695d..33f3185f6 100644 --- a/lib/components/fabro-dump/src/lib.rs +++ b/lib/components/fabro-dump/src/lib.rs @@ -4,6 +4,7 @@ )] use std::collections::HashMap; +use std::collections::hash_map::Entry; #[expect( clippy::disallowed_types, reason = "in-memory Vec::write_all for jsonl serialization; no filesystem or network I/O" @@ -233,12 +234,23 @@ impl RunDump { let Some(blob_hash) = parse_blob_ref(text) else { continue; }; - let blob = read_blob(blob_hash) - .await? - .with_context(|| format!("blob {blob_hash:?} is missing from the store"))?; - *text = serde_json::from_slice::(&blob).with_context(|| { - format!("blob {blob_hash:?} is not a JSON string text log") - })?; + let hydrated = match cache.entry(blob_hash) { + Entry::Occupied(entry) => entry.into_mut(), + Entry::Vacant(entry) => { + 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) + } + }; + *text = hydrated + .as_str() + .with_context(|| { + format!("blob {blob_hash:?} is not a JSON string text log") + })? + .to_string(); } RunDumpContents::Bytes(_) => {} } @@ -751,4 +763,43 @@ mod tests { }; assert_eq!(value["stdout"], legacy_ref); } + + #[test] + fn hydrate_referenced_blobs_fetches_shared_blobs_once() { + let blob = serde_json::to_vec("offloaded response text").unwrap(); + 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!({ "response": blob_ref })), + RunDumpEntry::text("stages/001-demo@1/response.md", blob_ref.clone()), + ], + stage_ranks: HashMap::new(), + dump_log_index: None, + }; + + let reads = std::cell::Cell::new(0); + executor::block_on(async { + dump.hydrate_referenced_blobs_with_reader(|read_blob_hash| { + reads.set(reads.get() + 1); + let blob = blob.clone(); + Box::pin(async move { + assert_eq!(read_blob_hash, blob_hash); + Ok(Some(bytes::Bytes::from(blob))) + }) + }) + .await + }) + .unwrap(); + + assert_eq!(reads.get(), 1, "shared blob should be fetched once"); + let RunDumpContents::Json(value) = &dump.entries[0].contents else { + panic!("entry should be JSON"); + }; + assert_eq!(value["response"], "offloaded response text"); + let RunDumpContents::Text(text) = &dump.entries[1].contents else { + panic!("entry should be text"); + }; + assert_eq!(text, "offloaded response text"); + } }