diff --git a/lib/crates/fabro-workflows/src/engine.rs b/lib/crates/fabro-workflows/src/engine.rs index 2348f9d8a..1de1c8e78 100644 --- a/lib/crates/fabro-workflows/src/engine.rs +++ b/lib/crates/fabro-workflows/src/engine.rs @@ -253,9 +253,9 @@ pub fn resolve_fidelity( // --- Thread ID resolution (spec 5.4) --- -/// Resolve the thread ID for a node, following the precedence (spec lines 1196-1204): -/// 1. Target node `thread_id` attribute -/// 2. Incoming edge `thread_id` attribute +/// Resolve the thread ID for a node, following the precedence: +/// 1. Incoming edge `thread_id` attribute +/// 2. Target node `thread_id` attribute /// 3. Graph-level default thread /// 4. Derived class from enclosing subgraph (first class from the node's classes list) /// 5. Fallback to previous node ID @@ -266,16 +266,16 @@ pub fn resolve_thread_id( graph: &Graph, previous_node_id: Option<&str>, ) -> Option { - // Step 1: Node thread_id - if let Some(tid) = node.thread_id() { - return Some(tid.to_string()); - } - // Step 2: Edge thread_id + // Step 1: Edge thread_id if let Some(edge) = incoming_edge { if let Some(tid) = edge.thread_id() { return Some(tid.to_string()); } } + // Step 2: Node thread_id + if let Some(tid) = node.thread_id() { + return Some(tid.to_string()); + } // Step 3: Graph-level default thread if let Some(tid) = graph.default_thread() { return Some(tid.to_string()); @@ -3691,7 +3691,25 @@ mod tests { } #[test] - fn thread_id_node_overrides_edge() { + fn thread_id_node_used_when_no_edge_thread() { + // When the edge has no thread_id, the node's thread_id is used. + let mut node = Node::new("work"); + node.attrs.insert( + "thread_id".to_string(), + AttrValue::String("node-thread".to_string()), + ); + let edge = Edge::new("prev", "work"); + let graph = Graph::new("test"); + assert_eq!( + resolve_thread_id(Some(&edge), &node, &graph, Some("prev")), + Some("node-thread".to_string()) + ); + } + + #[test] + fn thread_id_edge_overrides_node() { + // Edge thread_id should take precedence over node thread_id, + // matching the fidelity precedence where edge > node. let mut node = Node::new("work"); node.attrs.insert( "thread_id".to_string(), @@ -3705,7 +3723,8 @@ mod tests { let graph = Graph::new("test"); assert_eq!( resolve_thread_id(Some(&edge), &node, &graph, Some("prev")), - Some("node-thread".to_string()) + Some("edge-thread".to_string()), + "edge thread_id should override node thread_id" ); } diff --git a/lib/crates/fabro-workflows/tests/integration.rs b/lib/crates/fabro-workflows/tests/integration.rs index e1a23116b..6db83366c 100644 --- a/lib/crates/fabro-workflows/tests/integration.rs +++ b/lib/crates/fabro-workflows/tests/integration.rs @@ -5843,7 +5843,7 @@ async fn fidelity_checkpoint_roundtrip_preserves_fidelity() { #[tokio::test] async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { - // When both node and edge have thread_id, the node's takes precedence (spec step 1 > step 2). + // When both node and edge have thread_id, the edge's takes precedence (step 1 > step 2). let mut graph = make_graph_with_start_exit("NodeOverridesEdgeThreadTest"); let mut work = Node::new("work"); work.attrs.insert( @@ -5902,8 +5902,8 @@ async fn fidelity_node_thread_id_overrides_edge_thread_id_in_pipeline() { assert_eq!(thread_ids[0].0, "work"); assert_eq!( thread_ids[0].1, - Some("node-thread".to_string()), - "node thread_id should take precedence over edge thread_id" + Some("edge-thread".to_string()), + "edge thread_id should take precedence over node thread_id" ); }