Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 2 additions & 0 deletions CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
59 changes: 47 additions & 12 deletions src/conductor/ai/agents/runtime/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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
Expand Down Expand Up @@ -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 (
Expand All @@ -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:
Expand Down Expand Up @@ -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_<name> 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

Expand All @@ -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")
Expand Down Expand Up @@ -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_<name> workers for swarm agents.

Expand All @@ -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

Expand All @@ -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.
Expand Down
11 changes: 9 additions & 2 deletions src/conductor/ai/agents/runtime/tool_registry.py
Original file line number Diff line number Diff line change
Expand Up @@ -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.

Expand Down Expand Up @@ -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
Expand Down
90 changes: 90 additions & 0 deletions tests/unit/ai/test_hybrid_transfer_workers.py
Original file line number Diff line number Diff line change
Expand Up @@ -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"]
78 changes: 78 additions & 0 deletions tests/unit/ai/test_swarm_handoff_check.py
Original file line number Diff line number Diff line change
Expand Up @@ -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",
}
53 changes: 53 additions & 0 deletions tests/unit/ai/test_tool_registry_domain.py
Original file line number Diff line number Diff line change
@@ -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
Loading