From 06480db4170b6ce12eae2294668262f23181f8a1 Mon Sep 17 00:00:00 2001 From: Bryan Helmkamp Date: Mon, 23 Feb 2026 16:04:16 -0500 Subject: [PATCH] Register parallel handler to prevent dry-run hang on parallel pipelines The parallel handler was never registered in default_registry() due to a circular dependency (ParallelHandler::new needs Arc). Break the cycle by storing ParallelHandler as a separate field on PipelineEngine, created after Arc-wrapping the registry and emitter. Co-Authored-By: Claude Opus 4.6 --- crates/attractor/src/engine.rs | 27 ++++++++++++++++++++++++--- 1 file changed, 24 insertions(+), 3 deletions(-) diff --git a/crates/attractor/src/engine.rs b/crates/attractor/src/engine.rs index 68316bc88..15950065a 100644 --- a/crates/attractor/src/engine.rs +++ b/crates/attractor/src/engine.rs @@ -15,6 +15,7 @@ use crate::context::Context; use crate::error::{AttractorError, Result}; use crate::event::{EventEmitter, PipelineEvent}; use crate::graph::{Edge, Graph, Node}; +use crate::handler::parallel::ParallelHandler; use crate::handler::HandlerRegistry; use crate::interviewer::Interviewer; use crate::outcome::{Outcome, StageStatus}; @@ -451,17 +452,23 @@ pub struct RunConfig { /// The pipeline execution engine. pub struct PipelineEngine { - pub registry: HandlerRegistry, - pub emitter: EventEmitter, + registry: Arc, + emitter: Arc, + parallel_handler: ParallelHandler, pub interviewer: Option>, } impl PipelineEngine { #[must_use] pub fn new(registry: HandlerRegistry, emitter: EventEmitter) -> Self { + let registry = Arc::new(registry); + let emitter = Arc::new(emitter); + let parallel_handler = + ParallelHandler::new(Arc::clone(®istry), Arc::clone(&emitter)); Self { registry, emitter, + parallel_handler, interviewer: None, } } @@ -473,13 +480,27 @@ impl PipelineEngine { emitter: EventEmitter, interviewer: Arc, ) -> Self { + let registry = Arc::new(registry); + let emitter = Arc::new(emitter); + let parallel_handler = + ParallelHandler::new(Arc::clone(®istry), Arc::clone(&emitter)); Self { registry, emitter, + parallel_handler, interviewer: Some(interviewer), } } + /// Resolve the handler for a node, returning the parallel handler for + /// parallel nodes and delegating to the registry for everything else. + fn resolve_handler(&self, node: &Node) -> &dyn crate::handler::Handler { + if node.handler_type() == Some("parallel") { + return &self.parallel_handler; + } + self.registry.resolve(node) + } + /// Call inform on the interviewer, if one is configured. fn inform(&self, message: &str, stage: &str) { if let Some(ref interviewer) = self.interviewer { @@ -517,7 +538,7 @@ impl PipelineEngine { policy: &RetryPolicy, stage_index: usize, ) -> Result<(Outcome, u32)> { - let handler = self.registry.resolve(node); + let handler = self.resolve_handler(node); let node_timeout = node.timeout();