diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs index 20f5d98a0..880338894 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/event.rs @@ -511,6 +511,8 @@ mod tests { name: "Plan".into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, duration_ms: 5000, status: "success".into(), preferred_label: None, @@ -555,6 +557,8 @@ mod tests { }, session_id: None, parent_session_id: None, + parallel_group_id: None, + parallel_branch_id: None, }; let stored = to_run_event(&fixtures::RUN_1, &event); diff --git a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs index 63507433e..5c5ced36d 100644 --- a/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs +++ b/lib/crates/fabro-cli/src/commands/run/run_progress/mod.rs @@ -479,6 +479,8 @@ mod tests { event, session_id: None, parent_session_id: None, + parallel_group_id: None, + parallel_branch_id: None, } } @@ -488,6 +490,8 @@ mod tests { name: name.into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, handler_type: String::new(), attempt: 1, max_attempts: 1, @@ -512,6 +516,8 @@ mod tests { name: name.into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, duration_ms: 5000, status: "success".into(), preferred_label: None, @@ -709,6 +715,8 @@ mod tests { name: "Code".into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, attempt: 2, max_attempts: 3, delay_ms: 1500, @@ -954,6 +962,8 @@ mod tests { name: "Code".into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, attempt: 2, max_attempts: 3, delay_ms: 1500, @@ -1159,6 +1169,8 @@ mod tests { name: "Code".into(), index: 0, visit: 1, + parallel_group_id: None, + parallel_branch_id: None, handler_type: "agent".into(), attempt: 1, max_attempts: 1, diff --git a/lib/crates/fabro-cli/src/commands/store/dump.rs b/lib/crates/fabro-cli/src/commands/store/dump.rs index 3f25253a5..929deb4c1 100644 --- a/lib/crates/fabro-cli/src/commands/store/dump.rs +++ b/lib/crates/fabro-cli/src/commands/store/dump.rs @@ -600,6 +600,8 @@ mod tests { name: "Code".to_string(), index: 1, visit: 2, + parallel_group_id: None, + parallel_branch_id: None, duration_ms: 250, status: "partial_success".to_string(), preferred_label: None, diff --git a/lib/crates/fabro-cli/tests/it/cmd/support.rs b/lib/crates/fabro-cli/tests/it/cmd/support.rs index 62e7cf38e..d93851add 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/support.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/support.rs @@ -14,6 +14,7 @@ use fabro_config::Storage; use fabro_server::bind::Bind; use fabro_store::EventEnvelope; use fabro_test::TestContext; +use fabro_types::RunId; use serde_json::Value; use shlex::try_quote; diff --git a/lib/crates/fabro-server/src/server.rs b/lib/crates/fabro-server/src/server.rs index 1672da3d2..acf1fce6a 100644 --- a/lib/crates/fabro-server/src/server.rs +++ b/lib/crates/fabro-server/src/server.rs @@ -1480,8 +1480,8 @@ fn event_matches_run_filter(event: &EventEnvelope, run_filter: Option<&HashSet Option { - let event = api_event_envelope_from_store(event).ok()?; - let data = serde_json::to_string(&event).ok()?; + let wire = event.to_wire_value().ok()?; + let data = serde_json::to_string(&wire).ok()?; let data = redact_jsonl_line(&data); Some(Event::default().data(data)) } @@ -2381,20 +2381,15 @@ fn octet_stream_response(bytes: Bytes) -> Response { #[allow(clippy::result_large_err)] fn api_event_envelope_from_store(event: &EventEnvelope) -> Result { - let value = event.to_wire_value().map_err(|err| { + fn serialize_error(err: impl std::fmt::Display) -> Response { ApiError::new( StatusCode::INTERNAL_SERVER_ERROR, format!("Failed to serialize stored event: {err}"), ) .into_response() - })?; - serde_json::from_value(value).map_err(|err| { - ApiError::new( - StatusCode::INTERNAL_SERVER_ERROR, - format!("Failed to serialize stored event: {err}"), - ) - .into_response() - }) + } + let value = event.to_wire_value().map_err(serialize_error)?; + serde_json::from_value(value).map_err(serialize_error) } fn clear_live_run_state(run: &mut ManagedRun) { diff --git a/lib/crates/fabro-types/src/run_event/mod.rs b/lib/crates/fabro-types/src/run_event/mod.rs index 56f771698..d8a1bc1db 100644 --- a/lib/crates/fabro-types/src/run_event/mod.rs +++ b/lib/crates/fabro-types/src/run_event/mod.rs @@ -595,6 +595,7 @@ impl RunEvent { let obj = value.as_object().ok_or_else(|| { ::custom("run event must be a JSON object") })?; + let opt_str = |key: &str| obj.get(key).and_then(Value::as_str).map(str::to_string); let id = obj.get("id").and_then(Value::as_str).ok_or_else(|| { ::custom("missing or non-string field: id") })?; @@ -621,38 +622,14 @@ impl RunEvent { id: id.to_string(), ts, run_id, - node_id: obj - .get("node_id") - .and_then(Value::as_str) - .map(str::to_string), - node_label: obj - .get("node_label") - .and_then(Value::as_str) - .map(str::to_string), - stage_id: obj - .get("stage_id") - .and_then(Value::as_str) - .map(str::to_string), - 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), + node_id: opt_str("node_id"), + node_label: opt_str("node_label"), + stage_id: opt_str("stage_id"), + parallel_group_id: opt_str("parallel_group_id"), + parallel_branch_id: opt_str("parallel_branch_id"), + session_id: opt_str("session_id"), + parent_session_id: opt_str("parent_session_id"), + tool_call_id: opt_str("tool_call_id"), actor, event, properties: &properties, diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index eddf0f2d7..b018e97f1 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -5,8 +5,8 @@ use std::sync::atomic::{AtomicI64, Ordering}; use ::fabro_types::run_event as fabro_types; use ::fabro_types::{ - BilledTokenCounts, RunBlobId, RunControlAction, RunEvent, RunId, StageId, StageStatus, - StatusReason, + ActorKind, ActorRef, BilledTokenCounts, RunBlobId, RunControlAction, RunEvent, RunId, + RunProvenance, StageId, StageStatus, StatusReason, }; use anyhow::{Context, Result}; use chrono::Utc; @@ -55,7 +55,7 @@ pub enum Event { #[serde(default, skip_serializing_if = "Option::is_none")] db_prefix: Option, #[serde(default, skip_serializing_if = "Option::is_none")] - provenance: Option<::fabro_types::RunProvenance>, + provenance: Option, #[serde(default, skip_serializing_if = "Option::is_none")] manifest_blob: Option, }, @@ -1289,13 +1289,22 @@ struct StoredEventFields { parallel_group_id: Option, parallel_branch_id: Option, tool_call_id: Option, - actor: Option<::fabro_types::ActorRef>, + actor: Option, } fn default_node_label(node_id: Option<&String>, node_label: Option) -> Option { node_label.or_else(|| node_id.cloned()) } +fn node_stored_fields(node_id: Option) -> StoredEventFields { + let node_label = default_node_label(node_id.as_ref(), None); + StoredEventFields { + node_id, + node_label, + ..StoredEventFields::default() + } +} + fn billed_token_counts_from_llm(usage: &LlmTokenCounts) -> BilledTokenCounts { BilledTokenCounts { input_tokens: usage.input_tokens, @@ -1383,15 +1392,7 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { | 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 { - node_id, - node_label, - ..StoredEventFields::default() - } - } + | Event::AgentCliCompleted { node_id, .. } => node_stored_fields(Some(node_id.clone())), Event::Agent { stage, visit, @@ -1416,17 +1417,9 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { parallel_branch_id: parallel_branch_id.clone(), tool_call_id, actor, - ..StoredEventFields::default() - } - } - Event::GitCommit { node_id, .. } => { - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - node_id: node_id.clone(), - node_label, - ..StoredEventFields::default() } } + Event::GitCommit { node_id, .. } => node_stored_fields(node_id.clone()), Event::ParallelBranchStarted { parallel_group_id, parallel_branch_id, @@ -1453,35 +1446,16 @@ fn stored_event_fields(event: &Event) -> StoredEventFields { | Event::InterviewStarted { stage, .. } | Event::InterviewTimeout { stage, .. } | Event::InterviewInterrupted { stage, .. } - | Event::Failover { stage, .. } => { - let node_id = Some(stage.clone()); - let node_label = default_node_label(node_id.as_ref(), None); - StoredEventFields { - 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 { - node_id, - node_label, - ..StoredEventFields::default() - } - } + | Event::Failover { stage, .. } => node_stored_fields(Some(stage.clone())), + Event::StallWatchdogTimeout { node, .. } => node_stored_fields(Some(node.clone())), _ => StoredEventFields::default(), } } -fn actor_from_provenance( - provenance: &::fabro_types::RunProvenance, -) -> Option<::fabro_types::ActorRef> { - let subject = provenance.subject.as_ref()?; - let login = subject.login.clone()?; - Some(::fabro_types::ActorRef { - kind: ::fabro_types::ActorKind::User, +fn actor_from_provenance(provenance: &RunProvenance) -> Option { + let login = provenance.subject.as_ref()?.login.clone()?; + Some(ActorRef { + kind: ActorKind::User, id: Some(login.clone()), display: Some(login), }) @@ -1495,13 +1469,10 @@ fn agent_tool_call_id(event: &AgentEvent) -> Option<&str> { } } -fn agent_actor_for_event( - event: &AgentEvent, - session_id: Option<&str>, -) -> Option<::fabro_types::ActorRef> { +fn agent_actor_for_event(event: &AgentEvent, session_id: Option<&str>) -> Option { match event { - AgentEvent::AssistantMessage { model, .. } => Some(::fabro_types::ActorRef { - kind: ::fabro_types::ActorKind::Agent, + AgentEvent::AssistantMessage { model, .. } => Some(ActorRef { + kind: ActorKind::Agent, id: session_id.map(str::to_string), display: Some(model.clone()), }), @@ -3281,16 +3252,14 @@ mod tests { }, ); let actor = stored.actor.as_ref().expect("actor set"); - assert_eq!(actor.kind, ::fabro_types::ActorKind::Agent); + assert_eq!(actor.kind, ActorKind::Agent); assert_eq!(actor.id.as_deref(), Some("ses_agent")); assert_eq!(actor.display.as_deref(), Some("claude-sonnet")); } #[test] fn run_created_populates_user_actor_from_provenance() { - use ::fabro_types::{ - Graph, RunAuthMethod, RunProvenance, RunSubjectProvenance, Settings, fixtures, - }; + use ::fabro_types::{Graph, RunAuthMethod, RunSubjectProvenance, Settings, fixtures}; let provenance = RunProvenance { server: None, @@ -3322,7 +3291,7 @@ mod tests { }, ); let actor = stored.actor.as_ref().expect("actor set"); - assert_eq!(actor.kind, ::fabro_types::ActorKind::User); + assert_eq!(actor.kind, ActorKind::User); assert_eq!(actor.id.as_deref(), Some("alice")); assert_eq!(actor.display.as_deref(), Some("alice")); } diff --git a/lib/crates/fabro-workflow/src/handler/llm/api.rs b/lib/crates/fabro-workflow/src/handler/llm/api.rs index d8f7a934f..a95b0a64d 100644 --- a/lib/crates/fabro-workflow/src/handler/llm/api.rs +++ b/lib/crates/fabro-workflow/src/handler/llm/api.rs @@ -37,10 +37,6 @@ fn build_profile(model: &str, provider: Provider) -> Box { } } -fn current_visit(context: &Context) -> u32 { - u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX) -} - #[derive(Clone)] struct StageEventScope { visit: u32, @@ -50,7 +46,7 @@ struct StageEventScope { fn current_stage_event_scope(context: &Context) -> StageEventScope { StageEventScope { - visit: current_visit(context), + visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX), parallel_group_id: context.parallel_group_id(), parallel_branch_id: context.parallel_branch_id(), }