mirror of
https://github.com/fabro-sh/fabro.git
synced 2026-10-10 03:30:59 +00:00
Rename progress.jsonl fields for clarity
Add node_id to Stage* enum variants so both the programmatic ID and display label are available. Add rename_fields() post-processing in flatten_event() to give flattened JSONL fields self-describing names: - timestamp → ts (save space) - name → node_label (Stage*), workflow_name, snapshot_name - index → stage_index, branch_index, command_index - stage → node_id (Agent.*, Interview*, Prompt) - branch → node_id (ParallelBranch*) - node → node_id (StallWatchdogTimeout) - from_node/to_node → from_node_id/to_node_id - start_node → start_node_id - provider → sandbox_provider (Sandbox.*) - text → prompt_text (Prompt) - Insert node_label defaulting to node_id where only an id exists Co-Authored-By: Claude Opus 4.6 <noreply@anthropic.com>
This commit is contained in:
parent
b4f1495982
commit
c560c608e2
6 changed files with 323 additions and 28 deletions
|
|
@ -528,11 +528,11 @@ fn dry_run_writes_jsonl_and_live_json() {
|
|||
"progress.jsonl should have at least one line"
|
||||
);
|
||||
|
||||
// Every line must be valid JSON with timestamp, run_id, and event keys
|
||||
// Every line must be valid JSON with ts, run_id, and event keys
|
||||
let first_line: serde_json::Value = serde_json::from_str(lines[0]).unwrap();
|
||||
assert!(
|
||||
first_line.get("timestamp").is_some(),
|
||||
"line should have timestamp"
|
||||
first_line.get("ts").is_some(),
|
||||
"line should have ts"
|
||||
);
|
||||
assert!(
|
||||
first_line.get("run_id").is_some(),
|
||||
|
|
@ -560,7 +560,7 @@ fn dry_run_writes_jsonl_and_live_json() {
|
|||
assert!(live_path.exists(), "live.json should exist");
|
||||
let live_content: serde_json::Value =
|
||||
serde_json::from_str(&std::fs::read_to_string(&live_path).unwrap()).unwrap();
|
||||
assert!(live_content.get("timestamp").is_some());
|
||||
assert!(live_content.get("ts").is_some());
|
||||
assert!(live_content.get("run_id").is_some());
|
||||
assert!(live_content.get("event").is_some());
|
||||
}
|
||||
|
|
|
|||
|
|
@ -217,13 +217,14 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
format!("[WORKFLOW_RUN_FAILED] error=\"{error}\" duration={duration_ms}ms")
|
||||
}
|
||||
WorkflowRunEvent::StageStarted {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
handler_type,
|
||||
attempt,
|
||||
max_attempts,
|
||||
} => {
|
||||
let mut s = format!("[STAGE_STARTED] name={name} index={index}");
|
||||
let mut s = format!("[STAGE_STARTED] node_id={node_id} name={name} index={index}");
|
||||
if let Some(ht) = handler_type {
|
||||
s.push_str(&format!(" handler_type={ht}"));
|
||||
}
|
||||
|
|
@ -231,6 +232,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
s
|
||||
}
|
||||
WorkflowRunEvent::StageCompleted {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
duration_ms,
|
||||
|
|
@ -245,7 +247,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
max_attempts,
|
||||
failure_class,
|
||||
} => {
|
||||
let mut s = format!("[STAGE_COMPLETED] name={name} index={index} duration={duration_ms}ms status={status}");
|
||||
let mut s = format!("[STAGE_COMPLETED] node_id={node_id} name={name} index={index} duration={duration_ms}ms status={status}");
|
||||
if let Some(label) = preferred_label {
|
||||
s.push_str(&format!(" preferred_label=\"{label}\""));
|
||||
}
|
||||
|
|
@ -280,6 +282,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
s
|
||||
}
|
||||
WorkflowRunEvent::StageFailed {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
error,
|
||||
|
|
@ -288,7 +291,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
failure_class,
|
||||
} => {
|
||||
let mut s = format!(
|
||||
"[STAGE_FAILED] name={name} index={index} error=\"{error}\" will_retry={will_retry}"
|
||||
"[STAGE_FAILED] node_id={node_id} name={name} index={index} error=\"{error}\" will_retry={will_retry}"
|
||||
);
|
||||
if let Some(reason) = failure_reason {
|
||||
s.push_str(&format!(" failure_reason=\"{reason}\""));
|
||||
|
|
@ -299,6 +302,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
s
|
||||
}
|
||||
WorkflowRunEvent::StageRetrying {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
attempt,
|
||||
|
|
@ -306,7 +310,7 @@ pub fn format_event_summary(event: &WorkflowRunEvent, styles: &Styles) -> String
|
|||
delay_ms,
|
||||
} => {
|
||||
format!(
|
||||
"[STAGE_RETRYING] name={name} index={index} attempt={attempt}/{max_attempts} delay={delay_ms}ms"
|
||||
"[STAGE_RETRYING] node_id={node_id} name={name} index={index} attempt={attempt}/{max_attempts} delay={delay_ms}ms"
|
||||
)
|
||||
}
|
||||
WorkflowRunEvent::ParallelStarted {
|
||||
|
|
|
|||
|
|
@ -208,7 +208,7 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu
|
|||
let (event_name, event_fields) = crate::event::flatten_event(event);
|
||||
let mut envelope = serde_json::Map::new();
|
||||
envelope.insert(
|
||||
"timestamp".to_string(),
|
||||
"ts".to_string(),
|
||||
serde_json::Value::String(
|
||||
Utc::now()
|
||||
.to_rfc3339_opts(chrono::SecondsFormat::Millis, true),
|
||||
|
|
@ -223,7 +223,7 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu
|
|||
serde_json::Value::String(event_name),
|
||||
);
|
||||
for (k, v) in event_fields {
|
||||
if k != "timestamp" && k != "run_id" && k != "event" {
|
||||
if k != "ts" && k != "run_id" && k != "event" {
|
||||
envelope.insert(k, v);
|
||||
}
|
||||
}
|
||||
|
|
@ -259,7 +259,7 @@ pub async fn run_command(args: RunArgs, styles: &'static Styles) -> anyhow::Resu
|
|||
duration_ms,
|
||||
status,
|
||||
usage,
|
||||
..
|
||||
.. // node_id and other fields
|
||||
} => {
|
||||
let mut line = format!(
|
||||
"{dim}Stage \"{name}\" completed ({status}) in {duration}",
|
||||
|
|
|
|||
|
|
@ -885,6 +885,7 @@ impl WorkflowRunEngine {
|
|||
if attempt < policy.max_attempts && handler.should_retry(&e) {
|
||||
let delay = policy.backoff.delay_for_attempt(attempt);
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageFailed {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
error: e.to_string(),
|
||||
|
|
@ -893,6 +894,7 @@ impl WorkflowRunEngine {
|
|||
failure_class: Some(e.failure_class().to_string()),
|
||||
});
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageRetrying {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
attempt: usize::try_from(attempt).unwrap_or(usize::MAX),
|
||||
|
|
@ -918,6 +920,7 @@ impl WorkflowRunEngine {
|
|||
if attempt < policy.max_attempts {
|
||||
let delay = policy.backoff.delay_for_attempt(attempt);
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageRetrying {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
attempt: usize::try_from(attempt).unwrap_or(usize::MAX),
|
||||
|
|
@ -1274,6 +1277,7 @@ impl WorkflowRunEngine {
|
|||
let retry_policy = build_retry_policy(node, graph);
|
||||
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageStarted {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
handler_type: node.handler_type().map(String::from),
|
||||
|
|
@ -1362,6 +1366,7 @@ impl WorkflowRunEngine {
|
|||
|
||||
if outcome.status == StageStatus::Fail {
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageFailed {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
error: outcome
|
||||
|
|
@ -1375,6 +1380,7 @@ impl WorkflowRunEngine {
|
|||
});
|
||||
} else {
|
||||
self.services.emitter.emit(&WorkflowRunEvent::StageCompleted {
|
||||
node_id: node.id.clone(),
|
||||
name: node.label().to_string(),
|
||||
index: stage_index,
|
||||
duration_ms: stage_duration_ms,
|
||||
|
|
|
|||
|
|
@ -33,6 +33,7 @@ pub enum WorkflowRunEvent {
|
|||
git_commit_sha: Option<String>,
|
||||
},
|
||||
StageStarted {
|
||||
node_id: String,
|
||||
name: String,
|
||||
index: usize,
|
||||
handler_type: Option<String>,
|
||||
|
|
@ -40,6 +41,7 @@ pub enum WorkflowRunEvent {
|
|||
max_attempts: usize,
|
||||
},
|
||||
StageCompleted {
|
||||
node_id: String,
|
||||
name: String,
|
||||
index: usize,
|
||||
duration_ms: u64,
|
||||
|
|
@ -55,6 +57,7 @@ pub enum WorkflowRunEvent {
|
|||
failure_class: Option<String>,
|
||||
},
|
||||
StageFailed {
|
||||
node_id: String,
|
||||
name: String,
|
||||
index: usize,
|
||||
error: String,
|
||||
|
|
@ -63,6 +66,7 @@ pub enum WorkflowRunEvent {
|
|||
failure_class: Option<String>,
|
||||
},
|
||||
StageRetrying {
|
||||
node_id: String,
|
||||
name: String,
|
||||
index: usize,
|
||||
attempt: usize,
|
||||
|
|
@ -199,6 +203,7 @@ impl WorkflowRunEvent {
|
|||
error!(error, duration_ms, "Workflow run failed");
|
||||
}
|
||||
Self::StageStarted {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
handler_type,
|
||||
|
|
@ -206,6 +211,7 @@ impl WorkflowRunEvent {
|
|||
max_attempts,
|
||||
} => {
|
||||
debug!(
|
||||
node_id,
|
||||
stage = name.as_str(),
|
||||
index,
|
||||
handler_type = handler_type.as_deref().unwrap_or(""),
|
||||
|
|
@ -215,6 +221,7 @@ impl WorkflowRunEvent {
|
|||
);
|
||||
}
|
||||
Self::StageCompleted {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
duration_ms,
|
||||
|
|
@ -224,6 +231,7 @@ impl WorkflowRunEvent {
|
|||
..
|
||||
} => {
|
||||
debug!(
|
||||
node_id,
|
||||
stage = name.as_str(),
|
||||
index,
|
||||
duration_ms,
|
||||
|
|
@ -234,6 +242,7 @@ impl WorkflowRunEvent {
|
|||
);
|
||||
}
|
||||
Self::StageFailed {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
error,
|
||||
|
|
@ -242,6 +251,7 @@ impl WorkflowRunEvent {
|
|||
} => {
|
||||
if *will_retry {
|
||||
warn!(
|
||||
node_id,
|
||||
stage = name.as_str(),
|
||||
index,
|
||||
error,
|
||||
|
|
@ -250,6 +260,7 @@ impl WorkflowRunEvent {
|
|||
);
|
||||
} else {
|
||||
error!(
|
||||
node_id,
|
||||
stage = name.as_str(),
|
||||
index,
|
||||
error,
|
||||
|
|
@ -259,6 +270,7 @@ impl WorkflowRunEvent {
|
|||
}
|
||||
}
|
||||
Self::StageRetrying {
|
||||
node_id,
|
||||
name,
|
||||
index,
|
||||
attempt,
|
||||
|
|
@ -266,6 +278,7 @@ impl WorkflowRunEvent {
|
|||
delay_ms,
|
||||
} => {
|
||||
warn!(
|
||||
node_id,
|
||||
stage = name.as_str(),
|
||||
index,
|
||||
attempt,
|
||||
|
|
@ -450,7 +463,7 @@ pub fn flatten_event(
|
|||
event: &WorkflowRunEvent,
|
||||
) -> (String, serde_json::Map<String, serde_json::Value>) {
|
||||
let value = serde_json::to_value(event).expect("WorkflowRunEvent must serialize");
|
||||
match value {
|
||||
let (event_name, mut fields) = match value {
|
||||
serde_json::Value::Object(map) => {
|
||||
// Externally-tagged enum: { "VariantName": { fields } }
|
||||
let (variant_name, inner) = map.into_iter().next().expect("enum must have one key");
|
||||
|
|
@ -469,7 +482,9 @@ pub fn flatten_event(
|
|||
// Unit variants serialize as strings
|
||||
serde_json::Value::String(name) => (name, serde_json::Map::new()),
|
||||
_ => ("Unknown".to_string(), serde_json::Map::new()),
|
||||
}
|
||||
};
|
||||
rename_fields(&event_name, &mut fields);
|
||||
(event_name, fields)
|
||||
}
|
||||
|
||||
fn flatten_agent(inner: serde_json::Value) -> (String, serde_json::Map<String, serde_json::Value>) {
|
||||
|
|
@ -580,6 +595,74 @@ fn flatten_sub_agent_event(
|
|||
(event_name, fields)
|
||||
}
|
||||
|
||||
/// Rename flattened event fields for clarity in progress.jsonl output.
|
||||
///
|
||||
/// Applied as a post-processing step after `flatten_event` serialization to
|
||||
/// give fields self-describing names without changing the Rust enum.
|
||||
fn rename_fields(event_name: &str, fields: &mut serde_json::Map<String, serde_json::Value>) {
|
||||
/// Move a key from `old` to `new` if present.
|
||||
fn rename(fields: &mut serde_json::Map<String, serde_json::Value>, old: &str, new: &str) {
|
||||
if let Some(v) = fields.remove(old) {
|
||||
fields.insert(new.to_string(), v);
|
||||
}
|
||||
}
|
||||
|
||||
/// Insert `node_label` defaulting to the value of `node_id`, if not already present.
|
||||
fn default_node_label(fields: &mut serde_json::Map<String, serde_json::Value>) {
|
||||
if !fields.contains_key("node_label") {
|
||||
if let Some(id) = fields.get("node_id").cloned() {
|
||||
fields.insert("node_label".to_string(), id);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if event_name.starts_with("Stage") {
|
||||
// name → node_label, index → stage_index, node_id stays
|
||||
rename(fields, "name", "node_label");
|
||||
rename(fields, "index", "stage_index");
|
||||
// node_id already present from Rust enum
|
||||
} else if event_name == "WorkflowRunStarted" {
|
||||
rename(fields, "name", "workflow_name");
|
||||
} else if event_name.starts_with("Agent.") || event_name == "Agent" {
|
||||
rename(fields, "stage", "node_id");
|
||||
default_node_label(fields);
|
||||
} else if event_name.starts_with("Sandbox.Snapshot") {
|
||||
// Must check before generic Sandbox.* to catch Snapshot* first
|
||||
rename(fields, "name", "snapshot_name");
|
||||
rename(fields, "provider", "sandbox_provider");
|
||||
} else if event_name.starts_with("Sandbox.") {
|
||||
rename(fields, "provider", "sandbox_provider");
|
||||
} else if event_name.starts_with("ParallelBranch") {
|
||||
rename(fields, "branch", "node_id");
|
||||
default_node_label(fields);
|
||||
rename(fields, "index", "branch_index");
|
||||
} else if event_name.starts_with("SetupCommand") || event_name == "SetupFailed" {
|
||||
rename(fields, "index", "command_index");
|
||||
} else if event_name == "EdgeSelected" || event_name == "LoopRestart" {
|
||||
rename(fields, "from_node", "from_node_id");
|
||||
rename(fields, "to_node", "to_node_id");
|
||||
} else if event_name == "StallWatchdogTimeout" {
|
||||
rename(fields, "node", "node_id");
|
||||
default_node_label(fields);
|
||||
} else if event_name == "Prompt" {
|
||||
rename(fields, "stage", "node_id");
|
||||
default_node_label(fields);
|
||||
rename(fields, "text", "prompt_text");
|
||||
} else if event_name.starts_with("Interview") && event_name != "InterviewCompleted" {
|
||||
// InterviewStarted, InterviewTimeout have `stage`
|
||||
rename(fields, "stage", "node_id");
|
||||
default_node_label(fields);
|
||||
} else if event_name == "SubgraphStarted" {
|
||||
default_node_label(fields);
|
||||
rename(fields, "start_node", "start_node_id");
|
||||
} else if event_name == "SubgraphCompleted"
|
||||
|| event_name == "CheckpointSaved"
|
||||
|| event_name == "GitCheckpoint"
|
||||
{
|
||||
default_node_label(fields);
|
||||
}
|
||||
}
|
||||
|
||||
/// Current time as epoch milliseconds.
|
||||
fn epoch_millis() -> i64 {
|
||||
std::time::SystemTime::now()
|
||||
|
|
@ -685,6 +768,7 @@ mod tests {
|
|||
#[test]
|
||||
fn workflow_run_event_serialization() {
|
||||
let event = WorkflowRunEvent::StageStarted {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
handler_type: Some("codergen".to_string()),
|
||||
|
|
@ -700,6 +784,7 @@ mod tests {
|
|||
|
||||
// None handler_type serializes as null
|
||||
let event_none = WorkflowRunEvent::StageStarted {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
handler_type: None,
|
||||
|
|
@ -800,6 +885,7 @@ mod tests {
|
|||
#[test]
|
||||
fn stage_completed_event_serialization_with_new_fields() {
|
||||
let event = WorkflowRunEvent::StageCompleted {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
duration_ms: 1500,
|
||||
|
|
@ -823,6 +909,7 @@ mod tests {
|
|||
assert!(json.contains("\"failure_class\":null"));
|
||||
|
||||
let event_none = WorkflowRunEvent::StageCompleted {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
duration_ms: 1500,
|
||||
|
|
@ -845,6 +932,7 @@ mod tests {
|
|||
#[test]
|
||||
fn stage_failed_event_serialization() {
|
||||
let event = WorkflowRunEvent::StageFailed {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
error: "timeout".to_string(),
|
||||
|
|
@ -862,6 +950,7 @@ mod tests {
|
|||
);
|
||||
|
||||
let event_none = WorkflowRunEvent::StageFailed {
|
||||
node_id: "plan".to_string(),
|
||||
name: "plan".to_string(),
|
||||
index: 0,
|
||||
error: "timeout".to_string(),
|
||||
|
|
@ -1006,6 +1095,7 @@ mod tests {
|
|||
#[test]
|
||||
fn stage_retrying_event_serialization() {
|
||||
let event = WorkflowRunEvent::StageRetrying {
|
||||
node_id: "lint".to_string(),
|
||||
name: "lint".to_string(),
|
||||
index: 2,
|
||||
attempt: 3,
|
||||
|
|
@ -1177,7 +1267,8 @@ mod tests {
|
|||
#[test]
|
||||
fn flatten_event_simple_variant() {
|
||||
let event = WorkflowRunEvent::StageStarted {
|
||||
name: "plan".to_string(),
|
||||
node_id: "plan".to_string(),
|
||||
name: "Plan Stage".to_string(),
|
||||
index: 0,
|
||||
handler_type: Some("codergen".to_string()),
|
||||
attempt: 1,
|
||||
|
|
@ -1185,11 +1276,15 @@ mod tests {
|
|||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "StageStarted");
|
||||
assert_eq!(fields["name"], "plan");
|
||||
assert_eq!(fields["index"], 0);
|
||||
assert_eq!(fields["node_id"], "plan");
|
||||
assert_eq!(fields["node_label"], "Plan Stage");
|
||||
assert_eq!(fields["stage_index"], 0);
|
||||
assert_eq!(fields["handler_type"], "codergen");
|
||||
assert_eq!(fields["attempt"], 1);
|
||||
assert_eq!(fields["max_attempts"], 3);
|
||||
// Old keys should not be present
|
||||
assert!(!fields.contains_key("name"));
|
||||
assert!(!fields.contains_key("index"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -1204,9 +1299,11 @@ mod tests {
|
|||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Agent.ToolCallStarted");
|
||||
assert_eq!(fields["stage"], "code");
|
||||
assert_eq!(fields["node_id"], "code");
|
||||
assert_eq!(fields["node_label"], "code");
|
||||
assert_eq!(fields["tool_name"], "read_file");
|
||||
assert_eq!(fields["tool_call_id"], "call_1");
|
||||
assert!(!fields.contains_key("stage"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -1218,7 +1315,8 @@ mod tests {
|
|||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Sandbox.Initializing");
|
||||
assert_eq!(fields["provider"], "docker");
|
||||
assert_eq!(fields["sandbox_provider"], "docker");
|
||||
assert!(!fields.contains_key("provider"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
@ -1237,9 +1335,11 @@ mod tests {
|
|||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Agent.SubAgentEvent.ToolCallStarted");
|
||||
assert_eq!(fields["stage"], "code");
|
||||
assert_eq!(fields["node_id"], "code");
|
||||
assert_eq!(fields["node_label"], "code");
|
||||
assert_eq!(fields["agent_id"], "sub_1");
|
||||
assert_eq!(fields["depth"], 1);
|
||||
assert!(!fields.contains_key("stage"));
|
||||
// Inner event preserved as nested_event JSON (not flattened)
|
||||
let nested = fields["nested_event"].as_object().unwrap();
|
||||
let tool_call = nested["ToolCallStarted"].as_object().unwrap();
|
||||
|
|
@ -1269,7 +1369,9 @@ mod tests {
|
|||
// Outer SubAgentEvent fields at top level
|
||||
assert_eq!(fields["agent_id"], "sub_1");
|
||||
assert_eq!(fields["depth"], 1);
|
||||
assert_eq!(fields["stage"], "code");
|
||||
assert_eq!(fields["node_id"], "code");
|
||||
assert_eq!(fields["node_label"], "code");
|
||||
assert!(!fields.contains_key("stage"));
|
||||
// Inner SubAgentEvent preserved in nested_event with all data intact
|
||||
let nested = fields["nested_event"].as_object().unwrap();
|
||||
let inner_sub = nested["SubAgentEvent"].as_object().unwrap();
|
||||
|
|
@ -1288,7 +1390,188 @@ mod tests {
|
|||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Agent.SessionStarted");
|
||||
assert_eq!(fields["stage"], "plan");
|
||||
assert_eq!(fields["node_id"], "plan");
|
||||
assert_eq!(fields["node_label"], "plan");
|
||||
assert!(!fields.contains_key("stage"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_workflow_run_started() {
|
||||
let event = WorkflowRunEvent::WorkflowRunStarted {
|
||||
name: "my_pipeline".to_string(),
|
||||
run_id: "r1".to_string(),
|
||||
base_sha: None,
|
||||
run_branch: None,
|
||||
worktree_dir: None,
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "WorkflowRunStarted");
|
||||
assert_eq!(fields["workflow_name"], "my_pipeline");
|
||||
assert!(!fields.contains_key("name"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_parallel_branch_started() {
|
||||
let event = WorkflowRunEvent::ParallelBranchStarted {
|
||||
branch: "lint".to_string(),
|
||||
index: 0,
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "ParallelBranchStarted");
|
||||
assert_eq!(fields["node_id"], "lint");
|
||||
assert_eq!(fields["node_label"], "lint");
|
||||
assert_eq!(fields["branch_index"], 0);
|
||||
assert!(!fields.contains_key("branch"));
|
||||
assert!(!fields.contains_key("index"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_parallel_branch_completed() {
|
||||
let event = WorkflowRunEvent::ParallelBranchCompleted {
|
||||
branch: "lint".to_string(),
|
||||
index: 0,
|
||||
duration_ms: 1000,
|
||||
status: "success".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "ParallelBranchCompleted");
|
||||
assert_eq!(fields["node_id"], "lint");
|
||||
assert_eq!(fields["node_label"], "lint");
|
||||
assert_eq!(fields["branch_index"], 0);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_setup_command_started() {
|
||||
let event = WorkflowRunEvent::SetupCommandStarted {
|
||||
command: "npm install".to_string(),
|
||||
index: 2,
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "SetupCommandStarted");
|
||||
assert_eq!(fields["command_index"], 2);
|
||||
assert!(!fields.contains_key("index"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_setup_failed() {
|
||||
let event = WorkflowRunEvent::SetupFailed {
|
||||
command: "npm test".to_string(),
|
||||
index: 1,
|
||||
exit_code: 1,
|
||||
stderr: "fail".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "SetupFailed");
|
||||
assert_eq!(fields["command_index"], 1);
|
||||
assert!(!fields.contains_key("index"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_edge_selected() {
|
||||
let event = WorkflowRunEvent::EdgeSelected {
|
||||
from_node: "plan".to_string(),
|
||||
to_node: "code".to_string(),
|
||||
label: Some("success".to_string()),
|
||||
condition: None,
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "EdgeSelected");
|
||||
assert_eq!(fields["from_node_id"], "plan");
|
||||
assert_eq!(fields["to_node_id"], "code");
|
||||
assert!(!fields.contains_key("from_node"));
|
||||
assert!(!fields.contains_key("to_node"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_loop_restart() {
|
||||
let event = WorkflowRunEvent::LoopRestart {
|
||||
from_node: "review".to_string(),
|
||||
to_node: "code".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "LoopRestart");
|
||||
assert_eq!(fields["from_node_id"], "review");
|
||||
assert_eq!(fields["to_node_id"], "code");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_stall_watchdog_timeout() {
|
||||
let event = WorkflowRunEvent::StallWatchdogTimeout {
|
||||
node: "work".to_string(),
|
||||
idle_seconds: 600,
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "StallWatchdogTimeout");
|
||||
assert_eq!(fields["node_id"], "work");
|
||||
assert_eq!(fields["node_label"], "work");
|
||||
assert!(!fields.contains_key("node"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_prompt() {
|
||||
let event = WorkflowRunEvent::Prompt {
|
||||
stage: "gate".to_string(),
|
||||
text: "Approve?".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Prompt");
|
||||
assert_eq!(fields["node_id"], "gate");
|
||||
assert_eq!(fields["node_label"], "gate");
|
||||
assert_eq!(fields["prompt_text"], "Approve?");
|
||||
assert!(!fields.contains_key("stage"));
|
||||
assert!(!fields.contains_key("text"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_interview_started() {
|
||||
let event = WorkflowRunEvent::InterviewStarted {
|
||||
question: "OK?".to_string(),
|
||||
stage: "gate".to_string(),
|
||||
question_type: "yes_no".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "InterviewStarted");
|
||||
assert_eq!(fields["node_id"], "gate");
|
||||
assert_eq!(fields["node_label"], "gate");
|
||||
assert!(!fields.contains_key("stage"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_subgraph_started() {
|
||||
let event = WorkflowRunEvent::SubgraphStarted {
|
||||
node_id: "sub_1".to_string(),
|
||||
start_node: "start".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "SubgraphStarted");
|
||||
assert_eq!(fields["node_id"], "sub_1");
|
||||
assert_eq!(fields["node_label"], "sub_1");
|
||||
assert_eq!(fields["start_node_id"], "start");
|
||||
assert!(!fields.contains_key("start_node"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_checkpoint_saved() {
|
||||
let event = WorkflowRunEvent::CheckpointSaved {
|
||||
node_id: "plan".to_string(),
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "CheckpointSaved");
|
||||
assert_eq!(fields["node_id"], "plan");
|
||||
assert_eq!(fields["node_label"], "plan");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rename_fields_sandbox_snapshot_pulling() {
|
||||
let event = WorkflowRunEvent::Sandbox {
|
||||
event: SandboxEvent::SnapshotPulling {
|
||||
name: "base-image".into(),
|
||||
},
|
||||
};
|
||||
let (name, fields) = flatten_event(&event);
|
||||
assert_eq!(name, "Sandbox.SnapshotPulling");
|
||||
assert_eq!(fields["snapshot_name"], "base-image");
|
||||
assert!(!fields.contains_key("name"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
|
|
|||
|
|
@ -185,7 +185,7 @@ pub fn extract_stage_durations(logs_root: &Path) -> HashMap<String, u64> {
|
|||
if envelope.get("event").and_then(|v| v.as_str()) != Some("StageCompleted") {
|
||||
continue;
|
||||
}
|
||||
let Some(name) = envelope.get("name").and_then(|v| v.as_str()) else {
|
||||
let Some(name) = envelope.get("node_label").and_then(|v| v.as_str()) else {
|
||||
continue;
|
||||
};
|
||||
let Some(duration_ms) = envelope.get("duration_ms").and_then(|v| v.as_u64()) else {
|
||||
|
|
@ -519,11 +519,12 @@ mod tests {
|
|||
let jsonl = dir.path().join("progress.jsonl");
|
||||
|
||||
let event1 = serde_json::json!({
|
||||
"timestamp": "2025-01-01T00:00:00.000Z",
|
||||
"ts": "2025-01-01T00:00:00.000Z",
|
||||
"run_id": "r1",
|
||||
"event": "StageCompleted",
|
||||
"name": "plan",
|
||||
"index": 0,
|
||||
"node_id": "plan",
|
||||
"node_label": "plan",
|
||||
"stage_index": 0,
|
||||
"duration_ms": 5000,
|
||||
"status": "success",
|
||||
"preferred_label": null,
|
||||
|
|
@ -537,11 +538,12 @@ mod tests {
|
|||
"failure_class": null
|
||||
});
|
||||
let event2 = serde_json::json!({
|
||||
"timestamp": "2025-01-01T00:00:05.000Z",
|
||||
"ts": "2025-01-01T00:00:05.000Z",
|
||||
"run_id": "r1",
|
||||
"event": "StageCompleted",
|
||||
"name": "code",
|
||||
"index": 1,
|
||||
"node_id": "code",
|
||||
"node_label": "code",
|
||||
"stage_index": 1,
|
||||
"duration_ms": 15000,
|
||||
"status": "success",
|
||||
"preferred_label": null,
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue