From a00f7fd60e4526e36b5bb184801fdfe848b8b558 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Wed, 25 Feb 2026 11:40:19 -0500 Subject: [PATCH] Add missing data fields to PipelineEvent variants Thread existing data through to pipeline events to match kilroy's progress.ndjson schema: handler_type on StageStarted, failure_reason on StageFailed/StageCompleted, notes on StageCompleted, join_policy and error_policy on ParallelStarted, status (replacing success bool) on ParallelBranchCompleted, and question_type on InterviewStarted. Adds Display impls for QuestionType, JoinPolicy, and ErrorPolicy. Co-Authored-By: Claude Opus 4.6 (1M context) Entire-Checkpoint: 2d4a98c3b5f0 --- crates/attractor/src/cli/mod.rs | 74 +++++++++---- crates/attractor/src/engine.rs | 5 + crates/attractor/src/event.rs | 121 ++++++++++++++++++++- crates/attractor/src/handler/parallel.rs | 52 +++++++-- crates/attractor/src/handler/wait_human.rs | 1 + crates/attractor/src/interviewer/mod.rs | 21 ++++ 6 files changed, 245 insertions(+), 29 deletions(-) diff --git a/crates/attractor/src/cli/mod.rs b/crates/attractor/src/cli/mod.rs index baff2e601..2179c1c8d 100644 --- a/crates/attractor/src/cli/mod.rs +++ b/crates/attractor/src/cli/mod.rs @@ -178,8 +178,12 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { PipelineEvent::PipelineFailed { error, duration_ms } => { format!("[PIPELINE_FAILED] error=\"{error}\" duration={duration_ms}ms") } - PipelineEvent::StageStarted { name, index } => { - format!("[STAGE_STARTED] name={name} index={index}") + PipelineEvent::StageStarted { name, index, handler_type } => { + let mut s = format!("[STAGE_STARTED] name={name} index={index}"); + if let Some(ht) = handler_type { + s.push_str(&format!(" handler_type={ht}")); + } + s } PipelineEvent::StageCompleted { name, @@ -189,6 +193,8 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { preferred_label, suggested_next_ids, usage, + failure_reason, + notes, } => { let mut s = format!("[STAGE_COMPLETED] name={name} index={index} duration={duration_ms}ms status={status}"); if let Some(label) = preferred_label { @@ -206,6 +212,12 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { s.push_str(&format!(" tokens={tokens_str}")); } } + if let Some(reason) = failure_reason { + s.push_str(&format!(" failure_reason=\"{reason}\"")); + } + if let Some(n) = notes { + s.push_str(&format!(" notes=\"{n}\"")); + } s } PipelineEvent::StageFailed { @@ -213,10 +225,15 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { index, error, will_retry, + failure_reason, } => { - format!( + let mut s = format!( "[STAGE_FAILED] name={name} index={index} error=\"{error}\" will_retry={will_retry}" - ) + ); + if let Some(reason) = failure_reason { + s.push_str(&format!(" failure_reason=\"{reason}\"")); + } + s } PipelineEvent::StageRetrying { name, @@ -228,8 +245,8 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { "[STAGE_RETRYING] name={name} index={index} attempt={attempt} delay={delay_ms}ms" ) } - PipelineEvent::ParallelStarted { branch_count } => { - format!("[PARALLEL_STARTED] branches={branch_count}") + PipelineEvent::ParallelStarted { branch_count, join_policy, error_policy } => { + format!("[PARALLEL_STARTED] branches={branch_count} join_policy={join_policy} error_policy={error_policy}") } PipelineEvent::ParallelBranchStarted { branch, index } => { format!("[PARALLEL_BRANCH_STARTED] branch={branch} index={index}") @@ -238,9 +255,9 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { branch, index, duration_ms, - success, + status, } => { - format!("[PARALLEL_BRANCH_COMPLETED] branch={branch} index={index} duration={duration_ms}ms success={success}") + format!("[PARALLEL_BRANCH_COMPLETED] branch={branch} index={index} duration={duration_ms}ms status={status}") } PipelineEvent::ParallelCompleted { duration_ms, @@ -249,8 +266,8 @@ pub fn format_event_summary(event: &PipelineEvent, styles: &Styles) -> String { } => { format!("[PARALLEL_COMPLETED] duration={duration_ms}ms succeeded={success_count} failed={failure_count}") } - PipelineEvent::InterviewStarted { question, stage } => { - format!("[INTERVIEW_STARTED] stage={stage} question=\"{question}\"") + PipelineEvent::InterviewStarted { question, stage, question_type } => { + format!("[INTERVIEW_STARTED] stage={stage} question=\"{question}\" question_type={question_type}") } PipelineEvent::InterviewCompleted { question, @@ -357,10 +374,14 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { PipelineEvent::PipelineFailed { error, duration_ms } => { format!("{d}── PIPELINE_FAILED ──────────────────────────{r}\n {d}error:{r} {error}\n {d}duration_ms:{r} {duration_ms}\n") } - PipelineEvent::StageStarted { name, index } => { - format!( + PipelineEvent::StageStarted { name, index, handler_type } => { + let mut s = format!( "{d}── STAGE_STARTED ────────────────────────────{r}\n {d}name:{r} {name}\n {d}index:{r} {index}\n" - ) + ); + if let Some(ht) = handler_type { + s.push_str(&format!(" {d}handler_type:{r} {ht}\n")); + } + s } PipelineEvent::StageCompleted { name, @@ -370,6 +391,8 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { preferred_label, suggested_next_ids, usage, + failure_reason, + notes, } => { let mut s = format!("{d}── STAGE_COMPLETED ──────────────────────────{r}\n {d}name:{r} {name}\n {d}index:{r} {index}\n {d}duration_ms:{r} {duration_ms}\n {d}status:{r} {status}\n"); if let Some(label) = preferred_label { @@ -390,6 +413,12 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { s.push_str(&format!(" {d}cost:{r} {}\n", format_cost(cost))); } } + if let Some(reason) = failure_reason { + s.push_str(&format!(" {d}failure_reason:{r} {reason}\n")); + } + if let Some(n) = notes { + s.push_str(&format!(" {d}notes:{r} {n}\n")); + } s } PipelineEvent::StageFailed { @@ -397,8 +426,13 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { index, error, will_retry, + failure_reason, } => { - format!("{d}── STAGE_FAILED ─────────────────────────────{r}\n {d}name:{r} {name}\n {d}index:{r} {index}\n {d}error:{r} {error}\n {d}will_retry:{r} {will_retry}\n") + let mut s = format!("{d}── STAGE_FAILED ─────────────────────────────{r}\n {d}name:{r} {name}\n {d}index:{r} {index}\n {d}error:{r} {error}\n {d}will_retry:{r} {will_retry}\n"); + if let Some(reason) = failure_reason { + s.push_str(&format!(" {d}failure_reason:{r} {reason}\n")); + } + s } PipelineEvent::StageRetrying { name, @@ -408,8 +442,8 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { } => { format!("{d}── STAGE_RETRYING ───────────────────────────{r}\n {d}name:{r} {name}\n {d}index:{r} {index}\n {d}attempt:{r} {attempt}\n {d}delay_ms:{r} {delay_ms}\n") } - PipelineEvent::ParallelStarted { branch_count } => { - format!("{d}── PARALLEL_STARTED ─────────────────────────{r}\n {d}branch_count:{r} {branch_count}\n") + PipelineEvent::ParallelStarted { branch_count, join_policy, error_policy } => { + format!("{d}── PARALLEL_STARTED ─────────────────────────{r}\n {d}branch_count:{r} {branch_count}\n {d}join_policy:{r} {join_policy}\n {d}error_policy:{r} {error_policy}\n") } PipelineEvent::ParallelBranchStarted { branch, index } => { format!("{d}── PARALLEL_BRANCH_STARTED ──────────────────{r}\n {d}branch:{r} {branch}\n {d}index:{r} {index}\n") @@ -418,9 +452,9 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { branch, index, duration_ms, - success, + status, } => { - format!("{d}── PARALLEL_BRANCH_COMPLETED ────────────────{r}\n {d}branch:{r} {branch}\n {d}index:{r} {index}\n {d}duration_ms:{r} {duration_ms}\n {d}success:{r} {success}\n") + format!("{d}── PARALLEL_BRANCH_COMPLETED ────────────────{r}\n {d}branch:{r} {branch}\n {d}index:{r} {index}\n {d}duration_ms:{r} {duration_ms}\n {d}status:{r} {status}\n") } PipelineEvent::ParallelCompleted { duration_ms, @@ -429,8 +463,8 @@ pub fn format_event_detail(event: &PipelineEvent, styles: &Styles) -> String { } => { format!("{d}── PARALLEL_COMPLETED ───────────────────────{r}\n {d}duration_ms:{r} {duration_ms}\n {d}success_count:{r} {success_count}\n {d}failure_count:{r} {failure_count}\n") } - PipelineEvent::InterviewStarted { question, stage } => { - format!("{d}── INTERVIEW_STARTED ────────────────────────{r}\n {d}stage:{r} {stage}\n {d}question:{r} {question}\n") + PipelineEvent::InterviewStarted { question, stage, question_type } => { + format!("{d}── INTERVIEW_STARTED ────────────────────────{r}\n {d}stage:{r} {stage}\n {d}question:{r} {question}\n {d}question_type:{r} {question_type}\n") } PipelineEvent::InterviewCompleted { question, diff --git a/crates/attractor/src/engine.rs b/crates/attractor/src/engine.rs index 6de16222e..0732fd64c 100644 --- a/crates/attractor/src/engine.rs +++ b/crates/attractor/src/engine.rs @@ -566,6 +566,7 @@ impl PipelineEngine { index: stage_index, error: e.to_string(), will_retry: true, + failure_reason: None, }); self.services.emitter.emit(&PipelineEvent::StageRetrying { name: node.label().to_string(), @@ -809,6 +810,7 @@ impl PipelineEngine { self.services.emitter.emit(&PipelineEvent::StageStarted { name: node.label().to_string(), index: stage_index, + handler_type: node.handler_type().map(String::from), }); if node.handler_type() != Some("wait.human") { self.inform( @@ -850,6 +852,7 @@ impl PipelineEngine { .unwrap_or("unknown") .to_string(), will_retry: false, + failure_reason: outcome.failure_reason.clone(), }); } else { self.services.emitter.emit(&PipelineEvent::StageCompleted { @@ -860,6 +863,8 @@ impl PipelineEngine { preferred_label: outcome.preferred_label.clone(), suggested_next_ids: outcome.suggested_next_ids.clone(), usage: outcome.usage.clone(), + failure_reason: outcome.failure_reason.clone(), + notes: outcome.notes.clone(), }); self.inform( &format!("Stage completed: {}", node.label()), diff --git a/crates/attractor/src/event.rs b/crates/attractor/src/event.rs index ba231b64f..ec8ffdab9 100644 --- a/crates/attractor/src/event.rs +++ b/crates/attractor/src/event.rs @@ -20,6 +20,7 @@ pub enum PipelineEvent { StageStarted { name: String, index: usize, + handler_type: Option, }, StageCompleted { name: String, @@ -29,12 +30,15 @@ pub enum PipelineEvent { preferred_label: Option, suggested_next_ids: Vec, usage: Option, + failure_reason: Option, + notes: Option, }, StageFailed { name: String, index: usize, error: String, will_retry: bool, + failure_reason: Option, }, StageRetrying { name: String, @@ -44,6 +48,8 @@ pub enum PipelineEvent { }, ParallelStarted { branch_count: usize, + join_policy: String, + error_policy: String, }, ParallelBranchStarted { branch: String, @@ -53,7 +59,7 @@ pub enum PipelineEvent { branch: String, index: usize, duration_ms: u64, - success: bool, + status: String, }, ParallelCompleted { duration_ms: u64, @@ -63,6 +69,7 @@ pub enum PipelineEvent { InterviewStarted { question: String, stage: String, + question_type: String, }, InterviewCompleted { question: String, @@ -209,10 +216,21 @@ mod tests { let event = PipelineEvent::StageStarted { name: "plan".to_string(), index: 0, + handler_type: Some("codergen".to_string()), }; let json = serde_json::to_string(&event).unwrap(); assert!(json.contains("StageStarted")); assert!(json.contains("plan")); + assert!(json.contains("\"handler_type\":\"codergen\"")); + + // None handler_type serializes as null + let event_none = PipelineEvent::StageStarted { + name: "plan".to_string(), + index: 0, + handler_type: None, + }; + let json_none = serde_json::to_string(&event_none).unwrap(); + assert!(json_none.contains("\"handler_type\":null")); } #[test] @@ -254,6 +272,107 @@ mod tests { assert!(json.contains("claude-opus-4-6")); } + #[test] + fn stage_completed_event_serialization_with_new_fields() { + let event = PipelineEvent::StageCompleted { + name: "plan".to_string(), + index: 0, + duration_ms: 1500, + status: "partial_success".to_string(), + preferred_label: None, + suggested_next_ids: vec![], + usage: None, + failure_reason: Some("lint errors remain".to_string()), + notes: Some("fixed 3 of 5 issues".to_string()), + }; + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"failure_reason\":\"lint errors remain\"")); + assert!(json.contains("\"notes\":\"fixed 3 of 5 issues\"")); + + let event_none = PipelineEvent::StageCompleted { + name: "plan".to_string(), + index: 0, + duration_ms: 1500, + status: "success".to_string(), + preferred_label: None, + suggested_next_ids: vec![], + usage: None, + failure_reason: None, + notes: None, + }; + let json_none = serde_json::to_string(&event_none).unwrap(); + assert!(json_none.contains("\"failure_reason\":null")); + assert!(json_none.contains("\"notes\":null")); + } + + #[test] + fn stage_failed_event_serialization() { + let event = PipelineEvent::StageFailed { + name: "plan".to_string(), + index: 0, + error: "timeout".to_string(), + will_retry: true, + failure_reason: Some("LLM request timed out".to_string()), + }; + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"failure_reason\":\"LLM request timed out\"")); + + let event_none = PipelineEvent::StageFailed { + name: "plan".to_string(), + index: 0, + error: "timeout".to_string(), + will_retry: false, + failure_reason: None, + }; + let json_none = serde_json::to_string(&event_none).unwrap(); + assert!(json_none.contains("\"failure_reason\":null")); + } + + #[test] + fn parallel_branch_completed_event_serialization() { + let event = PipelineEvent::ParallelBranchCompleted { + branch: "branch_a".to_string(), + index: 0, + duration_ms: 1500, + status: "success".to_string(), + }; + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"status\":\"success\"")); + assert!(!json.contains("\"success\":")); + + let deserialized: PipelineEvent = serde_json::from_str(&json).unwrap(); + assert!(matches!(deserialized, PipelineEvent::ParallelBranchCompleted { status, .. } if status == "success")); + } + + #[test] + fn parallel_started_event_serialization() { + let event = PipelineEvent::ParallelStarted { + branch_count: 3, + join_policy: "wait_all".to_string(), + error_policy: "continue".to_string(), + }; + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"join_policy\":\"wait_all\"")); + assert!(json.contains("\"error_policy\":\"continue\"")); + + let deserialized: PipelineEvent = serde_json::from_str(&json).unwrap(); + assert!(matches!(deserialized, PipelineEvent::ParallelStarted { join_policy, error_policy, .. } if join_policy == "wait_all" && error_policy == "continue")); + } + + #[test] + fn interview_started_event_serialization() { + let event = PipelineEvent::InterviewStarted { + question: "Review changes?".to_string(), + stage: "gate".to_string(), + question_type: "multiple_choice".to_string(), + }; + let json = serde_json::to_string(&event).unwrap(); + assert!(json.contains("\"question_type\":\"multiple_choice\"")); + + let deserialized: PipelineEvent = serde_json::from_str(&json).unwrap(); + assert!(matches!(deserialized, PipelineEvent::InterviewStarted { question_type, .. } if question_type == "multiple_choice")); + } + #[test] fn compaction_pipeline_event_serialization() { let started = PipelineEvent::CompactionStarted { diff --git a/crates/attractor/src/handler/parallel.rs b/crates/attractor/src/handler/parallel.rs index 817064099..ec6fcea01 100644 --- a/crates/attractor/src/handler/parallel.rs +++ b/crates/attractor/src/handler/parallel.rs @@ -31,6 +31,17 @@ enum JoinPolicy { Quorum(f64), } +impl std::fmt::Display for JoinPolicy { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::WaitAll => write!(f, "wait_all"), + Self::FirstSuccess => write!(f, "first_success"), + Self::KOfN(k) => write!(f, "k_of_n({k})"), + Self::Quorum(frac) => write!(f, "quorum({frac})"), + } + } +} + fn parse_join_policy(raw: &str) -> JoinPolicy { if raw == "first_success" { return JoinPolicy::FirstSuccess; @@ -56,6 +67,16 @@ enum ErrorPolicy { Ignore, } +impl std::fmt::Display for ErrorPolicy { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::Continue => write!(f, "continue"), + Self::FailFast => write!(f, "fail_fast"), + Self::Ignore => write!(f, "ignore"), + } + } +} + fn parse_error_policy(raw: &str) -> ErrorPolicy { match raw { "fail_fast" => ErrorPolicy::FailFast, @@ -85,10 +106,6 @@ impl Handler for ParallelHandler { return Ok(Outcome::fail("No branches for parallel node")); } - services.emitter.emit(&PipelineEvent::ParallelStarted { - branch_count: branches.len(), - }); - let join_policy = parse_join_policy( node.attrs .get("join_policy") @@ -101,6 +118,12 @@ impl Handler for ParallelHandler { .and_then(|v| v.as_str()) .unwrap_or("continue"), ); + + services.emitter.emit(&PipelineEvent::ParallelStarted { + branch_count: branches.len(), + join_policy: join_policy.to_string(), + error_policy: error_policy.to_string(), + }); let max_parallel = node .attrs .get("max_parallel") @@ -138,7 +161,7 @@ impl Handler for ParallelHandler { branch: target_id.clone(), index: branch_index, duration_ms: millis_u64(branch_start.elapsed()), - success: false, + status: "fail".to_string(), }); return Ok(BranchResult { id: target_id.clone(), @@ -155,13 +178,11 @@ impl Handler for ParallelHandler { .execute(target_node, &branch_context, &graph, &logs_root, &branch_services) .await?; - let success = outcome.status == StageStatus::Success - || outcome.status == StageStatus::PartialSuccess; emitter.emit(&PipelineEvent::ParallelBranchCompleted { branch: target_id.clone(), index: branch_index, duration_ms: millis_u64(branch_start.elapsed()), - success, + status: outcome.status.to_string(), }); Ok::(BranchResult { @@ -427,6 +448,21 @@ mod tests { assert_eq!(outcome.status, StageStatus::Success); } + #[test] + fn join_policy_display() { + assert_eq!(JoinPolicy::WaitAll.to_string(), "wait_all"); + assert_eq!(JoinPolicy::FirstSuccess.to_string(), "first_success"); + assert_eq!(JoinPolicy::KOfN(3).to_string(), "k_of_n(3)"); + assert_eq!(JoinPolicy::Quorum(0.5).to_string(), "quorum(0.5)"); + } + + #[test] + fn error_policy_display() { + assert_eq!(ErrorPolicy::Continue.to_string(), "continue"); + assert_eq!(ErrorPolicy::FailFast.to_string(), "fail_fast"); + assert_eq!(ErrorPolicy::Ignore.to_string(), "ignore"); + } + #[test] fn parse_join_policy_variants() { assert!(matches!(parse_join_policy("wait_all"), JoinPolicy::WaitAll)); diff --git a/crates/attractor/src/handler/wait_human.rs b/crates/attractor/src/handler/wait_human.rs index c9e9048ea..99fa91ef9 100644 --- a/crates/attractor/src/handler/wait_human.rs +++ b/crates/attractor/src/handler/wait_human.rs @@ -155,6 +155,7 @@ impl Handler for WaitHumanHandler { self.emit(&PipelineEvent::InterviewStarted { question: question_text.clone(), stage: node.id.clone(), + question_type: question.question_type.to_string(), }); let interview_start = Instant::now(); let answer = self.interviewer.ask(question).await; diff --git a/crates/attractor/src/interviewer/mod.rs b/crates/attractor/src/interviewer/mod.rs index 0c3f84d40..8cd3d2df9 100644 --- a/crates/attractor/src/interviewer/mod.rs +++ b/crates/attractor/src/interviewer/mod.rs @@ -21,6 +21,18 @@ pub enum QuestionType { Confirmation, } +impl std::fmt::Display for QuestionType { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + match self { + Self::YesNo => write!(f, "yes_no"), + Self::MultipleChoice => write!(f, "multiple_choice"), + Self::MultiSelect => write!(f, "multi_select"), + Self::Freeform => write!(f, "freeform"), + Self::Confirmation => write!(f, "confirmation"), + } + } +} + /// An option presented to the user for multiple-choice questions. #[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)] pub struct QuestionOption { @@ -173,6 +185,15 @@ pub trait Interviewer: Send + Sync { mod tests { use super::*; + #[test] + fn question_type_display() { + assert_eq!(QuestionType::YesNo.to_string(), "yes_no"); + assert_eq!(QuestionType::MultipleChoice.to_string(), "multiple_choice"); + assert_eq!(QuestionType::MultiSelect.to_string(), "multi_select"); + assert_eq!(QuestionType::Freeform.to_string(), "freeform"); + assert_eq!(QuestionType::Confirmation.to_string(), "confirmation"); + } + #[test] fn question_new() { let q = Question::new("Do you approve?", QuestionType::YesNo);