Skip to content
Closed
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
16 changes: 14 additions & 2 deletions e2e/test_suite14_stateful_domain.py
Original file line number Diff line number Diff line change
Expand Up @@ -133,8 +133,20 @@ def _get_task_to_domain(execution_id):


def _find_tasks_by_type(tasks, task_def_name):
"""Find tasks matching a taskDefName (or containing it)."""
return [t for t in tasks if task_def_name in t.get("taskDefName", "")]
"""Find tasks matching a taskDefName or referenceTaskName (or containing it).

referenceTaskName is checked too because compiler-generated system tasks carry
their identity there, not in taskDefName: handoff_check is emitted as an INLINE
task, so its taskDefName is literally "INLINE" and only the reference name says
what it is. Matching taskDefName alone made assertions on those tasks fail even
when the tasks existed and had COMPLETED.
"""
return [
t
for t in tasks
if task_def_name in t.get("taskDefName", "")
or task_def_name in t.get("referenceTaskName", "")
]


def _find_scheduled_tasks(tasks):
Expand Down
14 changes: 11 additions & 3 deletions src/conductor/ai/agents/config_serializer.py
Original file line number Diff line number Diff line change
Expand Up @@ -35,7 +35,7 @@ def serialize(self, agent: "Agent") -> dict:
"""
return self._serialize_agent(agent)

def _serialize_agent(self, agent: "Agent") -> dict:
def _serialize_agent(self, agent: "Agent", *, inherited_stateful: bool = False) -> dict:
from conductor.ai.agents.agent import PromptTemplate

# Skill agents — emit the raw skill config so the server's
Expand Down Expand Up @@ -106,14 +106,22 @@ def _serialize_agent(self, agent: "Agent") -> dict:

# Tools
if agent.tools:
agent_stateful = getattr(agent, "stateful", False)
# Statefulness belongs to the composite, not only the agent that declares it:
# a member of a stateful swarm/team shares the parent's session, so its tools
# are stateful too. Reading only `agent.stateful` left member tools unmarked,
# which disagreed with how the server routes their tasks. Mirrors the same
# inheritance in AgentRuntime._register_workers.
agent_stateful = bool(getattr(agent, "stateful", False)) or inherited_stateful
config["tools"] = [
self._serialize_tool(t, agent_stateful=agent_stateful) for t in agent.tools
]

# Sub-agents (recursive)
if agent.agents:
config["agents"] = [self._serialize_agent(a) for a in agent.agents]
_stateful = bool(getattr(agent, "stateful", False)) or inherited_stateful
config["agents"] = [
self._serialize_agent(a, inherited_stateful=_stateful) for a in agent.agents
]

# Router
if agent.router is not None:
Expand Down
30 changes: 26 additions & 4 deletions src/conductor/ai/agents/runtime/runtime.py
Original file line number Diff line number Diff line change
Expand Up @@ -1032,7 +1032,12 @@ def _collect_worker_names(self, agent: Agent, *, required_workers: Optional[set]
return names

def _register_workers(
self, agent: Agent, *, required_workers: Optional[set] = None, domain: Optional[str] = None
self,
agent: Agent,
*,
required_workers: Optional[set] = None,
domain: Optional[str] = None,
inherited_stateful: bool = False,
) -> None:
"""Register all workers needed for SDK-side execution.

Expand Down Expand Up @@ -1084,14 +1089,22 @@ def _server_needs(task_name: str) -> bool:
self._register_passthrough_worker(worker)
return

# Statefulness is a property of the composite, not just the agent that declares it.
# A member of a stateful swarm/team shares the parent's session, so the server
# domain-routes ITS tool tasks too. Reading only `agent.stateful` here left member
# tool workers polling the default queue while their tasks sat in the domain queue:
# pollCount stayed 0, the fork's JOIN never satisfied, and the workflow timed out.
# Inherit downward so worker registration matches how the server routes.
effective_stateful = bool(getattr(agent, "stateful", False)) or inherited_stateful

# 1. Tools (and tool-level guardrails) — always registered
if agent.tools:
tc = ToolRegistry()
tc.register_tool_workers(
agent.tools,
agent.name,
domain=domain,
agent_stateful=getattr(agent, "stateful", False),
agent_stateful=effective_stateful,
)
for t in agent.tools:
from conductor.ai.agents.tool import get_tool_def
Expand All @@ -1101,7 +1114,11 @@ def _server_needs(task_name: str) -> bool:
if td.tool_type == "agent_tool" and td.config and "agent" in td.config:
nested_agent = td.config["agent"]
if not getattr(nested_agent, "external", False):
self._register_workers(nested_agent, required_workers=required_workers)
self._register_workers(
nested_agent,
required_workers=required_workers,
inherited_stateful=effective_stateful,
)
# Register tool-level guardrail workers
tool_guardrails = [
g
Expand Down Expand Up @@ -1240,7 +1257,12 @@ def _server_needs(task_name: str) -> bool:
)
self._register_passthrough_worker(worker)
elif not sub.external:
self._register_workers(sub, required_workers=required_workers, domain=domain)
self._register_workers(
sub,
required_workers=required_workers,
domain=domain,
inherited_stateful=effective_stateful,
)

# ── Worker registration helpers ────────────────────────────────

Expand Down
Loading