mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-08 03:10:26 +00:00
fabro(01KWFGXZ5P42QRWBYAPVEAXMX6): simplify_opus (succeeded)
Fabro-Run: 01KWFGXZ5P42QRWBYAPVEAXMX6 Fabro-Completed: 6 ⚒️ Generated with [Fabro](https://fabro.sh)
This commit is contained in:
parent
f8992b2f09
commit
88047477db
6 changed files with 75 additions and 77 deletions
|
|
@ -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()
|
||||
};
|
||||
|
|
|
|||
|
|
@ -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<dyn Sandbox>,
|
||||
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(),
|
||||
|
|
|
|||
|
|
@ -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<String> {
|
||||
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]
|
||||
|
|
|
|||
|
|
@ -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,
|
||||
|
|
|
|||
|
|
@ -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<RunEvent> {
|
||||
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<String> {
|
||||
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);
|
||||
}
|
||||
|
|
|
|||
|
|
@ -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?;
|
||||
}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue