From 93a7f19383a488806c6fa794b6208c32b6696a30 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Sat, 4 Apr 2026 10:51:15 -0400 Subject: [PATCH] refactor: remove legacy stored event compatibility --- .../src/commands/run/run_progress/event.rs | 557 ++++++++---------- .../src/commands/run/run_progress/mod.rs | 65 +- lib/crates/fabro-cli/tests/it/cmd/attach.rs | 1 - lib/crates/fabro-cli/tests/it/cmd/logs.rs | 4 +- lib/crates/fabro-cli/tests/it/cmd/run.rs | 1 - lib/crates/fabro-workflow/src/event.rs | 220 +++---- .../fabro-workflow/src/operations/create.rs | 6 +- 7 files changed, 380 insertions(+), 474 deletions(-) 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 1172612f7..95d54be96 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 @@ -1,10 +1,10 @@ use std::convert::TryFrom; use chrono::{DateTime, Utc}; -use fabro_types::StoredEvent; +use fabro_types::{EventBody, StageUsage, StoredEvent}; use fabro_workflow::event::RunNoticeLevel; -use fabro_workflow::outcome::{StageUsage, compute_stage_cost}; -use serde_json::{Map, Value}; +use fabro_workflow::outcome::compute_stage_cost; +use serde_json::Value; #[derive(Debug, Clone)] pub(super) struct ProgressUsage { @@ -16,17 +16,13 @@ pub(super) struct ProgressUsage { } impl ProgressUsage { - pub(super) fn from_value(value: &Value) -> Option { - let Value::Object(fields) = value else { - return None; - }; - + pub(super) fn from_stage_usage(usage: &StageUsage) -> Option { Some(Self { - model: string_field(fields, "model"), - input_tokens: u64_field(fields, "input_tokens"), - output_tokens: u64_field(fields, "output_tokens"), - speed: string_field(fields, "speed"), - cost: f64_field(fields, "cost"), + model: Some(usage.model.clone()), + input_tokens: u64::try_from(usage.input_tokens).ok()?, + output_tokens: u64::try_from(usage.output_tokens).ok()?, + speed: usage.speed.clone(), + cost: usage.cost, }) } @@ -240,268 +236,234 @@ pub(super) enum ProgressEvent { }, } -#[allow(clippy::needless_pass_by_value)] -fn from_envelope_fields(event_name: &str, fields: &Map) -> Option { - match event_name { - "run.started" => Some(ProgressEvent::WorkflowStarted { - worktree_dir: prop_string_field(fields, "worktree_dir"), - base_branch: prop_string_field(fields, "base_branch"), - base_sha: prop_string_field(fields, "base_sha"), +pub(super) fn from_stored_event(stored: &StoredEvent) -> Option { + let node_id = stored.node_id.clone().unwrap_or_else(|| "?".to_string()); + let node_label = stored.node_label.clone().unwrap_or_else(|| node_id.clone()); + + match &stored.body { + EventBody::RunStarted(props) => Some(ProgressEvent::WorkflowStarted { + worktree_dir: props.worktree_dir.clone(), + base_branch: props.base_branch.clone(), + base_sha: props.base_sha.clone(), }), - "sandbox.initialized" => Some(ProgressEvent::WorkingDirectorySet { - working_directory: prop_string_field(fields, "working_directory")?, + EventBody::SandboxInitialized(props) => Some(ProgressEvent::WorkingDirectorySet { + working_directory: props.working_directory.clone(), }), - "sandbox.initializing" => Some(ProgressEvent::SandboxInitializing { - provider: prop_string_field(fields, "provider") - .unwrap_or_else(|| "unknown".to_string()), + EventBody::SandboxInitializing(props) => Some(ProgressEvent::SandboxInitializing { + provider: props.provider.clone(), }), - "sandbox.ready" => Some(ProgressEvent::SandboxReady { - provider: prop_string_field(fields, "provider") - .unwrap_or_else(|| "unknown".to_string()), - duration_ms: prop_u64_field(fields, "duration_ms"), - name: prop_string_field(fields, "name"), - cpu: prop_f64_field(fields, "cpu"), - memory: prop_f64_field(fields, "memory"), - url: prop_string_field(fields, "url"), + EventBody::SandboxReady(props) => Some(ProgressEvent::SandboxReady { + provider: props.provider.clone(), + duration_ms: props.duration_ms, + name: props.name.clone(), + cpu: props.cpu, + memory: props.memory, + url: props.url.clone(), }), - "ssh.ready" => Some(ProgressEvent::SshAccessReady { - ssh_command: prop_string_field(fields, "ssh_command")?, + EventBody::SshAccessReady(props) => Some(ProgressEvent::SshAccessReady { + ssh_command: props.ssh_command.clone(), }), - "setup.started" => Some(ProgressEvent::SetupStarted { - command_count: prop_u64_field(fields, "command_count"), + EventBody::SetupStarted(props) => Some(ProgressEvent::SetupStarted { + command_count: props.command_count as u64, }), - "setup.completed" => Some(ProgressEvent::SetupCompleted { - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::SetupCompleted(props) => Some(ProgressEvent::SetupCompleted { + duration_ms: props.duration_ms, }), - "setup.command.completed" => Some(ProgressEvent::SetupCommandCompleted { - command: prop_string_field(fields, "command").unwrap_or_else(|| "?".to_string()), - command_index: prop_u64_field(fields, "index"), - exit_code: prop_i64_field(fields, "exit_code"), - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::SetupCommandCompleted(props) => Some(ProgressEvent::SetupCommandCompleted { + command: props.command.clone(), + command_index: props.index as u64, + exit_code: i64::from(props.exit_code), + duration_ms: props.duration_ms, }), - "cli.ensure.started" => Some(ProgressEvent::CliEnsureStarted { - cli_name: prop_string_field(fields, "cli_name").unwrap_or_else(|| "?".to_string()), + EventBody::CliEnsureStarted(props) => Some(ProgressEvent::CliEnsureStarted { + cli_name: props.cli_name.clone(), }), - "cli.ensure.completed" => Some(ProgressEvent::CliEnsureCompleted { - cli_name: prop_string_field(fields, "cli_name").unwrap_or_else(|| "?".to_string()), - already_installed: prop_bool_field(fields, "already_installed"), - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::CliEnsureCompleted(props) => Some(ProgressEvent::CliEnsureCompleted { + cli_name: props.cli_name.clone(), + already_installed: props.already_installed, + duration_ms: props.duration_ms, }), - "cli.ensure.failed" => Some(ProgressEvent::CliEnsureFailed { - cli_name: prop_string_field(fields, "cli_name").unwrap_or_else(|| "?".to_string()), + EventBody::CliEnsureFailed(props) => Some(ProgressEvent::CliEnsureFailed { + cli_name: props.cli_name.clone(), }), - "devcontainer.resolved" => Some(ProgressEvent::DevcontainerResolved { - dockerfile_lines: prop_u64_field(fields, "dockerfile_lines"), - environment_count: prop_u64_field(fields, "environment_count"), - lifecycle_command_count: prop_u64_field(fields, "lifecycle_command_count"), - workspace_folder: prop_string_field(fields, "workspace_folder") - .unwrap_or_else(|| "?".to_string()), + EventBody::DevcontainerResolved(props) => Some(ProgressEvent::DevcontainerResolved { + dockerfile_lines: props.dockerfile_lines as u64, + environment_count: props.environment_count as u64, + lifecycle_command_count: props.lifecycle_command_count as u64, + workspace_folder: props.workspace_folder.clone(), }), - "devcontainer.lifecycle.started" => Some(ProgressEvent::DevcontainerLifecycleStarted { - phase: prop_string_field(fields, "phase").unwrap_or_else(|| "?".to_string()), - command_count: prop_u64_field(fields, "command_count"), - }), - "devcontainer.lifecycle.completed" => Some(ProgressEvent::DevcontainerLifecycleCompleted { - phase: prop_string_field(fields, "phase").unwrap_or_else(|| "?".to_string()), - duration_ms: prop_u64_field(fields, "duration_ms"), - }), - "devcontainer.lifecycle.failed" => Some(ProgressEvent::DevcontainerLifecycleFailed { - phase: prop_string_field(fields, "phase").unwrap_or_else(|| "?".to_string()), - command: prop_string_field(fields, "command").unwrap_or_else(|| "?".to_string()), - exit_code: prop_i64_field(fields, "exit_code"), - stderr: prop_display_field(fields, "stderr").unwrap_or_default(), - }), - "devcontainer.lifecycle.command.completed" => { - Some(ProgressEvent::DevcontainerLifecycleCommandCompleted { - command: prop_string_field(fields, "command").unwrap_or_else(|| "?".to_string()), - command_index: prop_u64_field(fields, "index"), - exit_code: prop_i64_field(fields, "exit_code"), - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::DevcontainerLifecycleStarted(props) => { + Some(ProgressEvent::DevcontainerLifecycleStarted { + phase: props.phase.clone(), + command_count: props.command_count as u64, }) } - "stage.started" => Some(ProgressEvent::StageStarted { - node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - name: string_field(fields, "node_label").unwrap_or_else(|| "?".to_string()), - script: prop_string_field(fields, "script"), + EventBody::DevcontainerLifecycleCompleted(props) => { + Some(ProgressEvent::DevcontainerLifecycleCompleted { + phase: props.phase.clone(), + duration_ms: props.duration_ms, + }) + } + EventBody::DevcontainerLifecycleFailed(props) => { + Some(ProgressEvent::DevcontainerLifecycleFailed { + phase: props.phase.clone(), + command: props.command.clone(), + exit_code: i64::from(props.exit_code), + stderr: props.stderr.clone(), + }) + } + EventBody::DevcontainerLifecycleCommandCompleted(props) => { + Some(ProgressEvent::DevcontainerLifecycleCommandCompleted { + command: props.command.clone(), + command_index: props.index as u64, + exit_code: i64::from(props.exit_code), + duration_ms: props.duration_ms, + }) + } + EventBody::StageStarted(_) => Some(ProgressEvent::StageStarted { + node_id, + name: node_label, + script: None, }), - "stage.completed" => Some(ProgressEvent::StageCompleted { - node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - name: string_field(fields, "node_label").unwrap_or_else(|| "?".to_string()), - duration_ms: prop_u64_field(fields, "duration_ms"), - status: prop_string_field(fields, "status").unwrap_or_else(|| "success".to_string()), - usage: prop_value(fields, "usage").and_then(ProgressUsage::from_value), + EventBody::StageCompleted(props) => Some(ProgressEvent::StageCompleted { + node_id, + name: node_label, + duration_ms: props.duration_ms, + status: props.status.to_string(), + usage: props + .usage + .as_ref() + .and_then(ProgressUsage::from_stage_usage), }), - "stage.failed" => Some(ProgressEvent::StageFailed { - node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - name: string_field(fields, "node_label").unwrap_or_else(|| "?".to_string()), - error: prop_display_field(fields, "error") + EventBody::StageFailed(props) => Some(ProgressEvent::StageFailed { + node_id, + name: node_label, + error: props + .failure + .as_ref() + .map(|failure| failure.message.clone()) .unwrap_or_else(|| "unknown error".to_string()), }), - "stage.retrying" => Some(ProgressEvent::StageRetrying { - name: string_field(fields, "node_label").unwrap_or_else(|| "?".to_string()), - attempt: prop_u64_field(fields, "attempt"), - max_attempts: prop_u64_field(fields, "max_attempts"), - delay_ms: prop_u64_field(fields, "delay_ms"), + EventBody::StageRetrying(props) => Some(ProgressEvent::StageRetrying { + name: node_label, + attempt: props.attempt as u64, + max_attempts: props.max_attempts as u64, + delay_ms: props.delay_ms, }), - "parallel.started" => Some(ProgressEvent::ParallelStarted), - "parallel.branch.started" => Some(ProgressEvent::ParallelBranchStarted { - branch: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), + EventBody::ParallelStarted(_) => Some(ProgressEvent::ParallelStarted), + EventBody::ParallelBranchStarted(_) => { + Some(ProgressEvent::ParallelBranchStarted { branch: node_id }) + } + EventBody::ParallelBranchCompleted(props) => Some(ProgressEvent::ParallelBranchCompleted { + branch: node_id, + duration_ms: props.duration_ms, + status: props.status.clone(), }), - "parallel.branch.completed" => Some(ProgressEvent::ParallelBranchCompleted { - branch: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - duration_ms: prop_u64_field(fields, "duration_ms"), - status: prop_string_field(fields, "status").unwrap_or_else(|| "success".to_string()), + EventBody::ParallelCompleted(_) => Some(ProgressEvent::ParallelCompleted), + EventBody::AgentMessage(props) => Some(ProgressEvent::AssistantMessage { + stage_node_id: node_id, + model: props.model.clone(), }), - "parallel.completed" => Some(ProgressEvent::ParallelCompleted), - "agent.message" => Some(ProgressEvent::AssistantMessage { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - model: prop_string_field(fields, "model").unwrap_or_else(|| "?".to_string()), + EventBody::AgentToolStarted(props) => Some(ProgressEvent::ToolCallStarted { + stage_node_id: node_id, + tool_name: props.tool_name.clone(), + tool_call_id: props.tool_call_id.clone(), + arguments: props.arguments.clone(), + timestamp: Some(stored.ts), }), - "agent.tool.started" => Some(ProgressEvent::ToolCallStarted { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - tool_name: prop_string_field(fields, "tool_name").unwrap_or_else(|| "?".to_string()), - tool_call_id: prop_string_field(fields, "tool_call_id") - .unwrap_or_else(|| "?".to_string()), - arguments: prop_value(fields, "arguments") - .cloned() - .unwrap_or_else(|| Value::Object(Map::new())), - timestamp: timestamp_field(fields, "ts"), + EventBody::AgentToolCompleted(props) => Some(ProgressEvent::ToolCallCompleted { + stage_node_id: node_id, + tool_call_id: props.tool_call_id.clone(), + is_error: props.is_error, + duration_ms: None, + timestamp: Some(stored.ts), }), - "agent.tool.completed" => Some(ProgressEvent::ToolCallCompleted { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - tool_call_id: prop_string_field(fields, "tool_call_id") - .unwrap_or_else(|| "?".to_string()), - is_error: prop_bool_field(fields, "is_error"), - duration_ms: prop_optional_u64_field(fields, "duration_ms"), - timestamp: timestamp_field(fields, "ts"), - }), - "agent.warning" - if prop_string_field(fields, "kind").as_deref() == Some("context_window") => - { - let usage_percent = prop_value(fields, "details") - .and_then(Value::as_object) + EventBody::AgentWarning(props) if props.kind == "context_window" => { + let usage_percent = props + .details + .as_object() .and_then(|details| details.get("usage_percent")) .and_then(Value::as_u64) .unwrap_or(0); Some(ProgressEvent::ContextWindowWarning { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), + stage_node_id: node_id, usage_percent, }) } - "agent.compaction.started" => Some(ProgressEvent::CompactionStarted { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), + EventBody::AgentCompactionStarted(_) => Some(ProgressEvent::CompactionStarted { + stage_node_id: node_id, }), - "agent.compaction.completed" => Some(ProgressEvent::CompactionCompleted { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - original_turn_count: prop_u64_field(fields, "original_turn_count"), - preserved_turn_count: prop_u64_field(fields, "preserved_turn_count"), - tracked_file_count: prop_u64_field(fields, "tracked_file_count"), + EventBody::AgentCompactionCompleted(props) => Some(ProgressEvent::CompactionCompleted { + stage_node_id: node_id, + original_turn_count: props.original_turn_count as u64, + preserved_turn_count: props.preserved_turn_count as u64, + tracked_file_count: props.tracked_file_count as u64, }), - "agent.llm.retry" => { - let delay_secs = prop_f64_field(fields, "delay_secs").unwrap_or(0.0); + EventBody::AgentLlmRetry(props) => { #[allow(clippy::cast_possible_truncation, clippy::cast_sign_loss)] - let delay_ms = (delay_secs * 1000.0) as u64; + let delay_ms = (props.delay_secs * 1000.0) as u64; Some(ProgressEvent::LlmRetry { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - model: prop_string_field(fields, "model").unwrap_or_else(|| "?".to_string()), - attempt: prop_u64_field(fields, "attempt"), + stage_node_id: node_id, + model: props.model.clone(), + attempt: props.attempt as u64, delay_ms, - error: prop_display_field(fields, "error") - .unwrap_or_else(|| "unknown error".to_string()), + error: display_value(&props.error).unwrap_or_else(|| "unknown error".to_string()), }) } - "agent.sub.spawned" => Some(ProgressEvent::SubagentSpawned { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - agent_id: prop_string_field(fields, "agent_id").unwrap_or_else(|| "?".to_string()), - task: prop_string_field(fields, "task").unwrap_or_default(), + EventBody::AgentSubSpawned(props) => Some(ProgressEvent::SubagentSpawned { + stage_node_id: node_id, + agent_id: props.agent_id.clone(), + task: props.task.clone(), }), - "agent.sub.completed" => Some(ProgressEvent::SubagentCompleted { - stage_node_id: string_field(fields, "node_id").unwrap_or_else(|| "?".to_string()), - agent_id: prop_string_field(fields, "agent_id").unwrap_or_else(|| "?".to_string()), - success: prop_bool_field(fields, "success"), - turns_used: prop_u64_field(fields, "turns_used"), + EventBody::AgentSubCompleted(props) => Some(ProgressEvent::SubagentCompleted { + stage_node_id: node_id, + agent_id: props.agent_id.clone(), + success: props.success, + turns_used: props.turns_used as u64, }), - "edge.selected" => Some(ProgressEvent::EdgeSelected { - from_node: prop_string_field(fields, "from_node").unwrap_or_else(|| "?".to_string()), - to_node: prop_string_field(fields, "to_node").unwrap_or_else(|| "?".to_string()), - label: prop_string_field(fields, "label"), - condition: prop_string_field(fields, "condition"), + EventBody::EdgeSelected(props) => Some(ProgressEvent::EdgeSelected { + from_node: props.from_node.clone(), + to_node: props.to_node.clone(), + label: props.label.clone(), + condition: props.condition.clone(), }), - "loop.restart" => Some(ProgressEvent::LoopRestart { - from_node: prop_string_field(fields, "from_node").unwrap_or_else(|| "?".to_string()), - to_node: prop_string_field(fields, "to_node").unwrap_or_else(|| "?".to_string()), + EventBody::LoopRestart(props) => Some(ProgressEvent::LoopRestart { + from_node: props.from_node.clone(), + to_node: props.to_node.clone(), }), - "retro.started" => Some(ProgressEvent::RetroStarted), - "retro.completed" => Some(ProgressEvent::RetroCompleted { - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::RetroStarted(_) => Some(ProgressEvent::RetroStarted), + EventBody::RetroCompleted(props) => Some(ProgressEvent::RetroCompleted { + duration_ms: props.duration_ms, }), - "retro.failed" => Some(ProgressEvent::RetroFailed { - duration_ms: prop_u64_field(fields, "duration_ms"), + EventBody::RetroFailed(props) => Some(ProgressEvent::RetroFailed { + duration_ms: props.duration_ms, }), - "run.notice" => Some(ProgressEvent::RunNotice { - level: parse_run_notice_level(prop_string_field(fields, "level").as_deref()), - code: prop_string_field(fields, "code").unwrap_or_default(), - message: prop_string_field(fields, "message").unwrap_or_default(), + EventBody::RunNotice(props) => Some(ProgressEvent::RunNotice { + level: match props.level { + fabro_types::RunNoticeLevel::Info => RunNoticeLevel::Info, + fabro_types::RunNoticeLevel::Warn => RunNoticeLevel::Warn, + fabro_types::RunNoticeLevel::Error => RunNoticeLevel::Error, + }, + code: props.code.clone(), + message: props.message.clone(), }), - "pull_request.created" => Some(ProgressEvent::PullRequestCreated { - pr_url: prop_string_field(fields, "pr_url").unwrap_or_else(|| "?".to_string()), - draft: prop_bool_field(fields, "draft"), + EventBody::PullRequestCreated(props) => Some(ProgressEvent::PullRequestCreated { + pr_url: props.pr_url.clone(), + draft: props.draft, }), - "pull_request.failed" => Some(ProgressEvent::PullRequestFailed { - error: prop_display_field(fields, "error") - .unwrap_or_else(|| "unknown error".to_string()), + EventBody::PullRequestFailed(props) => Some(ProgressEvent::PullRequestFailed { + error: props.error.clone(), }), _ => None, } } -pub(super) fn from_stored_event(stored: &StoredEvent) -> Option { - let Value::Object(fields) = stored.to_value().ok()? else { - return None; - }; - let event_name = fields.get("event")?.as_str()?; - from_envelope_fields(event_name, &fields) -} - pub(super) fn from_json_line(line: &str) -> Option { - if let Ok(stored) = StoredEvent::from_json_str(line) { - return from_stored_event(&stored); - } - - let Value::Object(fields) = serde_json::from_str(line).ok()? else { - return None; - }; - let event_name = fields.get("event")?.as_str()?; - from_envelope_fields(event_name, &fields) + let stored = StoredEvent::from_json_str(line).ok()?; + from_stored_event(&stored) } -fn parse_run_notice_level(level: Option<&str>) -> RunNoticeLevel { - match level.unwrap_or("info") { - "warn" => RunNoticeLevel::Warn, - "error" => RunNoticeLevel::Error, - _ => RunNoticeLevel::Info, - } -} - -fn string_field(fields: &Map, key: &str) -> Option { - fields.get(key).and_then(Value::as_str).map(str::to_owned) -} - -fn prop_value<'a>(fields: &'a Map, key: &str) -> Option<&'a Value> { - fields - .get("properties") - .and_then(Value::as_object) - .and_then(|properties| properties.get(key)) -} - -fn prop_string_field(fields: &Map, key: &str) -> Option { - prop_value(fields, key) - .and_then(Value::as_str) - .map(str::to_owned) -} - -fn prop_display_field(fields: &Map, key: &str) -> Option { - let value = prop_value(fields, key)?; +fn display_value(value: &Value) -> Option { match value { Value::Null => None, Value::String(value) => Some(value.clone()), @@ -529,73 +491,32 @@ fn prop_display_field(fields: &Map, key: &str) -> Option } } -fn u64_field(fields: &Map, key: &str) -> u64 { - fields.get(key).and_then(Value::as_u64).unwrap_or(0) -} - -fn prop_u64_field(fields: &Map, key: &str) -> u64 { - prop_value(fields, key).and_then(Value::as_u64).unwrap_or(0) -} - -fn prop_optional_u64_field(fields: &Map, key: &str) -> Option { - prop_value(fields, key).and_then(Value::as_u64) -} - -fn prop_i64_field(fields: &Map, key: &str) -> i64 { - prop_value(fields, key).and_then(Value::as_i64).unwrap_or(0) -} - -fn f64_field(fields: &Map, key: &str) -> Option { - fields.get(key).and_then(Value::as_f64) -} - -fn prop_f64_field(fields: &Map, key: &str) -> Option { - prop_value(fields, key).and_then(Value::as_f64) -} - -fn prop_bool_field(fields: &Map, key: &str) -> bool { - prop_value(fields, key) - .and_then(Value::as_bool) - .unwrap_or(false) -} - -fn timestamp_field(fields: &Map, key: &str) -> Option> { - let value = fields.get(key)?.as_str()?; - DateTime::parse_from_rfc3339(value) - .ok() - .map(|timestamp| timestamp.with_timezone(&Utc)) -} - #[cfg(test)] mod tests { use fabro_agent::AgentEvent; use fabro_types::fixtures; - use fabro_workflow::event::{WorkflowRunEvent, canonicalize_event}; + use fabro_workflow::event::{WorkflowRunEvent, to_stored_event}; use super::*; - fn json_map(value: Value) -> Map { - value.as_object().cloned().expect("json object") - } - - fn canonical_fields(event: &WorkflowRunEvent) -> (String, Map) { - let envelope = canonicalize_event(&fixtures::RUN_1, event); - let event_name = envelope.event.clone(); - let fields = json_map(serde_json::to_value(envelope).expect("serializable envelope")); - (event_name, fields) - } - #[test] fn parse_edge_selected() { - let fields = json_map(serde_json::json!({ - "properties": { - "from_node": "a", - "to_node": "b", - "label": "yes" - } - })); + let stored = to_stored_event( + &fixtures::RUN_1, + &WorkflowRunEvent::EdgeSelected { + from_node: "a".into(), + to_node: "b".into(), + label: Some("yes".into()), + condition: None, + reason: "condition".into(), + preferred_label: None, + suggested_next_ids: Vec::new(), + stage_status: "success".into(), + is_jump: false, + }, + ); - let event = from_envelope_fields("edge.selected", &fields).unwrap(); + let event = from_stored_event(&stored).unwrap(); assert!(matches!( event, ProgressEvent::EdgeSelected { @@ -632,8 +553,8 @@ mod tests { max_attempts: 1, }; - let (name, fields) = canonical_fields(&event); - let parsed = from_envelope_fields(&name, &fields).unwrap(); + let stored = to_stored_event(&fixtures::RUN_1, &event); + let parsed = from_stored_event(&stored).unwrap(); assert!(matches!( parsed, ProgressEvent::StageCompleted { @@ -659,8 +580,8 @@ mod tests { parent_session_id: None, }; - let (name, fields) = canonical_fields(&event); - let parsed = from_envelope_fields(&name, &fields).unwrap(); + let stored = to_stored_event(&fixtures::RUN_1, &event); + let parsed = from_stored_event(&stored).unwrap(); assert!(matches!( parsed, ProgressEvent::ToolCallStarted { @@ -673,28 +594,44 @@ mod tests { } #[test] - fn parse_tool_call_timestamps_from_jsonl_envelope() { - let started_fields = json_map(serde_json::json!({ - "ts": "2026-03-30T12:00:00.000Z", - "node_id": "code", - "properties": { - "tool_name": "read_file", - "tool_call_id": "tc1", - "arguments": {"path": "src/main.rs"} - } - })); - let completed_fields = json_map(serde_json::json!({ - "ts": "2026-03-30T12:00:00.500Z", - "node_id": "code", - "properties": { - "tool_call_id": "tc1", - "is_error": false, - "duration_ms": 500 - } - })); - - let started = from_envelope_fields("agent.tool.started", &started_fields).unwrap(); - let completed = from_envelope_fields("agent.tool.completed", &completed_fields).unwrap(); + fn parse_tool_call_timestamps_from_jsonl() { + let started = from_json_line( + &serde_json::json!({ + "id": "evt_1", + "ts": "2026-03-30T12:00:00.000Z", + "run_id": fixtures::RUN_1.to_string(), + "event": "agent.tool.started", + "node_id": "code", + "node_label": "code", + "properties": { + "tool_name": "read_file", + "tool_call_id": "tc1", + "arguments": {"path": "src/main.rs"}, + "visit": 1 + } + }) + .to_string(), + ) + .unwrap(); + let completed = from_json_line( + &serde_json::json!({ + "id": "evt_2", + "ts": "2026-03-30T12:00:00.500Z", + "run_id": fixtures::RUN_1.to_string(), + "event": "agent.tool.completed", + "node_id": "code", + "node_label": "code", + "properties": { + "tool_name": "read_file", + "tool_call_id": "tc1", + "output": {"ok": true}, + "is_error": false, + "visit": 1 + } + }) + .to_string(), + ) + .unwrap(); assert!(matches!( started, @@ -708,7 +645,7 @@ mod tests { assert!(matches!( completed, ProgressEvent::ToolCallCompleted { - duration_ms: Some(500), + duration_ms: None, timestamp: Some(timestamp), .. } if timestamp == DateTime::parse_from_rfc3339("2026-03-30T12:00:00.500Z") @@ -730,8 +667,8 @@ mod tests { }, }; - let (name, fields) = canonical_fields(&event); - let parsed = from_envelope_fields(&name, &fields).unwrap(); + let stored = to_stored_event(&fixtures::RUN_1, &event); + let parsed = from_stored_event(&stored).unwrap(); assert!(matches!( parsed, ProgressEvent::SandboxReady { @@ -751,8 +688,8 @@ mod tests { message: "sandbox cleanup failed".into(), }; - let (name, fields) = canonical_fields(&event); - let parsed = from_envelope_fields(&name, &fields).unwrap(); + let stored = to_stored_event(&fixtures::RUN_1, &event); + let parsed = from_stored_event(&stored).unwrap(); assert!(matches!( parsed, ProgressEvent::RunNotice { 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 be40ac95b..dface5807 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 @@ -414,11 +414,12 @@ mod tests { use std::io::{self, Write}; use std::sync::{Arc, Mutex}; + use chrono::{DateTime, Utc}; use fabro_agent::{AgentEvent, SandboxEvent}; use fabro_llm::types::Usage; use fabro_types::fixtures; use fabro_workflow::event::{ - RunNoticeLevel, WorkflowRunEvent, canonicalize_event, to_stored_event, + RunNoticeLevel, WorkflowRunEvent, to_stored_event, to_stored_event_at, }; use fabro_workflow::outcome::StageUsage; @@ -772,7 +773,7 @@ mod tests { let (mut json_ui, json_buffer) = capture_ui(true); for event in &events { - let line = serde_json::to_string(&canonicalize_event(&fixtures::RUN_1, event)).unwrap(); + let line = serde_json::to_string(&to_stored_event(&fixtures::RUN_1, event)).unwrap(); json_ui.handle_json_line(&line); } @@ -1137,15 +1138,57 @@ mod tests { fn tty_tool_call_completion_uses_jsonl_timestamps() { let mut ui = ProgressUI::new(true, false); - ui.handle_json_line( - r#"{"ts":"2026-03-30T12:00:00.000Z","event":"stage.started","node_id":"code","node_label":"Code","properties":{"attempt":1,"max_attempts":1}}"#, - ); - ui.handle_json_line( - r#"{"ts":"2026-03-30T12:00:00.000Z","event":"agent.tool.started","node_id":"code","properties":{"tool_name":"read_file","tool_call_id":"tc1","arguments":{"path":"src/main.rs"}}}"#, - ); - ui.handle_json_line( - r#"{"ts":"2026-03-30T12:00:00.500Z","event":"agent.tool.completed","node_id":"code","properties":{"tool_call_id":"tc1","is_error":false}}"#, - ); + let started_ts = DateTime::parse_from_rfc3339("2026-03-30T12:00:00.000Z") + .unwrap() + .with_timezone(&Utc); + let completed_ts = DateTime::parse_from_rfc3339("2026-03-30T12:00:00.500Z") + .unwrap() + .with_timezone(&Utc); + + let stage_started = serde_json::to_string(&to_stored_event_at( + &fixtures::RUN_1, + &WorkflowRunEvent::StageStarted { + node_id: "code".into(), + name: "Code".into(), + index: 0, + handler_type: "agent".into(), + attempt: 1, + max_attempts: 1, + }, + started_ts, + )) + .unwrap(); + let tool_started = serde_json::to_string(&to_stored_event_at( + &fixtures::RUN_1, + &agent_event( + "code", + AgentEvent::ToolCallStarted { + tool_name: "read_file".into(), + tool_call_id: "tc1".into(), + arguments: serde_json::json!({"path": "src/main.rs"}), + }, + ), + started_ts, + )) + .unwrap(); + let tool_completed = serde_json::to_string(&to_stored_event_at( + &fixtures::RUN_1, + &agent_event( + "code", + AgentEvent::ToolCallCompleted { + tool_name: "read_file".into(), + tool_call_id: "tc1".into(), + output: serde_json::json!({"ok": true}), + is_error: false, + }, + ), + completed_ts, + )) + .unwrap(); + + ui.handle_json_line(&stage_started); + ui.handle_json_line(&tool_started); + ui.handle_json_line(&tool_completed); let stage = &ui.stage.active_stages["code"]; assert_eq!(stage.tool_calls[0].bar.prefix(), "500ms"); diff --git a/lib/crates/fabro-cli/tests/it/cmd/attach.rs b/lib/crates/fabro-cli/tests/it/cmd/attach.rs index 64566dec2..8f20830b4 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/attach.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/attach.rs @@ -396,7 +396,6 @@ fn attach_json_errors_without_prompting_for_human_input() { } }, "host_repo_path": "[TEMP_DIR]", - "labels": {}, "run_dir": "[RUN_DIR]", "settings": { "goal": "Wait for approval", diff --git a/lib/crates/fabro-cli/tests/it/cmd/logs.rs b/lib/crates/fabro-cli/tests/it/cmd/logs.rs index 837d534c8..ace2a0fa6 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/logs.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/logs.rs @@ -63,7 +63,7 @@ fn logs_completed_run_outputs_raw_ndjson() { success: true exit_code: 0 ----- stdout ----- - {"event":"run.created","id":"[EVENT_ID]","properties":{"graph":{"attrs":{"goal":{"String":"Run tests and report results"},"rankdir":{"String":"LR"}},"edges":[{"attrs":{},"from":"start","to":"run_tests"},{"attrs":{},"from":"run_tests","to":"report"},{"attrs":{},"from":"report","to":"exit"}],"name":"Simple","nodes":{"exit":{"attrs":{"label":{"String":"Exit"},"shape":{"String":"Msquare"}},"id":"exit"},"report":{"attrs":{"label":{"String":"Report"},"prompt":{"String":"Summarize the test results"}},"id":"report"},"run_tests":{"attrs":{"label":{"String":"Run Tests"},"prompt":{"String":"Run the test suite and report results"}},"id":"run_tests"},"start":{"attrs":{"label":{"String":"Start"},"shape":{"String":"Mdiamond"}},"id":"start"}}},"host_repo_path":"[TEMP_DIR]","labels":{},"run_dir":"[STORAGE_DIR]/runs/20260404-[ULID]","settings":{"auto_approve":true,"dry_run":true,"fabro":{"root":"fabro/"},"features":{"retros":false,"session_sandboxes":false},"goal":"Run tests and report results","hooks":[{"blocking":true,"command":"cargo fmt","event":"post_tool_use","matcher":"write_file|edit_file|apply_patch","name":"cargo-fmt","sandbox":null,"timeout_ms":null}],"llm":{"fallbacks":null,"model":"claude-sonnet-4-6","provider":"anthropic"},"mode":"standalone","no_retro":true,"pull_request":{"auto_merge":false,"draft":false,"enabled":true,"merge_strategy":"squash"},"sandbox":{"daytona":{"auto_stop_interval":30,"labels":{"repo":"fabro-sh/fabro"},"network":null,"skip_clone":false,"snapshot":{"cpu":4,"disk":20,"dockerfile":"FROM ubuntu:24.04/n/nRUN apt-get update && apt-get install -y --no-install-recommends curl git ca-certificates build-essential pkg-config libssl-dev unzip python3 && rm -rf /var/lib/apt/lists/*/n/n# GitHub CLI/nRUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg && echo \"deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main\" | tee /etc/apt/sources.list.d/github-cli.list > /dev/null && apt-get update && apt-get install -y --no-install-recommends gh && rm -rf /var/lib/apt/lists/*/n/n# Rust/nRUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y/nENV PATH=\"/root/.cargo/bin:${PATH}\"/nRUN cargo install cargo-nextest --locked/nENV CARGO_INCREMENTAL=0/n/n# Bun/nRUN curl -fsSL https://bun.sh/install | bash/nENV PATH=\"/root/.bun/bin:${PATH}\"/n/nWORKDIR /root/n","memory":8,"name":"fabro-v6"}},"devcontainer":null,"env":null,"local":null,"preserve":null,"provider":"local"},"storage_dir":"[STORAGE_DIR]","version":1},"workflow_slug":"simple","workflow_source":"digraph Simple {/n graph [goal=\"Run tests and report results\"]/n rankdir=LR/n/n start [shape=Mdiamond, label=\"Start\"]/n exit [shape=Msquare, label=\"Exit\"]/n/n run_tests [label=\"Run Tests\", prompt=\"Run the test suite and report results\"]/n report [label=\"Report\", prompt=\"Summarize the test results\"]/n/n start -> run_tests -> report -> exit/n}/n","working_directory":"[TEMP_DIR]"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} + {"event":"run.created","id":"[EVENT_ID]","properties":{"graph":{"attrs":{"goal":{"String":"Run tests and report results"},"rankdir":{"String":"LR"}},"edges":[{"attrs":{},"from":"start","to":"run_tests"},{"attrs":{},"from":"run_tests","to":"report"},{"attrs":{},"from":"report","to":"exit"}],"name":"Simple","nodes":{"exit":{"attrs":{"label":{"String":"Exit"},"shape":{"String":"Msquare"}},"id":"exit"},"report":{"attrs":{"label":{"String":"Report"},"prompt":{"String":"Summarize the test results"}},"id":"report"},"run_tests":{"attrs":{"label":{"String":"Run Tests"},"prompt":{"String":"Run the test suite and report results"}},"id":"run_tests"},"start":{"attrs":{"label":{"String":"Start"},"shape":{"String":"Mdiamond"}},"id":"start"}}},"host_repo_path":"[TEMP_DIR]","run_dir":"[STORAGE_DIR]/runs/20260404-[ULID]","settings":{"auto_approve":true,"dry_run":true,"fabro":{"root":"fabro/"},"features":{"retros":false,"session_sandboxes":false},"goal":"Run tests and report results","hooks":[{"blocking":true,"command":"cargo fmt","event":"post_tool_use","matcher":"write_file|edit_file|apply_patch","name":"cargo-fmt","sandbox":null,"timeout_ms":null}],"llm":{"fallbacks":null,"model":"claude-sonnet-4-6","provider":"anthropic"},"mode":"standalone","no_retro":true,"pull_request":{"auto_merge":false,"draft":false,"enabled":true,"merge_strategy":"squash"},"sandbox":{"daytona":{"auto_stop_interval":30,"labels":{"repo":"fabro-sh/fabro"},"network":null,"skip_clone":false,"snapshot":{"cpu":4,"disk":20,"dockerfile":"FROM ubuntu:24.04/n/nRUN apt-get update && apt-get install -y --no-install-recommends curl git ca-certificates build-essential pkg-config libssl-dev unzip python3 && rm -rf /var/lib/apt/lists/*/n/n# GitHub CLI/nRUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg && echo \"deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main\" | tee /etc/apt/sources.list.d/github-cli.list > /dev/null && apt-get update && apt-get install -y --no-install-recommends gh && rm -rf /var/lib/apt/lists/*/n/n# Rust/nRUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y/nENV PATH=\"/root/.cargo/bin:${PATH}\"/nRUN cargo install cargo-nextest --locked/nENV CARGO_INCREMENTAL=0/n/n# Bun/nRUN curl -fsSL https://bun.sh/install | bash/nENV PATH=\"/root/.bun/bin:${PATH}\"/n/nWORKDIR /root/n","memory":8,"name":"fabro-v6"}},"devcontainer":null,"env":null,"local":null,"preserve":null,"provider":"local"},"storage_dir":"[STORAGE_DIR]","version":1},"workflow_slug":"simple","workflow_source":"digraph Simple {/n graph [goal=\"Run tests and report results\"]/n rankdir=LR/n/n start [shape=Mdiamond, label=\"Start\"]/n exit [shape=Msquare, label=\"Exit\"]/n/n run_tests [label=\"Run Tests\", prompt=\"Run the test suite and report results\"]/n report [label=\"Report\", prompt=\"Summarize the test results\"]/n/n start -> run_tests -> report -> exit/n}/n","working_directory":"[TEMP_DIR]"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"run.submitted","id":"[EVENT_ID]","properties":{},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"run.starting","id":"[EVENT_ID]","properties":{"reason":"sandbox_initializing"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"sandbox.initializing","id":"[EVENT_ID]","properties":{"provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} @@ -228,7 +228,7 @@ fn logs_follow_detached_run_streams_until_completion() { success: true exit_code: 0 ----- stdout ----- - {"event":"run.created","id":"[EVENT_ID]","properties":{"graph":{"attrs":{"goal":{"String":"Run tests and report results"},"rankdir":{"String":"LR"}},"edges":[{"attrs":{},"from":"start","to":"run_tests"},{"attrs":{},"from":"run_tests","to":"report"},{"attrs":{},"from":"report","to":"exit"}],"name":"Simple","nodes":{"exit":{"attrs":{"label":{"String":"Exit"},"shape":{"String":"Msquare"}},"id":"exit"},"report":{"attrs":{"label":{"String":"Report"},"prompt":{"String":"Summarize the test results"}},"id":"report"},"run_tests":{"attrs":{"label":{"String":"Run Tests"},"prompt":{"String":"Run the test suite and report results"}},"id":"run_tests"},"start":{"attrs":{"label":{"String":"Start"},"shape":{"String":"Mdiamond"}},"id":"start"}}},"host_repo_path":"[TEMP_DIR]","labels":{},"run_dir":"[STORAGE_DIR]/runs/20260404-[ULID]","settings":{"auto_approve":true,"dry_run":true,"fabro":{"root":"fabro/"},"features":{"retros":false,"session_sandboxes":false},"goal":"Run tests and report results","hooks":[{"blocking":true,"command":"cargo fmt","event":"post_tool_use","matcher":"write_file|edit_file|apply_patch","name":"cargo-fmt","sandbox":null,"timeout_ms":null}],"llm":{"fallbacks":null,"model":"claude-sonnet-4-6","provider":"anthropic"},"mode":"standalone","no_retro":true,"pull_request":{"auto_merge":false,"draft":false,"enabled":true,"merge_strategy":"squash"},"sandbox":{"daytona":{"auto_stop_interval":30,"labels":{"repo":"fabro-sh/fabro"},"network":null,"skip_clone":false,"snapshot":{"cpu":4,"disk":20,"dockerfile":"FROM ubuntu:24.04/n/nRUN apt-get update && apt-get install -y --no-install-recommends curl git ca-certificates build-essential pkg-config libssl-dev unzip python3 && rm -rf /var/lib/apt/lists/*/n/n# GitHub CLI/nRUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg && echo \"deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main\" | tee /etc/apt/sources.list.d/github-cli.list > /dev/null && apt-get update && apt-get install -y --no-install-recommends gh && rm -rf /var/lib/apt/lists/*/n/n# Rust/nRUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y/nENV PATH=\"/root/.cargo/bin:${PATH}\"/nRUN cargo install cargo-nextest --locked/nENV CARGO_INCREMENTAL=0/n/n# Bun/nRUN curl -fsSL https://bun.sh/install | bash/nENV PATH=\"/root/.bun/bin:${PATH}\"/n/nWORKDIR /root/n","memory":8,"name":"fabro-v6"}},"devcontainer":null,"env":null,"local":null,"preserve":null,"provider":"local"},"storage_dir":"[STORAGE_DIR]","version":1},"workflow_slug":"simple","workflow_source":"digraph Simple {/n graph [goal=\"Run tests and report results\"]/n rankdir=LR/n/n start [shape=Mdiamond, label=\"Start\"]/n exit [shape=Msquare, label=\"Exit\"]/n/n run_tests [label=\"Run Tests\", prompt=\"Run the test suite and report results\"]/n report [label=\"Report\", prompt=\"Summarize the test results\"]/n/n start -> run_tests -> report -> exit/n}/n","working_directory":"[TEMP_DIR]"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} + {"event":"run.created","id":"[EVENT_ID]","properties":{"graph":{"attrs":{"goal":{"String":"Run tests and report results"},"rankdir":{"String":"LR"}},"edges":[{"attrs":{},"from":"start","to":"run_tests"},{"attrs":{},"from":"run_tests","to":"report"},{"attrs":{},"from":"report","to":"exit"}],"name":"Simple","nodes":{"exit":{"attrs":{"label":{"String":"Exit"},"shape":{"String":"Msquare"}},"id":"exit"},"report":{"attrs":{"label":{"String":"Report"},"prompt":{"String":"Summarize the test results"}},"id":"report"},"run_tests":{"attrs":{"label":{"String":"Run Tests"},"prompt":{"String":"Run the test suite and report results"}},"id":"run_tests"},"start":{"attrs":{"label":{"String":"Start"},"shape":{"String":"Mdiamond"}},"id":"start"}}},"host_repo_path":"[TEMP_DIR]","run_dir":"[STORAGE_DIR]/runs/20260404-[ULID]","settings":{"auto_approve":true,"dry_run":true,"fabro":{"root":"fabro/"},"features":{"retros":false,"session_sandboxes":false},"goal":"Run tests and report results","hooks":[{"blocking":true,"command":"cargo fmt","event":"post_tool_use","matcher":"write_file|edit_file|apply_patch","name":"cargo-fmt","sandbox":null,"timeout_ms":null}],"llm":{"fallbacks":null,"model":"claude-sonnet-4-6","provider":"anthropic"},"mode":"standalone","no_retro":true,"pull_request":{"auto_merge":false,"draft":false,"enabled":true,"merge_strategy":"squash"},"sandbox":{"daytona":{"auto_stop_interval":30,"labels":{"repo":"fabro-sh/fabro"},"network":null,"skip_clone":false,"snapshot":{"cpu":4,"disk":20,"dockerfile":"FROM ubuntu:24.04/n/nRUN apt-get update && apt-get install -y --no-install-recommends curl git ca-certificates build-essential pkg-config libssl-dev unzip python3 && rm -rf /var/lib/apt/lists/*/n/n# GitHub CLI/nRUN curl -fsSL https://cli.github.com/packages/githubcli-archive-keyring.gpg | dd of=/usr/share/keyrings/githubcli-archive-keyring.gpg && echo \"deb [arch=$(dpkg --print-architecture) signed-by=/usr/share/keyrings/githubcli-archive-keyring.gpg] https://cli.github.com/packages stable main\" | tee /etc/apt/sources.list.d/github-cli.list > /dev/null && apt-get update && apt-get install -y --no-install-recommends gh && rm -rf /var/lib/apt/lists/*/n/n# Rust/nRUN curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh -s -- -y/nENV PATH=\"/root/.cargo/bin:${PATH}\"/nRUN cargo install cargo-nextest --locked/nENV CARGO_INCREMENTAL=0/n/n# Bun/nRUN curl -fsSL https://bun.sh/install | bash/nENV PATH=\"/root/.bun/bin:${PATH}\"/n/nWORKDIR /root/n","memory":8,"name":"fabro-v6"}},"devcontainer":null,"env":null,"local":null,"preserve":null,"provider":"local"},"storage_dir":"[STORAGE_DIR]","version":1},"workflow_slug":"simple","workflow_source":"digraph Simple {/n graph [goal=\"Run tests and report results\"]/n rankdir=LR/n/n start [shape=Mdiamond, label=\"Start\"]/n exit [shape=Msquare, label=\"Exit\"]/n/n run_tests [label=\"Run Tests\", prompt=\"Run the test suite and report results\"]/n report [label=\"Report\", prompt=\"Summarize the test results\"]/n/n start -> run_tests -> report -> exit/n}/n","working_directory":"[TEMP_DIR]"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"run.submitted","id":"[EVENT_ID]","properties":{},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"run.starting","id":"[EVENT_ID]","properties":{"reason":"sandbox_initializing"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} {"event":"sandbox.initializing","id":"[EVENT_ID]","properties":{"provider":"local"},"run_id":"[ULID]","ts":"[TIMESTAMP]"} diff --git a/lib/crates/fabro-cli/tests/it/cmd/run.rs b/lib/crates/fabro-cli/tests/it/cmd/run.rs index c4828cc62..83eea71fe 100644 --- a/lib/crates/fabro-cli/tests/it/cmd/run.rs +++ b/lib/crates/fabro-cli/tests/it/cmd/run.rs @@ -342,7 +342,6 @@ fn json_run_implies_auto_approve_for_human_gates() { } }, "host_repo_path": "[TEMP_DIR]", - "labels": {}, "run_dir": "[RUN_DIR]", "settings": { "auto_approve": true, diff --git a/lib/crates/fabro-workflow/src/event.rs b/lib/crates/fabro-workflow/src/event.rs index 3374f7f8d..20d143386 100644 --- a/lib/crates/fabro-workflow/src/event.rs +++ b/lib/crates/fabro-workflow/src/event.rs @@ -6,7 +6,7 @@ use chrono::{SecondsFormat, Utc}; use fabro_store::{EventPayload, SlateRunStore}; use fabro_types::{RunId, StoredEvent}; use serde::{Deserialize, Serialize}; -use serde_json::{Map, Value}; +use serde_json::{Map, Value, json}; use std::collections::BTreeMap; use tokio::sync::{mpsc, oneshot}; use uuid::Uuid; @@ -20,30 +20,6 @@ use fabro_util::redact::redact_jsonl_line; pub use fabro_types::{EventBody, RunNoticeLevel}; -#[derive(Debug, Clone, Serialize, Deserialize)] -pub struct RunEventEnvelope { - pub id: String, - pub ts: String, - pub run_id: String, - pub event: String, - #[serde(skip_serializing_if = "Option::is_none")] - pub session_id: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub parent_session_id: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub node_id: Option, - #[serde(skip_serializing_if = "Option::is_none")] - pub node_label: Option, - pub properties: serde_json::Value, -} - -impl From<&RunEventEnvelope> for StoredEvent { - fn from(value: &RunEventEnvelope) -> Self { - StoredEvent::from_value(serde_json::to_value(value).expect("event envelope serializes")) - .expect("event envelope converts to stored event") - } -} - /// Events emitted during workflow run execution for observability. #[derive(Debug, Clone, Serialize, Deserialize)] #[allow(clippy::large_enum_variant)] @@ -1201,7 +1177,7 @@ pub fn event_name(event: &WorkflowRunEvent) -> &'static str { } #[derive(Debug)] -struct EnvelopeFields { +struct StoredEventFields { session_id: Option, parent_session_id: Option, node_id: Option, @@ -1243,33 +1219,16 @@ fn remove_string(fields: &mut Map, key: &str) -> Option { } } -fn flatten_failure_detail(fields: &mut Map) { - let Some(Value::Object(failure)) = fields.remove("failure") else { - return; - }; - if let Some(message) = failure.get("message").cloned() { - fields.insert("error".to_string(), message); - } - if let Some(failure_class) = failure.get("failure_class").cloned() { - fields.insert("failure_class".to_string(), failure_class); - } - if let Some(failure_signature) = failure.get("failure_signature").cloned() { - if !failure_signature.is_null() { - fields.insert("failure_signature".to_string(), failure_signature); - } - } -} - fn default_node_label(node_id: Option<&String>, node_label: Option) -> Option { node_label.or_else(|| node_id.cloned()) } -fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { +fn extract_stored_event_fields(event: &WorkflowRunEvent) -> StoredEventFields { match event { WorkflowRunEvent::RunCreated { .. } | WorkflowRunEvent::WorkflowRunStarted { .. } => { let mut fields = tagged_variant_fields(event); fields.remove("run_id"); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id: None, @@ -1280,7 +1239,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { WorkflowRunEvent::WorkflowRunFailed { error, .. } => { let mut fields = tagged_variant_fields(event); fields.insert("error".to_string(), Value::String(error.to_string())); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id: None, @@ -1293,8 +1252,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { let node_id = remove_string(&mut fields, "node_id"); let node_label = default_node_label(node_id.as_ref(), remove_string(&mut fields, "name")); - flatten_failure_detail(&mut fields); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1320,7 +1278,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { let node_id = remove_string(&mut fields, "node_id"); let node_label = default_node_label(node_id.as_ref(), remove_string(&mut fields, "name")); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1346,7 +1304,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { if let (Some(visit), Value::Object(map)) = (visit, &mut properties) { map.insert("visit".to_string(), visit); } - EnvelopeFields { + StoredEventFields { session_id: session_id.clone(), parent_session_id: parent_session_id.clone(), node_id, @@ -1360,7 +1318,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { || Value::Object(Map::new()), |value| Value::Object(tagged_variant_fields_from_value(value)), ); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id: None, @@ -1372,7 +1330,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { 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); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1385,7 +1343,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { 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); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1400,7 +1358,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { 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); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1412,7 +1370,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { 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); - EnvelopeFields { + StoredEventFields { session_id: None, parent_session_id: None, node_id, @@ -1420,7 +1378,7 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { properties: Value::Object(fields), } } - _ => EnvelopeFields { + _ => StoredEventFields { session_id: None, parent_session_id: None, node_id: None, @@ -1430,29 +1388,6 @@ fn extract_envelope_fields(event: &WorkflowRunEvent) -> EnvelopeFields { } } -pub fn canonicalize_event(run_id: &RunId, event: &WorkflowRunEvent) -> RunEventEnvelope { - canonicalize_event_at(run_id, event, Utc::now()) -} - -pub fn canonicalize_event_at( - run_id: &RunId, - event: &WorkflowRunEvent, - ts: chrono::DateTime, -) -> RunEventEnvelope { - let fields = extract_envelope_fields(event); - RunEventEnvelope { - id: Uuid::now_v7().to_string(), - ts: ts.to_rfc3339_opts(SecondsFormat::Millis, true), - run_id: run_id.to_string(), - event: event_name(event).to_string(), - session_id: fields.session_id, - parent_session_id: fields.parent_session_id, - node_id: fields.node_id, - node_label: fields.node_label, - properties: fields.properties, - } -} - pub fn to_stored_event(run_id: &RunId, event: &WorkflowRunEvent) -> StoredEvent { to_stored_event_at(run_id, event, Utc::now()) } @@ -1462,21 +1397,19 @@ pub fn to_stored_event_at( event: &WorkflowRunEvent, ts: chrono::DateTime, ) -> StoredEvent { - let envelope = canonicalize_event_at(run_id, event, ts); - let mut stored = StoredEvent::from(&envelope); - - match (event, &mut stored.body) { - (WorkflowRunEvent::StageCompleted { failure, .. }, EventBody::StageCompleted(props)) => { - props.failure = failure.clone(); - } - (WorkflowRunEvent::StageFailed { failure, .. }, EventBody::StageFailed(props)) => { - props.failure = Some(failure.clone()); - } - _ => {} - } - - stored.refresh_cache(); - stored + let fields = extract_stored_event_fields(event); + StoredEvent::from_value(json!({ + "id": Uuid::now_v7().to_string(), + "ts": ts.to_rfc3339_opts(SecondsFormat::Millis, true), + "run_id": run_id.to_string(), + "event": event_name(event), + "session_id": fields.session_id, + "parent_session_id": fields.parent_session_id, + "node_id": fields.node_id, + "node_label": fields.node_label, + "properties": fields.properties, + })) + .expect("workflow event converts to stored event") } pub fn build_redacted_event_payload(event: &StoredEvent, run_id: &RunId) -> Result { @@ -1752,8 +1685,8 @@ mod tests { } #[test] - fn canonicalize_stage_completed_places_node_fields_in_envelope() { - let envelope = canonicalize_event( + fn stored_stage_completed_places_node_fields_in_header() { + let stored = to_stored_event( &fixtures::RUN_2, &WorkflowRunEvent::StageCompleted { node_id: "plan".to_string(), @@ -1779,18 +1712,18 @@ mod tests { }, ); - assert_eq!(envelope.event, "stage.completed"); - assert_eq!(envelope.run_id, fixtures::RUN_2.to_string()); - assert_eq!(envelope.node_id.as_deref(), Some("plan")); - assert_eq!(envelope.node_label.as_deref(), Some("Plan")); - assert_eq!(envelope.properties["duration_ms"], 5000); - assert_eq!(envelope.properties["status"], "success"); - assert!(envelope.session_id.is_none()); + assert_eq!(stored.event_name(), "stage.completed"); + assert_eq!(stored.run_id, fixtures::RUN_2); + assert_eq!(stored.node_id.as_deref(), Some("plan")); + assert_eq!(stored.node_label.as_deref(), Some("Plan")); + assert_eq!(stored.properties["duration_ms"], 5000); + assert_eq!(stored.properties["status"], "success"); + assert!(stored.session_id.is_none()); } #[test] - fn canonicalize_stage_completed_keeps_response_and_signature_snapshots() { - let envelope = canonicalize_event( + fn stored_stage_completed_keeps_response_and_signature_snapshots() { + let stored = to_stored_event( &fixtures::RUN_2, &WorkflowRunEvent::StageCompleted { node_id: "plan".to_string(), @@ -1816,17 +1749,14 @@ mod tests { }, ); - assert_eq!(envelope.properties["response"], "done"); - assert_eq!(envelope.properties["loop_failure_signatures"]["sig-a"], 2); - assert_eq!( - envelope.properties["restart_failure_signatures"]["sig-b"], - 1 - ); + assert_eq!(stored.properties["response"], "done"); + assert_eq!(stored.properties["loop_failure_signatures"]["sig-a"], 2); + assert_eq!(stored.properties["restart_failure_signatures"]["sig-b"], 1); } #[test] - fn canonicalize_stage_failure_flattens_failure_detail() { - let envelope = canonicalize_event( + fn stored_stage_failure_keeps_failure_detail() { + let stored = to_stored_event( &fixtures::RUN_3, &WorkflowRunEvent::StageFailed { node_id: "code".to_string(), @@ -1840,16 +1770,18 @@ mod tests { }, ); - assert_eq!(envelope.event, "stage.failed"); - assert_eq!(envelope.properties["error"], "lint failed"); - assert_eq!(envelope.properties["failure_class"], "deterministic"); - assert_eq!(envelope.properties["will_retry"], true); - assert!(envelope.properties.get("failure").is_none()); + assert_eq!(stored.event_name(), "stage.failed"); + assert_eq!(stored.properties["failure"]["message"], "lint failed"); + assert_eq!( + stored.properties["failure"]["failure_class"], + "deterministic" + ); + assert_eq!(stored.properties["will_retry"], true); } #[test] - fn canonicalize_agent_tool_started_moves_session_metadata_to_envelope() { - let envelope = canonicalize_event( + fn stored_agent_tool_started_moves_session_metadata_to_header() { + let stored = to_stored_event( &fixtures::RUN_4, &WorkflowRunEvent::Agent { stage: "code".to_string(), @@ -1864,19 +1796,19 @@ mod tests { }, ); - assert_eq!(envelope.event, "agent.tool.started"); - assert_eq!(envelope.node_id.as_deref(), Some("code")); - assert_eq!(envelope.node_label.as_deref(), Some("code")); - assert_eq!(envelope.session_id.as_deref(), Some("ses_child")); - assert_eq!(envelope.parent_session_id.as_deref(), Some("ses_parent")); - assert_eq!(envelope.properties["tool_name"], "read_file"); - assert_eq!(envelope.properties["tool_call_id"], "call_1"); - assert_eq!(envelope.properties["visit"], 2); + assert_eq!(stored.event_name(), "agent.tool.started"); + assert_eq!(stored.node_id.as_deref(), Some("code")); + assert_eq!(stored.node_label.as_deref(), Some("code")); + assert_eq!(stored.session_id.as_deref(), Some("ses_child")); + assert_eq!(stored.parent_session_id.as_deref(), Some("ses_parent")); + assert_eq!(stored.properties["tool_name"], "read_file"); + assert_eq!(stored.properties["tool_call_id"], "call_1"); + assert_eq!(stored.properties["visit"], 2); } #[test] - fn canonicalize_sandbox_event_keeps_properties_nested() { - let envelope = canonicalize_event( + fn stored_sandbox_event_keeps_properties_nested() { + let stored = to_stored_event( &fixtures::RUN_5, &WorkflowRunEvent::Sandbox { event: SandboxEvent::Ready { @@ -1890,15 +1822,15 @@ mod tests { }, ); - assert_eq!(envelope.event, "sandbox.ready"); - assert!(envelope.node_id.is_none()); - assert_eq!(envelope.properties["provider"], "daytona"); - assert_eq!(envelope.properties["duration_ms"], 2500); + assert_eq!(stored.event_name(), "sandbox.ready"); + assert!(stored.node_id.is_none()); + assert_eq!(stored.properties["provider"], "daytona"); + assert_eq!(stored.properties["duration_ms"], 2500); } #[test] - fn canonicalize_workflow_failure_flattens_error_display() { - let envelope = canonicalize_event( + fn stored_workflow_failure_uses_display_error() { + let stored = to_stored_event( &fixtures::RUN_6, &WorkflowRunEvent::WorkflowRunFailed { error: FabroError::handler("boom"), @@ -1908,9 +1840,9 @@ mod tests { }, ); - assert_eq!(envelope.event, "run.failed"); - assert_eq!(envelope.properties["error"], "Handler error: boom"); - assert_eq!(envelope.properties["duration_ms"], 900); + assert_eq!(stored.event_name(), "run.failed"); + assert_eq!(stored.properties["error"], "Handler error: boom"); + assert_eq!(stored.properties["duration_ms"], 900); } #[tokio::test] @@ -1921,7 +1853,7 @@ mod tests { std::time::Duration::from_millis(1), ); let run_store = store.create_run(&fixtures::RUN_7).await.unwrap(); - let envelope = canonicalize_event( + let stored = to_stored_event( &fixtures::RUN_7, &WorkflowRunEvent::RunNotice { level: RunNoticeLevel::Warn, @@ -1929,8 +1861,6 @@ mod tests { message: "notice".to_string(), }, ); - - let stored = StoredEvent::from(&envelope); let payload = build_redacted_event_payload(&stored, &fixtures::RUN_7).unwrap(); run_store.append_event(&payload).await.unwrap(); @@ -1947,7 +1877,7 @@ mod tests { #[test] fn build_redacted_event_payload_requires_id() { - let envelope = canonicalize_event( + let stored = to_stored_event( &fixtures::RUN_8, &WorkflowRunEvent::RetroStarted { prompt: Some("Analyze the run".to_string()), @@ -1955,10 +1885,8 @@ mod tests { model: None, }, ); - - let stored = StoredEvent::from(&envelope); let payload = build_redacted_event_payload(&stored, &fixtures::RUN_8).unwrap(); - assert_eq!(payload.as_value()["id"], envelope.id); + assert_eq!(payload.as_value()["id"], stored.id); assert_eq!(payload.as_value()["event"], "retro.started"); assert_eq!( payload.as_value()["properties"]["prompt"], diff --git a/lib/crates/fabro-workflow/src/operations/create.rs b/lib/crates/fabro-workflow/src/operations/create.rs index 029a69c4b..a8f13514d 100644 --- a/lib/crates/fabro-workflow/src/operations/create.rs +++ b/lib/crates/fabro-workflow/src/operations/create.rs @@ -18,7 +18,7 @@ use fabro_sandbox::daytona::detect_repo_info; use super::source::{ResolveWorkflowInput, WorkflowInput, resolve_workflow}; use crate::event::{ - WorkflowRunEvent, append_workflow_event, canonicalize_event_at, normalize_json_value, + WorkflowRunEvent, append_workflow_event, normalize_json_value, to_stored_event_at, }; #[derive(Clone, Debug)] @@ -136,7 +136,7 @@ async fn persist_created_run( .map_err(|_| FabroError::engine(err.to_string()))?, }; - let envelope = canonicalize_event_at( + let stored = to_stored_event_at( &record.run_id, &WorkflowRunEvent::RunCreated { run_id: record.run_id, @@ -165,7 +165,7 @@ async fn persist_created_run( record.run_id.created_at(), ); let payload = fabro_store::EventPayload::new( - serde_json::to_value(&envelope).map_err(|err| FabroError::engine(err.to_string()))?, + serde_json::to_value(&stored).map_err(|err| FabroError::engine(err.to_string()))?, &record.run_id, ) .map_err(store_error)?;