refactor(retro): share run dump hydration

Move RunDump into fabro-dump so CLI export and retro uploads share the same hydrated run layout. Drop the legacy artifact file-ref parser, add best-effort run.log retrieval for retro, and update retro prompts/docs to use events.jsonl and checkpoints.
This commit is contained in:
Bryan Helmkamp 2026-04-30 23:31:12 -04:00
parent 7cb120b96c
commit 93908de46c
No known key found for this signature in database
24 changed files with 225 additions and 356 deletions

17
Cargo.lock generated
View file

@ -1661,6 +1661,7 @@ dependencies = [
"fabro-client",
"fabro-config",
"fabro-devcontainer",
"fabro-dump",
"fabro-github",
"fabro-graphviz",
"fabro-hooks",
@ -1837,6 +1838,20 @@ dependencies = [
"tracing",
]
[[package]]
name = "fabro-dump"
version = "0.219.0-nightly.0"
dependencies = [
"anyhow",
"bytes",
"chrono",
"fabro-store",
"fabro-types",
"futures",
"serde",
"serde_json",
]
[[package]]
name = "fabro-github"
version = "0.219.0-nightly.0"
@ -2067,6 +2082,7 @@ dependencies = [
"anyhow",
"chrono",
"fabro-agent",
"fabro-dump",
"fabro-llm",
"fabro-store",
"fabro-types",
@ -2408,6 +2424,7 @@ dependencies = [
"fabro-config",
"fabro-core",
"fabro-devcontainer",
"fabro-dump",
"fabro-github",
"fabro-graphviz",
"fabro-hooks",

View file

@ -34,7 +34,7 @@ These names are still real, but they are no longer live scratch files by default
- Metadata branch files such as `run.json`, `start.json`, `checkpoint.json`, and `retro.json`
- `fabro dump` exports such as `run.json`, `start.json`, `status.json`, `checkpoint.json`, `conclusion.json`, `retro.json`, `events.jsonl`, and per-node prompt/response/status/stdout/stderr files
- Retro-agent temp uploads named `progress.jsonl`, `checkpoint.json`, `run.json`, and `start.json` inside the retro sandbox
- Retro-agent temp uploads named `events.jsonl`, `run.json`, `graph.fabro`, `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{node_id}@{visit}/...` inside the retro sandbox
## Notes

View file

@ -101,7 +101,7 @@ Retro generation happens in two phases after a run completes:
1. **Derive** — Fabro extracts stage durations from durable run events and builds a retro from the checkpoint data. This is deterministic, fast, and produces the quantitative layer.
2. **Narrate** — An LLM agent session analyzes the run data. The agent receives `progress.jsonl`, `run.json`, `graph.fabro`, and per-stage files under `stages/{node_id}@{visit}/...` inside its sandbox so it can grep and read the event stream, run snapshot, workflow source, and full stage payloads. The narrative fields are merged back into durable retro state.
2. **Narrate** — An LLM agent session analyzes the run data. The agent receives `events.jsonl`, `run.json`, `graph.fabro`, checkpoint snapshots under `checkpoints/{seq:04}.json`, `run.log` when available, and per-stage files under `stages/{node_id}@{visit}/...` inside its sandbox so it can grep and read the event stream, run snapshot, workflow source, checkpoints, logs, and full stage payloads. The narrative fields are merged back into durable retro state.
Both phases run automatically at the end of every CLI run. The API server derives the quantitative layer but does not currently run the narrative agent.

View file

@ -26,6 +26,7 @@ fabro-oauth = { path = "../fabro-oauth" }
fabro-github = { path = "../fabro-github" }
fabro-agent = { path = "../fabro-agent" }
fabro-devcontainer = { path = "../fabro-devcontainer" }
fabro-dump = { path = "../fabro-dump" }
fabro-hooks = { path = "../fabro-hooks" }
fabro-install = { path = "../fabro-install" }
fabro-interview = { path = "../fabro-interview" }

View file

@ -7,9 +7,9 @@ use std::io::ErrorKind;
use std::path::Path;
use anyhow::{Context, Result};
use fabro_dump::RunDump;
use fabro_store::{RunProjection, StageId};
use fabro_types::RunId;
use fabro_workflow::run_dump::RunDump;
use tokio::task::spawn_blocking;
use crate::args::DumpArgs;

View file

@ -5,9 +5,7 @@ use anyhow::{Context as _, Result};
use cli_table::format::{Border, Justify, Separator};
use cli_table::{Cell, CellStruct, Style, Table};
use fabro_api::types;
use fabro_types::{
PullRequestRecord, RunBlobId, RunId, parse_blob_ref, parse_legacy_blob_file_ref,
};
use fabro_types::{PullRequestRecord, RunBlobId, RunId, parse_blob_ref};
use fabro_util::check_report::{CheckDetail, CheckReport, CheckResult, CheckSection, CheckStatus};
use fabro_util::printer::Printer;
use fabro_util::terminal::Styles;
@ -321,7 +319,7 @@ async fn resolve_response_string(
}
fn blob_id_from_response(response: &str) -> Option<RunBlobId> {
parse_blob_ref(response).or_else(|| parse_legacy_blob_file_ref(response))
parse_blob_ref(response)
}
async fn list_artifact_display_entries_with_client(

View file

@ -433,6 +433,15 @@ impl RunStoreBackend for HttpRunStore {
})
.await
}
async fn read_run_log(&self) -> Result<Option<Vec<u8>>> {
self.with_retries("get run logs", || {
let client = self.client.clone_for_reuse();
let run_id = self.run_id;
async move { client.get_run_logs(&run_id).await }
})
.await
}
}
fn set_worker_title(run_id: &RunId, phase: WorkerTitlePhase) {

View file

@ -0,0 +1,25 @@
[package]
name = "fabro-dump"
edition.workspace = true
version.workspace = true
publish = false
license.workspace = true
[lib]
doctest = false
[lints]
workspace = true
[dependencies]
anyhow.workspace = true
bytes.workspace = true
fabro-store = { path = "../fabro-store" }
fabro-types = { path = "../fabro-types" }
futures.workspace = true
serde.workspace = true
serde_json.workspace = true
[dev-dependencies]
chrono = { workspace = true, features = ["serde"] }
fabro-types = { path = "../fabro-types", features = ["test-support"] }

View file

@ -14,9 +14,11 @@ use std::path::{Component, Path, PathBuf};
use anyhow::{Context, Result, bail};
use bytes::Bytes;
use fabro_store::{EventEnvelope, RunProjection, SerializableProjection, StageId};
use fabro_types::{RunBlobId, parse_blob_ref, parse_legacy_blob_file_ref};
use fabro_types::{RunBlobId, parse_blob_ref};
use futures::future::BoxFuture;
pub type BlobReader = Box<dyn FnMut(RunBlobId) -> BoxFuture<'static, Result<Option<Bytes>>> + Send>;
#[derive(Debug, Clone)]
pub struct RunDump {
entries: Vec<RunDumpEntry>,
@ -29,7 +31,7 @@ pub struct RunDumpEntry {
}
#[derive(Debug, Clone)]
pub enum RunDumpContents {
pub(crate) enum RunDumpContents {
Text(String),
Json(serde_json::Value),
Bytes(Vec<u8>),
@ -238,6 +240,14 @@ impl RunDump {
}
impl RunDumpEntry {
pub fn path(&self) -> &str {
&self.path
}
pub fn to_bytes(&self) -> Result<Vec<u8>> {
self.contents.to_bytes()
}
fn text(path: impl Into<String>, contents: String) -> Self {
Self {
path: path.into(),
@ -291,7 +301,7 @@ impl RunDumpEntry {
}
impl RunDumpContents {
fn to_bytes(&self) -> Result<Vec<u8>> {
pub(crate) fn to_bytes(&self) -> Result<Vec<u8>> {
match self {
Self::Text(value) => Ok(value.as_bytes().to_vec()),
Self::Json(value) => Ok(serde_json::to_vec_pretty(value)?),
@ -350,9 +360,7 @@ fn validate_relative_path(kind: &str, value: &str) -> Result<PathBuf> {
fn collect_blob_refs_in_value(value: &serde_json::Value, blob_ids: &mut Vec<RunBlobId>) {
match value {
serde_json::Value::String(current) => {
if let Some(blob_id) =
parse_blob_ref(current).or_else(|| parse_legacy_blob_file_ref(current))
{
if let Some(blob_id) = parse_blob_ref(current) {
blob_ids.push(blob_id);
}
}
@ -376,9 +384,7 @@ fn replace_blob_refs_in_value(
) -> Result<()> {
match value {
serde_json::Value::String(current) => {
let Some(blob_id) =
parse_blob_ref(current).or_else(|| parse_legacy_blob_file_ref(current))
else {
let Some(blob_id) = parse_blob_ref(current) else {
return Ok(());
};
let hydrated = cache
@ -431,9 +437,9 @@ mod tests {
Checkpoint, Conclusion, NodeStatusRecord, RunStatus, SandboxRecord, StageOutcome,
StartRecord, SuccessReason, WorkflowSettings, fixtures,
};
use futures::executor;
use super::RunDump;
use crate::run_dump::RunDumpContents;
use super::{RunDump, RunDumpContents, RunDumpEntry};
fn sample_run_spec() -> RunSpec {
RunSpec {
@ -600,4 +606,34 @@ mod tests {
Some(serde_json::json!({ "provider": "openai" }))
);
}
#[test]
fn hydrate_referenced_blobs_ignores_legacy_artifact_file_refs() {
let blob = serde_json::to_vec("hydrated legacy text").unwrap();
let blob_id = fabro_types::RunBlobId::new(&blob);
let legacy_ref = format!("file:///sandbox/.fabro/artifacts/{blob_id}.json");
let mut dump = RunDump {
entries: vec![RunDumpEntry::json(
"run.json",
serde_json::json!({ "stdout": legacy_ref }),
)],
};
executor::block_on(async {
dump.hydrate_referenced_blobs_with_reader(|read_blob_id| {
let blob = blob.clone();
Box::pin(async move {
assert_eq!(read_blob_id, blob_id);
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["stdout"], legacy_ref);
}
}

View file

@ -16,6 +16,7 @@ workspace = true
anyhow = "1"
chrono = { workspace = true, features = ["serde"] }
fabro-agent = { path = "../fabro-agent" }
fabro-dump = { path = "../fabro-dump" }
fabro-llm = { path = "../fabro-llm" }
fabro-store = { path = "../fabro-store" }
fabro-types = { path = "../fabro-types" }

View file

@ -1,6 +1,4 @@
use std::future::Future;
use std::path::{Path, PathBuf};
use std::pin::Pin;
use std::path::Path;
use std::sync::{Arc, Mutex};
use std::time::Duration;
@ -9,12 +7,11 @@ use fabro_agent::{
AgentProfile, AnthropicProfile, GeminiProfile, OpenAiProfile, Sandbox, Session, SessionEvent,
SessionOptions, Turn, shell_quote,
};
use fabro_dump::{BlobReader, RunDump};
use fabro_llm::client::Client;
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 fabro_store::{EventEnvelope, RunProjection};
use tokio::task::JoinHandle;
use crate::retro::{RetroNarrative, SmoothnessRating};
@ -22,9 +19,11 @@ use crate::retro::{RetroNarrative, SmoothnessRating};
const RETRO_SYSTEM_PROMPT: &str = r"You are a workflow run retrospective analyst. Your job is to analyze a completed workflow run and generate a structured retrospective.
You have access to the run's data files:
- `progress.jsonl` — the full event stream (stage starts/completions, agent tool calls, errors, retries)
- `events.jsonl` — the full event stream (stage starts/completions, agent tool calls, errors, retries)
- `run.json` — serialized run projection with the run spec, checkpoint state, conclusion, retro data, and other metadata
- `graph.fabro` — the workflow source for the run
- `checkpoints/{seq:04}.json` — zero-padded checkpoint snapshots captured during the run
- `run.log` — server/worker log output for the run when available
- `stages/{node_id}@{visit}/...` — per-stage prompt, response, status, diff, stdout/stderr, and tool metadata files
## Your task
@ -124,33 +123,30 @@ pub struct RetroAgentResult {
pub response: String,
}
pub type RetroBlobReader = Arc<
dyn Fn(RunBlobId) -> Pin<Box<dyn Future<Output = anyhow::Result<Option<Vec<u8>>>> + Send>>
+ Send
+ Sync,
>;
#[must_use]
pub fn build_retro_prompt(retro_data_dir: &str) -> String {
format!(
"Analyze the workflow run data at `{retro_data_dir}/` and generate a retrospective. \
The key file is `{retro_data_dir}/progress.jsonl` which contains the full event stream. \
The key file is `{retro_data_dir}/events.jsonl` which contains the full event stream. \
Use `{retro_data_dir}/run.json` for the run-level snapshot, `{retro_data_dir}/graph.fabro` \
for the workflow source, and `{retro_data_dir}/stages/` for full per-stage payloads. \
for the workflow source, `{retro_data_dir}/checkpoints/` for checkpoint snapshots, \
`{retro_data_dir}/run.log` for server logs when present, and `{retro_data_dir}/stages/` \
for full per-stage payloads. \
Use grep to search for interesting signals (failures, retries, errors, approach changes) \
rather than reading the entire file. When done, call the `submit_retro` tool with your analysis."
)
}
/// Run a retro agent session that analyzes workflow run data and produces
/// a structured narrative. The agent explores `progress.jsonl` and other
/// a structured narrative. The agent explores `events.jsonl` and other
/// files via tool access, then calls `submit_retro` with its analysis.
pub async fn run_retro_agent(
sandbox: &Arc<dyn Sandbox>,
state: &RunProjection,
events: &[EventEnvelope],
run_dir: &Path,
blob_reader: Option<RetroBlobReader>,
_run_dir: &Path,
run_log: Option<Vec<u8>>,
blob_reader: Option<BlobReader>,
llm_client: &Client,
provider: Provider,
model: &str,
@ -158,15 +154,7 @@ pub async fn run_retro_agent(
) -> anyhow::Result<RetroAgentResult> {
// Upload data files into sandbox (needed for Daytona; no-op effect for local
// since the agent can also read from the original paths via tools).
upload_data_files(
sandbox,
state,
events,
run_dir,
RETRO_DATA_DIR,
blob_reader.as_ref(),
)
.await?;
upload_data_files(sandbox, state, events, RETRO_DATA_DIR, run_log, blob_reader).await?;
// Build provider profile with the submit_retro tool
let captured: Arc<Mutex<Option<RetroNarrative>>> = Arc::new(Mutex::new(None));
@ -317,240 +305,31 @@ async fn upload_data_files(
sandbox: &Arc<dyn Sandbox>,
state: &RunProjection,
events: &[EventEnvelope],
_run_dir: &Path,
target_dir: &str,
blob_reader: Option<&RetroBlobReader>,
run_log: Option<Vec<u8>>,
blob_reader: Option<BlobReader>,
) -> anyhow::Result<()> {
let progress_content = (!events.is_empty()).then(|| {
let mut buf = String::new();
for env in events {
if let Ok(line) = serde_json::to_string(&env.event) {
buf.push_str(&line);
buf.push('\n');
}
}
buf
});
upload_file(
sandbox,
target_dir,
Path::new("progress.jsonl"),
progress_content,
)
.await?;
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,
Path::new("run.json"),
Some(run_content),
)
.await?;
upload_file(
sandbox,
target_dir,
Path::new("graph.fabro"),
state.graph_source.clone(),
)
.await?;
let mut stage_ids: Vec<_> = state
.iter_nodes()
.map(|(stage_id, _)| stage_id.clone())
.collect();
stage_ids.sort();
for stage_id in stage_ids {
let Some(node) = retro_state.node(&stage_id) else {
continue;
};
let base = PathBuf::from("stages").join(stage_id.to_string());
upload_file(
sandbox,
target_dir,
&base.join("prompt.md"),
node.prompt.clone(),
)
.await?;
upload_file(
sandbox,
target_dir,
&base.join("response.md"),
node.response.clone(),
)
.await?;
upload_json_file(
sandbox,
target_dir,
&base.join("status.json"),
node.status.as_ref(),
)
.await?;
upload_json_file(
sandbox,
target_dir,
&base.join("provider_used.json"),
node.provider_used.as_ref(),
)
.await?;
upload_file(
sandbox,
target_dir,
&base.join("diff.patch"),
node.diff.clone(),
)
.await?;
upload_json_file(
sandbox,
target_dir,
&base.join("script_invocation.json"),
node.script_invocation.as_ref(),
)
.await?;
upload_json_file(
sandbox,
target_dir,
&base.join("script_timing.json"),
node.script_timing.as_ref(),
)
.await?;
upload_json_file(
sandbox,
target_dir,
&base.join("parallel_results.json"),
node.parallel_results.as_ref(),
)
.await?;
upload_file(
sandbox,
target_dir,
&base.join("stdout.log"),
resolve_text_file_content(node.stdout.clone(), blob_reader).await?,
)
.await?;
upload_file(
sandbox,
target_dir,
&base.join("stderr.log"),
resolve_text_file_content(node.stderr.clone(), blob_reader).await?,
)
.await?;
let mut dump = RunDump::from_store_state_and_events(state, events)?;
if let Some(log) = run_log {
dump.add_file_bytes("run.log", log);
}
if let Some(reader) = blob_reader {
dump.hydrate_referenced_blobs_with_reader(reader).await?;
}
for entry in dump.entries() {
let remote_path = format!("{target_dir}/{}", entry.path());
ensure_remote_dir(sandbox, Path::new(&remote_path)).await?;
let bytes = entry.to_bytes()?;
let text = String::from_utf8(bytes)
.map_err(|err| anyhow::anyhow!("non-UTF8 retro entry {}: {err}", entry.path()))?;
sandbox
.write_file(&remote_path, &text)
.await
.map_err(|err| anyhow::anyhow!("Failed to upload {}: {err}", entry.path()))?;
}
Ok(())
}
async fn hydrate_retro_projection(
state: &RunProjection,
blob_reader: Option<&RetroBlobReader>,
) -> anyhow::Result<RunProjection> {
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<Option<Value>> {
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<String>,
blob_reader: Option<&RetroBlobReader>,
) -> anyhow::Result<Option<String>> {
let Some(content) = content else {
return Ok(None);
};
let Some(blob_id) = parse_blob_ref(&content) else {
return Ok(Some(content));
};
let Some(blob_reader) = blob_reader else {
return Ok(Some(content));
};
let bytes = blob_reader(blob_id)
.await?
.ok_or_else(|| anyhow::anyhow!("text blob missing: {blob_id}"))?;
let text = serde_json::from_slice::<String>(&bytes)
.map_err(|err| anyhow::anyhow!("text blob was not a JSON string: {err}"))?;
Ok(Some(text))
}
async fn upload_file(
sandbox: &Arc<dyn Sandbox>,
target_dir: &str,
relative: &Path,
content: Option<String>,
) -> anyhow::Result<()> {
let Some(content) = content else {
return Ok(());
};
let path = Path::new(target_dir).join(relative);
let remote_path = path.to_string_lossy().into_owned();
ensure_remote_dir(sandbox, &path).await?;
sandbox
.write_file(&remote_path, &content)
.await
.map_err(|e| anyhow::anyhow!("Failed to upload {}: {e}", relative.display()))?;
Ok(())
}
async fn upload_json_file<T>(
sandbox: &Arc<dyn Sandbox>,
target_dir: &str,
relative: &Path,
value: Option<&T>,
) -> anyhow::Result<()>
where
T: serde::Serialize,
{
let content = value.map(serde_json::to_string_pretty).transpose()?;
upload_file(sandbox, target_dir, relative, content).await
}
async fn ensure_remote_dir(sandbox: &Arc<dyn Sandbox>, path: &Path) -> anyhow::Result<()> {
let parent = path
.parent()
@ -682,8 +461,8 @@ mod tests {
&sandbox,
&state,
&[],
output_dir.path(),
&target_dir_str,
Some(b"server log\n".to_vec()),
None,
)
.await
@ -723,13 +502,25 @@ mod tests {
.expect("stdout file should exist"),
"stdout"
);
assert_eq!(
fs::read_to_string(target_dir.join("events.jsonl"))
.await
.expect("events file should exist"),
""
);
assert_eq!(
fs::read_to_string(target_dir.join("run.log"))
.await
.expect("run.log should exist"),
"server log\n"
);
assert!(
target_dir.join("stages/build@2/status.json").exists(),
"status file should exist"
);
assert!(
!target_dir.join("progress.jsonl").exists(),
"progress file should be omitted when there are no events"
"legacy progress file should not be emitted"
);
}
@ -767,30 +558,23 @@ mod tests {
..NodeState::default()
});
let reader: RetroBlobReader = Arc::new(move |blob_id| {
let reader: BlobReader = Box::new(move |blob_id| {
let stdout_blob = stdout_blob.clone();
let stderr_blob = stderr_blob.clone();
Box::pin(async move {
if blob_id == stdout_id {
Ok(Some(stdout_blob))
Ok(Some(stdout_blob.into()))
} else if blob_id == stderr_id {
Ok(Some(stderr_blob))
Ok(Some(stderr_blob.into()))
} else {
Ok(None)
}
})
});
upload_data_files(
&sandbox,
&state,
&[],
output_dir.path(),
&target_dir_str,
Some(&reader),
)
.await
.expect("retro files should upload");
upload_data_files(&sandbox, &state, &[], &target_dir_str, None, Some(reader))
.await
.expect("retro files should upload");
assert_eq!(
fs::read_to_string(target_dir.join("stages/build@1/stdout.log"))
@ -820,14 +604,8 @@ mod tests {
.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)
);
assert_eq!(script_invocation["stdout"], "resolved stdout");
assert_eq!(script_invocation["stderr"], "resolved stderr");
let run_json: serde_json::Value = serde_json::from_str(
&fs::read_to_string(target_dir.join("run.json"))
@ -845,7 +623,9 @@ mod tests {
);
assert_eq!(
run_json["nodes"]["build@1"]["script_invocation"]["stdout"],
fabro_types::format_blob_ref(&stdout_id)
"resolved stdout"
);
assert!(run_json["nodes"]["build@1"]["stdout"].is_null());
assert!(run_json["nodes"]["build@1"]["stderr"].is_null());
}
}

View file

@ -14,18 +14,6 @@ pub fn parse_blob_ref(value: &str) -> Option<RunBlobId> {
value.strip_prefix(BLOB_REF_PREFIX)?.parse().ok()
}
#[must_use]
pub fn parse_legacy_blob_file_ref(value: &str) -> Option<RunBlobId> {
let path = value.strip_prefix("file://")?;
let blob_id = parse_blob_file_name(path)?;
if has_path_suffix(path, &[".fabro", "artifacts"]) {
Some(blob_id)
} else {
None
}
}
#[must_use]
pub fn parse_managed_blob_file_ref(value: &str) -> Option<RunBlobId> {
let path = value.strip_prefix("file://")?;
@ -57,9 +45,7 @@ fn has_path_suffix(path: &str, suffix: &[&str]) -> bool {
#[cfg(test)]
mod tests {
use super::{
format_blob_ref, parse_blob_ref, parse_legacy_blob_file_ref, parse_managed_blob_file_ref,
};
use super::{format_blob_ref, parse_blob_ref, parse_managed_blob_file_ref};
use crate::RunBlobId;
#[test]
@ -70,14 +56,6 @@ mod tests {
assert_eq!(parse_blob_ref(&formatted), Some(blob_id));
}
#[test]
fn legacy_remote_blob_file_ref_is_recognized() {
let blob_id = RunBlobId::new(b"hello");
let value = format!("file:///sandbox/.fabro/artifacts/{blob_id}.json");
assert_eq!(parse_legacy_blob_file_ref(&value), Some(blob_id));
}
#[test]
fn managed_local_blob_file_ref_is_recognized() {
let blob_id = RunBlobId::new(b"hello");
@ -96,7 +74,6 @@ mod tests {
#[test]
fn ordinary_file_refs_are_not_treated_as_blob_refs() {
assert_eq!(parse_legacy_blob_file_ref("file:///tmp/report.json"), None);
assert_eq!(parse_managed_blob_file_ref("file:///tmp/report.json"), None);
}
}

View file

@ -39,9 +39,7 @@ pub use billing::{
ModelBillingFacts, ModelBillingInput, ModelPricing, ModelPricingPolicy, ModelRef, ModelUsage,
OpenAiBillingFacts, OpenAiModelPricing, PricePerMTok, Speed, TokenCounts, UsdMicros,
};
pub use blob_ref::{
format_blob_ref, parse_blob_ref, parse_legacy_blob_file_ref, parse_managed_blob_file_ref,
};
pub use blob_ref::{format_blob_ref, parse_blob_ref, parse_managed_blob_file_ref};
pub use checkpoint::Checkpoint;
pub use command_output::{CommandOutputStream, CommandTermination};
pub use conclusion::{Conclusion, StageSummary};

View file

@ -25,6 +25,7 @@ fabro-graphviz = { path = "../fabro-graphviz" }
fabro-hooks = { path = "../fabro-hooks" }
fabro-validate = { path = "../fabro-validate" }
fabro-devcontainer = { path = "../fabro-devcontainer" }
fabro-dump = { path = "../fabro-dump" }
fabro-sandbox = { path = "../fabro-sandbox", features = ["daytona"] }
fabro-mcp = { path = "../fabro-mcp" }
fabro-github = { path = "../fabro-github" }

View file

@ -3,10 +3,7 @@ use std::path::{Path, PathBuf};
use fabro_agent::Sandbox;
use fabro_config::RunScratch;
use fabro_types::{
RunBlobId, format_blob_ref, parse_blob_ref, parse_legacy_blob_file_ref,
parse_managed_blob_file_ref,
};
use fabro_types::{RunBlobId, format_blob_ref, parse_blob_ref, parse_managed_blob_file_ref};
use futures::future::BoxFuture;
use serde_json::Value;
use tokio::fs;
@ -234,9 +231,7 @@ pub async fn sync_artifacts_to_env(
fn normalize_durable_value(value: &mut Value) {
match value {
Value::String(current) => {
if let Some(blob_id) =
parse_legacy_blob_file_ref(current).or_else(|| parse_managed_blob_file_ref(current))
{
if let Some(blob_id) = parse_managed_blob_file_ref(current) {
*current = format_blob_ref(&blob_id);
}
}
@ -283,9 +278,7 @@ fn resolve_execution_value<'a>(
Some(context::keys::COMMAND_OUTPUT | context::keys::COMMAND_STDERR)
) {
*current = resolve_text_or_blob_ref_str(current, run_store).await?;
} else if let Some(blob_id) =
parse_blob_ref(current).or_else(|| parse_legacy_blob_file_ref(current))
{
} else if let Some(blob_id) = parse_blob_ref(current) {
*current = materialize_blob_ref(&blob_id, run_store, env, run_dir).await?;
} else if current.starts_with(ARTIFACT_POINTER_PREFIX)
&& parse_managed_blob_file_ref(current).is_none()
@ -525,8 +518,8 @@ mod tests {
}
#[test]
fn normalize_checkpoint_for_resume_converts_legacy_blob_file_refs_and_drops_preamble() {
let blob_id = fabro_types::RunBlobId::new(b"legacy");
fn normalize_checkpoint_for_resume_converts_managed_blob_file_refs_and_drops_preamble() {
let blob_id = fabro_types::RunBlobId::new(b"managed");
let mut checkpoint = crate::records::Checkpoint {
timestamp: chrono::Utc::now(),
current_node: "work".to_string(),
@ -539,7 +532,7 @@ mod tests {
),
(
"response.work".to_string(),
serde_json::json!(format!("file:///sandbox/.fabro/artifacts/{blob_id}.json")),
serde_json::json!(format!("file:///sandbox/.fabro/blobs/{blob_id}.json")),
),
]),
node_outcomes: HashMap::from([(
@ -547,9 +540,7 @@ mod tests {
crate::outcome::Outcome {
context_updates: HashMap::from([(
"response.work".to_string(),
serde_json::json!(format!(
"file:///sandbox/.fabro/artifacts/{blob_id}.json"
)),
serde_json::json!(format!("file:///sandbox/.fabro/blobs/{blob_id}.json")),
)]),
..crate::outcome::Outcome::success()
},

View file

@ -342,12 +342,12 @@ mod tests {
use std::sync::Arc;
use std::time::Duration;
use fabro_dump::RunDump;
use fabro_store::Database;
use fabro_types::{CommandTermination, fixtures};
use object_store::memory::InMemory;
use super::*;
use crate::run_dump::RunDump;
/// Create a temporary git repo with an initial commit.
#[expect(

View file

@ -283,6 +283,10 @@ mod tests {
async fn read_blob(&self, id: &fabro_types::RunBlobId) -> anyhow::Result<Option<Bytes>> {
Ok(self.blobs.lock().await.get(id).cloned())
}
async fn read_run_log(&self) -> anyhow::Result<Option<Vec<u8>>> {
Ok(None)
}
}
fn make_services() -> EngineServices {

View file

@ -140,7 +140,6 @@ pub mod records;
mod retry;
pub mod run_control;
pub(crate) mod run_dir;
pub mod run_dump;
pub mod run_lookup;
pub use error::{Error, FailureCategory, FailureSignature, FailureSignatureExt, Result};

View file

@ -7,6 +7,7 @@ use fabro_core::graph::NodeSpec;
use fabro_core::lifecycle::RunLifecycle;
use fabro_core::outcome::NodeResult;
use fabro_core::state::ExecutionState;
use fabro_dump::RunDump;
use fabro_types::RunId;
use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase};
use fabro_util::error::collect_causes;
@ -17,7 +18,6 @@ use crate::event::{Emitter, Event, RunNoticeLevel, StageScope};
use crate::graph::{WorkflowGraph, WorkflowNode};
use crate::lifecycle::event::stage_scope_for;
use crate::outcome::BilledModelUsage;
use crate::run_dump::RunDump;
use crate::run_options::RunOptions;
use crate::runtime_store::RunStoreHandle;
use crate::sandbox_git::{checked_git_checkpoint, git_diff};
@ -965,5 +965,9 @@ mod tests {
async fn read_blob(&self, _id: &RunBlobId) -> Result<Option<Bytes>> {
Ok(None)
}
async fn read_run_log(&self) -> Result<Option<Vec<u8>>> {
Ok(None)
}
}
}

View file

@ -1,5 +1,6 @@
use std::time::Instant;
use fabro_dump::RunDump;
use fabro_hooks::{HookContext, HookEvent};
use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase};
use fabro_types::{BilledTokenCounts, EventBody};
@ -11,7 +12,6 @@ use crate::error::Error;
use crate::event::{Event, RunNoticeLevel};
use crate::outcome::{Outcome, OutcomeExt, StageOutcome};
use crate::records::{Checkpoint, Conclusion, StageSummary};
use crate::run_dump::RunDump;
use crate::run_options::RunOptions;
use crate::run_status::{FailureReason, RunStatus, SuccessReason};
use crate::runtime_store::RunStoreHandle;
@ -901,5 +901,9 @@ mod tests {
async fn read_blob(&self, _id: &RunBlobId) -> Result<Option<Bytes>> {
Ok(None)
}
async fn read_run_log(&self) -> Result<Option<Vec<u8>>> {
Ok(None)
}
}
}

View file

@ -1,10 +1,11 @@
use std::sync::Arc;
use fabro_agent::SessionEvent;
use fabro_dump::BlobReader;
use fabro_llm::client::Client;
use fabro_retro::retro::{Retro, derive_retro};
use fabro_retro::retro_agent::{
RETRO_DATA_DIR, RetroBlobReader, build_retro_prompt, dry_run_narrative, run_retro_agent,
RETRO_DATA_DIR, build_retro_prompt, dry_run_narrative, run_retro_agent,
};
use super::types::{Executed, RetroOptions, Retroed};
@ -82,20 +83,26 @@ pub async fn run_retro(options: &RetroOptions, dry_run: bool) -> Option<Retro> {
}
});
let run_store = services.run_store.clone();
let blob_reader: RetroBlobReader = Arc::new(move |blob_id| {
let run_log = match services.run_store.read_run_log().await {
Ok(log) => log,
Err(err) => {
tracing::warn!(
error = %err,
"failed to fetch run.log for retro; continuing without it"
);
None
}
};
let blob_reader: BlobReader = Box::new(move |blob_id| {
let run_store = run_store.clone();
Box::pin(async move {
Ok(run_store
.read_blob(&blob_id)
.await?
.map(|bytes| bytes.to_vec()))
})
Box::pin(async move { run_store.read_blob(&blob_id).await })
});
run_retro_agent(
&services.sandbox,
&state,
&events,
&options.run_dir,
run_log,
Some(blob_reader),
&client,
services.provider,
@ -412,7 +419,7 @@ mod tests {
assert!(
retro_started_properties["prompt"]
.as_str()
.is_some_and(|prompt| prompt.contains("/tmp/retro_data/progress.jsonl"))
.is_some_and(|prompt| prompt.contains("/tmp/retro_data/events.jsonl"))
);
let retro_completed = seen

View file

@ -15,6 +15,7 @@ pub trait RunStoreBackend: Send + Sync {
async fn append_run_event(&self, event: &RunEvent) -> Result<()>;
async fn write_blob(&self, data: &[u8]) -> Result<RunBlobId>;
async fn read_blob(&self, id: &RunBlobId) -> Result<Option<Bytes>>;
async fn read_run_log(&self) -> Result<Option<Vec<u8>>>;
}
#[derive(Clone)]
@ -52,6 +53,10 @@ impl RunStoreHandle {
pub async fn read_blob(&self, id: &RunBlobId) -> Result<Option<Bytes>> {
self.backend.read_blob(id).await
}
pub async fn read_run_log(&self) -> Result<Option<Vec<u8>>> {
self.backend.read_run_log().await
}
}
impl From<RunDatabase> for RunStoreHandle {
@ -99,6 +104,10 @@ impl RunStoreBackend for LocalRunStoreBackend {
.await
.map_err(anyhow::Error::from)
}
async fn read_run_log(&self) -> Result<Option<Vec<u8>>> {
Ok(None)
}
}
#[cfg(test)]
@ -208,4 +217,12 @@ mod tests {
assert_eq!(events.len(), 1);
assert_eq!(blob.as_ref(), br#"{"ok":true}"#);
}
#[tokio::test]
async fn local_handle_returns_no_run_log() {
let run_store = test_run_store().await;
let handle = RunStoreHandle::local(run_store);
assert_eq!(handle.read_run_log().await.unwrap(), None);
}
}

View file

@ -1339,7 +1339,7 @@ mod tests {
fork_source_ref: None,
in_place: false,
});
let mut dump = crate::run_dump::RunDump::from_projection(&projection);
let mut dump = fabro_dump::RunDump::from_projection(&projection);
dump.add_file_bytes("binary/payload.bin", vec![0, 159, 146, 150]);
dump.add_file_bytes("path with spaces.txt", b"quoted path\n".to_vec());
@ -1501,7 +1501,7 @@ mod tests {
fork_source_ref: None,
in_place: false,
});
let dump = crate::run_dump::RunDump::from_projection(&projection);
let dump = fabro_dump::RunDump::from_projection(&projection);
let expected_entries = dump.git_entries().unwrap();
let expected_entry_count = expected_entries.len();
let expected_bytes = expected_entries

View file

@ -3,12 +3,12 @@ use std::fmt::Write as _;
use std::sync::atomic::{AtomicBool, Ordering};
use fabro_agent::Sandbox;
use fabro_dump::RunDump;
use fabro_sandbox::shell_quote;
use tokio::fs;
use tokio::sync::OnceCell;
use crate::git::{GitAuthor, META_BRANCH_PREFIX};
use crate::run_dump::RunDump;
use crate::sandbox_git::GIT_REMOTE;
#[derive(Debug, thiserror::Error)]