diff --git a/lib/crates/fabro-cli/src/commands/run/rewind.rs b/lib/crates/fabro-cli/src/commands/run/rewind.rs index f5d027abd..eff957b5e 100644 --- a/lib/crates/fabro-cli/src/commands/run/rewind.rs +++ b/lib/crates/fabro-cli/src/commands/run/rewind.rs @@ -198,8 +198,13 @@ fn run_event(run_id: fabro_types::RunId, node_id: Option, body: EventBod run_id, node_id, node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, session_id: None, parent_session_id: None, + tool_call_id: None, + actor: None, body, } } diff --git a/lib/crates/fabro-cli/tests/it/cmd/pr_view.rs b/lib/crates/fabro-cli/tests/it/cmd/pr_view.rs index 771101196..49fe73cf4 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/pr_view.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/pr_view.rs @@ -90,8 +90,13 @@ fn pr_view_reads_pull_request_from_store_without_pull_request_json() { run_id, node_id: None, node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, session_id: None, parent_session_id: None, + tool_call_id: None, + actor: None, body: EventBody::PullRequestCreated(PullRequestCreatedProps { pr_url: "https://github.com/fabro-sh/fabro/pull/123".to_string(), pr_number: 123, diff --git a/lib/crates/fabro-store/src/run_state.rs b/lib/crates/fabro-store/src/run_state.rs index 6b6660342..f7e39e556 100644 --- a/lib/crates/fabro-store/src/run_state.rs +++ b/lib/crates/fabro-store/src/run_state.rs @@ -611,8 +611,13 @@ mod tests { run_id: fixtures::RUN_1, node_id: node_id.map(ToOwned::to_owned), node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, session_id: None, parent_session_id: None, + tool_call_id: None, + actor: None, body, }; diff --git a/lib/crates/fabro-types/src/lib.rs b/lib/crates/fabro-types/src/lib.rs index 6749e5545..c6dda288f 100644 --- a/lib/crates/fabro-types/src/lib.rs +++ b/lib/crates/fabro-types/src/lib.rs @@ -49,7 +49,7 @@ pub use run::{ RunSubjectProvenance, }; pub use run_blob_id::RunBlobId; -pub use run_event::{EventBody, RunEvent, RunNoticeLevel}; +pub use run_event::{ActorKind, ActorRef, EventBody, RunEvent, RunNoticeLevel}; pub use run_id::RunId; pub use run_id::fixtures; pub use sandbox_record::SandboxRecord; diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index f4817e3ed..56f771698 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -27,6 +27,23 @@ pub enum RunNoticeLevel { Error, } +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum ActorKind { + User, + Agent, + System, +} + +#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] +pub struct ActorRef { + pub kind: ActorKind, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub id: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] + pub display: Option, +} + #[derive(Debug, Clone, PartialEq)] pub struct RunEvent { pub id: String, @@ -34,8 +51,13 @@ pub struct RunEvent { pub run_id: RunId, pub node_id: Option, pub node_label: Option, + pub stage_id: Option, + pub parallel_group_id: Option, + pub parallel_branch_id: Option, pub session_id: Option, pub parent_session_id: Option, + pub tool_call_id: Option, + pub actor: Option, pub body: EventBody, } @@ -271,9 +293,19 @@ struct RunEventRaw { #[serde(default)] node_label: Option, #[serde(default)] + stage_id: Option, + #[serde(default)] + parallel_group_id: Option, + #[serde(default)] + parallel_branch_id: Option, + #[serde(default)] session_id: Option, #[serde(default)] parent_session_id: Option, + #[serde(default)] + tool_call_id: Option, + #[serde(default)] + actor: Option, event: String, #[serde(default = "default_properties")] properties: Value, @@ -283,6 +315,23 @@ fn default_properties() -> Value { Value::Object(Map::new()) } +struct RunEventParts<'a> { + id: String, + ts: DateTime, + run_id: RunId, + node_id: Option, + node_label: Option, + stage_id: Option, + parallel_group_id: Option, + parallel_branch_id: Option, + session_id: Option, + parent_session_id: Option, + tool_call_id: Option, + actor: Option, + event: &'a str, + properties: &'a Value, +} + impl EventBody { pub fn event_name(&self) -> &str { match self { @@ -524,17 +573,22 @@ 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, - ) + Self::from_parts(RunEventParts { + id: raw.id, + ts: raw.ts, + run_id: raw.run_id, + node_id: raw.node_id, + node_label: raw.node_label, + stage_id: raw.stage_id, + parallel_group_id: raw.parallel_group_id, + parallel_branch_id: raw.parallel_branch_id, + session_id: raw.session_id, + parent_session_id: raw.parent_session_id, + tool_call_id: raw.tool_call_id, + actor: raw.actor, + event: &raw.event, + properties: &raw.properties, + }) } pub fn from_ref(value: &Value) -> serde_json::Result { @@ -559,58 +613,78 @@ impl RunEvent { .get("properties") .cloned() .unwrap_or_else(default_properties); - Self::from_parts( - id.to_string(), + let actor = match obj.get("actor") { + Some(value) if !value.is_null() => Some(ActorRef::deserialize(value)?), + _ => None, + }; + Self::from_parts(RunEventParts { + id: id.to_string(), ts, run_id, - obj.get("node_id") + node_id: obj + .get("node_id") .and_then(Value::as_str) .map(str::to_string), - obj.get("node_label") + node_label: obj + .get("node_label") .and_then(Value::as_str) .map(str::to_string), - obj.get("session_id") + stage_id: obj + .get("stage_id") .and_then(Value::as_str) .map(str::to_string), - obj.get("parent_session_id") + parallel_group_id: obj + .get("parallel_group_id") .and_then(Value::as_str) .map(str::to_string), + parallel_branch_id: obj + .get("parallel_branch_id") + .and_then(Value::as_str) + .map(str::to_string), + session_id: obj + .get("session_id") + .and_then(Value::as_str) + .map(str::to_string), + parent_session_id: obj + .get("parent_session_id") + .and_then(Value::as_str) + .map(str::to_string), + tool_call_id: obj + .get("tool_call_id") + .and_then(Value::as_str) + .map(str::to_string), + actor, event, - &properties, - ) + properties: &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: &str, - properties: &Value, - ) -> serde_json::Result { + fn from_parts(parts: RunEventParts<'_>) -> serde_json::Result { let body_payload = json!({ - "event": event, - "properties": properties, + "event": parts.event, + "properties": parts.properties, }); let body: EventBody = match serde_json::from_value(body_payload) { Ok(body) => body, - Err(err) if is_known_event_name(event) => return Err(err), + Err(err) if is_known_event_name(parts.event) => return Err(err), Err(_) => EventBody::Unknown { - name: event.to_string(), - properties: properties.clone(), + name: parts.event.to_string(), + properties: parts.properties.clone(), }, }; Ok(Self { - id, - ts, - run_id, - node_id, - node_label, - session_id, - parent_session_id, + id: parts.id, + ts: parts.ts, + run_id: parts.run_id, + node_id: parts.node_id, + node_label: parts.node_label, + stage_id: parts.stage_id, + parallel_group_id: parts.parallel_group_id, + parallel_branch_id: parts.parallel_branch_id, + session_id: parts.session_id, + parent_session_id: parts.parent_session_id, + tool_call_id: parts.tool_call_id, + actor: parts.actor, body, }) } @@ -643,6 +717,27 @@ impl RunEvent { if let Some(value) = &self.node_label { map.insert("node_label".to_string(), Value::String(value.clone())); } + if let Some(value) = &self.stage_id { + map.insert("stage_id".to_string(), Value::String(value.clone())); + } + if let Some(value) = &self.parallel_group_id { + map.insert( + "parallel_group_id".to_string(), + Value::String(value.clone()), + ); + } + if let Some(value) = &self.parallel_branch_id { + map.insert( + "parallel_branch_id".to_string(), + Value::String(value.clone()), + ); + } + if let Some(value) = &self.tool_call_id { + map.insert("tool_call_id".to_string(), Value::String(value.clone())); + } + if let Some(actor) = &self.actor { + map.insert("actor".to_string(), serde_json::to_value(actor)?); + } map.insert("properties".to_string(), self.body.properties_value()?); Ok(Value::Object(map)) } @@ -697,8 +792,13 @@ mod tests { run_id: fixtures::RUN_1, node_id: Some("build".to_string()), node_label: Some("Build".to_string()), + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, session_id: None, parent_session_id: None, + tool_call_id: None, + actor: None, body: EventBody::StageCompleted(StageCompletedProps { index: 1, duration_ms: 1234, @@ -873,4 +973,90 @@ mod tests { assert_eq!(serialized["event"], value["event"]); assert_eq!(serialized["properties"], value["properties"]); } + + #[test] + fn run_event_round_trips_new_envelope_fields() { + let value = json!({ + "id": "evt_envelope", + "ts": "2026-04-08T16:21:11.106Z", + "run_id": fixtures::RUN_1, + "event": "agent.tool.completed", + "stage_id": "code@1", + "node_id": "code", + "node_label": "Code", + "parallel_group_id": "code@1", + "parallel_branch_id": "code@1:0", + "session_id": "ses_child", + "parent_session_id": "ses_parent", + "tool_call_id": "call_1", + "actor": { + "kind": "agent", + "id": "ses_child", + "display": "claude-sonnet" + }, + "properties": { + "tool_name": "read_file", + "tool_call_id": "call_1", + "output": {"summary": "read"}, + "is_error": false, + "visit": 1 + } + }); + + let parsed = RunEvent::from_value(value.clone()).unwrap(); + assert_eq!(parsed.stage_id.as_deref(), Some("code@1")); + assert_eq!(parsed.parallel_group_id.as_deref(), Some("code@1")); + assert_eq!(parsed.parallel_branch_id.as_deref(), Some("code@1:0")); + assert_eq!(parsed.tool_call_id.as_deref(), Some("call_1")); + let actor = parsed.actor.as_ref().expect("actor present"); + assert_eq!(actor.kind, ActorKind::Agent); + assert_eq!(actor.id.as_deref(), Some("ses_child")); + assert_eq!(actor.display.as_deref(), Some("claude-sonnet")); + + let serialized = parsed.to_value().unwrap(); + assert_eq!(serialized["stage_id"], value["stage_id"]); + assert_eq!(serialized["parallel_group_id"], value["parallel_group_id"]); + assert_eq!( + serialized["parallel_branch_id"], + value["parallel_branch_id"] + ); + assert_eq!(serialized["tool_call_id"], value["tool_call_id"]); + assert_eq!(serialized["actor"], value["actor"]); + } + + #[test] + fn run_event_omits_absent_envelope_fields() { + let event = RunEvent { + id: "evt_bare".to_string(), + ts: DateTime::parse_from_rfc3339("2026-04-04T12:00:00.000Z") + .unwrap() + .with_timezone(&Utc), + run_id: fixtures::RUN_1, + node_id: None, + node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, + session_id: None, + parent_session_id: None, + tool_call_id: None, + actor: None, + body: EventBody::RunStarted(RunStartedProps { + name: "demo".to_string(), + base_branch: None, + base_sha: None, + run_branch: None, + worktree_dir: None, + goal: None, + }), + }; + + let serialized = event.to_value().unwrap(); + let obj = serialized.as_object().unwrap(); + assert!(!obj.contains_key("stage_id")); + assert!(!obj.contains_key("parallel_group_id")); + assert!(!obj.contains_key("parallel_branch_id")); + assert!(!obj.contains_key("tool_call_id")); + assert!(!obj.contains_key("actor")); + } } diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index d55809675..220d58ee0 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -1248,12 +1248,17 @@ pub fn event_name(event: &Event) -> &'static str { } } -#[derive(Debug)] +#[derive(Debug, Default)] struct StoredEventFields { session_id: Option, parent_session_id: Option, node_id: Option, node_label: Option, + stage_id: Option, + parallel_group_id: Option, + parallel_branch_id: Option, + tool_call_id: Option, + actor: Option, } fn default_node_label(node_id: Option<&String>, node_label: Option) -> Option { @@ -1285,10 +1290,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { 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, + ..StoredEventFields::default() } } Event::CheckpointCompleted { node_id, .. } @@ -1306,10 +1310,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { 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, + ..StoredEventFields::default() } } Event::Agent { @@ -1325,15 +1328,15 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { parent_session_id: parent_session_id.clone(), node_id, node_label, + ..StoredEventFields::default() } } 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, + ..StoredEventFields::default() } } Event::ParallelBranchStarted { branch, .. } @@ -1341,10 +1344,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { 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, + ..StoredEventFields::default() } } Event::Prompt { stage, .. } @@ -1355,28 +1357,21 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { 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, + ..StoredEventFields::default() } } 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::default() } } - _ => StoredEventFields { - session_id: None, - parent_session_id: None, - node_id: None, - node_label: None, - }, + _ => StoredEventFields::default(), } } @@ -2394,8 +2389,13 @@ pub fn to_run_event_at(run_id: &RunId, event: &Event, ts: chrono::DateTime) run_id: *run_id, node_id: fields.node_id, node_label: fields.node_label, + stage_id: fields.stage_id, + parallel_group_id: fields.parallel_group_id, + parallel_branch_id: fields.parallel_branch_id, session_id: fields.session_id, parent_session_id: fields.parent_session_id, + tool_call_id: fields.tool_call_id, + actor: fields.actor, body, } } diff --git a/lib/crates/fabro-workflow/src/runtime_store.rs b/lib/crates/fabro-workflow/src/runtime_store.rs index d4ff85107..6eb093010 100644 --- a/lib/crates/fabro-workflow/src/runtime_store.rs +++ b/lib/crates/fabro-workflow/src/runtime_store.rs @@ -193,8 +193,13 @@ mod tests { run_id: fixtures::RUN_1, node_id: None, node_label: None, + stage_id: None, + parallel_group_id: None, + parallel_branch_id: None, session_id: None, parent_session_id: None, + tool_call_id: None, + actor: None, body: EventBody::RunSubmitted(RunSubmittedProps { reason: None, definition_blob: None,