diff --git a/e2e/test_suite14_stateful_domain.py b/e2e/test_suite14_stateful_domain.py index 76e1aad3..b76a3a5d 100644 --- a/e2e/test_suite14_stateful_domain.py +++ b/e2e/test_suite14_stateful_domain.py @@ -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): diff --git a/src/conductor/ai/agents/config_serializer.py b/src/conductor/ai/agents/config_serializer.py index fcc8ed67..51c50dbf 100644 --- a/src/conductor/ai/agents/config_serializer.py +++ b/src/conductor/ai/agents/config_serializer.py @@ -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 @@ -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: diff --git a/src/conductor/ai/agents/runtime/runtime.py b/src/conductor/ai/agents/runtime/runtime.py index e4a1709c..0872ff0b 100644 --- a/src/conductor/ai/agents/runtime/runtime.py +++ b/src/conductor/ai/agents/runtime/runtime.py @@ -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. @@ -1084,6 +1089,14 @@ 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() @@ -1091,7 +1104,7 @@ def _server_needs(task_name: str) -> bool: 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 @@ -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 @@ -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 ────────────────────────────────