From 76e1da8b35c46bd909347358553c518489577113 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 4 Apr 2026 13:24:19 -0400 Subject: [PATCH] refactor: remove remaining event json indirection --- .../fabro-cli/src/commands/run/attach.rs | 18 +- lib/crates/fabro-cli/src/commands/run/logs.rs | 18 +- lib/crates/fabro-store/src/run_state.rs | 2 +- lib/crates/fabro-store/src/types.rs | 2 +- lib/crates/fabro-types/src/run_event/mod.rs | 91 ++++- lib/crates/fabro-util/src/json.rs | 19 ++ lib/crates/fabro-util/src/lib.rs | 1 + lib/crates/fabro-util/src/redact/jsonl.rs | 43 +++ lib/crates/fabro-util/src/redact/mod.rs | 1 + lib/crates/fabro-workflow/src/event.rs | 322 +++++++----------- .../fabro-workflow/src/operations/create.rs | 3 +- 11 files changed, 269 insertions(+), 251 deletions(-) create mode 100644 lib/crates/fabro-util/src/json.rs diff --git a/lib/crates/fabro-cli/src/commands/run/attach.rs b/lib/crates/fabro-cli/src/commands/run/attach.rs index 238795d56..2035c2e12 100644 --- a/lib/crates/fabro-cli/src/commands/run/attach.rs +++ b/lib/crates/fabro-cli/src/commands/run/attach.rs @@ -11,10 +11,10 @@ use futures::StreamExt; use fabro_interview::{AnswerValue, ConsoleInterviewer}; use fabro_store::{EventEnvelope, RuntimeState, SlateRunStore}; +use fabro_util::json::normalize_json_value; use fabro_util::terminal::Styles; use fabro_workflow::outcome::StageStatus; use fabro_workflow::run_status::RunStatus; -use serde_json::{Map, Value}; use tokio::signal::ctrl_c; use tokio::time::{self, sleep}; @@ -343,22 +343,6 @@ fn event_payload_line(event: &EventEnvelope) -> Result { .map_err(Into::into) } -fn normalize_json_value(value: Value) -> Value { - match value { - Value::Object(map) => Value::Object( - map.into_iter() - .map(|(key, value)| (key, normalize_json_value(value))) - .collect::>() - .into_iter() - .collect::>(), - ), - Value::Array(values) => { - Value::Array(values.into_iter().map(normalize_json_value).collect()) - } - other => other, - } -} - fn read_launcher_pid(run_dir: &Path) -> Option { super::launcher::active_launcher_record_for_run(run_dir).map(|record| record.pid) } diff --git a/lib/crates/fabro-cli/src/commands/run/logs.rs b/lib/crates/fabro-cli/src/commands/run/logs.rs index df31e853b..fa878a2c8 100644 --- a/lib/crates/fabro-cli/src/commands/run/logs.rs +++ b/lib/crates/fabro-cli/src/commands/run/logs.rs @@ -6,11 +6,11 @@ use std::time::Duration; use anyhow::{Context, Result, bail}; use chrono::{DateTime, Utc}; use fabro_store::SlateRunStore; +use fabro_util::json::normalize_json_value; use fabro_util::redact::redact_jsonl_line; use fabro_util::terminal::Styles; use fabro_workflow::run_lookup::{resolve_run_combined, runs_base}; use futures::StreamExt; -use serde_json::{Map, Value}; use tokio::time; use tracing::{debug, info}; @@ -229,22 +229,6 @@ fn event_payload_line(event: &fabro_store::EventEnvelope) -> Result { Ok(redact_jsonl_line(&line)) } -fn normalize_json_value(value: Value) -> Value { - match value { - Value::Object(map) => Value::Object( - map.into_iter() - .map(|(key, value)| (key, normalize_json_value(value))) - .collect::>() - .into_iter() - .collect::>(), - ), - Value::Array(values) => { - Value::Array(values.into_iter().map(normalize_json_value).collect()) - } - other => other, - } -} - fn render_indented_markdown(styles: &Styles, text: &str, indent: &str) -> String { let term_width = Styles::terminal_width(); let wrap_width = term_width.saturating_sub(indent.len()); diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 2f26f4304..145e9d32d 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -75,7 +75,7 @@ impl RunProjection { } pub(crate) fn apply_event(&mut self, event: &EventEnvelope) -> Result<()> { - let stored = RunEvent::from_value(event.payload.as_value().clone()) + let stored = RunEvent::from_ref(event.payload.as_value()) .map_err(|err| StoreError::InvalidEvent(format!("invalid stored event: {err}")))?; let ts = stored.ts; let run_id = stored.run_id; diff --git a/lib/crates/fabro-store/src/types.rs b/lib/crates/fabro-store/src/types.rs index 49fc3f900..a859ccdba 100644 --- a/lib/crates/fabro-store/src/types.rs +++ b/lib/crates/fabro-store/src/types.rs @@ -74,7 +74,7 @@ impl TryFrom<&EventPayload> for RunEvent { type Error = StoreError; fn try_from(value: &EventPayload) -> Result { - RunEvent::from_value(value.as_value().clone()) + RunEvent::from_ref(value.as_value()) .map_err(|err| StoreError::InvalidEvent(format!("invalid stored event: {err}"))) } } diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index d3ac0b661..216bbe788 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -521,26 +521,93 @@ fn is_known_event_name(event: &str) -> bool { impl RunEvent { pub fn from_value(value: Value) -> serde_json::Result { let raw: RunEventRaw = serde_json::from_value(value)?; + Self::from_parts( + raw.id, + raw.ts, + raw.run_id, + raw.node_id, + raw.node_label, + raw.session_id, + raw.parent_session_id, + raw.event, + raw.properties, + ) + } + + pub fn from_ref(value: &Value) -> serde_json::Result { + let obj = value.as_object().ok_or_else(|| { + ::custom("run event must be a JSON object") + })?; + let id = obj.get("id").and_then(Value::as_str).ok_or_else(|| { + ::custom("missing or non-string field: id") + })?; + let ts = obj + .get("ts") + .ok_or_else(|| ::custom("missing field: ts")) + .and_then(DateTime::::deserialize)?; + let run_id = obj + .get("run_id") + .ok_or_else(|| ::custom("missing field: run_id")) + .and_then(RunId::deserialize)?; + let event = obj.get("event").and_then(Value::as_str).ok_or_else(|| { + ::custom("missing or non-string field: event") + })?; + let properties = obj + .get("properties") + .cloned() + .unwrap_or_else(default_properties); + Self::from_parts( + id.to_string(), + ts, + run_id, + obj.get("node_id") + .and_then(Value::as_str) + .map(str::to_string), + obj.get("node_label") + .and_then(Value::as_str) + .map(str::to_string), + obj.get("session_id") + .and_then(Value::as_str) + .map(str::to_string), + obj.get("parent_session_id") + .and_then(Value::as_str) + .map(str::to_string), + event.to_string(), + properties, + ) + } + + fn from_parts( + id: String, + ts: DateTime, + run_id: RunId, + node_id: Option, + node_label: Option, + session_id: Option, + parent_session_id: Option, + event: String, + properties: Value, + ) -> serde_json::Result { let body_payload = json!({ - "event": raw.event, - "properties": raw.properties, + "event": event, + "properties": properties, }); let body: EventBody = match serde_json::from_value(body_payload) { Ok(body) => body, - Err(err) if is_known_event_name(&raw.event) => return Err(err), + Err(err) if is_known_event_name(&event) => return Err(err), Err(_) => EventBody::Unknown { - name: raw.event.clone(), - properties: raw.properties.clone(), + name: event.clone(), + properties: properties.clone(), }, }; Ok(Self { - id: raw.id, - ts: raw.ts, - run_id: raw.run_id, - node_id: raw.node_id, - node_label: raw.node_label, - session_id: raw.session_id, - parent_session_id: raw.parent_session_id, + id, + ts, + run_id, + node_id, + node_label, + session_id, + parent_session_id, body, }) } diff --git a/lib/crates/fabro-util/src/json.rs b/lib/crates/fabro-util/src/json.rs new file mode 100644 index 000000000..22a68f9ce --- /dev/null +++ b/lib/crates/fabro-util/src/json.rs @@ -0,0 +1,19 @@ +use std::collections::BTreeMap; + +use serde_json::{Map, Value}; + +pub fn normalize_json_value(value: Value) -> Value { + match value { + Value::Object(map) => Value::Object( + map.into_iter() + .map(|(key, value)| (key, normalize_json_value(value))) + .collect::>() + .into_iter() + .collect::>(), + ), + Value::Array(values) => { + Value::Array(values.into_iter().map(normalize_json_value).collect()) + } + other => other, + } +} diff --git a/lib/crates/fabro-util/src/lib.rs b/lib/crates/fabro-util/src/lib.rs index 5b7c7d14e..85d5d1675 100644 --- a/lib/crates/fabro-util/src/lib.rs +++ b/lib/crates/fabro-util/src/lib.rs @@ -1,6 +1,7 @@ pub mod backoff; pub mod check_report; pub mod env; +pub mod json; pub mod path; pub mod printer; pub mod redact; diff --git a/lib/crates/fabro-util/src/redact/jsonl.rs b/lib/crates/fabro-util/src/redact/jsonl.rs index 0eccf6556..7ba213f8c 100644 --- a/lib/crates/fabro-util/src/redact/jsonl.rs +++ b/lib/crates/fabro-util/src/redact/jsonl.rs @@ -70,6 +70,36 @@ fn collect_replacements(v: &Value) -> Vec<(String, String)> { repls } +fn redact_value_in_place(value: &mut Value, skip_field: bool) { + match value { + Value::Object(obj) => { + if should_skip_object(obj) { + return; + } + for (key, child) in obj { + redact_value_in_place(child, should_skip_field(key)); + } + } + Value::Array(arr) => { + for child in arr { + redact_value_in_place(child, false); + } + } + Value::String(text) if !skip_field => { + let redacted = super::redact_string(text); + if redacted != *text { + *text = redacted; + } + } + _ => {} + } +} + +pub fn redact_json_value(mut value: Value) -> Value { + redact_value_in_place(&mut value, false); + value +} + /// JSON-encode a string value (with quotes), without HTML escaping. fn json_encode_string(s: &str) -> String { serde_json::to_string(s).unwrap_or_else(|_| format!("\"{s}\"")) @@ -205,6 +235,19 @@ mod tests { assert_eq!(redact_jsonl_line(input), input); } + #[test] + fn redact_json_value_matches_jsonl_behavior() { + let input = serde_json::json!({ + "content": format!("key={HIGH_ENTROPY_SECRET}"), + "session_id": HIGH_ENTROPY_SECRET, + }); + + let redacted = redact_json_value(input); + + assert_eq!(redacted["content"], "REDACTED"); + assert_eq!(redacted["session_id"], HIGH_ENTROPY_SECRET); + } + #[test] fn redact_jsonl_line_with_secret_in_content() { let input = format!(r#"{{"type":"text","content":"key={HIGH_ENTROPY_SECRET}"}}"#); diff --git a/lib/crates/fabro-util/src/redact/mod.rs b/lib/crates/fabro-util/src/redact/mod.rs index f15fd6755..bb26903aa 100644 --- a/lib/crates/fabro-util/src/redact/mod.rs +++ b/lib/crates/fabro-util/src/redact/mod.rs @@ -2,6 +2,7 @@ mod entropy; mod gitleaks; mod jsonl; +pub use jsonl::redact_json_value; pub use jsonl::redact_jsonl_line; /// A byte range within a string that should be redacted. diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 49477226c..2e8543689 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -6,8 +6,9 @@ use ::fabro_types::{RunEvent, RunId, StageStatus, StatusReason}; use anyhow::{Context, Result}; use chrono::Utc; use fabro_store::{EventPayload, SlateRunStore}; +use fabro_util::json::normalize_json_value; use serde::{Deserialize, Serialize}; -use serde_json::{Map, Value}; +use serde_json::Value; use std::collections::BTreeMap; use tokio::sync::{mpsc, oneshot}; use uuid::Uuid; @@ -16,7 +17,7 @@ use crate::error::FabroError; use crate::outcome::{FailureDetail, Outcome, StageUsage}; use fabro_agent::{AgentEvent, SandboxEvent, WorktreeEvent, WorktreeEventCallback}; use fabro_llm::types::Usage as LlmUsage; -use fabro_util::redact::redact_jsonl_line; +use fabro_util::redact::redact_json_value; pub use fabro_types::{EventBody, RunNoticeLevel}; @@ -1185,40 +1186,6 @@ struct StoredEventFields { node_label: Option, } -fn tagged_variant_fields(value: &T) -> Map { - tagged_variant_fields_from_value(serde_json::to_value(value).expect("serializable event")) -} - -fn tagged_variant_fields_from_value(value: Value) -> Map { - match value { - Value::Object(map) => { - let (_, inner) = map.into_iter().next().expect("enum must have one variant"); - match inner { - Value::Object(fields) => fields, - Value::String(_) | Value::Null => Map::new(), - other => { - let mut fields = Map::new(); - fields.insert("value".to_string(), other); - fields - } - } - } - Value::String(_) | Value::Null => Map::new(), - other => { - let mut fields = Map::new(); - fields.insert("value".to_string(), other); - fields - } - } -} - -fn remove_string(fields: &mut Map, key: &str) -> Option { - match fields.remove(key) { - Some(Value::String(value)) => Some(value), - _ => None, - } -} - fn default_node_label(node_id: Option<&String>, node_label: Option) -> Option { node_label.or_else(|| node_id.cloned()) } @@ -1240,6 +1207,116 @@ fn stage_status_from_string(status: &str) -> StageStatus { serde_json::from_value(Value::String(status.to_string())).expect("valid stage status") } +fn stored_event_fields(event: &Event) -> StoredEventFields { + match event { + Event::StageCompleted { node_id, name, .. } | Event::StageFailed { node_id, name, .. } => { + let node_id = Some(node_id.clone()); + let node_label = default_node_label(node_id.as_ref(), Some(name.clone())); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + Event::StageStarted { node_id, name, .. } | Event::StageRetrying { node_id, name, .. } => { + let node_id = Some(node_id.clone()); + let node_label = default_node_label(node_id.as_ref(), Some(name.clone())); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + Event::CheckpointCompleted { node_id, .. } + | Event::CheckpointFailed { node_id, .. } + | Event::SubgraphStarted { node_id, .. } + | Event::SubgraphCompleted { node_id, .. } + | Event::ArtifactCaptured { node_id, .. } + | Event::PromptCompleted { node_id, .. } + | Event::ParallelStarted { node_id, .. } + | Event::ParallelCompleted { node_id, .. } + | Event::CommandStarted { node_id, .. } + | Event::CommandCompleted { node_id, .. } + | Event::AgentCliStarted { node_id, .. } + | Event::AgentCliCompleted { node_id, .. } => { + let node_id = Some(node_id.clone()); + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + Event::Agent { + stage, + session_id, + parent_session_id, + .. + } => { + let node_id = Some(stage.clone()); + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: session_id.clone(), + parent_session_id: parent_session_id.clone(), + node_id, + node_label, + } + } + Event::GitCommit { node_id, .. } => { + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id: node_id.clone(), + node_label, + } + } + Event::ParallelBranchStarted { branch, .. } + | Event::ParallelBranchCompleted { branch, .. } => { + let node_id = Some(branch.clone()); + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + Event::Prompt { stage, .. } + | Event::InterviewStarted { stage, .. } + | Event::InterviewTimeout { stage, .. } + | Event::Failover { stage, .. } => { + let node_id = Some(stage.clone()); + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + Event::StallWatchdogTimeout { node, .. } => { + let node_id = Some(node.clone()); + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + session_id: None, + parent_session_id: None, + node_id, + node_label, + } + } + _ => StoredEventFields { + session_id: None, + parent_session_id: None, + node_id: None, + node_label: None, + }, + } +} + fn event_body_from_event(event: &Event) -> EventBody { match event { Event::RunCreated { @@ -2193,156 +2270,12 @@ fn event_body_from_event(event: &Event) -> EventBody { } } -fn extract_run_event_fields(event: &Event) -> StoredEventFields { - match event { - Event::RunCreated { .. } | Event::WorkflowRunStarted { .. } => { - let mut fields = tagged_variant_fields(event); - fields.remove("run_id"); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id: None, - node_label: None, - } - } - Event::WorkflowRunFailed { error, .. } => { - let mut fields = tagged_variant_fields(event); - fields.insert("error".to_string(), Value::String(error.to_string())); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id: None, - node_label: None, - } - } - Event::StageCompleted { .. } | Event::StageFailed { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "node_id"); - let node_label = - default_node_label(node_id.as_ref(), remove_string(&mut fields, "name")); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - Event::StageStarted { .. } - | Event::StageRetrying { .. } - | Event::CheckpointCompleted { .. } - | Event::CheckpointFailed { .. } - | Event::SubgraphStarted { .. } - | Event::SubgraphCompleted { .. } - | Event::ArtifactCaptured { .. } - | Event::PromptCompleted { .. } - | Event::ParallelStarted { .. } - | Event::ParallelCompleted { .. } - | Event::CommandStarted { .. } - | Event::CommandCompleted { .. } - | Event::AgentCliStarted { .. } - | Event::AgentCliCompleted { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "node_id"); - let node_label = - default_node_label(node_id.as_ref(), remove_string(&mut fields, "name")); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - Event::Agent { - session_id, - parent_session_id, - .. - } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "stage"); - let node_label = default_node_label(node_id.as_ref(), None); - fields.remove("visit"); - fields.remove("session_id"); - fields.remove("parent_session_id"); - fields.remove("event"); - StoredEventFields { - session_id: session_id.clone(), - parent_session_id: parent_session_id.clone(), - node_id, - node_label, - } - } - Event::Sandbox { .. } => { - let mut fields = tagged_variant_fields(event); - fields.remove("event"); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id: None, - node_label: None, - } - } - Event::GitCommit { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "node_id"); - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - Event::ParallelBranchStarted { .. } | Event::ParallelBranchCompleted { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "branch"); - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - Event::Prompt { .. } - | Event::InterviewStarted { .. } - | Event::InterviewTimeout { .. } - | Event::Failover { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "stage"); - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - Event::StallWatchdogTimeout { .. } => { - let mut fields = tagged_variant_fields(event); - let node_id = remove_string(&mut fields, "node"); - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - session_id: None, - parent_session_id: None, - node_id, - node_label, - } - } - _ => StoredEventFields { - session_id: None, - parent_session_id: None, - node_id: None, - node_label: None, - }, - } -} - pub fn to_run_event(run_id: &RunId, event: &Event) -> RunEvent { to_run_event_at(run_id, event, Utc::now()) } pub fn to_run_event_at(run_id: &RunId, event: &Event, ts: chrono::DateTime) -> RunEvent { - let fields = extract_run_event_fields(event); + let fields = stored_event_fields(event); let body = event_body_from_event(event); RunEvent { id: Uuid::now_v7().to_string(), @@ -2357,13 +2290,12 @@ pub fn to_run_event_at(run_id: &RunId, event: &Event, ts: chrono::DateTime) } pub fn build_redacted_event_payload(event: &RunEvent, run_id: &RunId) -> Result { - let line = redacted_event_json(event)?; - event_payload_from_redacted_json(&line, run_id) + let value = redacted_event_value(event)?; + EventPayload::new(value, run_id).map_err(anyhow::Error::from) } pub fn redacted_event_json(event: &RunEvent) -> Result { - let line = serde_json::to_string(&normalized_event_value(event)?)?; - Ok(redact_jsonl_line(&line)) + serde_json::to_string(&redacted_event_value(event)?).map_err(anyhow::Error::from) } fn normalized_event_value(event: &RunEvent) -> Result { @@ -2371,26 +2303,12 @@ fn normalized_event_value(event: &RunEvent) -> Result { Ok(normalize_json_value(value)) } -pub(crate) fn normalize_json_value(value: Value) -> Value { - match value { - Value::Object(map) => Value::Object( - map.into_iter() - .map(|(key, value)| (key, normalize_json_value(value))) - .collect::>() - .into_iter() - .collect::>(), - ), - Value::Array(values) => { - Value::Array(values.into_iter().map(normalize_json_value).collect()) - } - other => other, - } +fn redacted_event_value(event: &RunEvent) -> Result { + Ok(redact_json_value(normalized_event_value(event)?)) } pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result { - let value = normalize_json_value( - serde_json::from_str(line).context("Failed to parse redacted event payload")?, - ); + let value = serde_json::from_str(line).context("Failed to parse redacted event payload")?; EventPayload::new(value, run_id).map_err(anyhow::Error::from) } diff --git a/lib/crates/fabro-workflow/src/operations/create.rs b/lib/crates/fabro-workflow/src/operations/create.rs index d5b13b1ab..a3e11d186 100644 --- a/lib/crates/fabro-workflow/src/operations/create.rs +++ b/lib/crates/fabro-workflow/src/operations/create.rs @@ -15,9 +15,10 @@ use crate::records::RunRecord; use crate::run_lookup::default_runs_base; use crate::transforms::{Transform, expand_vars}; use fabro_sandbox::daytona::detect_repo_info; +use fabro_util::json::normalize_json_value; use super::source::{ResolveWorkflowInput, WorkflowInput, resolve_workflow}; -use crate::event::{Event, append_event, normalize_json_value, to_run_event_at}; +use crate::event::{Event, append_event, to_run_event_at}; #[derive(Clone, Debug)] pub struct CreateRunInput {