refactor: remove legacy stored event compatibility

This commit is contained in:
Bryan Helmkamp 2026-04-04 10:51:15 -04:00
parent 89a55b37e3
commit 93a7f19383
No known key found for this signature in database
7 changed files with 380 additions and 474 deletions

View file

@ -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<Self> {
let Value::Object(fields) = value else {
return None;
};
pub(super) fn from_stage_usage(usage: &StageUsage) -> Option<Self> {
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<String, Value>) -> Option<ProgressEvent> {
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<ProgressEvent> {
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<ProgressEvent> {
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<ProgressEvent> {
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<String, Value>, key: &str) -> Option<String> {
fields.get(key).and_then(Value::as_str).map(str::to_owned)
}
fn prop_value<'a>(fields: &'a Map<String, Value>, 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<String, Value>, key: &str) -> Option<String> {
prop_value(fields, key)
.and_then(Value::as_str)
.map(str::to_owned)
}
fn prop_display_field(fields: &Map<String, Value>, key: &str) -> Option<String> {
let value = prop_value(fields, key)?;
fn display_value(value: &Value) -> Option<String> {
match value {
Value::Null => None,
Value::String(value) => Some(value.clone()),
@ -529,73 +491,32 @@ fn prop_display_field(fields: &Map<String, Value>, key: &str) -> Option<String>
}
}
fn u64_field(fields: &Map<String, Value>, key: &str) -> u64 {
fields.get(key).and_then(Value::as_u64).unwrap_or(0)
}
fn prop_u64_field(fields: &Map<String, Value>, key: &str) -> u64 {
prop_value(fields, key).and_then(Value::as_u64).unwrap_or(0)
}
fn prop_optional_u64_field(fields: &Map<String, Value>, key: &str) -> Option<u64> {
prop_value(fields, key).and_then(Value::as_u64)
}
fn prop_i64_field(fields: &Map<String, Value>, key: &str) -> i64 {
prop_value(fields, key).and_then(Value::as_i64).unwrap_or(0)
}
fn f64_field(fields: &Map<String, Value>, key: &str) -> Option<f64> {
fields.get(key).and_then(Value::as_f64)
}
fn prop_f64_field(fields: &Map<String, Value>, key: &str) -> Option<f64> {
prop_value(fields, key).and_then(Value::as_f64)
}
fn prop_bool_field(fields: &Map<String, Value>, key: &str) -> bool {
prop_value(fields, key)
.and_then(Value::as_bool)
.unwrap_or(false)
}
fn timestamp_field(fields: &Map<String, Value>, key: &str) -> Option<DateTime<Utc>> {
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<String, Value> {
value.as_object().cloned().expect("json object")
}
fn canonical_fields(event: &WorkflowRunEvent) -> (String, Map<String, Value>) {
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 {

View file

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

View file

@ -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",

View file

@ -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]"}

View file

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

View file

@ -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<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub parent_session_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub node_id: Option<String>,
#[serde(skip_serializing_if = "Option::is_none")]
pub node_label: Option<String>,
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<String>,
parent_session_id: Option<String>,
node_id: Option<String>,
@ -1243,33 +1219,16 @@ fn remove_string(fields: &mut Map<String, Value>, key: &str) -> Option<String> {
}
}
fn flatten_failure_detail(fields: &mut Map<String, Value>) {
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<String>) -> Option<String> {
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<Utc>,
) -> 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<Utc>,
) -> 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<EventPayload> {
@ -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"],

View file

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