From 2b095612c8eb315d35f80ef4ac6eb728b22ddd3a Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Thu, 20 Aug 2026 20:35:52 -0400 Subject: [PATCH] Address run spec persistence review findings --- .../fabro-workflow/src/operations/create.rs | 54 +++++++++-------- .../fabro-workflow/src/pipeline/persist.rs | 58 +++++++++++++------ lib/foundation/fabro-redact/src/entropy.rs | 17 ++++++ lib/foundation/fabro-test/src/lib.rs | 18 ++---- 4 files changed, 95 insertions(+), 52 deletions(-) diff --git a/lib/components/fabro-workflow/src/operations/create.rs b/lib/components/fabro-workflow/src/operations/create.rs index 3f23acd5b..a872f0186 100644 --- a/lib/components/fabro-workflow/src/operations/create.rs +++ b/lib/components/fabro-workflow/src/operations/create.rs @@ -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, 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) -> Error { + Error::engine_with_source("run store operation failed", err) } /// Parse, transform, and validate `dot_source`. diff --git a/lib/components/fabro-workflow/src/pipeline/persist.rs b/lib/components/fabro-workflow/src/pipeline/persist.rs index 8c224c700..303241ff5 100644 --- a/lib/components/fabro-workflow/src/pipeline/persist.rs +++ b/lib/components/fabro-workflow/src/pipeline/persist.rs @@ -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(); diff --git a/lib/foundation/fabro-redact/src/entropy.rs b/lib/foundation/fabro-redact/src/entropy.rs index 884aa3713..ec6029c3c 100644 --- a/lib/foundation/fabro-redact/src/entropy.rs +++ b/lib/foundation/fabro-redact/src/entropy.rs @@ -81,6 +81,10 @@ pub(super) fn find_entropy_regions(s: &str) -> Vec { /// qualify simply measures under the threshold; no length guard is needed. fn assignment_value_offset(token: &str) -> Option { 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 diff --git a/lib/foundation/fabro-test/src/lib.rs b/lib/foundation/fabro-test/src/lib.rs index eec273d32..5c53f9ba2 100644 --- a/lib/foundation/fabro-test/src/lib.rs +++ b/lib/foundation/fabro-test/src/lib.rs @@ -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(),