From f8992b2f0984249232ceb140e8519a2100306f48 Mon Sep 17 00:00:00 2001 From: Fabro Date: Wed, 1 Jul 2026 21:13:36 +0000 Subject: [PATCH] fabro(01KWFGXZ5P42QRWBYAPVEAXMX6): implement (succeeded) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fabro-Run: 01KWFGXZ5P42QRWBYAPVEAXMX6 Fabro-Completed: 5 ⚒️ Generated with [Fabro](https://fabro.sh) --- docs/public/agents/hooks.mdx | 31 +- docs/public/execution/run-configuration.mdx | 15 +- docs/public/workflows/variables.mdx | 2 +- .../fabro-cli/src/commands/run/runner.rs | 21 +- lib/crates/fabro-hooks/src/bridge.rs | 1 + lib/crates/fabro-hooks/src/executor.rs | 315 +++++++++++++-- lib/crates/fabro-hooks/src/lib.rs | 2 + lib/crates/fabro-hooks/src/runner.rs | 62 ++- lib/crates/fabro-hooks/src/secrets.rs | 155 ++++++++ lib/crates/fabro-sandbox/src/error.rs | 47 ++- lib/crates/fabro-sandbox/src/lib.rs | 7 +- lib/crates/fabro-sandbox/src/sandbox.rs | 78 +++- lib/crates/fabro-workflow/src/event.rs | 3 +- lib/crates/fabro-workflow/src/event/events.rs | 6 +- .../fabro-workflow/src/event/redaction.rs | 194 +++++++++- lib/crates/fabro-workflow/src/event/sink.rs | 85 +++- .../fabro-workflow/src/handler/parallel.rs | 10 +- .../fabro-workflow/src/lifecycle/git.rs | 29 +- .../fabro-workflow/src/lifecycle/mod.rs | 2 + .../fabro-workflow/src/operations/start.rs | 363 +++++++++++++++++- .../fabro-workflow/src/pipeline/execute.rs | 1 + .../src/pipeline/execute/tests.rs | 6 + .../fabro-workflow/src/pipeline/finalize.rs | 15 +- .../fabro-workflow/src/pipeline/initialize.rs | 22 +- .../fabro-workflow/src/pipeline/types.rs | 3 + .../fabro-workflow/src/runtime_store.rs | 31 +- lib/crates/fabro-workflow/src/services.rs | 5 + lib/crates/fabro-workflow/src/test_support.rs | 2 + 28 files changed, 1415 insertions(+), 98 deletions(-) create mode 100644 lib/crates/fabro-hooks/src/secrets.rs diff --git a/docs/public/agents/hooks.mdx b/docs/public/agents/hooks.mdx index aafdd356f..77874ce27 100644 --- a/docs/public/agents/hooks.mdx +++ b/docs/public/agents/hooks.mdx @@ -36,8 +36,8 @@ Authorization = "Bearer {{ env.API_KEY }}" | Field | Description | |---|---| -| `url` | The endpoint to POST to. Must use `https://` unless `tls = "off"`. Supports `{{ env.NAME }}` interpolation. | -| `headers` | Optional HTTP headers. Values support `{{ env.NAME }}` interpolation, scoped to the names in `allowed_env_vars`. A token for any other env var fails to resolve and the hook blocks (fail-closed). | +| `url` | The endpoint to POST to. Must use `https://` unless `tls = "off"`. Supports `{{ env.NAME }}` and `{{ secrets.NAME }}` interpolation. | +| `headers` | Optional HTTP headers. Values support `{{ env.NAME }}` interpolation, scoped to the names in `allowed_env_vars`. A token for any other env var fails to resolve and the hook blocks (fail-closed). `{{ secrets.NAME }}` is not allowed in headers. | | `allowed_env_vars` | Allowlist of environment variable names a header may read via `{{ env.NAME }}`. Empty (the default) means no env vars may be interpolated into headers. | | `tls` | TLS mode: `"verify"` (default), `"no_verify"`, or `"off"`. | @@ -134,6 +134,33 @@ sandbox = false | `timeout_ms` | Hook timeout in milliseconds. Default: `60000` (60s) for most types, `30000` (30s) for prompt hooks. | | `sandbox` | Run inside the sandbox (`true`, default) or on the host (`false`). | +## Interpolation and secrets + +Hook `command`, HTTP `url`, prompt `prompt`, and agent `prompt` fields can reference server-vault token secrets with `{{ secrets.NAME }}`. Add these values with `fabro secret set NAME ...`. Missing secrets, empty/non-token vault entries, and unsupported namespaces fail closed so the hook does not run with a partially resolved value. + +```toml +[[hooks]] +event = "sandbox_ready" +command = "test \"{{ secrets.DEPLOY_ENV }}\" = staging" +blocking = true + +[[hooks]] +event = "run_failed" +type = "http" +url = "https://hooks.example.com/{{ secrets.WEBHOOK_PATH }}" +tls = "verify" + +[[hooks]] +event = "stage_start" +type = "prompt" +prompt = "Block deployments to {{ secrets.RESTRICTED_ENV }}?" +blocking = true +``` + +HTTP hook headers intentionally support only allowlisted `{{ env.NAME }}` interpolation. A `{{ secrets.NAME }}` token in a header blocks the hook with guidance to use secret interpolation in a hook command, prompt, or URL instead; Fabro does not provide a separate outbound-header secret allowlist. + +Resolved hook secrets are registered with the run's secret redactor. Fabro redacts those exact values from worker-side structured run surfaces such as events, `progress.jsonl`, and setup-error messages. Command output still crosses the sandbox boundary as process output, so treat sandbox-reemitted plaintext as best-effort: content-based redaction always runs, and worker-side event surfaces apply exact-match redaction after the output is captured. + ## Blocking vs. non-blocking Blocking hooks can affect workflow execution. Non-blocking hooks run for side effects only — their decisions are ignored. diff --git a/docs/public/execution/run-configuration.mdx b/docs/public/execution/run-configuration.mdx index 9c3cd8487..5e6c8b97a 100644 --- a/docs/public/execution/run-configuration.mdx +++ b/docs/public/execution/run-configuration.mdx @@ -288,12 +288,14 @@ When `provider = "local"`, Fabro runs directly in the resolved working directory. If you want local isolation, create or enter a separate clone or Git worktree yourself. -Environment variable values can be literal strings or host environment -references using `{{ env.VARNAME }}` syntax: +Environment variable values can be literal strings, host environment references +using `{{ env.VARNAME }}` syntax, or server-vault token secrets using +`{{ secrets.NAME }}` syntax: ```toml title="run.toml" [environments.ci.env] API_KEY = "{{ env.MY_API_KEY }}" +DEPLOY_TOKEN = "{{ secrets.DEPLOY_TOKEN }}" NODE_ENV = "production" SERVICE_URL = "https://api.{{ env.REGION }}.example.com" ``` @@ -302,9 +304,14 @@ SERVICE_URL = "https://api.{{ env.REGION }}.example.com" |---|---| | `"literal"` | Static value passed as-is | | `"{{ env.VARNAME }}"` | Whole-value reference resolved from the host environment at consumption time | -| `"prefix-{{ env.X }}-suffix"` | Substring interpolation; multiple tokens per string are supported | +| `"{{ secrets.NAME }}"` | Reference resolved from the server vault as a token secret at the run boundary | +| `"prefix-{{ env.X }}-{{ secrets.Y }}-suffix"` | Substring interpolation; multiple env and secret tokens per string are supported | -Missing host variables produce a hard error pointing at the specific field and unresolved token. +Missing host variables or missing/non-token secrets produce a hard error pointing at the specific field and unresolved token. + +Resolved declared secrets are registered with a per-run redactor. On worker-side structured surfaces — run events, `progress.jsonl`, the run store, SSE payloads derived from stored events, and setup-command errors — Fabro applies content-based redaction and then exact-match redaction for those declared secret values, even when the value is low entropy and does not look like a credential. + +The boundary is the worker/sandbox handoff. If you put a secret into sandbox process environment or a command line, software inside the sandbox can still print it as plain text. Fabro applies content-based redaction to command output and exact-match redaction when captured output is embedded back into worker-side structured events, but avoid intentionally echoing secrets from sandbox commands. ### `[run.integrations.github.permissions]` diff --git a/docs/public/workflows/variables.mdx b/docs/public/workflows/variables.mdx index 18b0c00b1..3344e1ab0 100644 --- a/docs/public/workflows/variables.mdx +++ b/docs/public/workflows/variables.mdx @@ -14,7 +14,7 @@ Goal and prompt templates can reference: | `{{ goal }}` | The workflow goal | | `{{ inputs.name }}` | A value from `[run.inputs]`, optionally overridden by CLI input flags | -Environment variables are **not** available in goal or prompt templates. Use `{{ env.NAME }}` only in config strings and HTTP hook headers. +Environment variables and secrets are **not** available in workflow graph goal or prompt templates. Use `{{ env.NAME }}` and `{{ secrets.NAME }}` only in supported run configuration and hook fields. ## Run config inputs diff --git a/lib/crates/fabro-cli/src/commands/run/runner.rs b/lib/crates/fabro-cli/src/commands/run/runner.rs index 9eba9a3d5..b0e912438 100644 --- a/lib/crates/fabro-cli/src/commands/run/runner.rs +++ b/lib/crates/fabro-cli/src/commands/run/runner.rs @@ -16,6 +16,7 @@ use fabro_interview::{ WorkerControlMessage, }; use fabro_model::Catalog; +use fabro_redact::SecretRedactor; use fabro_server::run_tool_manifest; use fabro_store::{EventEnvelope, RunProjection, RunProjectionReducer}; use fabro_tool::fabro_client::ClientBackend; @@ -26,7 +27,7 @@ use fabro_types::{ }; use fabro_vault::Vault; use fabro_workflow::artifact_upload::{ArtifactSink, StageArtifactUploader}; -use fabro_workflow::event::{Emitter, RunEventSink}; +use fabro_workflow::event::{Emitter, RunEventSink, build_redacted_event_payload_with_redactor}; use fabro_workflow::operations::{self, StartServices}; use fabro_workflow::run_control::RunControlState; use fabro_workflow::runtime_store::{RunStoreBackend, RunStoreHandle}; @@ -1006,6 +1007,22 @@ impl RunStoreBackend for HttpRunStore { } async fn append_run_event(&self, event: &RunEvent) -> Result<()> { + self.append_run_event_with_redactor(event, None).await + } + + async fn append_run_event_with_redactor( + &self, + event: &RunEvent, + 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")? + } else { + event.clone() + }; let seq = Box::pin(self.with_retries("append run event", || { let client = self.client.clone_for_reuse(); let run_id = self.run_id; @@ -1013,7 +1030,7 @@ impl RunStoreBackend for HttpRunStore { async move { client.append_run_event(&run_id, &event).await } })) .await?; - self.apply_acknowledged_event(seq, event).await + self.apply_acknowledged_event(seq, &event).await } async fn write_blob(&self, data: &[u8]) -> Result { diff --git a/lib/crates/fabro-hooks/src/bridge.rs b/lib/crates/fabro-hooks/src/bridge.rs index 6aa7e60f6..a7bfa05ed 100644 --- a/lib/crates/fabro-hooks/src/bridge.rs +++ b/lib/crates/fabro-hooks/src/bridge.rs @@ -102,6 +102,7 @@ mod tests { execution_context: &HookExecutionContext, _llm_source: &dyn fabro_auth::CredentialSource, _catalog: Arc, + _secrets: &crate::ResolvedHookSecrets, ) -> HookResult { self.captured_contexts.lock().unwrap().push(context.clone()); self.captured_execution_contexts diff --git a/lib/crates/fabro-hooks/src/executor.rs b/lib/crates/fabro-hooks/src/executor.rs index 9fda710d5..2e40bfe46 100644 --- a/lib/crates/fabro-hooks/src/executor.rs +++ b/lib/crates/fabro-hooks/src/executor.rs @@ -13,7 +13,7 @@ use fabro_llm::generate::{GenerateParams, generate_object}; use fabro_llm::types::{Message, Request, ToolResult}; use fabro_model::Catalog; use fabro_redact::redacted_url_for_log; -use fabro_types::settings::interp::Namespace; +use fabro_types::settings::interp::{Namespace, ResolveCtx}; use fabro_types::settings::{InterpString, ResolveError}; use fabro_util::env::{Env, SystemEnv}; use tokio::process::Command as TokioCommand; @@ -21,6 +21,7 @@ use tokio::time::timeout as tokio_timeout; use tokio_util::sync::CancellationToken; use crate::config::{HookDefinition, HookType, TlsMode}; +use crate::secrets::ResolvedHookSecrets; use crate::types::{ HookContext, HookDecision, HookExecutionContext, HookResult, PromptHookResponse, }; @@ -54,16 +55,18 @@ pub trait HookExecutor: Send + Sync { execution_context: &HookExecutionContext, llm_source: &dyn CredentialSource, catalog: Arc, + secrets: &ResolvedHookSecrets, ) -> HookResult; } /// Resolve a typed [`InterpString`] hook segment at fire time, looking up -/// `{{ env.* }}` tokens against `env`. +/// `{{ env.* }}` tokens against `env` and `{{ secrets.* }}` tokens against the +/// run-scoped secret resolver. /// -/// Only the `env` namespace is wired here; `{{ secrets.* }}`, `{{ vars.* }}`, -/// and `{{ inputs.* }}` tokens have no lookup in this context and resolve as -/// `Unavailable`, which is a hard error — so a hook that references one fails -/// closed rather than firing with a half-resolved value. +/// Only `env` and `secrets` are wired here; `{{ vars.* }}` and `{{ inputs.* }}` +/// tokens have no lookup in this context and resolve as `Unavailable`, which is +/// a hard error — so a hook that references one fails closed rather than firing +/// with a half-resolved value. /// /// The value stays typed end-to-end: it is carried as an `InterpString` /// through the config resolve layer and resolved here from its segments — @@ -73,13 +76,18 @@ pub trait HookExecutor: Send + Sync { /// /// Returns the typed [`ResolveError`] so callers keep the source until the /// decision boundary renders it; do not flatten it to a `String` here. -fn resolve_interp(value: &InterpString, env: &E) -> Result +fn resolve_interp( + value: &InterpString, + env: &E, + secrets: &ResolvedHookSecrets, +) -> Result where E: Env + ?Sized, { - value - .resolve(|name| env.var(name).ok()) - .map(|resolved| resolved.value) + let mut ctx = ResolveCtx::new() + .with_env(|name| env.var(name).ok()) + .with_secrets(|name| secrets.lookup(name)); + value.resolve_with(&mut ctx).map(|resolved| resolved.value) } #[expect( @@ -95,6 +103,7 @@ fn safe_url_source_for_log(url: &InterpString) -> String { #[derive(Debug, Clone, PartialEq, Eq)] enum HeaderResolveError { NotAllowed { name: String }, + SecretNotAllowed { name: String }, Resolve(ResolveError), } @@ -106,6 +115,11 @@ impl fmt::Display for HeaderResolveError { "environment variable {name:?} referenced by an HTTP hook header is not listed in \ allowed_env_vars" ), + Self::SecretNotAllowed { name } => write!( + f, + "secret {name:?} referenced by an HTTP hook header is not allowed; use secret \ + interpolation in a hook command, prompt, or url instead" + ), Self::Resolve(error) => error.fmt(f), } } @@ -114,7 +128,7 @@ impl fmt::Display for HeaderResolveError { impl std::error::Error for HeaderResolveError { fn source(&self) -> Option<&(dyn std::error::Error + 'static)> { match self { - Self::NotAllowed { .. } => None, + Self::NotAllowed { .. } | Self::SecretNotAllowed { .. } => None, Self::Resolve(error) => Some(error), } } @@ -135,10 +149,17 @@ fn resolve_header( value: &InterpString, allowed_env_vars: &[String], env: &E, + secrets: &ResolvedHookSecrets, ) -> Result where E: Env + ?Sized, { + if let Some(name) = value.names(Namespace::Secrets).into_iter().next() { + return Err(HeaderResolveError::SecretNotAllowed { + name: name.to_string(), + }); + } + if let Some(name) = value.names(Namespace::Env).into_iter().find(|name| { !allowed_env_vars .iter() @@ -149,7 +170,7 @@ where }); } - resolve_interp(value, env).map_err(HeaderResolveError::Resolve) + resolve_interp(value, env, secrets).map_err(HeaderResolveError::Resolve) } /// Executes hooks via shell commands or HTTP POST. @@ -181,20 +202,24 @@ impl HookExecutorImpl { /// Resolve the prompt and optional model segments at fire time. /// - /// Fail-closed: only `{{ env.* }}` is wired here; a missing env token (or a - /// token in any other, unavailable namespace) is a hard error so the hook - /// never fires with a half-resolved value. The caller turns the error into - /// a `Block` decision, matching the command-hook behavior. + /// Fail-closed: only `{{ env.* }}` and `{{ secrets.* }}` are wired here; a + /// missing token (or a token in any other, unavailable namespace) is a hard + /// error so the hook never fires with a half-resolved value. The caller + /// turns the error into a `Block` decision, matching the command-hook + /// behavior. fn resolve_prompt_and_model( prompt: &InterpString, model: Option<&InterpString>, env: &E, + secrets: &ResolvedHookSecrets, ) -> Result<(String, Option), ResolveError> where E: Env + ?Sized, { - let prompt = resolve_interp(prompt, env)?; - let model = model.map(|model| resolve_interp(model, env)).transpose()?; + let prompt = resolve_interp(prompt, env, secrets)?; + let model = model + .map(|model| resolve_interp(model, env, secrets)) + .transpose()?; Ok((prompt, model)) } @@ -206,11 +231,12 @@ impl HookExecutorImpl { sandbox: &Arc, execution_context: &HookExecutionContext, env: &E, + secrets: &ResolvedHookSecrets, ) -> HookDecision where E: Env + ?Sized, { - let command = match resolve_interp(command, env) { + let command = match resolve_interp(command, env, secrets) { Ok(command) => command, Err(error) => { return HookDecision::Block { @@ -363,13 +389,14 @@ impl HookExecutorImpl { model: Option<&InterpString>, context: &HookContext, env: &E, + secrets: &ResolvedHookSecrets, llm_source: &dyn CredentialSource, catalog: Arc, ) -> HookDecision where E: Env + ?Sized, { - let (prompt, model) = match Self::resolve_prompt_and_model(prompt, model, env) { + let (prompt, model) = match Self::resolve_prompt_and_model(prompt, model, env, secrets) { Ok(resolved) => resolved, Err(error) => { tracing::error!(error = %error, "prompt hook env resolution failed, not firing"); @@ -432,13 +459,14 @@ impl HookExecutorImpl { context: &HookContext, sandbox: Arc, env: &E, + secrets: &ResolvedHookSecrets, llm_source: &dyn CredentialSource, catalog: Arc, ) -> HookDecision where E: Env + ?Sized, { - let (prompt, model) = match Self::resolve_prompt_and_model(prompt, model, env) { + let (prompt, model) = match Self::resolve_prompt_and_model(prompt, model, env, secrets) { Ok(resolved) => resolved, Err(error) => { tracing::error!(error = %error, "agent hook env resolution failed, not firing"); @@ -578,11 +606,12 @@ impl HookExecutorImpl { context: &HookContext, timeout: std::time::Duration, env: &E, + secrets: &ResolvedHookSecrets, ) -> HookDecision where E: Env + ?Sized, { - let resolved_url = match resolve_interp(url, env) { + let resolved_url = match resolve_interp(url, env, secrets) { Ok(url) => url, Err(error) => { tracing::error!( @@ -618,7 +647,7 @@ impl HookExecutorImpl { // `{{ env.NAME }}` not in `allowed_env_vars` blocks before any // lookup, while an allowlisted-but-unset name still fails as // missing. - let interpolated = match resolve_header(value, allowed_env_vars, env) { + let interpolated = match resolve_header(value, allowed_env_vars, env, secrets) { Ok(rendered) => rendered, Err(error) => { tracing::error!( @@ -641,7 +670,11 @@ impl HookExecutorImpl { Err(e) => { tracing::warn!( url_source = %safe_url_source_for_log(url), - error = %e, + error_timeout = e.is_timeout(), + error_connect = e.is_connect(), + error_request = e.is_request(), + error_body = e.is_body(), + error_decode = e.is_decode(), "HTTP hook request failed, proceeding" ); return HookDecision::Proceed; @@ -662,7 +695,11 @@ impl HookExecutorImpl { Err(e) => { tracing::warn!( url_source = %safe_url_source_for_log(url), - error = %e, + error_timeout = e.is_timeout(), + error_connect = e.is_connect(), + error_request = e.is_request(), + error_body = e.is_body(), + error_decode = e.is_decode(), "HTTP hook body read failed, proceeding" ); return HookDecision::Proceed; @@ -728,6 +765,7 @@ impl HookExecutor for HookExecutorImpl { execution_context: &HookExecutionContext, llm_source: &dyn CredentialSource, catalog: Arc, + secrets: &ResolvedHookSecrets, ) -> HookResult { use std::sync::OnceLock; static HTTP_CLIENTS: OnceLock = OnceLock::new(); @@ -747,6 +785,7 @@ impl HookExecutor for HookExecutorImpl { &sandbox, execution_context, &env, + secrets, ) .await } @@ -774,6 +813,7 @@ impl HookExecutor for HookExecutorImpl { context, definition.timeout(), &env, + secrets, ) .await } @@ -793,6 +833,7 @@ impl HookExecutor for HookExecutorImpl { model.as_ref(), context, &env, + secrets, llm_source, Arc::clone(&catalog), ) @@ -818,6 +859,7 @@ impl HookExecutor for HookExecutorImpl { context, sandbox, &env, + secrets, llm_source, Arc::clone(&catalog), ) @@ -869,6 +911,19 @@ mod tests { HookExecutorImpl::build_http_client(TlsMode::Off) } + fn empty_secrets() -> ResolvedHookSecrets { + ResolvedHookSecrets::default() + } + + fn test_secrets(vars: &[(&str, &str)]) -> ResolvedHookSecrets { + ResolvedHookSecrets::new( + vars.iter() + .map(|(name, value)| ((*name).to_string(), (*value).to_string())) + .collect(), + fabro_redact::SecretRedactor::default(), + ) + } + fn make_definition(command: &str) -> HookDefinition { HookDefinition { name: Some("test-hook".into()), @@ -946,6 +1001,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -954,6 +1010,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert_eq!(result.decision, HookDecision::Proceed); @@ -967,6 +1024,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -975,6 +1033,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert!(matches!(result.decision, HookDecision::Block { .. })); @@ -987,6 +1046,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -995,6 +1055,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert!(matches!(result.decision, HookDecision::Block { .. })); @@ -1007,6 +1068,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -1015,6 +1077,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert_eq!(result.decision, HookDecision::Skip { @@ -1031,6 +1094,7 @@ mod tests { ctx.node_id = Some("plan".into()); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -1039,6 +1103,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert_eq!(result.decision, HookDecision::Proceed); @@ -1060,6 +1125,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -1068,6 +1134,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; assert!(matches!(result.decision, HookDecision::Block { .. })); @@ -1188,6 +1255,7 @@ mod tests { &interp("Bearer {{ env.FABRO_TEST_KEY_1 }}"), &["FABRO_TEST_KEY_1".to_string()], &env, + &empty_secrets(), ) .unwrap(); assert_eq!(result, "Bearer secret123"); @@ -1203,6 +1271,7 @@ mod tests { &interp("prefix-{{ env.FABRO_TEST_KEY_3 }}-suffix"), &[], &env, + &empty_secrets(), ) .unwrap_err(); assert_eq!(err, HeaderResolveError::NotAllowed { @@ -1217,11 +1286,12 @@ mod tests { &interp("prefix-{{ env.FABRO_TEST_KEY_3 }}-suffix"), &["FABRO_TEST_KEY_3".to_string()], &env, + &empty_secrets(), ) .unwrap_err(); match err { HeaderResolveError::Resolve(error) => assert_eq!(error.name, "FABRO_TEST_KEY_3"), - HeaderResolveError::NotAllowed { .. } => { + HeaderResolveError::NotAllowed { .. } | HeaderResolveError::SecretNotAllowed { .. } => { panic!("expected missing token resolve error, got {err:?}") } } @@ -1233,14 +1303,19 @@ mod tests { fn resolve_interp_resolves_embedded_token_from_typed_value() { let env = test_env(&[("FABRO_TEST_KEY_2", "val")]); let value = interp("x{{ env.FABRO_TEST_KEY_2 }}y"); - let result = resolve_interp(&value, &env).unwrap(); + let result = resolve_interp(&value, &env, &empty_secrets()).unwrap(); assert_eq!(result, "xvaly"); } #[test] fn resolve_interp_errors_on_missing_var() { let env = test_env(&[]); - let err = resolve_interp(&interp("a{{ env.FABRO_TEST_NOEXIST }}-b"), &env).unwrap_err(); + let err = resolve_interp( + &interp("a{{ env.FABRO_TEST_NOEXIST }}-b"), + &env, + &empty_secrets(), + ) + .unwrap_err(); assert_eq!(err.name, "FABRO_TEST_NOEXIST"); } @@ -1248,11 +1323,63 @@ mod tests { fn resolve_interp_without_tokens_passes_through() { let env = test_env(&[]); assert_eq!( - resolve_interp(&interp("plain text"), &env).unwrap(), + resolve_interp(&interp("plain text"), &env, &empty_secrets()).unwrap(), "plain text" ); } + #[test] + fn resolve_interp_resolves_secret_and_registers_value() { + let env = test_env(&[]); + let redactor = fabro_redact::SecretRedactor::default(); + let secrets = ResolvedHookSecrets::new( + HashMap::from([("HOOK_TOKEN".to_string(), "staging".to_string())]), + redactor.clone(), + ); + + let resolved = + resolve_interp(&interp("deploy {{ secrets.HOOK_TOKEN }}"), &env, &secrets).unwrap(); + + assert_eq!(resolved, "deploy staging"); + assert_eq!(redactor.redact_into("deploy staging"), "deploy REDACTED"); + } + + #[test] + fn resolve_prompt_and_model_resolves_secret_tokens() { + let env = test_env(&[]); + let secrets = test_secrets(&[("PROMPT_TOKEN", "staging")]); + + let (prompt, model) = HookExecutorImpl::resolve_prompt_and_model( + &interp("check {{ secrets.PROMPT_TOKEN }}"), + Some(&interp("haiku")), + &env, + &secrets, + ) + .unwrap(); + + assert_eq!(prompt, "check staging"); + assert_eq!(model.as_deref(), Some("haiku")); + } + + #[test] + fn header_rejects_secret_token_with_guidance() { + let env = test_env(&[]); + let err = resolve_header( + &interp("Bearer {{ secrets.HOOK_TOKEN }}"), + &[], + &env, + &test_secrets(&[("HOOK_TOKEN", "staging")]), + ) + .unwrap_err(); + + assert_eq!(err, HeaderResolveError::SecretNotAllowed { + name: "HOOK_TOKEN".to_string(), + }); + let message = err.to_string(); + assert!(message.contains("HTTP hook header")); + assert!(message.contains("command, prompt, or url")); + } + // --- HTTP hook execution tests --- #[tokio::test] @@ -1278,6 +1405,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1307,6 +1435,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1334,6 +1463,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1353,6 +1483,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(1), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1388,6 +1519,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &env, + &empty_secrets(), ) .await; @@ -1425,6 +1557,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &env, + &empty_secrets(), ) .await; @@ -1442,6 +1575,46 @@ mod tests { } } + #[tokio::test] + async fn http_hook_secret_header_blocks_without_firing() { + let server = httpmock::MockServer::start_async().await; + let mock = server + .mock_async(|when, then| { + when.method("POST").path("/hook"); + then.status(200).body(""); + }) + .await; + + let headers = HashMap::from([( + "Authorization".to_string(), + interp("Bearer {{ secrets.FABRO_TEST_TOKEN }}"), + )]); + + let client = test_http_client(); + let decision = HookExecutorImpl::execute_http( + &client, + &interp(&server.url("/hook")), + Some(&headers), + &[], + &TlsMode::Off, + &make_context(), + std::time::Duration::from_secs(5), + &test_env(&[]), + &test_secrets(&[("FABRO_TEST_TOKEN", "staging")]), + ) + .await; + + assert_eq!(mock.calls_async().await, 0); + match decision { + HookDecision::Block { reason } => { + let reason = reason.unwrap_or_default(); + assert!(reason.contains("HTTP hook header")); + assert!(reason.contains("command, prompt, or url")); + } + other => panic!("expected Block on secret header token, got {other:?}"), + } + } + #[tokio::test] async fn http_hook_resolves_url_before_dispatch() { let server = httpmock::MockServer::start_async().await; @@ -1463,6 +1636,36 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &env, + &empty_secrets(), + ) + .await; + + mock.assert_async().await; + assert_eq!(decision, HookDecision::Proceed); + } + + #[tokio::test] + async fn http_hook_resolves_secret_url_before_dispatch() { + let server = httpmock::MockServer::start_async().await; + let mock = server + .mock_async(|when, then| { + when.method("POST").path("/hook"); + then.status(200).body(""); + }) + .await; + + let client = test_http_client(); + let secrets = test_secrets(&[("FABRO_TEST_URL", &server.url("/hook"))]); + let decision = HookExecutorImpl::execute_http( + &client, + &interp("{{ secrets.FABRO_TEST_URL }}"), + None, + &[], + &TlsMode::Off, + &make_context(), + std::time::Duration::from_secs(5), + &test_env(&[]), + &secrets, ) .await; @@ -1490,6 +1693,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1534,6 +1738,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1557,6 +1762,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1575,6 +1781,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1601,6 +1808,7 @@ mod tests { &make_context(), std::time::Duration::from_secs(5), &test_env(&[]), + &empty_secrets(), ) .await; @@ -1637,6 +1845,7 @@ mod tests { let ctx = make_context(); let sandbox = make_sandbox(); let source = test_llm_source(); + let secrets = empty_secrets(); let result = executor .execute( &def, @@ -1645,6 +1854,7 @@ mod tests { &HookExecutionContext::default(), source.as_ref(), test_catalog(), + &secrets, ) .await; @@ -1663,12 +1873,57 @@ mod tests { &sandbox, &HookExecutionContext::default(), &test_env(&[]), + &empty_secrets(), ) .await; assert!(matches!(decision, HookDecision::Block { .. })); } + #[tokio::test] + async fn command_hook_resolves_secret_token() { + let sandbox = make_sandbox(); + let decision = HookExecutorImpl::execute_command( + &make_definition(r#"test "{{ secrets.HOOK_TOKEN }}" = "staging""#), + &interp(r#"test "{{ secrets.HOOK_TOKEN }}" = "staging""#), + &make_context(), + &sandbox, + &HookExecutionContext::default(), + &test_env(&[]), + &test_secrets(&[("HOOK_TOKEN", "staging")]), + ) + .await; + + assert_eq!(decision, HookDecision::Proceed); + } + + #[tokio::test] + async fn command_hook_missing_secret_blocks() { + let sandbox = make_sandbox(); + let decision = HookExecutorImpl::execute_command( + &make_definition("echo {{ secrets.MISSING_HOOK_SECRET }}"), + &interp("echo {{ secrets.MISSING_HOOK_SECRET }}"), + &make_context(), + &sandbox, + &HookExecutionContext::default(), + &test_env(&[]), + &empty_secrets(), + ) + .await; + + match decision { + HookDecision::Block { reason } => { + assert!( + reason + .as_deref() + .is_some_and(|reason| reason.contains("MISSING_HOOK_SECRET")), + "block reason should name the missing secret, got: {reason:?}" + ); + } + other => panic!("expected Block on missing command secret, got {other:?}"), + } + } + // Fail-closed: a prompt hook with a missing token does not fire the LLM // call; it blocks with the resolution error, matching command hooks. #[tokio::test] @@ -1679,6 +1934,7 @@ mod tests { None, &make_context(), &test_env(&[]), + &empty_secrets(), test_llm_source().as_ref(), test_catalog(), ) @@ -1708,6 +1964,7 @@ mod tests { &make_context(), make_sandbox(), &test_env(&[]), + &empty_secrets(), test_llm_source().as_ref(), test_catalog(), ) diff --git a/lib/crates/fabro-hooks/src/lib.rs b/lib/crates/fabro-hooks/src/lib.rs index 79a77bb98..d0d800f73 100644 --- a/lib/crates/fabro-hooks/src/lib.rs +++ b/lib/crates/fabro-hooks/src/lib.rs @@ -2,6 +2,7 @@ pub mod bridge; pub mod config; pub mod executor; pub mod runner; +mod secrets; pub mod types; pub use bridge::WorkflowToolHookCallback; @@ -10,4 +11,5 @@ pub use config::{HookDefinition, HookSettings, HookType, TlsMode}; // `InterpString`; constructing a hook definition requires it. pub use fabro_types::settings::InterpString; pub use runner::HookRunner; +pub use secrets::{HookSecretResolver, ResolvedHookSecrets}; pub use types::{HookContext, HookDecision, HookEvent, HookExecutionContext}; diff --git a/lib/crates/fabro-hooks/src/runner.rs b/lib/crates/fabro-hooks/src/runner.rs index 9435dba9a..fd005e0a5 100644 --- a/lib/crates/fabro-hooks/src/runner.rs +++ b/lib/crates/fabro-hooks/src/runner.rs @@ -6,10 +6,12 @@ use fabro_auth::CredentialSource; #[cfg(test)] use fabro_auth::EnvCredentialSource; use fabro_model::Catalog; +use fabro_redact::SecretRedactor; use crate::config::{HookDefinition, HookSettings}; use crate::executor::{HookExecutor, HookExecutorImpl}; -use crate::types::{HookContext, HookDecision, HookExecutionContext}; +use crate::secrets::HookSecretResolver; +use crate::types::{HookContext, HookDecision, HookExecutionContext, HookResult}; /// Central orchestrator: filters matching hooks, executes them, merges /// decisions. @@ -18,16 +20,53 @@ pub struct HookRunner { executor: Arc, llm_source: Arc, catalog: Arc, + secrets: HookSecretResolver, /// Pre-compiled regexes keyed by matcher pattern string. compiled_matchers: HashMap, } +fn decision_label(decision: &HookDecision) -> &'static str { + match decision { + HookDecision::Proceed => "proceed", + HookDecision::Skip { .. } => "skip", + HookDecision::Block { .. } => "block", + HookDecision::Override { .. } => "override", + } +} + +fn redact_hook_result(mut result: HookResult, redactor: &SecretRedactor) -> HookResult { + result.decision = redact_hook_decision(result.decision, redactor); + result +} + +fn redact_hook_decision(decision: HookDecision, redactor: &SecretRedactor) -> HookDecision { + match decision { + HookDecision::Skip { reason } => HookDecision::Skip { + reason: reason.map(|reason| redactor.redact_into(&reason)), + }, + HookDecision::Block { reason } => HookDecision::Block { + reason: reason.map(|reason| redactor.redact_into(&reason)), + }, + HookDecision::Proceed | HookDecision::Override { .. } => decision, + } +} + impl HookRunner { #[must_use] pub fn new( config: HookSettings, llm_source: Arc, catalog: Arc, + ) -> Self { + Self::new_with_secrets(config, llm_source, catalog, HookSecretResolver::default()) + } + + #[must_use] + pub fn new_with_secrets( + config: HookSettings, + llm_source: Arc, + catalog: Arc, + secrets: HookSecretResolver, ) -> Self { let compiled_matchers = Self::compile_matchers(&config); Self { @@ -35,6 +74,7 @@ impl HookRunner { executor: Arc::new(HookExecutorImpl), llm_source, catalog, + secrets, compiled_matchers, } } @@ -48,6 +88,7 @@ impl HookRunner { executor, llm_source: Arc::new(EnvCredentialSource::new()), catalog: Arc::new(Catalog::from_builtin().expect("default catalog should build")), + secrets: HookSecretResolver::default(), compiled_matchers, } } @@ -100,7 +141,7 @@ impl HookRunner { tracing::info!( event = %context.event, - decision = ?decision, + decision = decision_label(&decision), "Hooks complete" ); @@ -151,6 +192,7 @@ impl HookRunner { event = %context.event, "Executing hook" ); + let secrets = self.secrets.resolve_for_definition(hook).await; let result = self .executor .execute( @@ -160,12 +202,14 @@ impl HookRunner { 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 = ?result.decision, + decision = decision_label(&result.decision), "Hook complete" ); @@ -176,7 +220,7 @@ impl HookRunner { tracing::error!( hook = %hook.effective_name(), event = %context.event, - decision = ?merged, + decision = decision_label(&merged), "Hook blocked execution" ); return merged; @@ -185,7 +229,7 @@ impl HookRunner { tracing::warn!( hook = %hook.effective_name(), event = %context.event, - decision = ?result.decision, + decision = decision_label(&result.decision), "Non-blocking hook returned non-proceed, ignoring" ); } @@ -206,6 +250,7 @@ impl HookRunner { event = %context.event, "Executing hook" ); + let secrets = self.secrets.resolve_for_definition(hook).await; let result = self .executor .execute( @@ -215,19 +260,21 @@ impl HookRunner { 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 = ?result.decision, + decision = decision_label(&result.decision), "Hook complete" ); if !result.decision.is_proceed() { tracing::warn!( hook = %hook.effective_name(), event = %context.event, - decision = ?result.decision, + decision = decision_label(&result.decision), "Non-blocking hook failed, continuing" ); } @@ -259,6 +306,7 @@ mod tests { _execution_context: &HookExecutionContext, _llm_source: &dyn CredentialSource, _catalog: Arc, + _secrets: &crate::ResolvedHookSecrets, ) -> HookResult { HookResult { hook_name: definition.name.clone(), diff --git a/lib/crates/fabro-hooks/src/secrets.rs b/lib/crates/fabro-hooks/src/secrets.rs new file mode 100644 index 000000000..399af49d2 --- /dev/null +++ b/lib/crates/fabro-hooks/src/secrets.rs @@ -0,0 +1,155 @@ +use std::collections::HashMap; +use std::future::Future; +use std::pin::Pin; +use std::sync::Arc; + +use fabro_redact::SecretRedactor; +use fabro_types::settings::interp::Namespace; + +use crate::config::{HookDefinition, HookType}; + +type SecretLookupFuture = Pin> + Send + 'static>>; +type SecretLookup = dyn Fn(String) -> SecretLookupFuture + Send + Sync + 'static; + +/// Per-run hook secret resolver. +/// +/// The resolver is cheap to clone and must be constructed per run. It returns +/// only token-shaped vault secrets supplied by the worker and shares the run's +/// [`SecretRedactor`] so values resolved by hooks join the same redaction +/// registry as run-boundary environment and prepare-step secrets. +#[derive(Clone)] +pub struct HookSecretResolver { + lookup: Option>, + redactor: SecretRedactor, +} + +impl HookSecretResolver { + #[must_use] + pub fn new(redactor: SecretRedactor) -> Self { + Self { + lookup: None, + redactor, + } + } + + #[must_use] + pub fn with_lookup(redactor: SecretRedactor, lookup: F) -> Self + where + F: Fn(String) -> Fut + Send + Sync + 'static, + Fut: Future> + Send + 'static, + { + Self { + lookup: Some(Arc::new(move |name| Box::pin(lookup(name)))), + redactor, + } + } + + #[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 + } + + async fn resolve_names(&self, names: I) -> ResolvedHookSecrets + where + I: IntoIterator, + S: AsRef, + { + let mut values = HashMap::new(); + let Some(lookup) = self.lookup.as_ref() else { + return ResolvedHookSecrets::new(values, self.redactor.clone()); + }; + + for name in names { + let name = name.as_ref(); + if values.contains_key(name) { + continue; + } + if let Some(value) = lookup(name.to_string()).await { + values.insert(name.to_string(), value); + } + } + + ResolvedHookSecrets::new(values, self.redactor.clone()) + } +} + +impl Default for HookSecretResolver { + fn default() -> Self { + Self::new(SecretRedactor::default()) + } +} + +/// Secrets resolved for one hook firing. +#[derive(Clone, Default)] +pub struct ResolvedHookSecrets { + values: HashMap, + redactor: SecretRedactor, +} + +impl ResolvedHookSecrets { + #[must_use] + pub fn new(values: HashMap, redactor: SecretRedactor) -> Self { + for value in values.values() { + redactor.register(value); + } + 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 + } + + #[must_use] + pub fn redactor(&self) -> &SecretRedactor { + &self.redactor + } +} + +fn secret_names_for_definition(definition: &HookDefinition) -> Vec { + let Some(hook_type) = definition.resolved_hook_type() else { + return Vec::new(); + }; + match hook_type.as_ref() { + HookType::Command { command } => command + .names(Namespace::Secrets) + .into_iter() + .map(str::to_string) + .collect(), + HookType::Http { url, .. } => url + .names(Namespace::Secrets) + .into_iter() + .map(str::to_string) + .collect(), + HookType::Prompt { prompt, model } | HookType::Agent { prompt, model, .. } => { + let mut names: Vec = prompt + .names(Namespace::Secrets) + .into_iter() + .map(str::to_string) + .collect(); + if let Some(model) = model { + names.extend( + model + .names(Namespace::Secrets) + .into_iter() + .map(str::to_string), + ); + } + names + } + } +} diff --git a/lib/crates/fabro-sandbox/src/error.rs b/lib/crates/fabro-sandbox/src/error.rs index 7b2377660..a3c376f3b 100644 --- a/lib/crates/fabro-sandbox/src/error.rs +++ b/lib/crates/fabro-sandbox/src/error.rs @@ -2,6 +2,7 @@ use std::fmt::Write as _; #[cfg(feature = "docker")] use bollard::errors::Error as BollardError; +use fabro_redact::SecretRedactor; use fabro_util::error::{collect_causes, render_with_causes}; use crate::ExecResult; @@ -79,6 +80,13 @@ impl Error { default_redacted_output_tail(self) } + pub fn default_redacted_output_tail_with_redactor( + &self, + redactor: &SecretRedactor, + ) -> Option { + default_redacted_output_tail_with_redactor(self, redactor) + } + #[cfg(feature = "docker")] pub fn docker_connect(source: BollardError) -> Self { Self::DockerConnect { source } @@ -163,11 +171,28 @@ pub type Result = std::result::Result; pub fn default_redacted_output_tail( err: &(dyn std::error::Error + 'static), +) -> Option { + default_redacted_output_tail_inner(err, None) +} + +pub fn default_redacted_output_tail_with_redactor( + err: &(dyn std::error::Error + 'static), + redactor: &SecretRedactor, +) -> Option { + default_redacted_output_tail_inner(err, Some(redactor)) +} + +fn default_redacted_output_tail_inner( + err: &(dyn std::error::Error + 'static), + redactor: Option<&SecretRedactor>, ) -> Option { let mut current = Some(err); while let Some(err) = current { if let Some(Error::Exec { result, .. }) = err.downcast_ref::() { - return result.default_redacted_output_tail(); + return match redactor { + Some(redactor) => result.default_redacted_output_tail_with_redactor(redactor), + None => result.default_redacted_output_tail(), + }; } current = err.source(); } @@ -175,8 +200,26 @@ pub fn default_redacted_output_tail( } pub fn display_for_log(err: &(dyn std::error::Error + 'static)) -> String { + display_for_log_inner(err, None) +} + +pub fn display_for_log_with_redactor( + err: &(dyn std::error::Error + 'static), + redactor: &SecretRedactor, +) -> String { + display_for_log_inner(err, Some(redactor)) +} + +fn display_for_log_inner( + err: &(dyn std::error::Error + 'static), + redactor: Option<&SecretRedactor>, +) -> String { let mut rendered = render_with_causes(&err.to_string(), &collect_causes(err)); - if let Some(tail) = default_redacted_output_tail(err) { + let tail = match redactor { + Some(redactor) => default_redacted_output_tail_with_redactor(err, redactor), + None => default_redacted_output_tail(err), + }; + if let Some(tail) = tail { append_tail_for_log( &mut rendered, "stderr", diff --git a/lib/crates/fabro-sandbox/src/lib.rs b/lib/crates/fabro-sandbox/src/lib.rs index e53b25756..37b1809a1 100644 --- a/lib/crates/fabro-sandbox/src/lib.rs +++ b/lib/crates/fabro-sandbox/src/lib.rs @@ -39,7 +39,10 @@ pub mod test_support; pub use details::sandbox_details; #[cfg(feature = "docker")] pub use docker::{DockerSandbox, DockerSandboxOptions}; -pub use error::{Error, Result, default_redacted_output_tail, display_for_log}; +pub use error::{ + Error, Result, default_redacted_output_tail, default_redacted_output_tail_with_redactor, + display_for_log, display_for_log_with_redactor, +}; pub use fabro_types::{RunSandboxInstance, SandboxProviderKind}; pub use local::LocalSandbox; #[cfg(feature = "daytona")] @@ -57,7 +60,7 @@ pub use sandbox::{ ExecStreamingResult, GitRunInfo, GitSetupIntent, GrepOptions, Sandbox, SandboxEvent, SandboxEventCallback, StderrCollector, StdioProcess, StdioProcessHandle, StdioProcessTermination, format_lines_numbered, git_push_via_exec, redacted_output_tail, - setup_git_via_exec, shell_quote, + redacted_output_tail_with_redactor, setup_git_via_exec, shell_quote, }; pub use sandbox_spec::SandboxSpec; pub use terminal::{TerminalSession, TerminalSize, open_terminal_for_run}; diff --git a/lib/crates/fabro-sandbox/src/sandbox.rs b/lib/crates/fabro-sandbox/src/sandbox.rs index 68801f8bb..3db0d846a 100644 --- a/lib/crates/fabro-sandbox/src/sandbox.rs +++ b/lib/crates/fabro-sandbox/src/sandbox.rs @@ -7,6 +7,7 @@ use std::sync::Arc; use std::time::Duration; use async_trait::async_trait; +use fabro_redact::SecretRedactor; use fabro_types::{CommandOutputStream, CommandTermination}; use fabro_util::shell; use serde::{Deserialize, Serialize}; @@ -572,10 +573,30 @@ impl ExecResult { redacted_output_tail(&self.stdout, &self.stderr, max_bytes_per_stream) } + pub fn redacted_output_tail_with_redactor( + &self, + max_bytes_per_stream: usize, + redactor: &SecretRedactor, + ) -> Option { + redacted_output_tail_with_redactor( + &self.stdout, + &self.stderr, + max_bytes_per_stream, + redactor, + ) + } + pub fn default_redacted_output_tail(&self) -> Option { self.redacted_output_tail(DEFAULT_EXEC_OUTPUT_TAIL_BYTES) } + pub fn default_redacted_output_tail_with_redactor( + &self, + redactor: &SecretRedactor, + ) -> Option { + self.redacted_output_tail_with_redactor(DEFAULT_EXEC_OUTPUT_TAIL_BYTES, redactor) + } + /// Converts host process output into the canonical full exec result. /// /// This stores raw stdout/stderr. Callers must not log these fields @@ -607,8 +628,29 @@ pub fn redacted_output_tail( stderr: &str, max_bytes_per_stream: usize, ) -> Option { - let (stdout, stdout_truncated) = redacted_tail(stdout, max_bytes_per_stream); - let (stderr, stderr_truncated) = redacted_tail(stderr, max_bytes_per_stream); + redacted_output_tail_inner(stdout, stderr, max_bytes_per_stream, None) +} + +/// Build a redacted `ExecOutputTail` with a run-scoped exact-match secret +/// redactor applied after the content-based redaction baseline. +#[must_use] +pub fn redacted_output_tail_with_redactor( + stdout: &str, + stderr: &str, + max_bytes_per_stream: usize, + redactor: &SecretRedactor, +) -> Option { + redacted_output_tail_inner(stdout, stderr, max_bytes_per_stream, Some(redactor)) +} + +fn redacted_output_tail_inner( + stdout: &str, + stderr: &str, + max_bytes_per_stream: usize, + redactor: Option<&SecretRedactor>, +) -> Option { + let (stdout, stdout_truncated) = redacted_tail(stdout, max_bytes_per_stream, redactor); + let (stderr, stderr_truncated) = redacted_tail(stderr, max_bytes_per_stream, redactor); let tail = fabro_types::ExecOutputTail { stdout, stderr, @@ -618,12 +660,19 @@ pub fn redacted_output_tail( (!tail.is_empty()).then_some(tail) } -fn redacted_tail(text: &str, max_bytes: usize) -> (Option, bool) { +fn redacted_tail( + text: &str, + max_bytes: usize, + redactor: Option<&SecretRedactor>, +) -> (Option, bool) { if text.is_empty() || max_bytes == 0 { return (None, !text.is_empty()); } - let redacted = fabro_redact::redact_string(text); + let mut redacted = fabro_redact::redact_string(text); + if let Some(redactor) = redactor { + redacted = redactor.redact_into(&redacted); + } let sanitized = sanitize_exec_output(&redacted); let truncated = sanitized.len() > max_bytes; let start = if truncated { @@ -1282,6 +1331,27 @@ mod tests { assert!(tail.stdout_truncated); } + #[test] + fn exec_result_redacts_registered_low_entropy_secret_before_taking_tail() { + let redactor = SecretRedactor::default(); + redactor.register("staging"); + let result = ExecResult { + stdout: format!("{} staging done", "context ".repeat(20)), + stderr: String::new(), + exit_code: Some(1), + termination: CommandTermination::Exited, + duration_ms: 1, + }; + + let tail = result + .redacted_output_tail_with_redactor(32, &redactor) + .expect("redacted output tail"); + let stdout = tail.stdout.expect("stdout tail"); + assert!(stdout.contains("REDACTED"), "{stdout}"); + assert!(!stdout.contains("staging"), "{stdout}"); + assert!(tail.stdout_truncated); + } + #[test] fn exec_result_tail_sanitizes_terminal_control_sequences() { let result = ExecResult { diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index e64b2ae30..8ebf998f5 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -15,7 +15,8 @@ pub use self::emitter::Emitter; pub use self::events::Event; pub use self::names::event_name; pub use self::redaction::{ - build_redacted_event_payload, event_payload_from_redacted_json, redacted_event_json, + build_redacted_event_payload, build_redacted_event_payload_with_redactor, + event_payload_from_redacted_json, redacted_event_json, redacted_event_json_with_redactor, }; pub use self::sink::{ RunEventLogger, RunEventSink, StoreProgressLogger, append_event, append_event_to_sink, diff --git a/lib/crates/fabro-workflow/src/event/events.rs b/lib/crates/fabro-workflow/src/event/events.rs index 0b5afb1c7..85707d0a2 100644 --- a/lib/crates/fabro-workflow/src/event/events.rs +++ b/lib/crates/fabro-workflow/src/event/events.rs @@ -1330,7 +1330,7 @@ impl Event { info!(command_count, "Setup started"); } Self::SetupCommandStarted { command, index } => { - debug!(command, index, "Setup command started"); + debug!(command_len = command.len(), index, "Setup command started"); } Self::SetupCommandCompleted { command, @@ -1339,7 +1339,7 @@ impl Event { duration_ms, } => { debug!( - command, + command_len = command.len(), index, exit_code, duration_ms, "Setup command completed" ); } @@ -1355,7 +1355,7 @@ impl Event { } => { let tail = fabro_types::ExecOutputTail::trace_summary(exec_output_tail.as_ref()); error!( - command, + command_len = command.len(), index, exit_code, exec_output_tail_present = tail.present, diff --git a/lib/crates/fabro-workflow/src/event/redaction.rs b/lib/crates/fabro-workflow/src/event/redaction.rs index 5ae44d111..858818df3 100644 --- a/lib/crates/fabro-workflow/src/event/redaction.rs +++ b/lib/crates/fabro-workflow/src/event/redaction.rs @@ -1,17 +1,32 @@ use ::fabro_types::{RunEvent, RunId}; use anyhow::{Context, Result}; -use fabro_redact::redact_json_value; +use fabro_redact::{SecretRedactor, redact_json_value}; use fabro_store::EventPayload; use fabro_util::json::normalize_json_value; use serde_json::Value; pub fn build_redacted_event_payload(event: &RunEvent, run_id: &RunId) -> Result { - let value = redacted_event_value(event)?; + build_redacted_event_payload_with_redactor(event, run_id, None) +} + +pub fn build_redacted_event_payload_with_redactor( + event: &RunEvent, + run_id: &RunId, + redactor: Option<&SecretRedactor>, +) -> Result { + let value = redacted_event_value(event, redactor)?; EventPayload::new(value, run_id).map_err(anyhow::Error::from) } pub fn redacted_event_json(event: &RunEvent) -> Result { - serde_json::to_string(&redacted_event_value(event)?).map_err(anyhow::Error::from) + redacted_event_json_with_redactor(event, None) +} + +pub fn redacted_event_json_with_redactor( + event: &RunEvent, + redactor: Option<&SecretRedactor>, +) -> Result { + serde_json::to_string(&redacted_event_value(event, redactor)?).map_err(anyhow::Error::from) } fn normalized_event_value(event: &RunEvent) -> Result { @@ -19,8 +34,90 @@ fn normalized_event_value(event: &RunEvent) -> Result { Ok(normalize_json_value(value)) } -fn redacted_event_value(event: &RunEvent) -> Result { - Ok(redact_json_value(normalized_event_value(event)?)) +fn redacted_event_value(event: &RunEvent, redactor: Option<&SecretRedactor>) -> Result { + let mut value = redact_json_value(normalized_event_value(event)?); + if let Some(redactor) = redactor { + redact_event_payload_secrets(&mut value, redactor); + } + Ok(value) +} + +fn redact_event_payload_secrets(value: &mut Value, redactor: &SecretRedactor) { + if let Some(properties) = value.get_mut("properties") { + redact_redactable_event_properties(properties, redactor); + } + if let Some(Value::String(label)) = value.get_mut("node_label") { + let redacted = redactor.redact_into(label); + if redacted != *label { + *label = redacted; + } + } +} + +fn redact_redactable_event_properties(value: &mut Value, redactor: &SecretRedactor) { + match value { + Value::Object(obj) => { + for (key, child) in obj { + if is_secret_redactable_event_property(key) { + *child = redactor.redact_json(std::mem::take(child)); + } else { + redact_redactable_event_properties(child, redactor); + } + } + } + Value::Array(items) => { + for item in items { + redact_redactable_event_properties(item, redactor); + } + } + Value::Null | Value::Bool(_) | Value::Number(_) | Value::String(_) => {} + } +} + +// Exact-match secret values may be intentionally low entropy ("staging", +// "pause", "running"). Redacting every string in an event can therefore corrupt +// structural fields that are validated enum values or IDs. Keep this list to +// free-form text/blob fields where replacing a matched substring preserves the +// event schema and projection semantics. +fn is_secret_redactable_event_property(key: &str) -> bool { + matches!( + key, + "active_form" + | "answer" + | "arguments" + | "causes" + | "command" + | "context_display" + | "delta" + | "description" + | "details" + | "diff" + | "error" + | "error_message" + | "exec_output_tail" + | "failure" + | "final_patch" + | "goal" + | "input" + | "message" + | "notes" + | "output" + | "preview" + | "prompt" + | "question" + | "reason" + | "response" + | "script" + | "stderr" + | "stdout" + | "subject" + | "text" + | "title" + | "tool_input" + | "tool_output" + | "workflow_config" + | "workflow_source" + ) } pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result { @@ -30,7 +127,7 @@ pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result, inner: Box, }, + RedactSecrets { + redactor: SecretRedactor, + inner: Box, + }, Composite(Vec), } @@ -99,27 +107,60 @@ impl RunEventSink { } } + #[must_use] + pub fn with_secret_redactor(self, redactor: SecretRedactor) -> Self { + Self::RedactSecrets { + redactor, + inner: Box::new(self), + } + } + pub async fn write_run_event(&self, event: &RunEvent) -> Result<()> { - let mut pending = vec![(self, event.clone())]; - while let Some((sink, event)) = pending.pop() { + let mut pending = vec![(self, event.clone(), None::)]; + while let Some((sink, event, redactor)) = pending.pop() { match sink { Self::Store(run_store) => { - run_store.append_run_event(&event).await?; + run_store + .append_run_event_with_redactor(&event, redactor.as_ref()) + .await?; } Self::JsonLines(writer) => { - let line = redacted_event_json(&event)?; + let line = match redactor.as_ref() { + Some(redactor) => { + redacted_event_json_with_redactor(&event, Some(redactor))? + } + None => redacted_event_json(&event)?, + }; let mut writer = writer.lock().await; writer.write_all(line.as_bytes()).await?; writer.write_all(b"\n").await?; writer.flush().await?; } - Self::Callback(callback) => callback(event).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 + }; + callback(event).await?; + } Self::Map { transform, inner } => { - pending.push((inner.as_ref(), transform(event))); + pending.push((inner.as_ref(), transform(event), redactor)); + } + Self::RedactSecrets { + redactor: sink_redactor, + inner, + } => { + pending.push((inner.as_ref(), event, Some(sink_redactor.clone()))); } Self::Composite(sinks) => { for sink in sinks.iter().rev() { - pending.push((sink, event.clone())); + pending.push((sink, event.clone(), redactor.clone())); } } } @@ -334,6 +375,34 @@ mod tests { assert_eq!(second[0].actor, Some(user_principal("alice"))); } + #[tokio::test] + async fn run_event_sink_redacts_callback_events_with_run_secret_redactor() { + let captured = Arc::new(AsyncMutex::new(Vec::new())); + let captured_events = Arc::clone(&captured); + let redactor = fabro_redact::SecretRedactor::default(); + redactor.register("staging"); + let sink = RunEventSink::callback(move |event| { + let captured_events = Arc::clone(&captured_events); + async move { + captured_events.lock().await.push(event); + Ok(()) + } + }) + .with_secret_redactor(redactor); + let event = to_run_event(&fixtures::RUN_7, &Event::SetupCommandStarted { + command: "deploy staging".to_string(), + index: 0, + }); + + sink.write_run_event(&event).await.unwrap(); + + let captured = captured.lock().await; + assert_eq!(captured.len(), 1); + let captured_text = serde_json::to_string(&captured[0].to_value().unwrap()).unwrap(); + assert!(!captured_text.contains("staging")); + assert!(captured_text.contains("REDACTED")); + } + #[tokio::test] async fn run_event_logger_registers_emitter_events_to_json_lines() { use tokio::io::{AsyncBufReadExt, BufReader}; diff --git a/lib/crates/fabro-workflow/src/handler/parallel.rs b/lib/crates/fabro-workflow/src/handler/parallel.rs index 108a72369..b5f216100 100644 --- a/lib/crates/fabro-workflow/src/handler/parallel.rs +++ b/lib/crates/fabro-workflow/src/handler/parallel.rs @@ -205,14 +205,20 @@ impl Handler for ParallelHandler { } Err(e) => { tracing::warn!( - error = %fabro_sandbox::display_for_log(&e), + error = %fabro_sandbox::display_for_log_with_redactor( + &e, + &services.run.secret_redactor, + ), "parallel base checkpoint failed" ); services.run.emitter.notice_with_tail( RunNoticeLevel::Warn, RunNoticeCode::ParallelBaseCheckpointFailed, format!("Could not checkpoint base state before parallel branches: {e}"), - fabro_sandbox::default_redacted_output_tail(&e), + fabro_sandbox::default_redacted_output_tail_with_redactor( + &e, + &services.run.secret_redactor, + ), ); None } diff --git a/lib/crates/fabro-workflow/src/lifecycle/git.rs b/lib/crates/fabro-workflow/src/lifecycle/git.rs index 5133fa602..855e8de2f 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/git.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/git.rs @@ -8,6 +8,7 @@ use fabro_core::lifecycle::RunLifecycle; use fabro_core::outcome::NodeResult; use fabro_core::state::ExecutionState; use fabro_dump::RunDump; +use fabro_redact::SecretRedactor; use fabro_types::run_event::{MetadataSnapshotFailureKind, MetadataSnapshotPhase}; use fabro_types::{CheckpointRecord, DiffSummary, RunDiff, RunId}; use fabro_util::error::collect_causes; @@ -84,6 +85,7 @@ pub(crate) struct GitLifecycle { pub sandbox_git: Arc, pub metadata_runtime: Arc, pub metadata_writer: Option, + pub secret_redactor: SecretRedactor, pub start_node_id: Option, // Cross-lifecycle data (shared with EventLifecycle) pub checkpoint_git_result: Arc>>, @@ -311,10 +313,16 @@ impl RunLifecycle for GitLifecycle { Ok(()) => (true, None), Err(err) => { let exec_output_tail = - fabro_sandbox::default_redacted_output_tail(&err); + fabro_sandbox::default_redacted_output_tail_with_redactor( + &err, + &self.secret_redactor, + ); tracing::warn!( refspec = %refspec, - error = %fabro_sandbox::display_for_log(&err), + error = %fabro_sandbox::display_for_log_with_redactor( + &err, + &self.secret_redactor, + ), "git push from run lifecycle failed" ); self.emitter.notice_with_tail( @@ -365,7 +373,10 @@ impl RunLifecycle for GitLifecycle { Ok(_) => {} Err(err) => { let exec_output_tail = - fabro_sandbox::default_redacted_output_tail(&err); + fabro_sandbox::default_redacted_output_tail_with_redactor( + &err, + &self.secret_redactor, + ); self.emitter.notice_with_tail( RunNoticeLevel::Warn, RunNoticeCode::GitDiffFailed, @@ -380,7 +391,10 @@ impl RunLifecycle for GitLifecycle { } Some(Err(err)) => { let exec_output_tail = - fabro_sandbox::default_redacted_output_tail(&err); + fabro_sandbox::default_redacted_output_tail_with_redactor( + &err, + &self.secret_redactor, + ); self.emitter.notice_with_tail( RunNoticeLevel::Warn, RunNoticeCode::GitDiffFailed, @@ -399,7 +413,10 @@ impl RunLifecycle for GitLifecycle { .expect("git lifecycle mutex should not be poisoned: no code panics while holding this lock") = Some(git_result); } Err(e) => { - let exec_output_tail = fabro_sandbox::default_redacted_output_tail(&e); + let exec_output_tail = fabro_sandbox::default_redacted_output_tail_with_redactor( + &e, + &self.secret_redactor, + ); let error = e.to_string(); // Emit CheckpointFailed and return error let scope = stage_scope_for(state, node_id); @@ -797,6 +814,7 @@ mod tests { sandbox_git: Arc::new(SandboxGitRuntime::new()), metadata_runtime, metadata_writer, + secret_redactor: SecretRedactor::default(), start_node_id: Some("start".to_string()), checkpoint_git_result: Arc::new(Mutex::new(None)), last_git_sha: Arc::new(Mutex::new(None)), @@ -1260,6 +1278,7 @@ mod tests { "claude-sonnet-4-6".to_string(), Arc::new(fabro_auth::EnvCredentialSource::new()), Arc::new(Catalog::from_builtin().expect("default catalog should build")), + fabro_redact::SecretRedactor::default(), Arc::new(SandboxGitRuntime::new()), Arc::clone(&lifecycle.metadata_runtime), lifecycle.metadata_writer.clone(), diff --git a/lib/crates/fabro-workflow/src/lifecycle/mod.rs b/lib/crates/fabro-workflow/src/lifecycle/mod.rs index 2db1e8d75..97b46a5a8 100644 --- a/lib/crates/fabro-workflow/src/lifecycle/mod.rs +++ b/lib/crates/fabro-workflow/src/lifecycle/mod.rs @@ -94,6 +94,7 @@ impl WorkflowLifecycle { sandbox_git: Arc, metadata_runtime: Arc, metadata_writer: Option, + secret_redactor: fabro_redact::SecretRedactor, is_resume: bool, on_node: crate::OnNodeCallback, run_control: Option>, @@ -161,6 +162,7 @@ impl WorkflowLifecycle { sandbox_git, metadata_runtime, metadata_writer, + secret_redactor, start_node_id, checkpoint_git_result: Arc::clone(&checkpoint_git_result), last_git_sha, diff --git a/lib/crates/fabro-workflow/src/operations/start.rs b/lib/crates/fabro-workflow/src/operations/start.rs index 7fa47e938..9036e8032 100644 --- a/lib/crates/fabro-workflow/src/operations/start.rs +++ b/lib/crates/fabro-workflow/src/operations/start.rs @@ -8,6 +8,7 @@ use fabro_interview::{AutoApproveInterviewer, Interviewer}; use fabro_llm::client::Client as LlmClient; use fabro_mcp::config::McpServerSettings; use fabro_model::{Catalog, FallbackTarget, ProviderId}; +use fabro_redact::SecretRedactor; use fabro_sandbox::daytona::DaytonaConfig; use fabro_sandbox::from_environment::{ daytona_config_from_environment, docker_config_from_environment_with_secrets, @@ -79,6 +80,8 @@ struct RunSession { workflow_bundle: Option>, run_control: Option>, vault: Option>>, + secret_redactor: SecretRedactor, + hook_secrets: fabro_hooks::HookSecretResolver, catalog: Arc, fabro_run_tools: Option, } @@ -252,10 +255,13 @@ pub(super) async fn execute_persisted_run( }; bootstrap_guard.defuse(); + let terminal_event_sink = event_sink + .clone() + .with_secret_redactor(session.secret_redactor.clone()); let mut completion_guard = DetachedRunCompletionGuard::arm( run_id, run_store.clone(), - event_sink.clone(), + terminal_event_sink.clone(), cancel_token, ); let run_start = Instant::now(); @@ -270,7 +276,7 @@ pub(super) async fn execute_persisted_run( persist_terminal_engine_failure( run_id, &run_store, - &event_sink, + &terminal_event_sink, run_dir, &err, run_start.elapsed(), @@ -377,17 +383,21 @@ impl RunSession { Some(vault) => Some(vault.read().await), None => None, }; + let secret_redactor = SecretRedactor::default(); + let hook_secrets = hook_secret_resolver(services.vault.clone(), secret_redactor.clone()); // Token-only secrets lookup over the vault read guard, shared across // every run-boundary resolver. A missing or non-Token secret becomes // `None`, so resolution fails closed with a secret error. - let secret_lookup = |name: &str| vault_token_lookup(vault_guard.as_deref(), name); + let secret_lookup = |name: &str| { + registered_vault_token_lookup(vault_guard.as_deref(), &secret_redactor, name) + }; let mcp_servers = resolved .agent .mcps .iter() .map(|(key, entry)| match entry { ResolvedMcpEntry::Resolved(server) => { - runtime_mcp_server(server, process_env_var, secret_lookup) + runtime_mcp_server(server, process_env_var, |name| secret_lookup(name)) } // References must be resolved to concrete servers before the run // spec is persisted (server-side run-preparation pass). Reaching @@ -419,7 +429,7 @@ impl RunSession { SandboxSpec::Local { working_directory } } SandboxProviderKind::Docker => SandboxSpec::Docker { - config: resolve_docker_config(resolved, secret_lookup)?, + config: resolve_docker_config(resolved, |name| secret_lookup(name))?, github_app: services.github_app.clone(), run_id: Some(record.run_id), clone_origin_url: record.repo_origin_url().map(str::to_string), @@ -443,7 +453,7 @@ impl RunSession { let toml_env = resolved .environment - .resolve_env(process_env_var, secret_lookup) + .resolve_env(process_env_var, |name| secret_lookup(name)) .map_err(|err| Error::engine_with_source("failed to resolve run environment", err))?; let github_permissions: Option> = (!services.github_permissions.is_empty()).then(|| services.github_permissions.clone()); @@ -461,8 +471,9 @@ impl RunSession { }; let pr_config = resolved.pull_request.clone(); - let setup_commands = - runtime_setup_commands(&resolved.prepare, process_env_var, secret_lookup)?; + let setup_commands = runtime_setup_commands(&resolved.prepare, process_env_var, |name| { + secret_lookup(name) + })?; drop(vault_guard); Ok(Self { @@ -505,6 +516,8 @@ impl RunSession { workflow_path, workflow_bundle, vault: services.vault, + secret_redactor, + hook_secrets, catalog, fabro_run_tools: services.fabro_run_tools, }) @@ -566,6 +579,34 @@ fn vault_token_lookup(vault: Option<&Vault>, name: &str) -> Option { vault.and_then(|vault| fabro_auth::vault_get_token(vault, name).ok().flatten()) } +fn registered_vault_token_lookup( + vault: Option<&Vault>, + redactor: &SecretRedactor, + name: &str, +) -> Option { + let value = vault_token_lookup(vault, name); + if let Some(value) = value.as_deref() { + redactor.register(value); + } + value +} + +fn hook_secret_resolver( + vault: Option>>, + redactor: SecretRedactor, +) -> fabro_hooks::HookSecretResolver { + match vault { + Some(vault) => fabro_hooks::HookSecretResolver::with_lookup(redactor, move |name| { + let vault = Arc::clone(&vault); + async move { + let guard = vault.read().await; + vault_token_lookup(Some(&guard), &name) + } + }), + None => fabro_hooks::HookSecretResolver::new(redactor), + } +} + async fn load_accepted_run_definition( run_store: &RunStoreHandle, blob_id: fabro_types::RunBlobId, @@ -827,7 +868,11 @@ impl RunSession { }); } - let store_progress_logger = RunEventLogger::new(self.event_sink.clone()); + let store_progress_logger = RunEventLogger::new( + self.event_sink + .clone() + .with_secret_redactor(self.secret_redactor.clone()), + ); store_progress_logger.register(self.emitter.as_ref()); let init_options = InitOptions { @@ -845,8 +890,10 @@ impl RunSession { workflow_path: self.workflow_path, workflow_bundle: self.workflow_bundle, hooks: self.hooks, + hook_secrets: self.hook_secrets, sandbox_env: self.sandbox_env, vault: self.vault, + secret_redactor: self.secret_redactor, git: self.git, registry_override: self.registry_override, artifact_sink: self.artifact_sink, @@ -855,7 +902,13 @@ impl RunSession { seed_context: self.seed_context, fabro_run_tools: self.fabro_run_tools, }; - let mut initialized = Box::pin(pipeline::initialize(persisted, init_options)).await?; + let mut initialized = match Box::pin(pipeline::initialize(persisted, init_options)).await { + Ok(initialized) => initialized, + Err(err) => { + store_progress_logger.flush().await; + return Err(err); + } + }; initialized.on_node = on_node; let sandbox_for_cleanup = Arc::clone(&initialized.engine.run.sandbox); @@ -1125,8 +1178,8 @@ mod tests { }; use fabro_store::Database; use fabro_types::settings::run::{ - McpTransport as ResolvedMcpTransport, PreparedStep, PreparedStepRun, RunMode, - RunPrepareSettings, + HookDefinition, HookEvent, HookType, McpTransport as ResolvedMcpTransport, PreparedStep, + PreparedStepRun, RunMode, RunPrepareSettings, TlsMode, }; use fabro_types::settings::{InterpString, ModelRef}; use fabro_types::{ @@ -1568,6 +1621,49 @@ reasoning = false ); } + #[tokio::test] + async fn run_session_secret_redactors_are_isolated_between_runs() { + async fn session_with_secret(secret_value: &str) -> RunSession { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, _run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings.run.environment.env.insert( + "DEPLOY_ENV".to_string(), + InterpString::parse("{{ secrets.DEPLOY_ENV }}"), + ); + let (persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(token_vault("DEPLOY_ENV", secret_value))); + + RunSession::new(&persisted, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + .unwrap() + } + + let first = session_with_secret("alpha").await; + let second = session_with_secret("bravo").await; + + assert_eq!( + first.secret_redactor.redact_into("alpha bravo"), + "REDACTED bravo" + ); + assert_eq!( + second.secret_redactor.redact_into("alpha bravo"), + "alpha REDACTED" + ); + } + #[tokio::test] async fn run_session_new_missing_secret_fails_startup() { let temp = tempfile::tempdir().unwrap(); @@ -1605,6 +1701,218 @@ reasoning = false assert!(err.causes()[0].contains("DEPLOY_TOKEN")); } + #[tokio::test] + async fn setup_failure_redacts_low_entropy_secret_in_event_and_error() { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings.run.prepare = prepare_with_step(script_step( + "echo {{ secrets.DEPLOY_ENV }} >&2; exit 7", + HashMap::new(), + )); + let (_persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(token_vault("DEPLOY_ENV", "staging"))); + + let Err(err) = start(&run_dir, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + else { + panic!("setup failure should fail the run"); + }; + + let error_text = err.to_string(); + assert!(!error_text.contains("staging")); + assert!(error_text.contains("REDACTED")); + + let run_store = store.open_run(&fixtures::RUN_1).await.unwrap(); + let events = run_store.list_events().await.unwrap(); + let events_text = serde_json::to_string(&events).unwrap(); + assert!(!events_text.contains("staging")); + assert!(events_text.contains("REDACTED")); + assert!( + events + .iter() + .any(|event| event.event.event_name() == "setup.failed") + ); + } + + #[tokio::test] + async fn hook_command_resolves_secret_from_vault() { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings.run.hooks.push(sandbox_ready_command_hook( + "test \"{{ secrets.HOOK_TOKEN }}\" = staging", + )); + let (_persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(token_vault("HOOK_TOKEN", "staging"))); + + start(&run_dir, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + .expect("hook command should resolve secret and proceed"); + } + + #[tokio::test] + async fn hook_missing_secret_fails_closed() { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings + .run + .hooks + .push(sandbox_ready_command_hook("echo {{ secrets.HOOK_TOKEN }}")); + let (_persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(temp_vault(&[]))); + + let Err(err) = start(&run_dir, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + else { + panic!("missing hook secret should fail the run"); + }; + + assert!(err.to_string().contains("HOOK_TOKEN")); + } + + #[tokio::test] + async fn hook_http_url_resolves_secret_from_vault() { + let server = httpmock::MockServer::start_async().await; + let mock = server + .mock_async(|when, then| { + when.method("POST").path("/hook"); + then.status(200).body(""); + }) + .await; + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings + .run + .hooks + .push(sandbox_ready_http_hook("{{ secrets.HOOK_URL }}")); + let (_persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(token_vault( + "HOOK_URL", + &server.url("/hook"), + ))); + + start(&run_dir, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + .expect("HTTP hook URL should resolve secret and proceed"); + + mock.assert_async().await; + } + + #[tokio::test] + async fn hook_prompt_secret_resolver_resolves_from_vault_and_registers() { + let redactor = SecretRedactor::default(); + let vault = Arc::new(AsyncRwLock::new(token_vault("PROMPT_TOKEN", "staging"))); + let resolver = hook_secret_resolver(Some(vault), redactor.clone()); + let hook = HookDefinition { + name: Some("prompt-secret".to_string()), + event: HookEvent::StageStart, + command: None, + hook_type: Some(HookType::Prompt { + prompt: InterpString::parse("check {{ secrets.PROMPT_TOKEN }}"), + model: None, + }), + matcher: None, + blocking: Some(true), + timeout_ms: None, + sandbox: None, + }; + + let secrets = resolver.resolve_for_definition(&hook).await; + + assert_eq!(secrets.lookup("PROMPT_TOKEN").as_deref(), Some("staging")); + assert_eq!(redactor.redact_into("deploy staging"), "deploy REDACTED"); + } + + #[tokio::test] + async fn hook_block_reason_redacts_resolved_secret_in_error_and_events() { + let temp = tempfile::tempdir().unwrap(); + let (storage_root, run_dir) = storage_root_and_run_dir(&temp); + let mut settings = settings_from_run_layer(RunLayer { + execution: Some(RunExecutionLayer { + mode: Some(RunMode::DryRun), + ..RunExecutionLayer::default() + }), + ..RunLayer::default() + }); + settings.run.hooks.push(sandbox_ready_command_hook( + r#"printf '%s' '{"decision":"block","reason":"{{ secrets.HOOK_TOKEN }}"}'"#, + )); + let (_persisted, store) = + persisted_workflow_with_settings(MINIMAL_DOT, &storage_root, settings).await; + let emitter = Arc::new(Emitter::new(fixtures::RUN_1)); + let registry = Arc::new(test_registry()); + let vault = Arc::new(AsyncRwLock::new(token_vault("HOOK_TOKEN", "staging"))); + + let Err(err) = start(&run_dir, StartServices { + vault: Some(vault), + ..test_start_services(&store, &storage_root, emitter, registry).await + }) + .await + else { + panic!("blocking hook should fail the run"); + }; + + let error_text = err.to_string(); + assert!(!error_text.contains("staging")); + assert!(error_text.contains("REDACTED")); + + let run_store = store.open_run(&fixtures::RUN_1).await.unwrap(); + let events_text = serde_json::to_string(&run_store.list_events().await.unwrap()).unwrap(); + assert!(!events_text.contains("staging")); + assert!(events_text.contains("REDACTED")); + } + #[test] fn runtime_docker_config_maps_environment_hints() { let settings = settings_from_run_layer(RunLayer { @@ -1806,6 +2114,37 @@ reasoning = false } } + fn sandbox_ready_command_hook(command: &str) -> HookDefinition { + HookDefinition { + name: Some("sandbox-ready".to_string()), + event: HookEvent::SandboxReady, + command: Some(InterpString::parse(command)), + hook_type: None, + matcher: None, + blocking: Some(true), + timeout_ms: Some(5_000), + sandbox: Some(false), + } + } + + fn sandbox_ready_http_hook(url: &str) -> HookDefinition { + HookDefinition { + name: Some("sandbox-ready-http".to_string()), + event: HookEvent::SandboxReady, + command: None, + hook_type: Some(HookType::Http { + url: InterpString::parse(url), + headers: None, + allowed_env_vars: Vec::new(), + tls: TlsMode::Off, + }), + matcher: None, + blocking: Some(true), + timeout_ms: Some(5_000), + sandbox: Some(false), + } + } + use crate::test_support::{mark_run_running, test_usage}; async fn append_completed_stage( diff --git a/lib/crates/fabro-workflow/src/pipeline/execute.rs b/lib/crates/fabro-workflow/src/pipeline/execute.rs index 4a82a509c..f84b03e7c 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute.rs @@ -89,6 +89,7 @@ pub async fn execute(init: Initialized) -> Executed { Arc::clone(&engine.run.sandbox_git), Arc::clone(&engine.run.metadata_runtime), engine.run.metadata_writer.clone(), + engine.run.secret_redactor.clone(), checkpoint.is_some(), on_node, run_control, diff --git a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs index e2b96ed37..09b9a4f6d 100644 --- a/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs +++ b/lib/crates/fabro-workflow/src/pipeline/execute/tests.rs @@ -282,12 +282,14 @@ async fn execute_test_run_with_options( workflow_path: None, workflow_bundle: None, hooks: HookSettings { hooks: vec![] }, + hook_secrets: fabro_hooks::HookSecretResolver::default(), sandbox_env: SandboxEnvSpec { toml_env: HashMap::new(), github_permissions: None, origin_url: None, }, vault: None, + secret_redactor: fabro_redact::SecretRedactor::default(), git: git_options, run_control: None, registry_override, @@ -345,12 +347,14 @@ async fn execute_runs_start_to_exit_and_returns_final_context() { workflow_path: None, workflow_bundle: None, hooks: HookSettings { hooks: vec![] }, + hook_secrets: fabro_hooks::HookSecretResolver::default(), sandbox_env: SandboxEnvSpec { toml_env: HashMap::new(), github_permissions: None, origin_url: None, }, vault: None, + secret_redactor: fabro_redact::SecretRedactor::default(), git: None, run_control: None, registry_override: None, @@ -416,12 +420,14 @@ async fn run_with_lifecycle( workflow_path: None, workflow_bundle: None, hooks: HookSettings { hooks: vec![] }, + hook_secrets: fabro_hooks::HookSecretResolver::default(), sandbox_env: SandboxEnvSpec { toml_env: HashMap::new(), github_permissions: None, origin_url: None, }, vault: None, + secret_redactor: fabro_redact::SecretRedactor::default(), git: None, run_control: None, registry_override: Some(Arc::new(registry)), diff --git a/lib/crates/fabro-workflow/src/pipeline/finalize.rs b/lib/crates/fabro-workflow/src/pipeline/finalize.rs index 6e90991dc..59746ac4d 100644 --- a/lib/crates/fabro-workflow/src/pipeline/finalize.rs +++ b/lib/crates/fabro-workflow/src/pipeline/finalize.rs @@ -615,8 +615,17 @@ pub async fn finalize(executed: Executed, options: &FinalizeOptions) -> Result, pub workflow_bundle: Option>, pub hooks: fabro_hooks::HookSettings, + pub hook_secrets: fabro_hooks::HookSecretResolver, pub sandbox_env: SandboxEnvSpec, pub vault: Option>>, + pub secret_redactor: SecretRedactor, pub git: Option, pub registry_override: Option>, pub artifact_sink: Option, diff --git a/lib/crates/fabro-workflow/src/runtime_store.rs b/lib/crates/fabro-workflow/src/runtime_store.rs index 761a8c4fd..0216c20af 100644 --- a/lib/crates/fabro-workflow/src/runtime_store.rs +++ b/lib/crates/fabro-workflow/src/runtime_store.rs @@ -3,16 +3,25 @@ use std::sync::Arc; use anyhow::Result; use async_trait::async_trait; use bytes::Bytes; +use fabro_redact::SecretRedactor; use fabro_store::{EventEnvelope, RunDatabase, RunProjection}; use fabro_types::{RunBlobId, RunEvent}; -use crate::event::build_redacted_event_payload; +use crate::event::build_redacted_event_payload_with_redactor; #[async_trait] pub trait RunStoreBackend: Send + Sync { async fn load_state(&self) -> Result; async fn list_events(&self) -> Result>; async fn append_run_event(&self, event: &RunEvent) -> Result<()>; + async fn append_run_event_with_redactor( + &self, + event: &RunEvent, + redactor: Option<&SecretRedactor>, + ) -> Result<()> { + let _ = redactor; + self.append_run_event(event).await + } async fn write_blob(&self, data: &[u8]) -> Result; async fn read_blob(&self, id: &RunBlobId) -> Result>; async fn read_run_log(&self) -> Result>>; @@ -46,6 +55,16 @@ impl RunStoreHandle { self.backend.append_run_event(event).await } + pub async fn append_run_event_with_redactor( + &self, + event: &RunEvent, + redactor: Option<&SecretRedactor>, + ) -> Result<()> { + self.backend + .append_run_event_with_redactor(event, redactor) + .await + } + pub async fn write_blob(&self, data: &[u8]) -> Result { self.backend.write_blob(data).await } @@ -83,7 +102,15 @@ impl RunStoreBackend for LocalRunStoreBackend { } async fn append_run_event(&self, event: &RunEvent) -> Result<()> { - let payload = build_redacted_event_payload(event, &event.run_id)?; + self.append_run_event_with_redactor(event, None).await + } + + async fn append_run_event_with_redactor( + &self, + event: &RunEvent, + redactor: Option<&SecretRedactor>, + ) -> Result<()> { + let payload = build_redacted_event_payload_with_redactor(event, &event.run_id, redactor)?; self.run_store .append_event(&payload) .await diff --git a/lib/crates/fabro-workflow/src/services.rs b/lib/crates/fabro-workflow/src/services.rs index 44e0eca21..c4dbd6735 100644 --- a/lib/crates/fabro-workflow/src/services.rs +++ b/lib/crates/fabro-workflow/src/services.rs @@ -11,6 +11,7 @@ use fabro_auth::ResolvedCredentials; use fabro_hooks::{HookContext, HookDecision, HookExecutionContext, HookRunner}; use fabro_interview::Interviewer; use fabro_model::{Catalog, ProviderId}; +use fabro_redact::SecretRedactor; use fabro_types::{ManifestPath, RunId}; use tokio_util::sync::CancellationToken; @@ -103,6 +104,7 @@ pub struct RunServices { pub model: String, pub llm_source: Arc, pub catalog: Arc, + pub secret_redactor: SecretRedactor, pub(crate) sandbox_git: Arc, pub(crate) metadata_runtime: Arc, pub(crate) metadata_writer: Option, @@ -122,6 +124,7 @@ impl RunServices { model: String, llm_source: Arc, catalog: Arc, + secret_redactor: SecretRedactor, sandbox_git: Arc, metadata_runtime: Arc, metadata_writer: Option, @@ -137,6 +140,7 @@ impl RunServices { model, llm_source, catalog, + secret_redactor, sandbox_git, metadata_runtime, metadata_writer, @@ -337,6 +341,7 @@ impl EngineServices { "claude-sonnet-4-6".to_string(), Arc::new(StubCredentialSource), Arc::new(Catalog::from_builtin().expect("default catalog should build")), + SecretRedactor::default(), Arc::new(SandboxGitRuntime::new()), Arc::new(RunMetadataRuntime::new()), None, diff --git a/lib/crates/fabro-workflow/src/test_support.rs b/lib/crates/fabro-workflow/src/test_support.rs index 1acf410a4..c890667e6 100644 --- a/lib/crates/fabro-workflow/src/test_support.rs +++ b/lib/crates/fabro-workflow/src/test_support.rs @@ -9,6 +9,7 @@ use fabro_auth::{CredentialSource, EnvCredentialSource}; use fabro_graphviz::graph::Graph as GvGraph; use fabro_interview::AutoApproveInterviewer; use fabro_model::Catalog; +use fabro_redact::SecretRedactor; use fabro_store::{ArtifactStore, Database, RunProjection}; use object_store::local::LocalFileSystem; @@ -235,6 +236,7 @@ async fn initialized( .llm_source .unwrap_or_else(|| Arc::new(EnvCredentialSource::new())), Arc::new(Catalog::from_builtin().expect("default catalog should build")), + SecretRedactor::default(), Arc::new(SandboxGitRuntime::new()), Arc::new(RunMetadataRuntime::new()), None,