fix(dump): propagate serde errors from RunDump::from_projection

push_json_entry and push_json_entry_path silently dropped entries via if let Ok(...) on serde_json::to_value, hiding any future Serialize impl failure as missing files. They now return Result, RunDump::from_projection returns Result<Self>, and the three production callers (pipeline/finalize, lifecycle/git init + checkpoint) report failures via emit_metadata_snapshot_failed with MetadataSnapshotFailureKind::Write.

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-05-01 09:32:48 -04:00
parent becc2d0254
commit 2ef87240e7
No known key found for this signature in database
6 changed files with 109 additions and 42 deletions

2
Cargo.lock generated
View file

@ -1840,7 +1840,7 @@ dependencies = [
[[package]]
name = "fabro-dump"
version = "0.219.0-nightly.0"
version = "0.220.0-nightly.0"
dependencies = [
"anyhow",
"bytes",

View file

@ -38,11 +38,10 @@ pub(crate) enum RunDumpContents {
}
impl RunDump {
#[must_use]
pub fn from_projection(state: &RunProjection) -> Self {
pub fn from_projection(state: &RunProjection) -> Result<Self> {
let mut entries = Vec::new();
push_json_entry(&mut entries, "run.json", &SerializableProjection(state));
push_json_entry(&mut entries, "run.json", &SerializableProjection(state))?;
if let Some(graph_source) = state.graph_source.as_ref() {
entries.push(RunDumpEntry::text("graph.fabro", graph_source.clone()));
@ -73,7 +72,7 @@ impl RunDump {
));
}
if let Some(status) = node.status.as_ref() {
push_json_entry_path(&mut entries, &base.join("status.json"), status);
push_json_entry_path(&mut entries, &base.join("status.json"), status)?;
}
if let Some(provider_used) = node.provider_used.as_ref() {
entries.push(RunDumpEntry::json_path(
@ -129,14 +128,14 @@ impl RunDump {
));
}
Self { entries }
Ok(Self { entries })
}
pub fn from_store_state_and_events(
state: &RunProjection,
events: &[EventEnvelope],
) -> Result<Self> {
let mut dump = Self::from_projection(state);
let mut dump = Self::from_projection(state)?;
let mut events_jsonl = Vec::new();
for event in events {
@ -151,7 +150,7 @@ impl RunDump {
&mut dump.entries,
&PathBuf::from("checkpoints").join(format!("{seq:04}.json")),
checkpoint,
);
)?;
}
Ok(dump)
@ -310,22 +309,20 @@ impl RunDumpContents {
}
}
fn push_json_entry<T>(entries: &mut Vec<RunDumpEntry>, path: &str, value: &T)
fn push_json_entry<T>(entries: &mut Vec<RunDumpEntry>, path: &str, value: &T) -> Result<()>
where
T: serde::Serialize,
{
if let Ok(value) = serde_json::to_value(value) {
entries.push(RunDumpEntry::json(path, value));
}
entries.push(RunDumpEntry::json(path, serde_json::to_value(value)?));
Ok(())
}
fn push_json_entry_path<T>(entries: &mut Vec<RunDumpEntry>, path: &Path, value: &T)
fn push_json_entry_path<T>(entries: &mut Vec<RunDumpEntry>, path: &Path, value: &T) -> Result<()>
where
T: serde::Serialize,
{
if let Ok(value) = serde_json::to_value(value) {
entries.push(RunDumpEntry::json_path(path, value));
}
entries.push(RunDumpEntry::json_path(path, serde_json::to_value(value)?));
Ok(())
}
fn path_to_string(path: &Path) -> String {
@ -551,7 +548,7 @@ mod tests {
termination: None,
});
let dump = RunDump::from_projection(&projection);
let dump = RunDump::from_projection(&projection).unwrap();
let paths: Vec<&str> = dump
.entries()
.iter()

View file

@ -551,7 +551,10 @@ mod tests {
.unwrap();
let state = run.state().await.unwrap();
let files = RunDump::from_projection(&state).git_entries().unwrap();
let files = RunDump::from_projection(&state)
.unwrap()
.git_entries()
.unwrap();
let paths: Vec<&str> = files.iter().map(|(path, _)| path.as_str()).collect();
assert!(paths.contains(&"stages/work@2/prompt.md"));
assert!(paths.contains(&"stages/work@2/response.md"));

View file

@ -91,19 +91,36 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
let started = Instant::now();
self.emit_metadata_snapshot_started(phase, &meta_branch, None);
match self.run_store.state().await {
Ok(state) => {
let dump = RunDump::from_projection(&state);
let _ = self
.write_metadata_snapshot(
Ok(state) => match RunDump::from_projection(&state) {
Ok(dump) => {
let _ = self
.write_metadata_snapshot(
phase,
&meta_branch,
started,
&dump,
"init run",
None,
)
.await;
}
Err(err) => {
let message = format!("failed to build run dump for metadata init: {err}");
self.emit_metadata_snapshot_failed(
phase,
&meta_branch,
started,
&dump,
"init run",
MetadataSnapshotFailureKind::Write,
message.clone(),
collect_causes(err.as_ref()),
None,
)
.await;
}
None,
None,
None,
);
self.emit_metadata_warning("checkpoint_metadata_write_failed", message);
}
},
Err(err) => {
let message = format!("failed to load run state for metadata init: {err}");
self.emit_metadata_snapshot_failed(
@ -161,16 +178,41 @@ impl RunLifecycle<WorkflowGraph> for GitLifecycle {
match self.run_store.state().await {
Ok(mut projection) => {
projection.checkpoint = Some(checkpoint);
let dump = RunDump::from_projection(&projection);
self.write_metadata_snapshot(
phase,
&meta_branch,
started,
&dump,
"checkpoint",
Some(&scope),
)
.await
match RunDump::from_projection(&projection) {
Ok(dump) => {
self.write_metadata_snapshot(
phase,
&meta_branch,
started,
&dump,
"checkpoint",
Some(&scope),
)
.await
}
Err(err) => {
let message = format!(
"failed to build run dump for metadata checkpoint: {err}"
);
self.emit_metadata_snapshot_failed(
phase,
&meta_branch,
started,
MetadataSnapshotFailureKind::Write,
message.clone(),
collect_causes(err.as_ref()),
None,
None,
None,
Some(&scope),
);
self.emit_metadata_warning(
"checkpoint_metadata_write_failed",
message,
);
None
}
}
}
Err(err) => {
let message =
@ -624,7 +666,10 @@ mod tests {
let run_store = run_store(fixtures::RUN_1).await;
let handle = RunStoreHandle::local(run_store.clone());
let state = handle.state().await.unwrap();
let expected_entries = RunDump::from_projection(&state).git_entries().unwrap();
let expected_entries = RunDump::from_projection(&state)
.unwrap()
.git_entries()
.unwrap();
let expected_entry_count = expected_entries.len();
let expected_bytes = expected_entries
.iter()
@ -719,7 +764,10 @@ mod tests {
let run_store = run_store(fixtures::RUN_1).await;
let handle = RunStoreHandle::local(run_store.clone());
let state = handle.state().await.unwrap();
let expected_entries = RunDump::from_projection(&state).git_entries().unwrap();
let expected_entries = RunDump::from_projection(&state)
.unwrap()
.git_entries()
.unwrap();
let expected_entry_count = expected_entries.len();
let expected_bytes = expected_entries
.iter()

View file

@ -187,7 +187,26 @@ pub async fn write_finalize_commit(
}
};
projection.conclusion = Some(conclusion.clone());
let dump = RunDump::from_projection(&projection);
let dump = match RunDump::from_projection(&projection) {
Ok(dump) => dump,
Err(err) => {
let message = format!("failed to build run dump for final metadata snapshot: {err}");
emit_metadata_snapshot_failed(
services,
phase,
meta_branch,
started,
MetadataSnapshotFailureKind::Write,
message.clone(),
collect_causes(err.as_ref()),
None,
None,
None,
);
emit_metadata_warning(services, "checkpoint_metadata_write_failed", message);
return;
}
};
let run_id = run_options.run_id.to_string();
let writer = SandboxMetadataWriter::new(
&*services.sandbox,

View file

@ -1339,7 +1339,7 @@ mod tests {
fork_source_ref: None,
in_place: false,
});
let mut dump = fabro_dump::RunDump::from_projection(&projection);
let mut dump = fabro_dump::RunDump::from_projection(&projection).unwrap();
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 = fabro_dump::RunDump::from_projection(&projection);
let dump = fabro_dump::RunDump::from_projection(&projection).unwrap();
let expected_entries = dump.git_entries().unwrap();
let expected_entry_count = expected_entries.len();
let expected_bytes = expected_entries