feat(workflow): carry visit + parallel group/branch ids on stage events

Adds visit: u32 to Event::StageStarted/Completed/Failed/Retrying so
stored_event_fields() can derive stage_id = "{node_id}@{visit}".
Adds parallel_group_id/parallel_branch_id to ParallelBranchStarted/
Completed Events, computed once in handler/parallel.rs from the
parent parallel node id + visit_from_context + branch index.
Emission sites in lifecycle/event.rs populate visit from
state.node_visits via a new stage_visit helper.

Stored_event_fields() still leaves stage_id and parallel ids None
pending the extraction pass in the next commit.

Co-Authored-By: Claude Opus 4.6 (1M context) <noreply@anthropic.com>
This commit is contained in:
Bryan Helmkamp 2026-04-09 10:14:45 -04:00
parent 07f28dd9d5
commit 91d610162b
No known key found for this signature in database
9 changed files with 65 additions and 2 deletions

View file

@ -510,6 +510,7 @@ mod tests {
node_id: "plan".into(), node_id: "plan".into(),
name: "Plan".into(), name: "Plan".into(),
index: 0, index: 0,
visit: 1,
duration_ms: 5000, duration_ms: 5000,
status: "success".into(), status: "success".into(),
preferred_label: None, preferred_label: None,

View file

@ -487,6 +487,7 @@ mod tests {
node_id: node_id.into(), node_id: node_id.into(),
name: name.into(), name: name.into(),
index: 0, index: 0,
visit: 1,
handler_type: String::new(), handler_type: String::new(),
attempt: 1, attempt: 1,
max_attempts: 1, max_attempts: 1,
@ -510,6 +511,7 @@ mod tests {
node_id: node_id.into(), node_id: node_id.into(),
name: name.into(), name: name.into(),
index: 0, index: 0,
visit: 1,
duration_ms: 5000, duration_ms: 5000,
status: "success".into(), status: "success".into(),
preferred_label: None, preferred_label: None,
@ -561,6 +563,8 @@ mod tests {
emit( emit(
&mut ui, &mut ui,
Event::ParallelBranchStarted { Event::ParallelBranchStarted {
parallel_group_id: "fork1@1".into(),
parallel_branch_id: "fork1@1:0".into(),
branch: "security".into(), branch: "security".into(),
index: 0, index: 0,
}, },
@ -576,6 +580,8 @@ mod tests {
emit( emit(
&mut ui, &mut ui,
Event::ParallelBranchCompleted { Event::ParallelBranchCompleted {
parallel_group_id: "fork1@1".into(),
parallel_branch_id: "fork1@1:0".into(),
branch: "security".into(), branch: "security".into(),
index: 0, index: 0,
duration_ms: 2000, duration_ms: 2000,
@ -607,6 +613,8 @@ mod tests {
emit( emit(
&mut ui, &mut ui,
Event::ParallelBranchStarted { Event::ParallelBranchStarted {
parallel_group_id: "fork1@1".into(),
parallel_branch_id: "fork1@1:0".into(),
branch: "security".into(), branch: "security".into(),
index: 0, index: 0,
}, },
@ -700,6 +708,7 @@ mod tests {
node_id: "code".into(), node_id: "code".into(),
name: "Code".into(), name: "Code".into(),
index: 0, index: 0,
visit: 1,
attempt: 2, attempt: 2,
max_attempts: 3, max_attempts: 3,
delay_ms: 1500, delay_ms: 1500,
@ -944,6 +953,7 @@ mod tests {
node_id: "code".into(), node_id: "code".into(),
name: "Code".into(), name: "Code".into(),
index: 0, index: 0,
visit: 1,
attempt: 2, attempt: 2,
max_attempts: 3, max_attempts: 3,
delay_ms: 1500, delay_ms: 1500,
@ -1108,6 +1118,8 @@ mod tests {
emit( emit(
&mut ui, &mut ui,
Event::ParallelBranchStarted { Event::ParallelBranchStarted {
parallel_group_id: "fork1@1".into(),
parallel_branch_id: "fork1@1:0".into(),
branch: "security".into(), branch: "security".into(),
index: 0, index: 0,
}, },
@ -1115,6 +1127,8 @@ mod tests {
emit( emit(
&mut ui, &mut ui,
Event::ParallelBranchCompleted { Event::ParallelBranchCompleted {
parallel_group_id: "fork1@1".into(),
parallel_branch_id: "fork1@1:0".into(),
branch: "security".into(), branch: "security".into(),
index: 0, index: 0,
duration_ms: 500, duration_ms: 500,
@ -1144,6 +1158,7 @@ mod tests {
node_id: "code".into(), node_id: "code".into(),
name: "Code".into(), name: "Code".into(),
index: 0, index: 0,
visit: 1,
handler_type: "agent".into(), handler_type: "agent".into(),
attempt: 1, attempt: 1,
max_attempts: 1, max_attempts: 1,

View file

@ -599,6 +599,7 @@ mod tests {
node_id: "code".to_string(), node_id: "code".to_string(),
name: "Code".to_string(), name: "Code".to_string(),
index: 1, index: 1,
visit: 2,
duration_ms: 250, duration_ms: 250,
status: "partial_success".to_string(), status: "partial_success".to_string(),
preferred_label: None, preferred_label: None,

View file

@ -1692,6 +1692,7 @@ mod tests {
node_id: "code".into(), node_id: "code".into(),
name: "code".into(), name: "code".into(),
index: 0, index: 0,
visit: 1,
failure: failure.clone(), failure: failure.clone(),
will_retry: false, will_retry: false,
}; };

View file

@ -137,6 +137,7 @@ pub enum Event {
node_id: String, node_id: String,
name: String, name: String,
index: usize, index: usize,
visit: u32,
handler_type: String, handler_type: String,
attempt: usize, attempt: usize,
max_attempts: usize, max_attempts: usize,
@ -145,6 +146,7 @@ pub enum Event {
node_id: String, node_id: String,
name: String, name: String,
index: usize, index: usize,
visit: u32,
duration_ms: u64, duration_ms: u64,
status: String, status: String,
preferred_label: Option<String>, preferred_label: Option<String>,
@ -175,6 +177,7 @@ pub enum Event {
node_id: String, node_id: String,
name: String, name: String,
index: usize, index: usize,
visit: u32,
failure: FailureDetail, failure: FailureDetail,
will_retry: bool, will_retry: bool,
}, },
@ -182,6 +185,7 @@ pub enum Event {
node_id: String, node_id: String,
name: String, name: String,
index: usize, index: usize,
visit: u32,
attempt: usize, attempt: usize,
max_attempts: usize, max_attempts: usize,
delay_ms: u64, delay_ms: u64,
@ -193,10 +197,14 @@ pub enum Event {
join_policy: String, join_policy: String,
}, },
ParallelBranchStarted { ParallelBranchStarted {
parallel_group_id: String,
parallel_branch_id: String,
branch: String, branch: String,
index: usize, index: usize,
}, },
ParallelBranchCompleted { ParallelBranchCompleted {
parallel_group_id: String,
parallel_branch_id: String,
branch: String, branch: String,
index: usize, index: usize,
duration_ms: u64, duration_ms: u64,
@ -676,6 +684,7 @@ impl Event {
index, index,
failure, failure,
will_retry, will_retry,
..
} => { } => {
let error_msg = &failure.message; let error_msg = &failure.message;
if *will_retry { if *will_retry {
@ -705,6 +714,7 @@ impl Event {
attempt, attempt,
max_attempts, max_attempts,
delay_ms, delay_ms,
..
} => { } => {
warn!( warn!(
node_id, node_id,
@ -723,7 +733,7 @@ impl Event {
} => { } => {
debug!(branch_count, join_policy, "Parallel execution started"); debug!(branch_count, join_policy, "Parallel execution started");
} }
Self::ParallelBranchStarted { branch, index } => { Self::ParallelBranchStarted { branch, index, .. } => {
debug!(branch, index, "Parallel branch started"); debug!(branch, index, "Parallel branch started");
} }
Self::ParallelBranchCompleted { Self::ParallelBranchCompleted {
@ -2759,6 +2769,7 @@ mod tests {
node_id: "plan".to_string(), node_id: "plan".to_string(),
name: "Plan".to_string(), name: "Plan".to_string(),
index: 0, index: 0,
visit: 1,
duration_ms: 5000, duration_ms: 5000,
status: "success".to_string(), status: "success".to_string(),
preferred_label: None, preferred_label: None,
@ -2797,6 +2808,7 @@ mod tests {
node_id: "plan".to_string(), node_id: "plan".to_string(),
name: "Plan".to_string(), name: "Plan".to_string(),
index: 0, index: 0,
visit: 1,
duration_ms: 5000, duration_ms: 5000,
status: "success".to_string(), status: "success".to_string(),
preferred_label: None, preferred_label: None,
@ -2831,6 +2843,7 @@ mod tests {
node_id: "code".to_string(), node_id: "code".to_string(),
name: "Code".to_string(), name: "Code".to_string(),
index: 1, index: 1,
visit: 1,
failure: FailureDetail::new( failure: FailureDetail::new(
"lint failed", "lint failed",
crate::outcome::FailureCategory::Deterministic, crate::outcome::FailureCategory::Deterministic,
@ -3016,6 +3029,8 @@ mod tests {
); );
assert_eq!( assert_eq!(
event_name(&Event::ParallelBranchStarted { event_name(&Event::ParallelBranchStarted {
parallel_group_id: "plan@1".to_string(),
parallel_branch_id: "plan@1:0".to_string(),
branch: "fork".to_string(), branch: "fork".to_string(),
index: 0, index: 0,
}), }),

View file

@ -449,6 +449,7 @@ mod tests {
node_id: "work".into(), node_id: "work".into(),
name: "Work".into(), name: "Work".into(),
index: 2, index: 2,
visit: 2,
duration_ms: 100, duration_ms: 100,
status: "success".into(), status: "success".into(),
preferred_label: None, preferred_label: None,

View file

@ -132,6 +132,7 @@ impl Handler for ParallelHandler {
struct BranchSetup { struct BranchSetup {
target_id: String, target_id: String,
branch_index: usize, branch_index: usize,
parallel_branch_id: String,
branch_context: Context, branch_context: Context,
sandbox: Arc<dyn Sandbox>, sandbox: Arc<dyn Sandbox>,
worktree_path: Option<PathBuf>, worktree_path: Option<PathBuf>,
@ -150,9 +151,12 @@ impl Handler for ParallelHandler {
.unwrap_or("wait_all"), .unwrap_or("wait_all"),
); );
let parallel_visit = u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX);
let parallel_group_id = format!("{}@{}", node.id, parallel_visit);
services.emitter.emit(&Event::ParallelStarted { services.emitter.emit(&Event::ParallelStarted {
node_id: node.id.clone(), node_id: node.id.clone(),
visit: u32::try_from(visit_from_context(context)).unwrap_or(u32::MAX), visit: parallel_visit,
branch_count: branches.len(), branch_count: branches.len(),
join_policy: join_policy.to_string(), join_policy: join_policy.to_string(),
}); });
@ -253,9 +257,11 @@ impl Handler for ParallelHandler {
(Arc::clone(&services.sandbox), None) (Arc::clone(&services.sandbox), None)
}; };
let parallel_branch_id = format!("{parallel_group_id}:{branch_index}");
branch_setups.push(BranchSetup { branch_setups.push(BranchSetup {
target_id, target_id,
branch_index, branch_index,
parallel_branch_id,
branch_context, branch_context,
sandbox: branch_sandbox, sandbox: branch_sandbox,
worktree_path, worktree_path,
@ -283,6 +289,7 @@ impl Handler for ParallelHandler {
.as_ref() .as_ref()
.map(|gs| gs.git_author.clone()) .map(|gs| gs.git_author.clone())
.unwrap_or_default(); .unwrap_or_default();
let group_id = parallel_group_id.clone();
let handle = tokio::spawn(async move { let handle = tokio::spawn(async move {
let _permit = sem let _permit = sem
@ -291,6 +298,8 @@ impl Handler for ParallelHandler {
.map_err(|e| FabroError::handler(format!("semaphore error: {e}")))?; .map_err(|e| FabroError::handler(format!("semaphore error: {e}")))?;
emitter.emit(&Event::ParallelBranchStarted { emitter.emit(&Event::ParallelBranchStarted {
parallel_group_id: group_id.clone(),
parallel_branch_id: setup.parallel_branch_id.clone(),
branch: setup.target_id.clone(), branch: setup.target_id.clone(),
index: setup.branch_index, index: setup.branch_index,
}); });
@ -302,6 +311,8 @@ impl Handler for ParallelHandler {
setup.target_id setup.target_id
)); ));
emitter.emit(&Event::ParallelBranchCompleted { emitter.emit(&Event::ParallelBranchCompleted {
parallel_group_id: group_id.clone(),
parallel_branch_id: setup.parallel_branch_id.clone(),
branch: setup.target_id.clone(), branch: setup.target_id.clone(),
index: setup.branch_index, index: setup.branch_index,
duration_ms: millis_u64(branch_start.elapsed()), duration_ms: millis_u64(branch_start.elapsed()),
@ -386,6 +397,8 @@ impl Handler for ParallelHandler {
}; };
emitter.emit(&Event::ParallelBranchCompleted { emitter.emit(&Event::ParallelBranchCompleted {
parallel_group_id: group_id.clone(),
parallel_branch_id: setup.parallel_branch_id.clone(),
branch: setup.target_id.clone(), branch: setup.target_id.clone(),
index: setup.branch_index, index: setup.branch_index,
duration_ms: millis_u64(branch_start.elapsed()), duration_ms: millis_u64(branch_start.elapsed()),

View file

@ -79,6 +79,11 @@ fn response_from_outcome(node_id: &str, outcome: &Outcome) -> Option<String> {
.and_then(|value| value.as_str().map(ToOwned::to_owned)) .and_then(|value| value.as_str().map(ToOwned::to_owned))
} }
fn stage_visit(state: &WfRunState, node_id: &str) -> u32 {
let visits = state.node_visits.get(node_id).copied().unwrap_or(1);
u32::try_from(visits.max(1)).unwrap_or(u32::MAX)
}
#[async_trait] #[async_trait]
impl RunLifecycle<WorkflowGraph> for EventLifecycle { impl RunLifecycle<WorkflowGraph> for EventLifecycle {
async fn on_run_start(&self, _graph: &WorkflowGraph, _state: &WfRunState) -> CoreResult<()> { async fn on_run_start(&self, _graph: &WorkflowGraph, _state: &WfRunState) -> CoreResult<()> {
@ -120,12 +125,14 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
} }
let gv = node.inner(); let gv = node.inner();
let stage_index = state.stage_index; let stage_index = state.stage_index;
let visit = stage_visit(state, &gv.id);
let (loop_failure_signatures, restart_failure_signatures) = let (loop_failure_signatures, restart_failure_signatures) =
snapshot_failure_signatures(&self.circuit_breaker); snapshot_failure_signatures(&self.circuit_breaker);
self.emitter.emit(&Event::StageStarted { self.emitter.emit(&Event::StageStarted {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
handler_type: gv.handler_type().unwrap_or_default().to_string(), handler_type: gv.handler_type().unwrap_or_default().to_string(),
attempt: 1, attempt: 1,
max_attempts: 1, max_attempts: 1,
@ -134,6 +141,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
duration_ms: 0, duration_ms: 0,
status: StageStatus::Success.to_string(), status: StageStatus::Success.to_string(),
preferred_label: None, preferred_label: None,
@ -167,6 +175,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: state.stage_index, index: state.stage_index,
visit: stage_visit(state, &gv.id),
handler_type: gv.handler_type().unwrap_or_default().to_string(), handler_type: gv.handler_type().unwrap_or_default().to_string(),
attempt: ctx.attempt as usize, attempt: ctx.attempt as usize,
max_attempts: ctx.max_attempts as usize, max_attempts: ctx.max_attempts as usize,
@ -183,11 +192,13 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
let gv = ctx.node.inner(); let gv = ctx.node.inner();
let outcome = &ctx.result.outcome; let outcome = &ctx.result.outcome;
let stage_index = state.stage_index; let stage_index = state.stage_index;
let visit = stage_visit(state, &gv.id);
self.emitter.emit(&Event::StageFailed { self.emitter.emit(&Event::StageFailed {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
failure: outcome.failure.clone().unwrap_or_else(|| { failure: outcome.failure.clone().unwrap_or_else(|| {
FailureDetail::new("handler failed", FailureCategory::TransientInfra) FailureDetail::new("handler failed", FailureCategory::TransientInfra)
}), }),
@ -198,6 +209,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
attempt: ctx.attempt as usize, attempt: ctx.attempt as usize,
max_attempts: ctx.result.max_attempts as usize, max_attempts: ctx.result.max_attempts as usize,
delay_ms: ctx delay_ms: ctx
@ -221,6 +233,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
} }
let gv = node.inner(); let gv = node.inner();
let stage_index = state.stage_index; let stage_index = state.stage_index;
let visit = stage_visit(state, &gv.id);
let duration_ms = u64::try_from(result.duration.as_millis()).unwrap(); let duration_ms = u64::try_from(result.duration.as_millis()).unwrap();
let (loop_failure_signatures, restart_failure_signatures) = let (loop_failure_signatures, restart_failure_signatures) =
snapshot_failure_signatures(&self.circuit_breaker); snapshot_failure_signatures(&self.circuit_breaker);
@ -230,6 +243,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
failure: outcome.failure.clone().unwrap_or_else(|| { failure: outcome.failure.clone().unwrap_or_else(|| {
FailureDetail::new("handler failed", FailureCategory::Deterministic) FailureDetail::new("handler failed", FailureCategory::Deterministic)
}), }),
@ -240,6 +254,7 @@ impl RunLifecycle<WorkflowGraph> for EventLifecycle {
node_id: gv.id.clone(), node_id: gv.id.clone(),
name: gv.label().to_string(), name: gv.label().to_string(),
index: stage_index, index: stage_index,
visit,
duration_ms, duration_ms,
status: outcome.status.to_string(), status: outcome.status.to_string(),
preferred_label: outcome.preferred_label.clone(), preferred_label: outcome.preferred_label.clone(),

View file

@ -1197,6 +1197,7 @@ mod tests {
node_id: "plan".to_string(), node_id: "plan".to_string(),
name: "plan".to_string(), name: "plan".to_string(),
index: 0, index: 0,
visit: 1,
duration_ms: 1, duration_ms: 1,
status: "success".to_string(), status: "success".to_string(),
preferred_label: None, preferred_label: None,