diff --git a/CHANGELOG.md b/CHANGELOG.md index adcb49dd..c9de2b50 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -29,3 +29,5 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 - Worker processes killed by a signal now log a diagnostic hint (signal number, `PYTHONFAULTHANDLER=1` guidance) instead of restarting silently - Per-schedule pause/resume now work on both Conductor server families: the client sends `PUT` (the OSS Conductor dialect — the spec-generated `GET` failed there) and transparently falls back to `GET` on a 405 for Orkes servers - `Agent(model="claude-code/...")` registration no longer crashes with `SpawnSafetyError`/`PicklingError` under the `spawn` start method (top-level or nested as a sub-agent); the tool worker is now built through the same picklable passthrough path already used by other frameworks +- Stateful `Strategy.SWARM` handoffs no longer hang forever (conductor-oss #1363): SWARM/hybrid transfer-tool worker registration now respects the server's `requiredWorkers` list instead of always attempting to `PUT` a `TaskDef` for compiler-owned `{source}_transfer_to_{peer}` control-signal names, which 404s against a server that (correctly) never created one for them. Requires a paired conductor-oss server fix (maps the transfer tool's compiled type to `INLINE`) to fully close the gap. +- A stateful multi-agent run (SWARM, hybrid handoff, sequential, router, manual-selection) no longer hangs on its member agents' own `@tool` calls: tool workers now always poll the execution's resolved worker domain when one is set, instead of only when that specific sub-agent (or tool) was individually marked `stateful=True`. The server domain-routes every declared worker-tool name for the whole agent tree once any part of it is stateful, so the client's extra per-agent re-check was redundant and, for nested sub-agents, wrong — their tool worker polled the domain-less queue while its task sat in the domain-scoped queue forever. diff --git a/src/conductor/ai/agents/runtime/runtime.py b/src/conductor/ai/agents/runtime/runtime.py index e4a1709c..80631ed7 100644 --- a/src/conductor/ai/agents/runtime/runtime.py +++ b/src/conductor/ai/agents/runtime/runtime.py @@ -17,7 +17,17 @@ import threading import time import uuid -from typing import TYPE_CHECKING, Any, AsyncIterator, Dict, Iterator, List, Optional, Union +from typing import ( + TYPE_CHECKING, + Any, + AsyncIterator, + Callable, + Dict, + Iterator, + List, + Optional, + Union, +) if TYPE_CHECKING: from conductor.ai.agents.runtime.config import AgentConfig @@ -1091,7 +1101,6 @@ def _server_needs(task_name: str) -> bool: agent.tools, agent.name, domain=domain, - agent_stateful=getattr(agent, "stateful", False), ) for t in agent.tools: from conductor.ai.agents.tool import get_tool_def @@ -1177,9 +1186,12 @@ def _server_needs(task_name: str) -> bool: task_name = f"{agent.name}_check_transfer" if _server_needs(task_name): self._register_check_transfer_worker(agent.name, domain=domain) - # Always register transfer tool workers — same reasoning as swarm: - # collectSimpleTaskNames may not recurse into nested sub-workflows. - self._register_hybrid_transfer_workers(agent, domain=domain) + # Transfer tools are compiler-owned control signals (toolType="handoff", + # never dispatched as real tasks) — only register the ones the server + # actually lists as required. + self._register_hybrid_transfer_workers( + agent, domain=domain, server_needs=_server_needs + ) # 6. Function-based router if ( @@ -1201,11 +1213,12 @@ def _server_needs(task_name: str) -> bool: # 7b. Swarm transfer tools and check_transfer workers if agent.strategy == "swarm" and agent.agents: - # Always register transfer workers for swarm agents — the server's - # requiredWorkers may not include them when the swarm is a nested - # registered sub-workflow (collectSimpleTaskNames doesn't recurse - # into separately-stored sub-workflow definitions). - self._register_swarm_transfer_workers(agent, domain=domain) + # Transfer tools are compiler-owned control signals (toolType="handoff", + # never dispatched as real tasks) — only register the ones the server + # actually lists as required. + self._register_swarm_transfer_workers( + agent, domain=domain, server_needs=_server_needs + ) if _server_needs(f"{agent.name}_check_transfer"): self._register_check_transfer_worker(agent.name, domain=domain) # parent for sub in agent.agents: @@ -1507,12 +1520,21 @@ def _register_check_transfer_worker( )(check_transfer_worker) def _register_hybrid_transfer_workers( - self, agent: Agent, domain: "Optional[str]" = None + self, + agent: Agent, + domain: "Optional[str]" = None, + server_needs: "Optional[Callable[[str], bool]]" = None, ) -> None: """Register transfer_to_ no-op workers for hybrid agents (tools + sub-agents). The transfer tools are no-ops — the actual handoff is detected by check_transfer which inspects toolCalls output from the LLM task. + + These are compiler-owned control-signal tools (toolType="handoff"), + never dispatched as real tasks — only register the ones *server_needs* + (the server's requiredWorkers gate) actually lists. When + *server_needs* is None, register everything (fallback for older + servers), matching every other call site's convention. """ from conductor.client.worker.worker_task import worker_task @@ -1523,6 +1545,8 @@ def _register_hybrid_transfer_workers( for sub in agent.agents: tool_name = f"{agent.name}_transfer_to_{sub.name}" + if server_needs is not None and not server_needs(tool_name): + continue transfer_worker = TransferNoopEntry() transfer_worker.__annotations__ = {"message": str, "return": object} probe_spawn_safety(transfer_worker, tool_name, group="system") @@ -1603,7 +1627,10 @@ def _register_handoff_worker(self, agent: Agent, domain: "Optional[str]" = None) )(handoff_check_worker) def _register_swarm_transfer_workers( - self, agent: Agent, domain: "Optional[str]" = None + self, + agent: Agent, + domain: "Optional[str]" = None, + server_needs: "Optional[Callable[[str], bool]]" = None, ) -> None: """Register transfer_to_ workers for swarm agents. @@ -1614,6 +1641,12 @@ def _register_swarm_transfer_workers( When allowed_transitions is set, transfers to targets that no agent is allowed to reach return an error message so the LLM knows to try a different tool. + + These are compiler-owned control-signal tools (toolType="handoff"), + never dispatched as real tasks — only register the ones + *server_needs* (the server's requiredWorkers gate) actually lists. + When *server_needs* is None, register everything (fallback for older + servers), matching every other call site's convention. """ from conductor.client.worker.worker_task import worker_task @@ -1636,6 +1669,8 @@ def _register_swarm_transfer_workers( if tool_name in registered: continue registered.add(tool_name) + if server_needs is not None and not server_needs(tool_name): + continue # If this target is never reachable via allowed_transitions, # return an error message so the LLM knows to stop trying. diff --git a/src/conductor/ai/agents/runtime/tool_registry.py b/src/conductor/ai/agents/runtime/tool_registry.py index 595e2673..3f4496e8 100644 --- a/src/conductor/ai/agents/runtime/tool_registry.py +++ b/src/conductor/ai/agents/runtime/tool_registry.py @@ -30,7 +30,6 @@ def register_tool_workers( tools: List[Any], agent_name: str, domain: Optional[str] = None, - agent_stateful: bool = False, ) -> None: """Register tool functions as Conductor workers and populate global registries. @@ -84,7 +83,15 @@ def register_tool_workers( ), register_task_def=True, overwrite_task_def=True, - domain=domain if (agent_stateful or td.stateful) else None, + # *domain* is already None unless this execution is stateful + # (resolved once, tree-wide, by + # AgentRuntime._resolve_worker_domain / _has_stateful_tools) — + # re-deriving "is this stateful" from the *current* agent's own + # .stateful flag broke nested sub-agents (e.g. SWARM members) + # whose own .stateful is False even though the top-level + # orchestrator's is True and the whole compiled graph — this + # tool's TaskDef included — is domain-routed either way. + domain=domain, lease_extend_enabled=True, )(wrapper) _tool_task_names[td.name] = td.name diff --git a/tests/unit/ai/test_hybrid_transfer_workers.py b/tests/unit/ai/test_hybrid_transfer_workers.py index 1dcf6f14..cafa3489 100644 --- a/tests/unit/ai/test_hybrid_transfer_workers.py +++ b/tests/unit/ai/test_hybrid_transfer_workers.py @@ -175,3 +175,93 @@ def capture(agent, **kw): ) assert mgr in called_with, "_register_hybrid_transfer_workers was not called for hybrid agent" + + def test_skips_worker_when_not_server_needed(self): + """conductor-oss #1363: transfer tools are compiler-owned control signals (never + dispatched as real tasks) -- don't PUT a taskdef for one the server didn't ask for.""" + + @tool + def lookup(k: str) -> str: + """Look up.""" + return k + + mgr = Agent(name="manager", model="openai/gpt-4o", tools=[lookup]) + mgr.agents = [ + Agent(name="researcher", model="openai/gpt-4o"), + Agent(name="writer", model="openai/gpt-4o"), + ] + + from conductor.ai.agents.runtime.runtime import AgentRuntime + + rt = AgentRuntime.__new__(AgentRuntime) + + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_hybrid_transfer_workers( + mgr, domain=None, server_needs=lambda name: False + ) + + assert registered == [] + + def test_registers_only_server_needed_workers(self): + """Per-name gating: one transfer name needed, the other isn't.""" + + @tool + def lookup(k: str) -> str: + """Look up.""" + return k + + mgr = Agent(name="manager", model="openai/gpt-4o", tools=[lookup]) + mgr.agents = [ + Agent(name="researcher", model="openai/gpt-4o"), + Agent(name="writer", model="openai/gpt-4o"), + ] + + from conductor.ai.agents.runtime.runtime import AgentRuntime + + rt = AgentRuntime.__new__(AgentRuntime) + + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + needed = {"manager_transfer_to_researcher"} + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_hybrid_transfer_workers( + mgr, domain=None, server_needs=lambda name: name in needed + ) + + assert registered == ["manager_transfer_to_researcher"] + + def test_registers_all_when_server_needs_is_none(self): + """server_needs=None (older server / fallback) must preserve prior unconditional behavior.""" + + @tool + def lookup(k: str) -> str: + """Look up.""" + return k + + mgr = Agent(name="manager", model="openai/gpt-4o", tools=[lookup]) + mgr.agents = [Agent(name="researcher", model="openai/gpt-4o")] + + from conductor.ai.agents.runtime.runtime import AgentRuntime + + rt = AgentRuntime.__new__(AgentRuntime) + + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_hybrid_transfer_workers(mgr, domain=None, server_needs=None) + + assert registered == ["manager_transfer_to_researcher"] diff --git a/tests/unit/ai/test_swarm_handoff_check.py b/tests/unit/ai/test_swarm_handoff_check.py index 62d5013b..c2e343f9 100644 --- a/tests/unit/ai/test_swarm_handoff_check.py +++ b/tests/unit/ai/test_swarm_handoff_check.py @@ -686,3 +686,81 @@ def test_server_required_workers_controls_registration(self): patch.object(rt2, "_register_check_transfer_worker"): rt2._register_workers(swarm, required_workers=required_without, domain=None) mock_handoff2.assert_not_called() + + +class TestRegisterSwarmTransferWorkers: + """conductor-oss #1363: transfer tools are compiler-owned control signals + (toolType="handoff", never dispatched as real tasks) -- only register the + ones the server's requiredWorkers gate actually lists.""" + + def test_skips_worker_when_not_server_needed(self): + a = Agent(name="agent_a", model="openai/gpt-4o") + b = Agent(name="agent_b", model="openai/gpt-4o") + swarm = Agent( + name="swarm_parent", + model="openai/gpt-4o", + agents=[a, b], + strategy=Strategy.SWARM, + ) + + rt = AgentRuntime.__new__(AgentRuntime) + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_swarm_transfer_workers(swarm, domain=None, server_needs=lambda name: False) + + assert registered == [] + + def test_registers_only_server_needed_workers(self): + a = Agent(name="agent_a", model="openai/gpt-4o") + b = Agent(name="agent_b", model="openai/gpt-4o") + swarm = Agent( + name="swarm_parent", + model="openai/gpt-4o", + agents=[a, b], + strategy=Strategy.SWARM, + ) + + rt = AgentRuntime.__new__(AgentRuntime) + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + needed = {"agent_a_transfer_to_agent_b"} + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_swarm_transfer_workers( + swarm, domain=None, server_needs=lambda name: name in needed + ) + + assert registered == ["agent_a_transfer_to_agent_b"] + + def test_registers_all_when_server_needs_is_none(self): + """server_needs=None (older server / fallback) must preserve prior unconditional behavior.""" + a = Agent(name="agent_a", model="openai/gpt-4o") + swarm = Agent( + name="swarm_parent", + model="openai/gpt-4o", + agents=[a], + strategy=Strategy.SWARM, + ) + + rt = AgentRuntime.__new__(AgentRuntime) + registered = [] + + def fake_worker_task(**kwargs): + registered.append(kwargs["task_definition_name"]) + return lambda fn: fn + + with patch("conductor.client.worker.worker_task.worker_task", side_effect=fake_worker_task): + rt._register_swarm_transfer_workers(swarm, domain=None, server_needs=None) + + assert set(registered) == { + "swarm_parent_transfer_to_agent_a", + "agent_a_transfer_to_swarm_parent", + } diff --git a/tests/unit/ai/test_tool_registry_domain.py b/tests/unit/ai/test_tool_registry_domain.py new file mode 100644 index 00000000..584ea193 --- /dev/null +++ b/tests/unit/ai/test_tool_registry_domain.py @@ -0,0 +1,53 @@ +"""Tests for ToolRegistry.register_tool_workers domain propagation. + +A nested sub-agent (e.g. a SWARM member) is never itself marked +`stateful=True` — only the top-level orchestrator carries that flag, even +though the whole compiled graph (including the sub-agent's own tools) is +domain-routed once any agent in the tree is stateful. `domain` is already +resolved once, tree-wide, by `AgentRuntime._resolve_worker_domain` / +`_has_stateful_tools` before `register_tool_workers` is ever called, so it +must be applied whenever non-None — regardless of the immediate agent's own +`.stateful` attribute. Re-deriving statefulness locally left a SWARM +member's own tool worker polling the domain-less queue while its task sat in +the domain-scoped queue forever (conductor-oss #1363 follow-up). +""" + +from unittest.mock import patch + +from conductor.ai.agents.runtime.tool_registry import ToolRegistry +from conductor.ai.agents.tool import tool + + +@tool +def swarm_tool(task: str) -> str: + """Perform a task.""" + return f"done:{task}" + + +def _register_and_capture_domain(domain): + captured = {} + + def fake_worker_task(**kwargs): + def _decorator(fn): + captured["domain"] = kwargs.get("domain") + return fn + + return _decorator + + with patch( + "conductor.client.worker.worker_task.worker_task", + side_effect=fake_worker_task, + ): + ToolRegistry().register_tool_workers([swarm_tool], "swarm_agent_a", domain=domain) + + return captured + + +class TestRegisterToolWorkersDomain: + def test_domain_applied_regardless_of_agent_or_tool_stateful_flag(self): + captured = _register_and_capture_domain(domain="abc123") + assert captured["domain"] == "abc123" + + def test_domain_none_stays_none(self): + captured = _register_and_capture_domain(domain=None) + assert captured["domain"] is None