Address run spec persistence review findings

This commit is contained in:
Bryan Helmkamp 2026-08-20 20:35:52 -04:00
parent 3421c4f06f
commit 2b095612c8
No known key found for this signature in database
4 changed files with 95 additions and 52 deletions

View file

@ -13,10 +13,11 @@ use std::sync::Arc;
use fabro_config::Storage;
use fabro_graphviz::graph::{AttrValue, Graph};
use fabro_model::{Catalog, ProviderId};
use fabro_store::Database;
use fabro_store::{Database, RunDatabase};
use fabro_template::TemplateContext;
use fabro_types::{
AutomationRef, ForkSourceRef, GitContext, ManifestPath, RunId, RunProvenance, WorkflowSettings,
AutomationRef, ForkSourceRef, GitContext, ManifestPath, RunBlobId, RunId, RunProvenance,
WorkflowSettings,
};
use fabro_util::json::normalize_json_value;
use tokio::task::spawn_blocking;
@ -517,24 +518,17 @@ async fn persist_created_run(
.create_run(&record.run_id)
.await
.map_err(|err| Error::engine_with_source("failed to create run store", err))?;
let manifest_blob = match submitted_manifest_bytes {
Some(bytes) => Some(run_store.write_blob(bytes).await.map_err(store_error)?),
None => None,
};
let definition_blob = match accepted_definition {
Some(definition) => {
let bytes =
serde_json::to_vec(definition).map_err(|err| Error::engine(err.to_string()))?;
Some(run_store.write_blob(&bytes).await.map_err(store_error)?)
}
None => None,
};
// The spec on the run.created event is subject to secret redaction in
// stored copies; the blob keeps the exact bytes execution needs.
let spec_blob = {
let bytes = serde_json::to_vec(record).map_err(|err| Error::engine(err.to_string()))?;
Some(run_store.write_blob(&bytes).await.map_err(store_error)?)
};
let definition_bytes = accepted_definition
.map(serde_json::to_vec)
.transpose()
.map_err(|err| Error::engine_with_source("failed to serialize run definition", err))?;
let spec_bytes = serde_json::to_vec(record)
.map_err(|err| Error::engine_with_source("failed to serialize run spec", err))?;
let (manifest_blob, definition_blob, spec_blob) = tokio::try_join!(
write_optional_blob(&run_store, submitted_manifest_bytes),
write_optional_blob(&run_store, definition_bytes.as_deref()),
async { run_store.write_blob(&spec_bytes).await.map_err(store_error) },
)?;
let title = explicit_title.unwrap_or_else(|| fabro_types::infer_run_title(record.graph.goal()));
let stored = to_run_event_at(
@ -561,7 +555,7 @@ async fn persist_created_run(
automation: record.automation.clone(),
provenance: record.provenance.clone(),
manifest_blob,
spec_blob,
spec_blob: Some(spec_blob),
git: record.git.clone(),
fork_source_ref: record.fork_source_ref.clone(),
retried_from: None,
@ -588,8 +582,22 @@ async fn persist_created_run(
.map_err(store_error)
}
fn store_error(err: impl std::fmt::Display) -> Error {
Error::engine(err.to_string())
async fn write_optional_blob(
run_store: &RunDatabase,
bytes: Option<&[u8]>,
) -> Result<Option<RunBlobId>, Error> {
match bytes {
Some(bytes) => run_store
.write_blob(bytes)
.await
.map(Some)
.map_err(store_error),
None => Ok(None),
}
}
fn store_error(err: impl Into<anyhow::Error>) -> Error {
Error::engine_with_source("run store operation failed", err)
}
/// Parse, transform, and validate `dot_source`.

View file

@ -65,22 +65,23 @@ async fn executable_run_spec(
let bytes = run_store
.read_blob(&blob_id)
.await
.map_err(|err| Error::engine(err.to_string()))?
.map_err(|err| Error::engine_with_anyhow("failed to read run spec blob", err))?
.ok_or_else(|| {
Error::engine(format!(
"run spec blob is missing from the run store: {blob_id}"
))
})?;
let mut spec: RunSpec =
serde_json::from_slice(&bytes).map_err(|err| Error::Parse(err.to_string()))?;
// The event stream stays authoritative for run identity, for provenance
// (a retry rewrites it), for blob ids recorded on events after the spec
// blob was written, and for a graph source the blob does not carry.
let mut spec: RunSpec = serde_json::from_slice(&bytes)
.map_err(|err| Error::engine_with_source("run spec blob was not valid JSON", err))?;
// The event stream stays authoritative for run identity, provenance, and
// blob ids. Prefer the unredacted graph source from the blob, with the
// folded source as a compatibility fallback.
spec.run_id = folded.run_id;
spec.provenance = folded.provenance;
spec.manifest_blob = folded.manifest_blob;
spec.definition_blob = folded.definition_blob;
spec.spec_blob = folded.spec_blob;
spec.fork_source_ref = folded.fork_source_ref;
spec.graph_source = spec.graph_source.or(folded.graph_source);
Ok(spec)
}
@ -191,27 +192,24 @@ mod tests {
}
async fn seeded_store(record: &RunSpec, source: Option<&str>) -> RunDatabase {
seeded_store_with(record, source, true).await
seeded_store_with(record, source, Some(record)).await
}
async fn seeded_store_with(
record: &RunSpec,
source: Option<&str>,
write_spec_blob: bool,
blob_record: Option<&RunSpec>,
) -> RunDatabase {
let store = memory_store();
let run_store = store.create_run(&record.run_id).await.unwrap();
// Mirror the production producer: the unredacted spec rides a blob
// and the redacted event carries its id.
let spec_blob = if write_spec_blob {
Some(
let spec_blob = match blob_record {
Some(blob_record) => Some(
run_store
.write_blob(&serde_json::to_vec(record).unwrap())
.write_blob(&serde_json::to_vec(blob_record).unwrap())
.await
.unwrap(),
)
} else {
None
),
None => None,
};
append_event(&run_store, &record.run_id, &Event::RunCreated {
run_id: record.run_id,
@ -384,7 +382,7 @@ mod tests {
let mut record = sample_record(different_graph());
record.graph = graph;
let run_store = seeded_store_with(&record, Some(&source), false).await;
let run_store = seeded_store_with(&record, Some(&source), None).await;
let loaded = load_from_store(&run_store.clone().into(), &run_dir)
.await
.unwrap();
@ -393,6 +391,32 @@ mod tests {
assert_eq!(loaded.run_spec().spec_blob, None);
}
#[tokio::test]
async fn load_from_store_uses_fork_reference_from_event_fold() {
let temp = tempfile::tempdir().unwrap();
let run_dir = temp.path().join("run");
std::fs::create_dir_all(&run_dir).unwrap();
let (graph, source) = graph_and_source();
let source_record = sample_record(graph.clone());
let mut fork_record = source_record.clone();
fork_record.run_id = fixtures::RUN_7;
fork_record.fork_source_ref = Some(fabro_types::ForkSourceRef {
source_run_id: source_record.run_id,
checkpoint_sha: "checkpoint-sha".to_string(),
});
let run_store = seeded_store_with(&fork_record, Some(&source), Some(&source_record)).await;
let loaded = load_from_store(&run_store.clone().into(), &run_dir)
.await
.unwrap();
assert_eq!(loaded.run_spec().run_id, fork_record.run_id);
assert_eq!(
loaded.run_spec().fork_source_ref,
fork_record.fork_source_ref
);
}
#[test]
fn persist_returns_error_on_io_failure() {
let temp = tempfile::tempdir().unwrap();

View file

@ -81,6 +81,10 @@ pub(super) fn find_entropy_regions(s: &str) -> Vec<Region> {
/// qualify simply measures under the threshold; no length guard is needed.
fn assignment_value_offset(token: &str) -> Option<usize> {
let eq = token.find('=')?;
let value = &token[eq + 1..];
if value.is_empty() || value.starts_with('=') {
return None;
}
let name = &token[..eq];
let mut chars = name.chars();
let first = chars.next()?;
@ -135,6 +139,19 @@ mod tests {
assert_eq!(regions[0].end, input.len());
}
#[test]
fn regions_find_padded_base64_tokens() {
for input in [
"WxFhjC5EAnh30M0JIe0Wa58Xb1BYf8kedTTdKUbbd9Y=",
"AbCdEfGhIjKlMnOpQrStUvWxYz0123456789ABCDEF==",
] {
assert_eq!(find_entropy_regions(input), vec![Region {
start: 0,
end: input.len(),
}]);
}
}
#[test]
fn regions_empty_for_json_escape_sequence() {
// "controller.go\nmodel.go" — the regex could match across the \n boundary

View file

@ -1955,18 +1955,12 @@ pub fn json_snapshot_filters(mut filters: Vec<(String, String)>) -> Vec<(String,
r#""id": "[EVENT_ID]""#.to_string(),
));
filters = json_elapsed_ms_snapshot_filters(filters);
filters.push((
r#""manifest_blob":\s*"[0-9a-f]{64}""#.to_string(),
r#""manifest_blob": "[BLOB_ID]""#.to_string(),
));
filters.push((
r#""definition_blob":\s*"[0-9a-f]{64}""#.to_string(),
r#""definition_blob": "[BLOB_ID]""#.to_string(),
));
filters.push((
r#""spec_blob":\s*"[0-9a-f]{64}""#.to_string(),
r#""spec_blob": "[BLOB_ID]""#.to_string(),
));
for field in ["manifest_blob", "definition_blob", "spec_blob"] {
filters.push((
format!(r#""{field}":\s*"[0-9a-f]{{64}}""#),
format!(r#""{field}": "[BLOB_ID]""#),
));
}
filters.push((
r#""run_dir":\s*"\[STORAGE_DIR\]/scratch/\d{8}-\[ULID\]""#.to_string(),
r#""run_dir": "[RUN_DIR]""#.to_string(),