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) <noreply@anthropic.com>
Entire-Checkpoint: 2d4a98c3b5f0
This commit is contained in:
Bryan Helmkamp 2026-02-25 11:40:19 -05:00
parent 3d39e06c29
commit a00f7fd60e
6 changed files with 245 additions and 29 deletions

View file

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

View file

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

View file

@ -20,6 +20,7 @@ pub enum PipelineEvent {
StageStarted {
name: String,
index: usize,
handler_type: Option<String>,
},
StageCompleted {
name: String,
@ -29,12 +30,15 @@ pub enum PipelineEvent {
preferred_label: Option<String>,
suggested_next_ids: Vec<String>,
usage: Option<StageUsage>,
failure_reason: Option<String>,
notes: Option<String>,
},
StageFailed {
name: String,
index: usize,
error: String,
will_retry: bool,
failure_reason: Option<String>,
},
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 {

View file

@ -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, AttractorError>(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));

View file

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

View file

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