diff --git a/docs/public/api-reference/fabro-api.yaml b/docs/public/api-reference/fabro-api.yaml index 8d1c8fa14..ba092f6a6 100644 --- a/docs/public/api-reference/fabro-api.yaml +++ b/docs/public/api-reference/fabro-api.yaml @@ -11441,6 +11441,8 @@ components: type: ["string", "null"] definition_blob: type: ["string", "null"] + spec_blob: + type: ["string", "null"] git: oneOf: - $ref: "#/components/schemas/GitContext" diff --git a/lib/apps/fabro-cli/src/commands/run/attach.rs b/lib/apps/fabro-cli/src/commands/run/attach.rs index cd356ffa0..7c1cf93b7 100644 --- a/lib/apps/fabro-cli/src/commands/run/attach.rs +++ b/lib/apps/fabro-cli/src/commands/run/attach.rs @@ -849,6 +849,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }; diff --git a/lib/apps/fabro-cli/tests/it/cmd/attach.rs b/lib/apps/fabro-cli/tests/it/cmd/attach.rs index 15709012d..9f388c14b 100644 --- a/lib/apps/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/apps/fabro-cli/tests/it/cmd/attach.rs @@ -1012,6 +1012,7 @@ fn attach_json_errors_without_prompting_for_human_input() { } }, "source_directory": "[TEMP_DIR]", + "spec_blob": "[BLOB_HASH]", "title": "Wait for approval", "web_url": "http://localhost:3000/runs/[ULID]", "workflow_slug": "human-gate", diff --git a/lib/apps/fabro-cli/tests/it/support/mod.rs b/lib/apps/fabro-cli/tests/it/support/mod.rs index 2f3557044..57e4a98e5 100644 --- a/lib/apps/fabro-cli/tests/it/support/mod.rs +++ b/lib/apps/fabro-cli/tests/it/support/mod.rs @@ -53,6 +53,7 @@ pub(crate) fn run_projection_json(run_id: &str, status: &serde_json::Value) -> s provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }; diff --git a/lib/apps/fabro-server/src/run_files.rs b/lib/apps/fabro-server/src/run_files.rs index 920460f79..2343933b6 100644 --- a/lib/apps/fabro-server/src/run_files.rs +++ b/lib/apps/fabro-server/src/run_files.rs @@ -2387,6 +2387,7 @@ index 1111111..2222222 160000 provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, diff --git a/lib/apps/fabro-server/src/server/handler/events.rs b/lib/apps/fabro-server/src/server/handler/events.rs index 1bfb7f889..5fd655f1d 100644 --- a/lib/apps/fabro-server/src/server/handler/events.rs +++ b/lib/apps/fabro-server/src/server/handler/events.rs @@ -627,6 +627,7 @@ mod stage_events_tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/apps/fabro-server/src/server/handler/pair.rs b/lib/apps/fabro-server/src/server/handler/pair.rs index ec31d2e54..6e244bcac 100644 --- a/lib/apps/fabro-server/src/server/handler/pair.rs +++ b/lib/apps/fabro-server/src/server/handler/pair.rs @@ -1027,6 +1027,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/apps/fabro-server/src/server/handler/sessions.rs b/lib/apps/fabro-server/src/server/handler/sessions.rs index c2ee94b32..e9bacdfff 100644 --- a/lib/apps/fabro-server/src/server/handler/sessions.rs +++ b/lib/apps/fabro-server/src/server/handler/sessions.rs @@ -1923,6 +1923,7 @@ reasoning = false provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }; diff --git a/lib/apps/fabro-server/src/server/tests.rs b/lib/apps/fabro-server/src/server/tests.rs index c6562fc3f..b63923e79 100644 --- a/lib/apps/fabro-server/src/server/tests.rs +++ b/lib/apps/fabro-server/src/server/tests.rs @@ -4650,6 +4650,7 @@ async fn append_default_run_created(run_store: &fabro_store::RunDatabase, run_id automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -4701,6 +4702,7 @@ async fn create_slack_notification_run( automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -5774,6 +5776,7 @@ async fn list_run_stages_distinguishes_visits() { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -5910,6 +5913,7 @@ async fn list_run_stages_exposes_execution_identity_for_resumed_stage() { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -7097,6 +7101,7 @@ async fn create_completed_run_ready_for_pull_request( provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }; @@ -7113,6 +7118,7 @@ async fn create_completed_run_ready_for_pull_request( automation: None, provenance: run_spec.provenance.clone(), manifest_blob: None, + spec_blob: None, git, fork_source_ref: None, retried_from: None, @@ -14085,6 +14091,7 @@ async fn create_preserved_local_sandbox_run(state: &Arc, run_id: RunId automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -14834,6 +14841,7 @@ async fn delete_run_retry_after_missing_provider_resource_removes_metadata() { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/apps/fabro-server/tests/it/api/run_files.rs b/lib/apps/fabro-server/tests/it/api/run_files.rs index 6d286cbab..307c4a574 100644 --- a/lib/apps/fabro-server/tests/it/api/run_files.rs +++ b/lib/apps/fabro-server/tests/it/api/run_files.rs @@ -68,6 +68,7 @@ async fn append_completed_run_with_final_patch( automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-store/src/run_state.rs b/lib/components/fabro-store/src/run_state.rs index 4e2c35621..0c59af57a 100644 --- a/lib/components/fabro-store/src/run_state.rs +++ b/lib/components/fabro-store/src/run_state.rs @@ -1046,6 +1046,7 @@ fn projection_from_created(event: &EventEnvelope) -> Result { provenance: props.provenance.clone(), manifest_blob: props.manifest_blob, definition_blob: None, + spec_blob: props.spec_blob, git: props.git.clone(), fork_source_ref: props.fork_source_ref.clone(), }; diff --git a/lib/components/fabro-store/src/run_summary_store.rs b/lib/components/fabro-store/src/run_summary_store.rs index 5db1a49b1..dcbec847e 100644 --- a/lib/components/fabro-store/src/run_summary_store.rs +++ b/lib/components/fabro-store/src/run_summary_store.rs @@ -601,6 +601,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, diff --git a/lib/components/fabro-store/src/slate/mod.rs b/lib/components/fabro-store/src/slate/mod.rs index efed28796..6158e18a4 100644 --- a/lib/components/fabro-store/src/slate/mod.rs +++ b/lib/components/fabro-store/src/slate/mod.rs @@ -593,6 +593,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: Some(fabro_types::GitContext { origin_url: "https://github.com/fabro-sh/fabro".to_string(), branch: "main".to_string(), diff --git a/lib/components/fabro-workflow/src/event/convert.rs b/lib/components/fabro-workflow/src/event/convert.rs index 6403d8933..ef255d5b8 100644 --- a/lib/components/fabro-workflow/src/event/convert.rs +++ b/lib/components/fabro-workflow/src/event/convert.rs @@ -36,6 +36,7 @@ fn event_body_from_event(event: &Event) -> EventBody { automation, provenance, manifest_blob, + spec_blob, git, fork_source_ref, retried_from, @@ -54,6 +55,7 @@ fn event_body_from_event(event: &Event) -> EventBody { automation: automation.clone(), provenance: provenance.clone(), manifest_blob: *manifest_blob, + spec_blob: *spec_blob, git: git.clone(), fork_source_ref: fork_source_ref.clone(), retried_from: *retried_from, @@ -2669,6 +2671,7 @@ mod tests { automation: Some(automation.clone()), provenance, manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/event/events.rs b/lib/components/fabro-workflow/src/event/events.rs index a38cefbd0..1ee3faf9f 100644 --- a/lib/components/fabro-workflow/src/event/events.rs +++ b/lib/components/fabro-workflow/src/event/events.rs @@ -41,6 +41,8 @@ pub enum Event { #[serde(default, skip_serializing_if = "Option::is_none")] manifest_blob: Option, #[serde(default, skip_serializing_if = "Option::is_none")] + spec_blob: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] git: Option, #[serde(default, skip_serializing_if = "Option::is_none")] fork_source_ref: Option, diff --git a/lib/components/fabro-workflow/src/event/sink.rs b/lib/components/fabro-workflow/src/event/sink.rs index f6abb1724..150c69f72 100644 --- a/lib/components/fabro-workflow/src/event/sink.rs +++ b/lib/components/fabro-workflow/src/event/sink.rs @@ -290,6 +290,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/git.rs b/lib/components/fabro-workflow/src/git.rs index de144b6a1..e7ca4feda 100644 --- a/lib/components/fabro-workflow/src/git.rs +++ b/lib/components/fabro-workflow/src/git.rs @@ -364,6 +364,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/handler/agent.rs b/lib/components/fabro-workflow/src/handler/agent.rs index 48b234a8b..99d20431b 100644 --- a/lib/components/fabro-workflow/src/handler/agent.rs +++ b/lib/components/fabro-workflow/src/handler/agent.rs @@ -501,6 +501,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/handler/command.rs b/lib/components/fabro-workflow/src/handler/command.rs index 82593c59b..db9682a1c 100644 --- a/lib/components/fabro-workflow/src/handler/command.rs +++ b/lib/components/fabro-workflow/src/handler/command.rs @@ -374,6 +374,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, @@ -474,6 +475,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/handler/parallel.rs b/lib/components/fabro-workflow/src/handler/parallel.rs index 9da654d1d..82773470b 100644 --- a/lib/components/fabro-workflow/src/handler/parallel.rs +++ b/lib/components/fabro-workflow/src/handler/parallel.rs @@ -956,6 +956,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/handler/prompt.rs b/lib/components/fabro-workflow/src/handler/prompt.rs index d586ca6b7..630332878 100644 --- a/lib/components/fabro-workflow/src/handler/prompt.rs +++ b/lib/components/fabro-workflow/src/handler/prompt.rs @@ -279,6 +279,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/lifecycle/git.rs b/lib/components/fabro-workflow/src/lifecycle/git.rs index 18f4139bf..2cb9fd2dd 100644 --- a/lib/components/fabro-workflow/src/lifecycle/git.rs +++ b/lib/components/fabro-workflow/src/lifecycle/git.rs @@ -750,6 +750,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/operations/archive.rs b/lib/components/fabro-workflow/src/operations/archive.rs index c44ff0006..92d9124cf 100644 --- a/lib/components/fabro-workflow/src/operations/archive.rs +++ b/lib/components/fabro-workflow/src/operations/archive.rs @@ -227,6 +227,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/operations/create.rs b/lib/components/fabro-workflow/src/operations/create.rs index 59f424e2d..7783a9b59 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, BlobHash, ForkSourceRef, GitContext, ManifestPath, RunId, RunProvenance, + WorkflowSettings, }; use fabro_util::json::normalize_json_value; use tokio::task::spawn_blocking; @@ -470,6 +471,7 @@ pub async fn persist_create_run( provenance, manifest_blob: None, definition_blob: None, + spec_blob: None, git, fork_source_ref, }; @@ -516,18 +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, - }; + 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( @@ -554,6 +555,7 @@ async fn persist_created_run( automation: record.automation.clone(), provenance: record.provenance.clone(), manifest_blob, + spec_blob: Some(spec_blob), git: record.git.clone(), fork_source_ref: record.fork_source_ref.clone(), retried_from: None, @@ -580,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/operations/fork.rs b/lib/components/fabro-workflow/src/operations/fork.rs index 7e5e65516..3f728211d 100644 --- a/lib/components/fabro-workflow/src/operations/fork.rs +++ b/lib/components/fabro-workflow/src/operations/fork.rs @@ -162,6 +162,9 @@ async fn persist_forked_run( automation: spec.automation.clone(), provenance: spec.provenance.clone(), manifest_blob: spec.manifest_blob, + // Content-addressed, so the forked run reads the source run's + // unredacted spec bytes through the same id. + spec_blob: spec.spec_blob, git: spec.git.clone(), fork_source_ref: spec.fork_source_ref.clone(), retried_from: None, @@ -381,6 +384,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: Some(fabro_types::GitContext { origin_url: "https://github.com/example/repo.git".to_string(), branch: "main".to_string(), diff --git a/lib/components/fabro-workflow/src/operations/retry.rs b/lib/components/fabro-workflow/src/operations/retry.rs index 8163efa55..010c08f3f 100644 --- a/lib/components/fabro-workflow/src/operations/retry.rs +++ b/lib/components/fabro-workflow/src/operations/retry.rs @@ -54,6 +54,7 @@ pub async fn retry_run( provenance: _, manifest_blob, definition_blob, + spec_blob, git, fork_source_ref, } = source.spec; @@ -78,6 +79,9 @@ pub async fn retry_run( automation, provenance: input.provenance.clone(), manifest_blob, + // Blobs are content-addressed, so the retried run reads the source + // run's unredacted spec bytes through the same id. + spec_blob, git, fork_source_ref, retried_from: Some(source_run_id), @@ -185,6 +189,7 @@ mod tests { automation: None, provenance: provenance("source-user"), manifest_blob, + spec_blob: None, git: Some(git_context()), fork_source_ref, retried_from: None, diff --git a/lib/components/fabro-workflow/src/operations/timeline.rs b/lib/components/fabro-workflow/src/operations/timeline.rs index b96d10aa4..83319745e 100644 --- a/lib/components/fabro-workflow/src/operations/timeline.rs +++ b/lib/components/fabro-workflow/src/operations/timeline.rs @@ -252,6 +252,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, diff --git a/lib/components/fabro-workflow/src/pipeline/execute/tests.rs b/lib/components/fabro-workflow/src/pipeline/execute/tests.rs index 5aadd6018..7bc1d6dff 100644 --- a/lib/components/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/components/fabro-workflow/src/pipeline/execute/tests.rs @@ -173,6 +173,7 @@ fn persisted_workflow(graph: Graph, source: String, run_dir: &Path, run_id: RunI provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }, ) @@ -218,6 +219,7 @@ async fn seed_created_and_starting( automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: run_options.pre_run_git.clone(), fork_source_ref: run_options.fork_source_ref.clone(), retried_from: None, diff --git a/lib/components/fabro-workflow/src/pipeline/finalize.rs b/lib/components/fabro-workflow/src/pipeline/finalize.rs index b1e7847ee..bcc520fcb 100644 --- a/lib/components/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/components/fabro-workflow/src/pipeline/finalize.rs @@ -788,6 +788,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -906,6 +907,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, diff --git a/lib/components/fabro-workflow/src/pipeline/initialize.rs b/lib/components/fabro-workflow/src/pipeline/initialize.rs index 4d4d529c4..2d5477434 100644 --- a/lib/components/fabro-workflow/src/pipeline/initialize.rs +++ b/lib/components/fabro-workflow/src/pipeline/initialize.rs @@ -866,6 +866,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref, }, ) @@ -1049,6 +1050,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: run_options.fork_source_ref.clone(), retried_from: None, diff --git a/lib/components/fabro-workflow/src/pipeline/persist.rs b/lib/components/fabro-workflow/src/pipeline/persist.rs index cc12c3cba..303241ff5 100644 --- a/lib/components/fabro-workflow/src/pipeline/persist.rs +++ b/lib/components/fabro-workflow/src/pipeline/persist.rs @@ -2,6 +2,7 @@ use std::path::Path; use super::types::{PersistOptions, Persisted, Validated}; use crate::error::Error; +use crate::records::RunSpec; use crate::runtime_store::RunStoreHandle; /// PERSIST phase: create the run directory and return durable metadata for @@ -37,7 +38,7 @@ pub(crate) async fn load_from_store( .state() .await .map_err(|err| Error::engine(err.to_string()))?; - let run_spec = state.spec; + let run_spec = executable_run_spec(run_store, state.spec).await?; let graph = run_spec.graph.clone(); let source = run_spec.graph_source.clone().unwrap_or_default(); @@ -50,6 +51,41 @@ pub(crate) async fn load_from_store( )) } +/// Replace the event-folded spec content with the exact bytes from the spec +/// blob. Stored events pass through secret redaction, so the folded spec is +/// display data; the blob written at creation is what execution must see. +/// Runs created before the blob existed fall back to the folded spec. +async fn executable_run_spec( + run_store: &RunStoreHandle, + folded: RunSpec, +) -> Result { + let Some(blob_id) = folded.spec_blob else { + return Ok(folded); + }; + let bytes = run_store + .read_blob(&blob_id) + .await + .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::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) +} + #[cfg(test)] #[expect(clippy::disallowed_methods, reason = "tests stage pipeline fixtures")] mod tests { @@ -150,30 +186,49 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, } } async fn seeded_store(record: &RunSpec, source: Option<&str>) -> RunDatabase { + seeded_store_with(record, source, Some(record)).await + } + + async fn seeded_store_with( + record: &RunSpec, + source: Option<&str>, + blob_record: Option<&RunSpec>, + ) -> RunDatabase { let store = memory_store(); let run_store = store.create_run(&record.run_id).await.unwrap(); + let spec_blob = match blob_record { + Some(blob_record) => Some( + run_store + .write_blob(&serde_json::to_vec(blob_record).unwrap()) + .await + .unwrap(), + ), + None => None, + }; append_event(&run_store, &record.run_id, &Event::RunCreated { - run_id: record.run_id, - title: None, - settings: serde_json::to_value(&record.settings).unwrap(), - graph: serde_json::to_value(&record.graph).unwrap(), - workflow_source: source.map(ToOwned::to_owned), - labels: record.labels.clone().into_iter().collect(), + run_id: record.run_id, + title: None, + settings: serde_json::to_value(&record.settings).unwrap(), + graph: serde_json::to_value(&record.graph).unwrap(), + workflow_source: source.map(ToOwned::to_owned), + labels: record.labels.clone().into_iter().collect(), source_directory: record.source_directory.clone(), - workflow_slug: record.workflow_slug.clone(), - automation: record.automation.clone(), - provenance: record.provenance.clone(), - manifest_blob: None, - git: record.git.clone(), - fork_source_ref: record.fork_source_ref.clone(), - retried_from: None, - parent_id: None, - web_url: None, + workflow_slug: record.workflow_slug.clone(), + automation: record.automation.clone(), + provenance: record.provenance.clone(), + manifest_blob: None, + spec_blob, + git: record.git.clone(), + fork_source_ref: record.fork_source_ref.clone(), + retried_from: None, + parent_id: None, + web_url: None, }) .await .unwrap(); @@ -272,6 +327,96 @@ mod tests { assert!(loaded.diagnostics().is_empty()); } + #[tokio::test] + async fn load_from_store_preserves_high_entropy_dockerfile_content() { + // The spec the worker executes must survive the store byte-identical. + // Event redaction is a storage/display concern; when it reaches the + // spec that `load_from_store` rehydrates, the sandbox builds a + // corrupted Dockerfile: `ARG NAME=` pairs come back as + // `ARG REDACTED`, the build's `set -eu` step fails on the unset + // variable, and the environment's snapshot identity silently changes. + 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(); + + // Two shapes that must both survive: the hex pins that triggered the + // production failure, and a token high-entropy enough that any + // detector will keep flagging it in stored events. The second keeps + // this test red until execution stops reading redacted content, + // independent of how the entropy heuristic evolves. + let dockerfile = "FROM buildpack-deps:noble\n\ + ARG DOCKER_INSTALL_COMMIT=5ce20f2eef3615d08fea941eda5a109e949e8ebf\n\ + ARG DOCKER_INSTALL_SHA256=b991f2806186f7287bb9e53362060c382e906d154599b2fb0982f34246bacfd4\n\ + ENV CACHE_SALT=xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6p\n\ + RUN install-docker \"${DOCKER_INSTALL_COMMIT}\" \"${DOCKER_INSTALL_SHA256}\"\n"; + + let mut record = sample_record(different_graph()); + record.graph = graph; + record.settings.run.environment.image.dockerfile = Some( + fabro_types::settings::run::DockerfileSource::Inline(dockerfile.to_string()), + ); + + let run_store = seeded_store(&record, Some(&source)).await; + let loaded = load_from_store(&run_store.clone().into(), &run_dir) + .await + .unwrap(); + + assert_eq!( + loaded.run_spec().settings.run.environment.image.dockerfile, + Some(fabro_types::settings::run::DockerfileSource::Inline( + dockerfile.to_string() + )), + "the executable run spec must round-trip through the store unredacted" + ); + } + + #[tokio::test] + async fn load_from_store_falls_back_to_folded_spec_without_spec_blob() { + // Runs created before the spec blob existed carry no spec_blob on + // run.created; the folded spec is their only copy. + 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 mut record = sample_record(different_graph()); + record.graph = graph; + + let run_store = seeded_store_with(&record, Some(&source), None).await; + let loaded = load_from_store(&run_store.clone().into(), &run_dir) + .await + .unwrap(); + + assert_eq!(loaded.run_spec().settings, record.settings); + 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/components/fabro-workflow/src/pipeline/pull_request.rs b/lib/components/fabro-workflow/src/pipeline/pull_request.rs index 1ab76a24b..f3f7e4a78 100644 --- a/lib/components/fabro-workflow/src/pipeline/pull_request.rs +++ b/lib/components/fabro-workflow/src/pipeline/pull_request.rs @@ -830,6 +830,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }, @@ -1112,6 +1113,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }; append_event(&run_store, &fixtures::RUN_1, &Event::RunCreated { @@ -1126,6 +1128,7 @@ mod tests { automation: None, provenance: run_spec.provenance.clone(), manifest_blob: None, + spec_blob: None, git: run_spec.git.clone(), fork_source_ref: None, retried_from: None, @@ -1179,6 +1182,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }; append_event(&run_store, &fixtures::RUN_1, &Event::RunCreated { @@ -1193,6 +1197,7 @@ mod tests { automation: None, provenance: run_spec.provenance.clone(), manifest_blob: None, + spec_blob: None, git: run_spec.git.clone(), fork_source_ref: None, retried_from: None, @@ -1596,6 +1601,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }; append_event(&run_store, &fixtures::RUN_1, &Event::RunCreated { @@ -1610,6 +1616,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, @@ -1813,6 +1820,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }; append_event(&run_store, &fixtures::RUN_1, &Event::RunCreated { @@ -1827,6 +1835,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/run_lookup.rs b/lib/components/fabro-workflow/src/run_lookup.rs index 1edd36290..8a66dbf91 100644 --- a/lib/components/fabro-workflow/src/run_lookup.rs +++ b/lib/components/fabro-workflow/src/run_lookup.rs @@ -501,6 +501,7 @@ mod tests { automation: None, provenance: run_spec.provenance.clone(), manifest_blob: None, + spec_blob: None, git: run_spec.git.clone(), fork_source_ref: run_spec.fork_source_ref.clone(), retried_from: None, diff --git a/lib/components/fabro-workflow/src/run_metadata.rs b/lib/components/fabro-workflow/src/run_metadata.rs index eb9ae4b71..20d54fb9a 100644 --- a/lib/components/fabro-workflow/src/run_metadata.rs +++ b/lib/components/fabro-workflow/src/run_metadata.rs @@ -620,6 +620,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, fork_source_ref: None, }, chrono::Utc::now(), diff --git a/lib/components/fabro-workflow/src/runtime_store.rs b/lib/components/fabro-workflow/src/runtime_store.rs index 63d55ca64..2e506400c 100644 --- a/lib/components/fabro-workflow/src/runtime_store.rs +++ b/lib/components/fabro-workflow/src/runtime_store.rs @@ -157,6 +157,7 @@ mod tests { automation: None, provenance: test_support::test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/components/fabro-workflow/src/stage_execution.rs b/lib/components/fabro-workflow/src/stage_execution.rs index c05a3547e..ec755127f 100644 --- a/lib/components/fabro-workflow/src/stage_execution.rs +++ b/lib/components/fabro-workflow/src/stage_execution.rs @@ -209,6 +209,7 @@ mod tests { provenance: test_support::test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, }; diff --git a/lib/components/fabro-workflow/src/test_support.rs b/lib/components/fabro-workflow/src/test_support.rs index 19390762d..c7407b7b9 100644 --- a/lib/components/fabro-workflow/src/test_support.rs +++ b/lib/components/fabro-workflow/src/test_support.rs @@ -204,6 +204,7 @@ async fn initialized( }, }, manifest_blob: None, + spec_blob: None, git: run_options.pre_run_git.clone(), fork_source_ref: run_options.fork_source_ref.clone(), retried_from: None, diff --git a/lib/foundation/fabro-redact/src/entropy.rs b/lib/foundation/fabro-redact/src/entropy.rs index dc744cf39..ec6029c3c 100644 --- a/lib/foundation/fabro-redact/src/entropy.rs +++ b/lib/foundation/fabro-redact/src/entropy.rs @@ -36,6 +36,12 @@ pub(super) fn shannon_entropy(s: &str) -> f64 { /// Returns regions where tokens match `[A-Za-z0-9+_=-]{10,}` and have /// Shannon entropy above the threshold (4.5 bits). Protects against /// consuming characters from JSON escape sequences. +/// +/// An assignment token (`NAME=value`) is measured and redacted by its +/// value alone. Measuring the pair merges the name's charset into the +/// value's and pushes innocuous values (a pure-hex git SHA can never +/// exceed 4.0 bits by itself) over the threshold, and redacting the +/// pair destroys the name that says what was redacted. pub(super) fn find_entropy_regions(s: &str) -> Vec { let mut regions = Vec::new(); for m in SECRET_PATTERN.find_iter(s) { @@ -58,6 +64,10 @@ pub(super) fn find_entropy_regions(s: &str) -> Vec { } } + if let Some(offset) = assignment_value_offset(&s[start..end]) { + start += offset; + } + if shannon_entropy(&s[start..end]) > ENTROPY_THRESHOLD { regions.push(Region { start, end }); } @@ -65,6 +75,28 @@ pub(super) fn find_entropy_regions(s: &str) -> Vec { regions } +/// For an assignment token (`NAME=value` with an identifier-shaped name), +/// return the byte offset where the value begins. Entropy above the 4.5-bit +/// threshold needs at least 23 distinct characters, so a value too short to +/// 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()?; + if !(first.is_ascii_alphabetic() || first == '_') { + return None; + } + if !chars.all(|c| c.is_ascii_alphanumeric() || c == '_') { + return None; + } + Some(eq + 1) +} + #[cfg(test)] mod tests { use super::*; @@ -98,14 +130,28 @@ mod tests { #[test] fn regions_finds_high_entropy_token() { - // `=` is in the regex pattern, so "key=xK9..." matches as one token + // "key=xK9..." matches as one token, but only the value is + // measured and flagged; the name survives redaction. let input = "key=xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6p"; let regions = find_entropy_regions(input); assert_eq!(regions.len(), 1); - assert_eq!(regions[0].start, 0); + assert_eq!(regions[0].start, "key=".len()); 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-redact/src/jsonl.rs b/lib/foundation/fabro-redact/src/jsonl.rs index e466f7002..f0f50f6fa 100644 --- a/lib/foundation/fabro-redact/src/jsonl.rs +++ b/lib/foundation/fabro-redact/src/jsonl.rs @@ -209,7 +209,7 @@ mod tests { let redacted = redact_json_value(input); assert_eq!(redacted["name"], "fabro-01KQR3V9D4VPFFWMNTVH09J48G"); - assert_eq!(redacted["content"], "REDACTED"); + assert_eq!(redacted["content"], "token=REDACTED"); } #[test] @@ -262,7 +262,7 @@ mod tests { let redacted = redact_json_value(input); - assert_eq!(redacted["content"], "REDACTED"); + assert_eq!(redacted["content"], "key=REDACTED"); assert_eq!(redacted["session_id"], HIGH_ENTROPY_SECRET); } diff --git a/lib/foundation/fabro-redact/src/lib.rs b/lib/foundation/fabro-redact/src/lib.rs index 8ed562e53..4a7efa45b 100644 --- a/lib/foundation/fabro-redact/src/lib.rs +++ b/lib/foundation/fabro-redact/src/lib.rs @@ -111,6 +111,27 @@ mod tests { assert_eq!(result, "key=REDACTED"); } + #[test] + fn redact_string_keeps_assignment_with_low_entropy_value() { + // A pinned git SHA is pure hex, so the value alone can never exceed + // 4.0 bits of entropy. Only the merged NAME=value token crosses the + // 4.5-bit threshold, because the uppercase name widens the charset. + // Measuring the name together with the value redacts innocuous + // pins; the pair must survive. + let input = "ARG DOCKER_INSTALL_COMMIT=5ce20f2eef3615d08fea941eda5a109e949e8ebf"; + assert_eq!(redact_string(input), input); + } + + #[test] + fn redact_string_keeps_assignment_key_for_high_entropy_value() { + // The value alone is above the entropy threshold, so it is + // redacted either way — but the name says which setting was + // redacted and must survive, as the gitleaks layer already + // does for `key=REDACTED`. + let result = redact_string("BUILD_STAMP=xK9mZ2vL8nQ5rT1wY4bC7dF0gH3jE6p"); + assert_eq!(result, "BUILD_STAMP=REDACTED"); + } + #[test] fn redact_string_overlapping_detections_produce_single_redacted() { // A high-entropy string that also matches a gitleaks pattern diff --git a/lib/foundation/fabro-test/src/lib.rs b/lib/foundation/fabro-test/src/lib.rs index 00af8bf4c..5d442b411 100644 --- a/lib/foundation/fabro-test/src/lib.rs +++ b/lib/foundation/fabro-test/src/lib.rs @@ -1955,7 +1955,7 @@ 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); - for field in ["manifest_blob", "definition_blob"] { + for field in ["manifest_blob", "definition_blob", "spec_blob"] { filters.push(( format!(r#""{field}":\s*"[0-9a-f]{{64}}""#), format!(r#""{field}": "[BLOB_HASH]""#), diff --git a/lib/foundation/fabro-types/src/run.rs b/lib/foundation/fabro-types/src/run.rs index 287f5b85f..86f1910d7 100644 --- a/lib/foundation/fabro-types/src/run.rs +++ b/lib/foundation/fabro-types/src/run.rs @@ -76,6 +76,11 @@ pub struct RunSpec { pub manifest_blob: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub definition_blob: Option, + /// Unredacted copy of this spec in the blob store. Stored events pass + /// through secret redaction, so the spec folded from them is display + /// data; execution must load the spec from this blob. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub spec_blob: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub git: Option, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/lib/foundation/fabro-types/src/run_event/run.rs b/lib/foundation/fabro-types/src/run_event/run.rs index 68a994dc9..c7a17d990 100644 --- a/lib/foundation/fabro-types/src/run_event/run.rs +++ b/lib/foundation/fabro-types/src/run_event/run.rs @@ -28,6 +28,11 @@ pub struct RunCreatedProps { pub provenance: RunProvenance, #[serde(default, skip_serializing_if = "Option::is_none")] pub manifest_blob: Option, + /// Unredacted copy of the run spec in the blob store. The settings and + /// graph on this event are redacted at the sink; execution loads the + /// spec from this blob instead. + #[serde(default, skip_serializing_if = "Option::is_none")] + pub spec_blob: Option, #[serde(default, skip_serializing_if = "Option::is_none")] pub git: Option, #[serde(default, skip_serializing_if = "Option::is_none")] diff --git a/lib/foundation/fabro-types/src/test_support.rs b/lib/foundation/fabro-types/src/test_support.rs index b00e792b8..41af963f1 100644 --- a/lib/foundation/fabro-types/src/test_support.rs +++ b/lib/foundation/fabro-types/src/test_support.rs @@ -49,6 +49,7 @@ pub fn test_run_spec() -> RunSpec { provenance: test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: None, fork_source_ref: None, } diff --git a/lib/foundation/fabro-types/tests/run_event_serde.rs b/lib/foundation/fabro-types/tests/run_event_serde.rs index 633abab04..f287a6acc 100644 --- a/lib/foundation/fabro-types/tests/run_event_serde.rs +++ b/lib/foundation/fabro-types/tests/run_event_serde.rs @@ -32,6 +32,7 @@ fn run_created_props_round_trip_templated_settings() { }), provenance: test_run_provenance(), manifest_blob: None, + spec_blob: None, git: Some(GitContext { origin_url: "https://github.com/fabro-sh/fabro.git".to_string(), branch: "main".to_string(), @@ -93,6 +94,7 @@ fn run_created_props_omits_web_url_when_absent() { automation: None, provenance: test_run_provenance(), manifest_blob: None, + spec_blob: None, git: None, fork_source_ref: None, retried_from: None, diff --git a/lib/foundation/fabro-types/tests/run_spec_serde.rs b/lib/foundation/fabro-types/tests/run_spec_serde.rs index 97529bc93..ec7a7ef40 100644 --- a/lib/foundation/fabro-types/tests/run_spec_serde.rs +++ b/lib/foundation/fabro-types/tests/run_spec_serde.rs @@ -31,6 +31,7 @@ fn run_spec_round_trips_templated_settings() { provenance: test_run_provenance(), manifest_blob: None, definition_blob: None, + spec_blob: None, git: Some(GitContext { origin_url: "https://github.com/fabro-sh/fabro.git".to_string(), branch: "main".to_string(), diff --git a/lib/packages/fabro-api-client/src/models/run-spec.ts b/lib/packages/fabro-api-client/src/models/run-spec.ts index 4797dc2ad..e83be609b 100644 --- a/lib/packages/fabro-api-client/src/models/run-spec.ts +++ b/lib/packages/fabro-api-client/src/models/run-spec.ts @@ -44,6 +44,7 @@ export interface RunSpec { 'provenance': RunProvenance; 'manifest_blob'?: string | null; 'definition_blob'?: string | null; + 'spec_blob'?: string | null; 'git'?: GitContext | null; 'fork_source_ref'?: ForkSourceRef | null; }