refactor: remove remaining event json indirection

This commit is contained in:
Bryan Helmkamp 2026-04-04 13:24:19 -04:00
parent 6c9877cc73
commit 76e1da8b35
No known key found for this signature in database
11 changed files with 269 additions and 251 deletions

View file

@ -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<String> {
.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::<std::collections::BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
Value::Array(values) => {
Value::Array(values.into_iter().map(normalize_json_value).collect())
}
other => other,
}
}
fn read_launcher_pid(run_dir: &Path) -> Option<u32> {
super::launcher::active_launcher_record_for_run(run_dir).map(|record| record.pid)
}

View file

@ -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<String> {
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::<std::collections::BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
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());

View file

@ -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;

View file

@ -74,7 +74,7 @@ impl TryFrom<&EventPayload> for RunEvent {
type Error = StoreError;
fn try_from(value: &EventPayload) -> Result<Self> {
RunEvent::from_value(value.as_value().clone())
RunEvent::from_ref(value.as_value())
.map_err(|err| StoreError::InvalidEvent(format!("invalid stored event: {err}")))
}
}

View file

@ -521,26 +521,93 @@ fn is_known_event_name(event: &str) -> bool {
impl RunEvent {
pub fn from_value(value: Value) -> serde_json::Result<Self> {
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<Self> {
let obj = value.as_object().ok_or_else(|| {
<serde_json::Error as DeError>::custom("run event must be a JSON object")
})?;
let id = obj.get("id").and_then(Value::as_str).ok_or_else(|| {
<serde_json::Error as DeError>::custom("missing or non-string field: id")
})?;
let ts = obj
.get("ts")
.ok_or_else(|| <serde_json::Error as DeError>::custom("missing field: ts"))
.and_then(DateTime::<Utc>::deserialize)?;
let run_id = obj
.get("run_id")
.ok_or_else(|| <serde_json::Error as DeError>::custom("missing field: run_id"))
.and_then(RunId::deserialize)?;
let event = obj.get("event").and_then(Value::as_str).ok_or_else(|| {
<serde_json::Error as DeError>::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<Utc>,
run_id: RunId,
node_id: Option<String>,
node_label: Option<String>,
session_id: Option<String>,
parent_session_id: Option<String>,
event: String,
properties: Value,
) -> serde_json::Result<Self> {
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,
})
}

View file

@ -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::<BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
Value::Array(values) => {
Value::Array(values.into_iter().map(normalize_json_value).collect())
}
other => other,
}
}

View file

@ -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;

View file

@ -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}"}}"#);

View file

@ -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.

View file

@ -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<String>,
}
fn tagged_variant_fields<T: Serialize>(value: &T) -> Map<String, Value> {
tagged_variant_fields_from_value(serde_json::to_value(value).expect("serializable event"))
}
fn tagged_variant_fields_from_value(value: Value) -> Map<String, Value> {
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<String, Value>, key: &str) -> Option<String> {
match fields.remove(key) {
Some(Value::String(value)) => Some(value),
_ => None,
}
}
fn default_node_label(node_id: Option<&String>, node_label: Option<String>) -> Option<String> {
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<Utc>) -> 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<Utc>)
}
pub fn build_redacted_event_payload(event: &RunEvent, run_id: &RunId) -> Result<EventPayload> {
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<String> {
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<Value> {
@ -2371,26 +2303,12 @@ fn normalized_event_value(event: &RunEvent) -> Result<Value> {
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::<BTreeMap<_, _>>()
.into_iter()
.collect::<Map<_, _>>(),
),
Value::Array(values) => {
Value::Array(values.into_iter().map(normalize_json_value).collect())
}
other => other,
}
fn redacted_event_value(event: &RunEvent) -> Result<Value> {
Ok(redact_json_value(normalized_event_value(event)?))
}
pub fn event_payload_from_redacted_json(line: &str, run_id: &RunId) -> Result<EventPayload> {
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)
}

View file

@ -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 {