diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs index b0e912438..b6fbd76c0 100644 --- a/lib/crates/fabro-cli/src/commands/run/runner.rs +++ b/lib/crates/fabro-cli/src/commands/run/runner.rs @@ -27,7 +27,7 @@ use fabro_types::{ }; use fabro_vault::Vault; use fabro_workflow::artifact_upload::{ArtifactSink, StageArtifactUploader}; -use fabro_workflow::event::{Emitter, RunEventSink, build_redacted_event_payload_with_redactor}; +use fabro_workflow::event::{Emitter, RunEventSink, redacted_run_event}; use fabro_workflow::operations::{self, StartServices}; use fabro_workflow::run_control::RunControlState; use fabro_workflow::runtime_store::{RunStoreBackend, RunStoreHandle}; @@ -1016,10 +1016,8 @@ impl RunStoreBackend for HttpRunStore { redactor: Option<&SecretRedactor>, ) -> Result<()> { let event = if let Some(redactor) = redactor { - let payload = - build_redacted_event_payload_with_redactor(event, &self.run_id, Some(redactor)) - .context("failed to build redacted run event payload")?; - RunEvent::try_from(&payload).context("redacted run event payload is invalid")? + redacted_run_event(event, &self.run_id, Some(redactor)) + .context("failed to build redacted run event payload")? } else { event.clone() }; diff --git a/lib/crates/fabro-hooks/src/runner.rs b/lib/crates/fabro-hooks/src/runner.rs index fd005e0a5..07cb1c8f4 100644 --- a/lib/crates/fabro-hooks/src/runner.rs +++ b/lib/crates/fabro-hooks/src/runner.rs @@ -47,6 +47,8 @@ fn redact_hook_decision(decision: HookDecision, redactor: &SecretRedactor) -> Ho HookDecision::Block { reason } => HookDecision::Block { reason: reason.map(|reason| redactor.redact_into(&reason)), }, + // `Override.edge_to` is a structural graph edge id, not free-form text, + // so it is intentionally left unredacted. HookDecision::Proceed | HookDecision::Override { .. } => decision, } } @@ -178,6 +180,43 @@ impl HookRunner { .any(|field| field.is_some_and(|v| re.is_match(v))) } + /// Resolve secrets, run a single hook through the executor, and redact its + /// result. Shared by the blocking and non-blocking loops. + async fn execute_one( + &self, + hook: &HookDefinition, + context: &HookContext, + sandbox: Arc, + execution_context: &HookExecutionContext, + ) -> HookResult { + tracing::debug!( + hook = %hook.effective_name(), + event = %context.event, + "Executing hook" + ); + let secrets = self.secrets.resolve_for_definition(hook).await; + let result = self + .executor + .execute( + hook, + context, + sandbox, + execution_context, + self.llm_source.as_ref(), + Arc::clone(&self.catalog), + &secrets, + ) + .await; + let result = redact_hook_result(result, secrets.redactor()); + tracing::debug!( + hook = %hook.effective_name(), + duration_ms = result.duration_ms, + decision = decision_label(&result.decision), + "Hook complete" + ); + result + } + async fn run_sequential( &self, hooks: &[&HookDefinition], @@ -187,31 +226,9 @@ impl HookRunner { ) -> HookDecision { let mut merged = HookDecision::Proceed; for hook in hooks { - tracing::debug!( - hook = %hook.effective_name(), - event = %context.event, - "Executing hook" - ); - let secrets = self.secrets.resolve_for_definition(hook).await; let result = self - .executor - .execute( - hook, - context, - sandbox.clone(), - execution_context, - self.llm_source.as_ref(), - Arc::clone(&self.catalog), - &secrets, - ) + .execute_one(hook, context, sandbox.clone(), execution_context) .await; - let result = redact_hook_result(result, secrets.redactor()); - tracing::debug!( - hook = %hook.effective_name(), - duration_ms = result.duration_ms, - decision = decision_label(&result.decision), - "Hook complete" - ); if hook.is_blocking() { merged = merged.merge(result.decision); @@ -245,31 +262,9 @@ impl HookRunner { execution_context: &HookExecutionContext, ) -> HookDecision { for hook in hooks { - tracing::debug!( - hook = %hook.effective_name(), - event = %context.event, - "Executing hook" - ); - let secrets = self.secrets.resolve_for_definition(hook).await; let result = self - .executor - .execute( - hook, - context, - sandbox.clone(), - execution_context, - self.llm_source.as_ref(), - Arc::clone(&self.catalog), - &secrets, - ) + .execute_one(hook, context, sandbox.clone(), execution_context) .await; - let result = redact_hook_result(result, secrets.redactor()); - tracing::debug!( - hook = %hook.effective_name(), - duration_ms = result.duration_ms, - decision = decision_label(&result.decision), - "Hook complete" - ); if !result.decision.is_proceed() { tracing::warn!( hook = %hook.effective_name(), diff --git a/lib/crates/fabro-hooks/src/secrets.rs b/lib/crates/fabro-hooks/src/secrets.rs index 399af49d2..3ffcd9b27 100644 --- a/lib/crates/fabro-hooks/src/secrets.rs +++ b/lib/crates/fabro-hooks/src/secrets.rs @@ -44,11 +44,6 @@ impl HookSecretResolver { } } - #[must_use] - pub fn redactor(&self) -> &SecretRedactor { - &self.redactor - } - pub async fn resolve_for_definition(&self, definition: &HookDefinition) -> ResolvedHookSecrets { self.resolve_names(secret_names_for_definition(definition)) .await @@ -100,18 +95,11 @@ impl ResolvedHookSecrets { Self { values, redactor } } - #[must_use] - pub fn empty_with_redactor(redactor: SecretRedactor) -> Self { - Self::new(HashMap::new(), redactor) - } - #[must_use] pub fn lookup(&self, name: &str) -> Option { - let value = self.values.get(name).cloned(); - if let Some(value) = value.as_deref() { - self.redactor.register(value); - } - value + // Values were registered with the redactor at construction time in + // `new`, so this is a pure read. + self.values.get(name).cloned() } #[must_use] diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 8ebf998f5..30be2bfee 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -17,6 +17,7 @@ pub use self::names::event_name; pub use self::redaction::{ build_redacted_event_payload, build_redacted_event_payload_with_redactor, event_payload_from_redacted_json, redacted_event_json, redacted_event_json_with_redactor, + redacted_run_event, }; pub use self::sink::{ RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_to_sink, diff --git a/lib/crates/fabro-workflow/src/event/redaction.rs b/lib/crates/fabro-workflow/src/event/redaction.rs index 858818df3..eba6a209d 100644 --- a/lib/crates/fabro-workflow/src/event/redaction.rs +++ b/lib/crates/fabro-workflow/src/event/redaction.rs @@ -18,6 +18,21 @@ pub fn build_redacted_event_payload_with_redactor( EventPayload::new(value, run_id).map_err(anyhow::Error::from) } +/// Redact an event and reconstruct it as a `RunEvent`. +/// +/// This runs the event through the content-based pass plus the optional +/// per-run [`SecretRedactor`], then reparses the redacted payload back into a +/// `RunEvent` so downstream sinks that require a typed event never see the raw +/// value. +pub fn redacted_run_event( + event: &RunEvent, + run_id: &RunId, + redactor: Option<&SecretRedactor>, +) -> Result { + let payload = build_redacted_event_payload_with_redactor(event, run_id, redactor)?; + RunEvent::try_from(&payload).map_err(anyhow::Error::from) +} + pub fn redacted_event_json(event: &RunEvent) -> Result { redacted_event_json_with_redactor(event, None) } @@ -43,6 +58,11 @@ fn redacted_event_value(event: &RunEvent, redactor: Option<&SecretRedactor>) -> } fn redact_event_payload_secrets(value: &mut Value, redactor: &SecretRedactor) { + // No declared secrets (the common case): skip the recursive property walk + // entirely. Content-based redaction already ran in `redacted_event_value`. + if redactor.is_empty() { + return; + } if let Some(properties) = value.get_mut("properties") { redact_redactable_event_properties(properties, redactor); } diff --git a/lib/crates/fabro-workflow/src/event/sink.rs b/lib/crates/fabro-workflow/src/event/sink.rs index 14e98d6d1..4473fae37 100644 --- a/lib/crates/fabro-workflow/src/event/sink.rs +++ b/lib/crates/fabro-workflow/src/event/sink.rs @@ -11,8 +11,8 @@ use tokio::sync::{Mutex as AsyncMutex, mpsc, oneshot}; use super::emitter::Emitter; use super::redaction::{ - build_redacted_event_payload, build_redacted_event_payload_with_redactor, redacted_event_json, - redacted_event_json_with_redactor, + build_redacted_event_payload, redacted_event_json, redacted_event_json_with_redactor, + redacted_run_event, }; use super::{Event, to_run_event}; use crate::runtime_store::RunStoreHandle; @@ -137,15 +137,11 @@ impl RunEventSink { writer.flush().await?; } Self::Callback(callback) => { - let event = if let Some(redactor) = redactor.as_ref() { - let payload = build_redacted_event_payload_with_redactor( - &event, - &event.run_id, - Some(redactor), - )?; - RunEvent::try_from(&payload)? - } else { - event + let event = match redactor.as_ref() { + Some(redactor) => { + redacted_run_event(&event, &event.run_id, Some(redactor))? + } + None => event, }; callback(event).await?; }