diff --git a/docs/pyagentspec/source/agentspec/language_spec_nightly.rst b/docs/pyagentspec/source/agentspec/language_spec_nightly.rst index 99de2c3a..fc3d9dad 100644 --- a/docs/pyagentspec/source/agentspec/language_spec_nightly.rst +++ b/docs/pyagentspec/source/agentspec/language_spec_nightly.rst @@ -2270,6 +2270,9 @@ The ManagerWorkers has two main parameters: - Workers cannot interact with the end user directly. - When invoked, each worker can leverage its equipped tools to complete the assigned task and report the result back to the group manager. +The ``ManagerWorkers`` input and output schemas must match those of its ``group_manager``. +In particular, the two components must declare the same input property names and the same output property names, +and each corresponding property must have the same type. Datastores ~~~~~~~~~~ diff --git a/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.json b/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.json index 1159e7f9..221088a6 100644 --- a/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.json +++ b/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.json @@ -4,7 +4,16 @@ "name": "managerworkers", "description": null, "metadata": {}, - "inputs": [], + "inputs": [ + { + "title": "customer_id", + "type": "string" + }, + { + "title": "company_policy_info", + "type": "string" + } + ], "outputs": [], "group_manager": { "component_type": "Agent", @@ -208,5 +217,5 @@ "model_id": "llama-4-maverick" } }, - "agentspec_version": "26.1.0" + "agentspec_version": "26.4.0" } diff --git a/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.yaml b/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.yaml index fb74f087..01bb633f 100644 --- a/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.yaml +++ b/docs/pyagentspec/source/agentspec_config_examples/howto_managerworkers.yaml @@ -9,7 +9,11 @@ id: 248045cb-ca6f-4d6f-9e22-d28b452a25da name: managerworkers description: null metadata: {} -inputs: [] +inputs: +- title: company_policy_info + type: string +- title: customer_id + type: string outputs: [] group_manager: component_type: Agent @@ -225,4 +229,4 @@ $referenced_components: default_generation_parameters: null url: http://url.to.my.vllm.server/llama4mav model_id: llama-4-maverick -agentspec_version: 26.1.0 +agentspec_version: 26.4.0 diff --git a/docs/pyagentspec/source/changelog.rst b/docs/pyagentspec/source/changelog.rst index 60795a13..c9ceb2b5 100644 --- a/docs/pyagentspec/source/changelog.rst +++ b/docs/pyagentspec/source/changelog.rst @@ -71,6 +71,25 @@ New features Added ``DbmsVectorChainLlmConfig`` for configuring LLM requests executed through Oracle Database ``DBMS_VECTOR_CHAIN``. +* **ManagerWorkers I/O update** + + The ManagerWorkers language specification now requires its input and output + schemas to use the same property names and types as its group manager. + This allows exposing inputs required by the manager agent (e.g., the placeholders + in its system prompt). + + We thank @spichen for the contribution! + +* **ManagerWorkers support in the LangGraph adapter** + + The LangGraph adapter now converts ``ManagerWorkers`` into hierarchical graphs: + the group manager delegates tasks to workers and receives their results before + producing a final response. Nested ``ManagerWorkers`` can be used as workers. + The adapter also supports ``ManagerWorkers`` in Flow ``AgentNode`` steps with + one string output only. + + We thank @spichen for the contribution! + * **MCP tool retry policies** Added ``retry_policy`` support to ``MCPTool`` and ``MCPToolBox`` so runtimes can diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/__init__.py b/pyagentspec/src/pyagentspec/adapters/langgraph/__init__.py index baeb855b..fb75ddeb 100644 --- a/pyagentspec/src/pyagentspec/adapters/langgraph/__init__.py +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/__init__.py @@ -6,10 +6,13 @@ """Agent Spec adapter for the LangGraph agentic framework.""" +from ._managerworkers import DELEGATE_TOOL_PREFIX, is_delegation_tool_name from .agentspecexporter import AgentSpecExporter from .agentspecloader import AgentSpecLoader __all__ = [ "AgentSpecLoader", "AgentSpecExporter", + "DELEGATE_TOOL_PREFIX", + "is_delegation_tool_name", ] diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/_execution_span.py b/pyagentspec/src/pyagentspec/adapters/langgraph/_execution_span.py new file mode 100644 index 00000000..1fe7e0bb --- /dev/null +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/_execution_span.py @@ -0,0 +1,87 @@ +# Copyright © 2025, 2026 Oracle and/or its affiliates. +# +# This software is under the Apache License 2.0 +# (LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0) or Universal Permissive License +# (UPL) 1.0 (LICENSE-UPL or https://oss.oracle.com/licenses/upl), at your option. + +"""Wrap a compiled graph's ``stream``/``astream`` in an Agent Spec execution span. + +Agent, Flow and ManagerWorkers graphs all need the same wrapper: open a span, emit a +Start event carrying the invocation inputs, yield the chunks the underlying stream +produces while remembering the last state chunk, then emit an End event built from +that final state. Only the span class and the two event payloads differ, so they come +in as factories. ``invoke``/``ainvoke`` need no patch; they use ``stream``/``astream`` +internally. +""" + +from typing import Any, AsyncGenerator, Callable, Dict, Generator + +from pyagentspec.adapters.langgraph._types import CompiledStateGraph + + +def _invocation_inputs(kwargs: Dict[str, Any]) -> Dict[str, Any]: + """The ``input=`` argument of the patched call, or ``{}`` when it isn't a dict.""" + inputs = kwargs.get("input", {}) + return inputs if isinstance(inputs, dict) else {} + + +def _final_state(chunk: Any, so_far: Any) -> Any: + """Fold one streamed chunk into the running "last state seen". + + State arrives as ``(namespace, state)`` tuples; other chunk shapes aren't + something to build the End event from, so they leave the fold untouched. + """ + return chunk[1] if isinstance(chunk, tuple) else so_far + + +async def _async_or_sync( + async_call: Callable[..., Any], sync_call: Callable[..., Any], *args: Any +) -> None: + """Await ``async_call``, falling back to ``sync_call`` for spans that don't + implement the async half of the tracing protocol.""" + try: + await async_call(*args) + except NotImplementedError: + sync_call(*args) + + +def patch_with_execution_span( + compiled_graph: CompiledStateGraph[Any, Any, Any], + make_span: Callable[[], Any], + make_start_event: Callable[[Dict[str, Any]], Any], + make_end_event: Callable[[Dict[str, Any]], Any], +) -> None: + """Monkey-patch ``compiled_graph.stream`` / ``.astream`` to run inside a span. + + ``make_start_event`` receives the invocation inputs; ``make_end_event`` + receives the final state chunk (``{}`` when the run produced none). + """ + original_stream = compiled_graph.stream + original_astream = compiled_graph.astream + + def patched_stream(*args: Any, **kwargs: Any) -> Generator[Any, Any, None]: + with make_span() as span: + span.add_event(make_start_event(_invocation_inputs(kwargs))) + state: Any = {} + for chunk in original_stream(*args, **kwargs): + yield chunk + state = _final_state(chunk, state) + span.add_event(make_end_event(state if isinstance(state, dict) else {})) + + async def patched_astream(*args: Any, **kwargs: Any) -> AsyncGenerator[Any, Any]: + span = make_span() + await _async_or_sync(span.start_async, span.start) + try: + start_event = make_start_event(_invocation_inputs(kwargs)) + await _async_or_sync(span.add_event_async, span.add_event, start_event) + state: Any = {} + async for chunk in original_astream(*args, **kwargs): + yield chunk + state = _final_state(chunk, state) + end_event = make_end_event(state if isinstance(state, dict) else {}) + await _async_or_sync(span.add_event_async, span.add_event, end_event) + finally: + await _async_or_sync(span.end_async, span.end) + + compiled_graph.stream = patched_stream # type: ignore[method-assign] + compiled_graph.astream = patched_astream # type: ignore[method-assign] diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/_langgraphconverter.py b/pyagentspec/src/pyagentspec/adapters/langgraph/_langgraphconverter.py index 549e24b2..0ca55cb0 100644 --- a/pyagentspec/src/pyagentspec/adapters/langgraph/_langgraphconverter.py +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/_langgraphconverter.py @@ -12,11 +12,10 @@ from typing import ( TYPE_CHECKING, Any, - AsyncGenerator, Awaitable, Callable, Dict, - Generator, + Hashable, List, Optional, Tuple, @@ -37,6 +36,16 @@ _build_type_from_schema, create_pydantic_model_from_properties, ) +from pyagentspec.adapters.langgraph._execution_span import patch_with_execution_span +from pyagentspec.adapters.langgraph._managerworkers import ( + _MANAGER_NODE_KEY, + _append_workers_roster, + _make_manager_router, + _make_worker_delegation_tool, + _patch_with_manager_workers_execution_span, + _safe_node_name, + _wrap_worker_for_subgraph, +) from pyagentspec.adapters.langgraph._node_execution import ( NodeExecutor, extract_outputs_from_invoke_result, @@ -70,6 +79,7 @@ AgentSpecToolCallbackHandler, ) from pyagentspec.agent import Agent as AgentSpecAgent +from pyagentspec.agenticcomponent import AgenticComponent as AgentSpecAgenticComponent from pyagentspec.flows.edges import ControlFlowEdge as AgentSpecControlFlowEdge from pyagentspec.flows.edges import DataFlowEdge as AgentSpecDataFlowEdge from pyagentspec.flows.flow import Flow as AgentSpecFlow @@ -103,6 +113,7 @@ ) from pyagentspec.llms.openaiconfig import OpenAiConfig from pyagentspec.llms.vllmconfig import VllmConfig +from pyagentspec.managerworkers import ManagerWorkers as AgentSpecManagerWorkers from pyagentspec.mcp.clienttransport import ClientTransport as AgentSpecClientTransport from pyagentspec.mcp.clienttransport import SSEmTLSTransport as AgentSpecSSEmTLSTransport from pyagentspec.mcp.clienttransport import SSETransport as AgentSpecSSETransport @@ -274,6 +285,15 @@ def _convert( config=config, middleware=middleware, ) + elif isinstance(agentspec_component, AgentSpecManagerWorkers): + return self._manager_workers_convert_to_langgraph( + agentspec_component, + tool_registry=tool_registry, + converted_components=converted_components, + checkpointer=checkpointer, + config=config, + middleware=middleware, + ) elif isinstance(agentspec_component, AgentSpecLlmConfig): return self._llm_convert_to_langgraph(agentspec_component, config=config) elif isinstance(agentspec_component, AgentSpecClientTransport): @@ -475,87 +495,18 @@ def _find_property(properties: List[AgentSpecProperty], name: str) -> AgentSpecP "Prefer invoke/stream or upgrade to Python 3.11+ for ainvoke/astream." ) - # To enable flow execution traces monkey patch all the functions that invoke the compiled graph - - original_stream = compiled_graph.stream - - def patch_with_flow_execution_span(*args: Any, **kwargs: Any) -> Generator[Any, Any, None]: - span_name = f"FlowExecution[{flow.name}]" - inputs = kwargs.get("input", {}) - if not isinstance(inputs, dict): - inputs = {} - with AgentSpecFlowExecutionSpan(name=span_name, flow=flow) as span: - span.add_event(AgentSpecFlowExecutionStart(flow=flow, inputs=inputs)) - original_result: dict[str, Any] | Any = {} - result: dict[str, Any] - # This is going to patch stream and astream, that return iterators and yield chunks - for chunk in original_stream(*args, **kwargs): - yield chunk - if isinstance(chunk, tuple): - original_result = chunk[1] - if not isinstance(original_result, dict): - result = {} - else: - result = original_result - span.add_event( - AgentSpecFlowExecutionEnd( - flow=flow, - outputs=result.get("outputs", {}), - branch_selected=result.get("node_execution_details", {}).get("branch", ""), - ) - ) - - original_astream = compiled_graph.astream - - async def patch_async_with_flow_execution_span( - *args: Any, **kwargs: Any - ) -> AsyncGenerator[Any, Any]: - span_name = f"FlowExecution[{flow.name}]" - inputs = kwargs.get("input", {}) - if not isinstance(inputs, dict): - inputs = {} - span = AgentSpecFlowExecutionSpan(name=span_name, flow=flow) - try: - await span.start_async() - except NotImplementedError: - span.start() - try: - try: - await span.add_event_async( - AgentSpecFlowExecutionStart(flow=flow, inputs=inputs) - ) - except NotImplementedError: - span.add_event(AgentSpecFlowExecutionStart(flow=flow, inputs=inputs)) - original_result: dict[str, Any] | Any = {} - result: dict[str, Any] - # This is going to patch stream and astream, that return iterators and yield chunks - async for chunk in original_astream(*args, **kwargs): - yield chunk - if isinstance(chunk, tuple): - original_result = chunk[1] - if not isinstance(original_result, dict): - result = {} - else: - result = original_result - span_end_event = AgentSpecFlowExecutionEnd( - flow=flow, - outputs=result.get("outputs", {}), - branch_selected=result.get("node_execution_details", {}).get("branch", ""), - ) - try: - await span.add_event_async(span_end_event) - except NotImplementedError: - span.add_event(span_end_event) - finally: - try: - await span.end_async() - except NotImplementedError: - span.end() - - # Monkey patch invocation functions to inject tracing - # No need to patch `(a)invoke` as the internally use `(a)stream` - compiled_graph.stream = patch_with_flow_execution_span # type: ignore - compiled_graph.astream = patch_async_with_flow_execution_span # type: ignore + patch_with_execution_span( + compiled_graph, + make_span=lambda: AgentSpecFlowExecutionSpan( + name=f"FlowExecution[{flow.name}]", flow=flow + ), + make_start_event=lambda inputs: AgentSpecFlowExecutionStart(flow=flow, inputs=inputs), + make_end_event=lambda result: AgentSpecFlowExecutionEnd( + flow=flow, + outputs=result.get("outputs", {}), + branch_selected=result.get("node_execution_details", {}).get("branch", ""), + ), + ) return compiled_graph def _node_convert_to_langgraph( @@ -762,9 +713,15 @@ def _agent_node_convert_to_langgraph( config: RunnableConfig, middleware: List[Any], ) -> "NodeExecutor": + from pyagentspec.adapters.langgraph._managerworkers_node import ManagerWorkersNodeExecutor from pyagentspec.adapters.langgraph._node_execution import AgentNodeExecutor - return AgentNodeExecutor( + executor_class = ( + ManagerWorkersNodeExecutor + if isinstance(agent_node.agent, AgentSpecManagerWorkers) + else AgentNodeExecutor + ) + return executor_class( agent_node, tool_registry=tool_registry, converted_components=converted_components, @@ -1123,6 +1080,108 @@ def _swarm_convert_to_langgraph( default_active_agent=agentspec_component.first_agent.name, ).compile(name=agentspec_component.name, checkpointer=checkpointer) + def _manager_workers_convert_to_langgraph( + self, + mw: AgentSpecManagerWorkers, + tool_registry: Dict[str, "LangGraphTool"], + converted_components: Dict[str, Any], + checkpointer: Optional[Checkpointer], + config: RunnableConfig, + middleware: List[Any], + ) -> CompiledStateGraph[Any, Any, Any]: + """Compile a ``ManagerWorkers`` into a hierarchical LangGraph. + + Topology:: + + ┌─ __delegate_to__w1 ─→ worker_1 ─┐ + START → manager ┤ ├→ manager (loop) + └─ __delegate_to__w2 ─→ worker_2 ─┘ + │ + └─ no tool_call ─→ END + + The manager is a react-agent holding one synthetic ``__delegate_to__`` + tool per worker. A conditional edge routes each delegation to its worker, + which runs in an isolated message context and answers with a ``ToolMessage`` + matched to the pending delegation id. Workers are converted recursively and + wired in as subgraph nodes, so ``astream_events`` still exposes the + parent/child boundary (``subgraph=True``) for tracing and streaming. + """ + if not isinstance(mw.group_manager, AgentSpecAgent): + # Delegation is routed off the manager's tool_calls, so the manager needs + # a chat-LLM; a Flow, Swarm or nested ManagerWorkers gives nothing to route on. + raise NotImplementedError( + f"ManagerWorkers.group_manager must be an Agent for LangGraph " + f"conversion; got {type(mw.group_manager).__name__}." + ) + + named_workers: List[Tuple[str, AgentSpecAgenticComponent]] = [ + (_safe_node_name(worker.name, fallback_id=worker.id), worker) for worker in mw.workers + ] + worker_node_names = [node_name for node_name, _ in named_workers] + if len(set(worker_node_names)) != len(worker_node_names): + raise ValueError( + "ManagerWorkers worker names collide after normalization: " + f"{worker_node_names}. Give each worker a unique name." + ) + + conversion_kwargs: Dict[str, Any] = { + "tool_registry": tool_registry, + "converted_components": converted_components, + "checkpointer": checkpointer, + "config": config, + "middleware": middleware, + } + + # The roster tells the LLM which delegation tool maps to which worker. + manager_agent = mw.group_manager + rendered_prompt = _append_workers_roster( + manager_agent.system_prompt, + [(node_name, worker.description or "") for node_name, worker in named_workers], + ) + + # The delegation tools execute inside the react loop: their Command(graph=PARENT) + # is how the call escapes the subgraph so the conditional edge below can route on it. + manager_graph = self._create_react_agent_with_given_info( + name=manager_agent.name, + system_prompt=rendered_prompt, + agent=manager_agent, + llm_config=manager_agent.llm_config, + tools=manager_agent.tools, + toolboxes=manager_agent.toolboxes, + inputs=manager_agent.inputs or [], + outputs=manager_agent.outputs or [], + additional_langgraph_tools=[ + _make_worker_delegation_tool(node_name) for node_name in worker_node_names + ], + **conversion_kwargs, + ) + + # Manager and workers all go in as compiled subgraph nodes, which is what + # makes LangGraph stream them with ``subgraph=True``. + builder = StateGraph(langgraph_graph.MessagesState) + builder.add_node(_MANAGER_NODE_KEY, manager_graph) + for node_name, worker in named_workers: + worker_graph = self.convert(worker, **conversion_kwargs) + builder.add_node(node_name, _wrap_worker_for_subgraph(worker_graph, node_name)) + builder.add_edge(node_name, _MANAGER_NODE_KEY) + + builder.add_edge(langgraph_graph.START, _MANAGER_NODE_KEY) + path_map: Dict[Hashable, str] = {} + for node_name in worker_node_names: + path_map[node_name] = node_name + path_map[langgraph_graph.END] = langgraph_graph.END + builder.add_conditional_edges( + _MANAGER_NODE_KEY, + _make_manager_router(worker_node_names), + # The path map covers every worker plus END, so langgraph can validate + # the routing statically. + path_map, + ) + + compiled_graph = builder.compile(checkpointer=checkpointer, name=mw.name) + _patch_with_manager_workers_execution_span(compiled_graph, mw) + return compiled_graph + def _create_react_agent_with_given_info( self, *, @@ -1213,81 +1272,19 @@ def _create_react_agent_with_given_info( **create_agent_kwargs ) - # To enable flow execution traces monkey patch all the functions that invoke the compiled graph - - original_stream = compiled_graph.stream - - def patch_with_agent_execution_span(*args: Any, **kwargs: Any) -> Generator[Any, Any, Any]: - span_name = f"AgentExecution[{agent.name}]" - inputs = kwargs.get("input", {}) - if not isinstance(inputs, dict): - inputs = {} - with AgentSpecAgentExecutionSpan(name=span_name, agent=agent) as span: - span.add_event(AgentSpecAgentExecutionStart(agent=agent, inputs=inputs)) - original_result: dict[str, Any] | Any = {} - result: dict[str, Any] - # This is going to patch stream and astream, that return iterators and yield chunks - for chunk in original_stream(*args, **kwargs): - yield chunk - if isinstance(chunk, tuple): - original_result = chunk[1] - if not isinstance(original_result, dict): - result = {} - else: - result = original_result - outputs = extract_outputs_from_invoke_result(result, agent.outputs or []) - span.add_event(AgentSpecAgentExecutionEnd(agent=agent, outputs=outputs)) - - original_astream = compiled_graph.astream - - async def patch_async_with_agent_execution_span( - *args: Any, **kwargs: Any - ) -> AsyncGenerator[Any, Any]: - span_name = f"AgentExecution[{agent.name}]" - inputs = kwargs.get("input", {}) - if not isinstance(inputs, dict): - inputs = {} - span = AgentSpecAgentExecutionSpan(name=span_name, agent=agent) - try: - await span.start_async() - except NotImplementedError: - span.start() - try: - try: - await span.add_event_async( - AgentSpecAgentExecutionStart(agent=agent, inputs=inputs) - ) - except NotImplementedError: - span.add_event(AgentSpecAgentExecutionStart(agent=agent, inputs=inputs)) - original_result: dict[str, Any] | Any = {} - result: dict[str, Any] - # This is going to patch stream and astream, that return iterators and yield chunks - async for chunk in original_astream(*args, **kwargs): - yield chunk - if isinstance(chunk, tuple): - original_result = chunk[1] - if not isinstance(original_result, dict): - result = {} - else: - result = original_result - - outputs = extract_outputs_from_invoke_result(result, agent.outputs or []) - try: - await span.add_event_async( - AgentSpecAgentExecutionEnd(agent=agent, outputs=outputs) - ) - except NotImplementedError: - span.add_event(AgentSpecAgentExecutionEnd(agent=agent, outputs=outputs)) - finally: - try: - await span.end_async() - except NotImplementedError: - span.end() - - # Monkey patch invocation functions to inject tracing - # No need to patch `(a)invoke` as they internally use `(a)stream` - compiled_graph.stream = patch_with_agent_execution_span # type: ignore - compiled_graph.astream = patch_async_with_agent_execution_span # type: ignore + patch_with_execution_span( + compiled_graph, + make_span=lambda: AgentSpecAgentExecutionSpan( + name=f"AgentExecution[{agent.name}]", agent=agent + ), + make_start_event=lambda inputs: AgentSpecAgentExecutionStart( + agent=agent, inputs=inputs + ), + make_end_event=lambda result: AgentSpecAgentExecutionEnd( + agent=agent, + outputs=extract_outputs_from_invoke_result(result, agent.outputs or []), + ), + ) return compiled_graph def _agent_convert_to_langgraph( diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers.py b/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers.py new file mode 100644 index 00000000..9d2298fb --- /dev/null +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers.py @@ -0,0 +1,232 @@ +# Copyright © 2025, 2026 Oracle and/or its affiliates. +# +# This software is under the Apache License 2.0 +# (LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0) or Universal Permissive License +# (UPL) 1.0 (LICENSE-UPL or https://oss.oracle.com/licenses/upl), at your option. + +"""Helpers for compiling a ``ManagerWorkers`` into LangGraph, orchestrated by +``AgentSpecToLangGraphConverter._manager_workers_convert_to_langgraph``. + +The delegation protocol is visible on purpose: ``__delegate_to__`` calls +stream like any other tool call. Consumers that would rather not render them can +filter on :func:`is_delegation_tool_name`. +""" + +import re +from typing import Annotated, Any, Dict, Iterable, List, Tuple + +from pyagentspec.adapters.langgraph._execution_span import patch_with_execution_span +from pyagentspec.adapters.langgraph._types import ( + CompiledStateGraph, + RunnableConfig, + langgraph_graph, +) +from pyagentspec.managerworkers import ManagerWorkers as AgentSpecManagerWorkers +from pyagentspec.tracing.events import ( + ManagerWorkersExecutionEnd as AgentSpecManagerWorkersExecutionEnd, +) +from pyagentspec.tracing.events import ( + ManagerWorkersExecutionStart as AgentSpecManagerWorkersExecutionStart, +) +from pyagentspec.tracing.spans import ( + ManagerWorkersExecutionSpan as AgentSpecManagerWorkersExecutionSpan, +) + +# Cannot collide with a worker node name: _normalize_identifier strips leading and +# trailing underscores, so no normalized name ever starts with one. +_MANAGER_NODE_KEY = "__manager__" + +#: Prefix of the synthetic ``__delegate_to__`` tool names the manager's LLM +#: uses to address a worker. The dunder prefix, like the delegation keys below, keeps +#: it from colliding with a real tool named ``delegate_to_``. Re-exported +#: from ``pyagentspec.adapters.langgraph`` so consumers can recognize the protocol. +DELEGATE_TOOL_PREFIX = "__delegate_to__" + +# Keys of the per-delegation ``Send`` payload: the task to run, and the tool_call_id +# the worker's reply must answer. Routing per delegation (instead of off shared state) +# lets one manager turn delegate to several workers at once. +_DELEGATE_TASK_KEY = "__delegate_task__" +_DELEGATE_CALL_ID_KEY = "__delegate_tool_call_id__" + +_WHITESPACE_RE = re.compile(r"\s+") + + +def is_delegation_tool_name(name: Any) -> bool: + """True for the synthetic ``__delegate_to__`` tool names a manager emits.""" + return isinstance(name, str) and name.startswith(DELEGATE_TOOL_PREFIX) + + +def _normalize_identifier(s: str) -> str: + """Lowercase, collapse non-alphanumerics to underscores, strip leading/trailing ones.""" + return re.sub(r"[^a-z0-9]+", "_", s.lower()).strip("_") + + +def _safe_node_name(name: str, fallback_id: str) -> str: + """Normalize a worker name into a LangGraph node identifier. + + The LLM has to emit ``__delegate_to__`` reliably as a tool name, so node + names stay ASCII identifiers. Falls back to the normalized component id when the + name slugifies to nothing. + """ + return _normalize_identifier(name) or _normalize_identifier(fallback_id) or "worker" + + +def _messages_of(state: Any) -> List[Any]: + """Read ``messages`` off a state, which langgraph injects as a dict or an object.""" + if isinstance(state, dict): + return list(state.get("messages") or []) + return list(getattr(state, "messages", []) or []) + + +def _append_workers_roster(system_prompt: str, entries: List[Tuple[str, str]]) -> str: + """Append an ``Available workers:`` block listing ``- : ``. + + Descriptions are flattened to one line each, since the LLM routes off the block's + one-line-per-worker shape. + """ + if not entries: + return system_prompt + lines = [ + f"- {name}: {_WHITESPACE_RE.sub(' ', description).strip()}" for name, description in entries + ] + roster = "Available workers:\n" + "\n".join(lines) + return f"{system_prompt}\n\n{roster}" if system_prompt else roster + + +def _make_worker_delegation_tool(worker_node_name: str) -> Any: + """Build the ``__delegate_to__`` tool the manager's LLM emits to route to + a worker. + + Executing the tool is only how the call escapes the react subgraph: its body + surfaces the subgraph messages to the parent with ``Command(graph=PARENT)`` and no + ``goto``. Routing stays in the edge built by :func:`_make_manager_router`; a ``goto`` here + would collapse several same-turn delegations into one parent Command and leave the + other ``tool_call_id``s unanswered. + """ + from langchain_core.tools import InjectedToolCallId, tool + from langgraph.prebuilt import InjectedState + from langgraph.types import Command + + tool_name = f"{DELEGATE_TOOL_PREFIX}{worker_node_name}" + description = ( + f"Delegate a task to the {worker_node_name} worker and receive its response. " + f"Use this when the task fits the worker's described capability." + ) + + @tool(tool_name, description=description) + def _delegate( + task: str, + state: Annotated[Any, InjectedState], + tool_call_id: Annotated[str, InjectedToolCallId], + ) -> Any: + # task and tool_call_id are declared for the LLM-facing schema; the routing + # edge recovers both off the surfaced AIMessage's tool_calls. The + # add_messages reducer dedupes by id, so re-surfacing messages is a no-op. + del task, tool_call_id + return Command(graph=Command.PARENT, update={"messages": _messages_of(state)}) + + return _delegate + + +def _make_manager_router(worker_node_names: Iterable[str]) -> Any: + """Build the conditional edge routing the parent graph off the manager's last + AIMessage: one ``Send`` per ``__delegate_to__`` tool call, or ``END`` + when it emitted none. + + Every delegation gets its own ``Send``, so each tool_call_id is answered + independently; an unanswered one breaks the manager's next-turn tool-call/result + sequence. Plain tool calls already ran inside the react loop; that includes a + real tool whose name merely starts with the prefix, which is why a suffix that + is not a worker node is not routed. + """ + known_workers = frozenset(worker_node_names) + + def _route_manager_to_worker_or_end(state: Dict[str, Any]) -> Any: + from langgraph.types import Send + + messages = state.get("messages") or [] + last = messages[-1] if messages else None + sends = [] + for tool_call in getattr(last, "tool_calls", None) or []: + name = tool_call.get("name") + if not is_delegation_tool_name(name): + continue + worker_node_name = name[len(DELEGATE_TOOL_PREFIX) :] + if worker_node_name not in known_workers: + continue + args = tool_call.get("args") or {} + sends.append( + Send( + worker_node_name, + { + _DELEGATE_TASK_KEY: args.get("task") or "", + _DELEGATE_CALL_ID_KEY: tool_call.get("id") or "", + }, + ) + ) + return sends or langgraph_graph.END + + return _route_manager_to_worker_or_end + + +def _worker_input(state: Dict[str, Any]) -> Dict[str, Any]: + """The single-message context a worker run starts from.""" + from langchain_core.messages import HumanMessage + + return {"messages": [HumanMessage(content=state.get(_DELEGATE_TASK_KEY) or "")]} + + +def _worker_reply(state: Dict[str, Any], result: Any) -> Dict[str, Any]: + """The worker's last message, as a ToolMessage answering this delegation.""" + from langchain_core.messages import ToolMessage + + messages = result.get("messages") if isinstance(result, dict) else None + content = (getattr(messages[-1], "content", "") if messages else "") or "" + return { + "messages": [ + ToolMessage(content=content, tool_call_id=state.get(_DELEGATE_CALL_ID_KEY) or "") + ] + } + + +def _wrap_worker_for_subgraph( + worker_graph: CompiledStateGraph[Any, Any, Any], + worker_node_name: str, +) -> Any: + """Wrap a worker subgraph as a node of the ManagerWorkers parent graph. + + Hierarchical rather than shared-state like a Swarm: each run is handed only the + manager's chosen task, and the worker's answer comes back as a ToolMessage so the + manager's react loop sees a well-formed tool response on its next turn. The worker + receives this node's ambient run config explicitly, which streams its token events + under the worker node's checkpoint namespace. Explicit propagation is necessary on + Python 3.10, where LangChain cannot preserve the callback context across tasks. + """ + from pyagentspec.adapters.langgraph._types import RunnableLambda + + def run(state: Dict[str, Any], config: RunnableConfig) -> Dict[str, Any]: + return _worker_reply(state, worker_graph.invoke(_worker_input(state), config=config)) + + async def arun(state: Dict[str, Any], config: RunnableConfig) -> Dict[str, Any]: + return _worker_reply(state, await worker_graph.ainvoke(_worker_input(state), config=config)) + + return RunnableLambda(func=run, afunc=arun, name=f"worker:{worker_node_name}") + + +def _patch_with_manager_workers_execution_span( + compiled_graph: CompiledStateGraph[Any, Any, Any], + mw: AgentSpecManagerWorkers, +) -> None: + """Wrap ``stream``/``astream`` so each run emits a ``ManagerWorkersExecutionSpan``.""" + patch_with_execution_span( + compiled_graph, + make_span=lambda: AgentSpecManagerWorkersExecutionSpan( + name=f"ManagerWorkersExecution[{mw.name}]", managerworkers=mw + ), + make_start_event=lambda inputs: AgentSpecManagerWorkersExecutionStart( + managerworkers=mw, inputs=inputs + ), + make_end_event=lambda result: AgentSpecManagerWorkersExecutionEnd( + managerworkers=mw, outputs={"messages": result.get("messages", [])} + ), + ) diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers_node.py b/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers_node.py new file mode 100644 index 00000000..74a0e700 --- /dev/null +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/_managerworkers_node.py @@ -0,0 +1,118 @@ +# Copyright © 2026 Oracle and/or its affiliates. +# +# This software is under the Apache License 2.0 +# (LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0) or Universal Permissive License +# (UPL) 1.0 (LICENSE-UPL or https://oss.oracle.com/licenses/upl), at your option. + +"""Runs a ``ManagerWorkers`` as a flow step. + +``AgentSpecToLangGraphConverter._agent_node_convert_to_langgraph`` selects +`ManagerWorkersNodeExecutor` when the node's agent is a ``ManagerWorkers``, +so `AgentNodeExecutor` keeps the plain-Agent behavior only. +""" + +from typing import Any, Dict, List, Optional, Tuple + +from pyagentspec.adapters._utils import render_template +from pyagentspec.adapters.langgraph._node_execution import AgentNodeExecutor +from pyagentspec.adapters.langgraph._types import ( + Checkpointer, + CompiledStateGraph, + ExecuteOutput, + LangGraphTool, + Messages, + NodeExecutionDetails, + RunnableConfig, +) +from pyagentspec.agent import Agent as AgentSpecAgent +from pyagentspec.flows.nodes import AgentNode as AgentSpecAgentNode +from pyagentspec.managerworkers import ManagerWorkers as AgentSpecManagerWorkers + + +class ManagerWorkersNodeExecutor(AgentNodeExecutor): + """Executes an ``AgentNode`` whose agent is a ``ManagerWorkers``. + + The hierarchical graph runs over ``MessagesState``, which can carry neither + structured inputs inward nor a ``structured_response`` outward. Inputs are + therefore rendered into the group-manager's system prompt before compiling, and + the manager's final message is the node's single string output. + """ + + def __init__( + self, + node: AgentSpecAgentNode, + tool_registry: Dict[str, "LangGraphTool"], + converted_components: Dict[str, Any], + checkpointer: Optional[Checkpointer], + config: RunnableConfig, + middleware: Optional[List[Any]] = None, + ) -> None: + super().__init__( + node, tool_registry, converted_components, checkpointer, config, middleware + ) + if not isinstance(node.agent, AgentSpecManagerWorkers): + raise TypeError( + "ManagerWorkersNodeExecutor requires an AgentNode holding a ManagerWorkers" + ) + self._manager_workers: AgentSpecManagerWorkers = node.agent + # Anything but a single string output cannot be honored (see class docstring); + # raising here fails at conversion time rather than mid-run. + outputs = node.outputs or [] + if outputs and (len(outputs) != 1 or outputs[0].type != "string"): + raise NotImplementedError( + "A ManagerWorkers flow step supports a single string output; " + f"node `{node.name}` declares {[o.title for o in outputs]}." + ) + + def _create_manager_workers_with_given_input_values( + self, inputs: Dict[str, Any] + ) -> CompiledStateGraph[Any, Any]: + """Compile the ``ManagerWorkers`` with the node inputs rendered into the + group-manager's ``system_prompt`` and the satisfied ports dropped. + + Cached by rendered prompt, the same key + :meth:`AgentNodeExecutor._create_react_agent_with_given_input_values` uses. + Calling the private converter entry point is deliberate, mirroring the react + path: the public ``convert`` caches by component id, which would collapse the + differently-rendered copies (all sharing the original's id) into one graph. + """ + from pyagentspec.adapters.langgraph._langgraphconverter import AgentSpecToLangGraphConverter + + converter = AgentSpecToLangGraphConverter() + component = self._manager_workers + entry_agent = component.group_manager + if not isinstance(entry_agent, AgentSpecAgent): + # Nothing to render or cache; the converter owns the error for this case. + return converter._manager_workers_convert_to_langgraph( + component, **self._conversion_kwargs() + ) + + system_prompt = render_template(entry_agent.system_prompt, inputs) + if system_prompt not in self._agents_cache: + rendered = component.model_copy( + update={ + "group_manager": entry_agent.model_copy( + update={"system_prompt": system_prompt, "inputs": []} + ), + "inputs": [], + } + ) + self._agents_cache[system_prompt] = converter._manager_workers_convert_to_langgraph( + rendered, **self._conversion_kwargs() + ) + return self._agents_cache[system_prompt] + + def _prepare_agent_and_inputs( + self, inputs: Dict[str, Any], messages: Messages + ) -> Tuple[CompiledStateGraph[Any, Any], Dict[str, Any]]: + # Inputs were baked into the group-manager's prompt, so this graph runs on + # messages alone rather than the react-agent's remaining_steps state. + graph = self._create_manager_workers_with_given_input_values(inputs) + return graph, {"messages": self._with_driving_message(messages)} + + def _format_agent_result(self, result: Dict[str, Any]) -> ExecuteOutput: + node_outputs = self.node.outputs + if not node_outputs: + return super()._format_agent_result(result) + # __init__ already rejected any shape but a single string output. + return {node_outputs[0].title: result["messages"][-1].content}, NodeExecutionDetails() diff --git a/pyagentspec/src/pyagentspec/adapters/langgraph/_node_execution.py b/pyagentspec/src/pyagentspec/adapters/langgraph/_node_execution.py index f1999aef..a1398df9 100644 --- a/pyagentspec/src/pyagentspec/adapters/langgraph/_node_execution.py +++ b/pyagentspec/src/pyagentspec/adapters/langgraph/_node_execution.py @@ -499,6 +499,16 @@ def __init__( self._middleware: List[Any] = list(middleware or []) self._agents_cache: Dict[str, CompiledStateGraph[Any, Any]] = {} + def _conversion_kwargs(self) -> Dict[str, Any]: + """The converter arguments every compile from this executor passes through.""" + return { + "tool_registry": self.tool_registry, + "converted_components": self.converted_components, + "checkpointer": self.checkpointer, + "config": self.config, + "middleware": self._middleware, + } + def _create_react_agent_with_given_input_values( self, inputs: Dict[str, Any] ) -> CompiledStateGraph[Any, Any]: @@ -521,27 +531,25 @@ def _create_react_agent_with_given_input_values( toolboxes=agentspec_component.toolboxes, inputs=agentspec_component.inputs or [], outputs=agentspec_component.outputs or [], - tool_registry=self.tool_registry, - converted_components=self.converted_components, - checkpointer=self.checkpointer, - config=self.config, - middleware=self._middleware, + **self._conversion_kwargs(), ) return self._agents_cache[system_prompt] + @staticmethod + def _with_driving_message(messages: Messages) -> Messages: + # LangGraph's agent expects at least one user message to drive execution. + # When an AgentNode is used with a templated system prompt and no messages are + # provided by the flow, the agent can crash. To avoid this, we artificially + # insert an empty user message when the message list is empty. + return messages if messages else cast(Messages, [{"role": "user", "content": ""}]) + def _prepare_agent_and_inputs( self, inputs: Dict[str, Any], messages: Messages ) -> Tuple[CompiledStateGraph[Any, Any], Dict[str, Any]]: agent = self._create_react_agent_with_given_input_values(inputs) - # LangGraph's agent expects at least one user message to drive execution. - # When an AgentNode is used with a templated system prompt and no messages are provided - # by the flow, the agent can crash. To avoid this, we artificially insert an empty - # user message when the message list is empty. - if not messages: - messages = cast(Messages, [{"role": "user", "content": ""}]) inputs |= { "remaining_steps": 20, # Get the right number of steps left - "messages": messages, + "messages": self._with_driving_message(messages), "structured_response": {}, } return agent, inputs diff --git a/pyagentspec/src/pyagentspec/managerworkers.py b/pyagentspec/src/pyagentspec/managerworkers.py index 26c2bdeb..11c9088c 100644 --- a/pyagentspec/src/pyagentspec/managerworkers.py +++ b/pyagentspec/src/pyagentspec/managerworkers.py @@ -13,6 +13,7 @@ from typing_extensions import Self from pyagentspec.agenticcomponent import AgenticComponent +from pyagentspec.property import Property, properties_have_same_type from pyagentspec.validation_helpers import model_validator_with_error_accumulation from pyagentspec.versioning import AgentSpecVersionEnum @@ -65,6 +66,61 @@ class ManagerWorkers(AgenticComponent): default=AgentSpecVersionEnum.v25_4_2, init=False, exclude=True ) + def _get_inferred_inputs(self) -> List[Property]: + # Per the language spec, the inputs of a ManagerWorkers are the inputs of its + # group manager (same name and type): the manager drives the conversation. + return ( + self.group_manager.inputs or [] + if getattr(self, "group_manager", None) + and self.min_agentspec_version >= AgentSpecVersionEnum.v26_4_0 + else [] + ) + + def _get_inferred_outputs(self) -> List[Property]: + # Symmetric with the inferred inputs: the group manager's outputs. + return ( + self.group_manager.outputs or [] + if getattr(self, "group_manager", None) + and self.min_agentspec_version >= AgentSpecVersionEnum.v26_4_0 + else [] + ) + + def _infer_min_agentspec_version_from_configuration(self) -> AgentSpecVersionEnum: + min_version = super()._infer_min_agentspec_version_from_configuration() + # ManagerWorkers I/O matching was introduced in 26.4.0. Omitted I/O + # inherits from the group manager; explicitly empty I/O is legacy-only. + manager = getattr(self, "group_manager", None) + inherits_manager_io = bool( + manager + and ( + ("inputs" not in self.model_fields_set and manager.inputs) + or ("outputs" not in self.model_fields_set and manager.outputs) + ) + ) + if inherits_manager_io or getattr(self, "inputs", []) or getattr(self, "outputs", []): + min_version = max(min_version, AgentSpecVersionEnum.v26_4_0) + return min_version + + def _infer_max_agentspec_version_from_configuration(self) -> AgentSpecVersionEnum: + max_version = super()._infer_max_agentspec_version_from_configuration() + # Before 26.4.0 a ManagerWorkers did not inherit its manager's I/O. + manager = getattr(self, "group_manager", None) + inherits_manager_io = bool( + manager + and ( + ("inputs" not in self.model_fields_set and manager.inputs) + or ("outputs" not in self.model_fields_set and manager.outputs) + ) + ) + if ( + manager + and (manager.inputs or manager.outputs) + and not inherits_manager_io + and not (self.inputs or self.outputs) + ): + max_version = min(max_version, AgentSpecVersionEnum.v26_3_0) + return max_version + @model_validator_with_error_accumulation def _validate_one_or_more_workers(self) -> Self: if len(self.workers) == 0: @@ -79,3 +135,28 @@ def _validate_group_manager_is_not_included_as_a_worker(self) -> Self: if any(self.group_manager is agent for agent in self.workers): raise ValueError("Group manager cannot be a worker.") return self + + @model_validator_with_error_accumulation + def _validate_ios_match_group_manager_ios(self) -> Self: + # Per the language spec, the I/Os of a ManagerWorkers must be the I/Os of its + # group manager, same name and type. The base ComponentWithIO validators + # already enforce matching titles; use the shared property helper here so + # nested JSON Schema types are compared correctly as well. + for kind, own_properties, manager_properties, explicitly_provided in ( + ("input", self.inputs or [], self.group_manager.inputs or [], "inputs"), + ("output", self.outputs or [], self.group_manager.outputs or [], "outputs"), + ): + if explicitly_provided not in self.model_fields_set: + continue + manager_property_by_title = {p.title: p for p in manager_properties} + for own_property in own_properties: + manager_property = manager_property_by_title.get(own_property.title) + if manager_property is not None and not properties_have_same_type( + own_property, manager_property + ): + raise ValueError( + f"The {kind}s of a `ManagerWorkers` must match the {kind}s of its " + f"group manager (same name and type), but {kind} " + f"`{own_property.title}` has a different type from the group manager." + ) + return self diff --git a/pyagentspec/tests/adapters/langgraph/test_managerworkers.py b/pyagentspec/tests/adapters/langgraph/test_managerworkers.py new file mode 100644 index 00000000..01a14389 --- /dev/null +++ b/pyagentspec/tests/adapters/langgraph/test_managerworkers.py @@ -0,0 +1,624 @@ +# Copyright © 2026 Oracle and/or its affiliates. +# +# This software is under the Apache License 2.0 +# (LICENSE-APACHE or http://www.apache.org/licenses/LICENSE-2.0) or Universal Permissive License +# (UPL) 1.0 (LICENSE-UPL or https://oss.oracle.com/licenses/upl), at your option. + +from typing import Any, List, Optional +from unittest.mock import patch + +import pytest + +from pyagentspec.agent import Agent +from pyagentspec.llms import OpenAiCompatibleConfig +from pyagentspec.managerworkers import ManagerWorkers +from pyagentspec.property import Property + + +def _agent( + name: str, + llm_name: str, + description: str = "", + system_prompt: str = ".", + outputs: Optional[List[Property]] = None, +) -> Agent: + return Agent( + name=name, + description=description, + system_prompt=system_prompt, + llm_config=OpenAiCompatibleConfig(name=llm_name, model_id="fake", url="null"), + outputs=outputs, + ) + + +def _fake_llm(*ai_responses: Any) -> Any: + """A FakeMessagesListChatModel subclassed under ChatOpenAI so the react-agent + treats it as an OpenAI-style chat model.""" + from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel + from langchain_openai import ChatOpenAI + + class _FakeModel(FakeMessagesListChatModel, ChatOpenAI): + pass + + return _FakeModel(responses=list(ai_responses)) + + +def _load_with_fake_llms(mw: Any, default: Any = None, **fakes_by_llm_name: Any) -> Any: + """Compile ``mw`` offline, answering each LLM config (keyed by ``llm_config.name``) + with a queued fake; ``default`` answers any config not named. + + ``bind_tools`` is stubbed to return the same fake because + ``FakeMessagesListChatModel`` inherits it from ``ChatOpenAI``, which calls OpenAI. + """ + from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel + from langgraph.checkpoint.memory import MemorySaver + + from pyagentspec.adapters.langgraph import AgentSpecLoader + from pyagentspec.adapters.langgraph._langgraphconverter import AgentSpecToLangGraphConverter + + def _dispatch(_self: Any, llm_config: Any, *args: Any, **kwargs: Any) -> Any: + fake = fakes_by_llm_name.get(llm_config.name, default) + if fake is None: + raise AssertionError(f"unexpected llm_config: {llm_config.name}") + return fake + + loader = AgentSpecLoader(tool_registry={}, checkpointer=MemorySaver()) + with patch.object( + AgentSpecToLangGraphConverter, + "_llm_convert_to_langgraph", + autospec=True, + side_effect=_dispatch, + ), patch.object( + FakeMessagesListChatModel, "bind_tools", new=lambda self_obj, *a, **kw: self_obj + ): + return loader.load_component(mw) + + +def test_safe_node_name_normalizes_and_falls_back() -> None: + from pyagentspec.adapters.langgraph._managerworkers import _safe_node_name + + assert _safe_node_name("Research Helper", "id-1") == "research_helper" + assert _safe_node_name("My-Worker!! v2", "id-1") == "my_worker_v2" + # Name slugifies to empty → normalized id; both empty → constant fallback. + assert _safe_node_name("!!!", "sub-1") == "sub_1" + assert _safe_node_name("", "") == "worker" + + +def test_append_workers_roster_renders_one_line_per_worker() -> None: + from pyagentspec.adapters.langgraph._managerworkers import _append_workers_roster + + out = _append_workers_roster( + "Coordinate the team.", + [("research_helper", "Handles research"), ("drafter", "Drafts text")], + ) + assert out == ( + "Coordinate the team.\n\n" + "Available workers:\n" + "- research_helper: Handles research\n" + "- drafter: Drafts text" + ) + # Multiline descriptions are flattened so the one-line-per-worker shape survives. + out = _append_workers_roster("", [("helper", "First line\nsecond line\n third line ")]) + assert out == "Available workers:\n- helper: First line second line third line" + + +def test_route_manager_to_worker_or_end_returns_end_when_no_delegation() -> None: + from langchain_core.messages import AIMessage + from langgraph.graph import END + + from pyagentspec.adapters.langgraph._managerworkers import _make_manager_router + + route = _make_manager_router(["drafter"]) + not_delegating = AIMessage(content="Done.", tool_calls=[]) + assert route({"messages": [not_delegating]}) == END + assert route({"messages": []}) == END + + +def test_route_manager_to_worker_or_end_fans_out_one_send_per_delegation() -> None: + from langchain_core.messages import AIMessage + from langgraph.types import Send + + from pyagentspec.adapters.langgraph._managerworkers import ( + _DELEGATE_CALL_ID_KEY, + _DELEGATE_TASK_KEY, + _make_manager_router, + ) + + route = _make_manager_router(["drafter", "research_helper"]) + msg = AIMessage( + content="", + tool_calls=[ + {"name": "some_other_tool", "args": {}, "id": "c0"}, + {"name": "__delegate_to__drafter", "args": {"task": "x"}, "id": "c1"}, + {"name": "__delegate_to__research_helper", "args": {"task": "y"}, "id": "c2"}, + ], + ) + sends = route({"messages": [msg]}) + # Every delegation gets its own Send carrying the task and the tool_call_id its + # reply must answer. The non-delegation tool call already ran inside the manager's + # react loop and is ignored by routing. + assert all(isinstance(s, Send) for s in sends) + assert [s.node for s in sends] == ["drafter", "research_helper"] + assert [s.arg[_DELEGATE_TASK_KEY] for s in sends] == ["x", "y"] + assert [s.arg[_DELEGATE_CALL_ID_KEY] for s in sends] == ["c1", "c2"] + + +def test_route_manager_ignores_prefixed_tool_whose_suffix_is_not_a_worker() -> None: + """A tool call that merely looks like a delegation must not be routed: its suffix + is not a worker node, so a Send would target a non-existing node. It already ran + as a plain tool inside the react loop.""" + from langchain_core.messages import AIMessage + from langgraph.graph import END + + from pyagentspec.adapters.langgraph._managerworkers import _make_manager_router + + route = _make_manager_router(["drafter"]) + msg = AIMessage( + content="", + tool_calls=[{"name": "__delegate_to__nobody", "args": {"task": "x"}, "id": "c1"}], + ) + assert route({"messages": [msg]}) == END + + +def test_manager_workers_compiles_to_hierarchical_graph_topology() -> None: + from langchain_core.messages import AIMessage + from langgraph.graph import START + + from pyagentspec.adapters.langgraph._managerworkers import _MANAGER_NODE_KEY + + mw = ManagerWorkers( + name="ResearchTeam", + group_manager=_agent("Coordinator", "manager_llm", system_prompt="Coordinate the team."), + workers=[ + _agent("Research Helper", "worker_a_llm", description="Handles research"), + _agent("Drafter", "worker_b_llm", description="Drafts text"), + ], + ) + + compiled = _load_with_fake_llms(mw, default=_fake_llm(AIMessage(content="Done."))) + + builder = compiled.builder + assert _MANAGER_NODE_KEY in builder.nodes + assert "research_helper" in builder.nodes + assert "drafter" in builder.nodes + + # START → manager; every worker → manager (loop). + edge_pairs = {(src, dst) for src, dst in builder.edges} + assert (START, _MANAGER_NODE_KEY) in edge_pairs + assert ("research_helper", _MANAGER_NODE_KEY) in edge_pairs + assert ("drafter", _MANAGER_NODE_KEY) in edge_pairs + + # Manager → worker is a conditional edge. + assert builder.branches.get(_MANAGER_NODE_KEY) + + +def test_manager_workers_registers_a_delegation_tool_per_worker() -> None: + from langchain_core.messages import AIMessage + + from pyagentspec.adapters.langgraph._managerworkers import _MANAGER_NODE_KEY + + mw = ManagerWorkers( + name="Team", + group_manager=_agent("Coordinator", "manager_llm"), + workers=[_agent("Research Helper", "worker_llm", description="Handles research tasks")], + ) + + compiled = _load_with_fake_llms(mw, default=_fake_llm(AIMessage(content="Done."))) + + # The delegation tool the roster advertises is registered on the manager + # react-agent's tools node, so the LLM has the matching contract. + manager_subgraph = compiled.builder.nodes[_MANAGER_NODE_KEY].runnable + tools_node = manager_subgraph.builder.nodes["tools"].runnable + assert "__delegate_to__research_helper" in tools_node.tools_by_name + + +def test_manager_workers_delegates_and_routes_back_with_tool_message() -> None: + """The manager delegates, the worker runs in an isolated message context and its + answer comes back as a ToolMessage matched to the pending tool_call_id, and the + manager's next turn terminates the graph.""" + from langchain_core.messages import AIMessage, HumanMessage + + mw = ManagerWorkers( + name="Team", + group_manager=_agent("Coordinator", "manager_llm", system_prompt="You coordinate."), + workers=[_agent("Research Helper", "worker_llm", description="Handles research")], + ) + + # Manager turn 1: delegate. Manager turn 2: final answer (no tool call → END). + manager_responses = [ + AIMessage( + content="", + tool_calls=[ + { + "name": "__delegate_to__research_helper", + "args": {"task": "Look up Saturn"}, + "id": "call_1", + } + ], + ), + AIMessage(content="The worker reports: Saturn has rings."), + ] + worker_responses = [AIMessage(content="Saturn has rings.")] + + compiled = _load_with_fake_llms( + mw, + manager_llm=_fake_llm(*manager_responses), + worker_llm=_fake_llm(*worker_responses), + ) + + # Sync invocation only: FakeMessagesListChatModel overrides ``_generate`` but not + # ``_agenerate``, so the async path would resolve to ``ChatOpenAI._agenerate`` + # and call OpenAI. + result = compiled.invoke( + {"messages": [HumanMessage(content="Tell me about Saturn.")]}, + {"configurable": {"thread_id": "mw-1"}}, + ) + messages = result["messages"] + + assert isinstance(messages[-1], AIMessage) + assert "Saturn has rings" in messages[-1].content + tool_msgs = [m for m in messages if type(m).__name__ == "ToolMessage"] + assert tool_msgs and tool_msgs[0].tool_call_id == "call_1" + assert "Saturn has rings" in tool_msgs[0].content + + +def test_manager_workers_answers_every_delegation_in_a_single_turn() -> None: + """When one manager turn emits several delegations, each must be answered by its + own ToolMessage matched to the originating tool_call_id; an unanswered one is an + invalid tool-call/result sequence the manager would hallucinate around.""" + from langchain_core.messages import AIMessage, HumanMessage + + mw = ManagerWorkers( + name="Team", + group_manager=_agent("Coordinator", "manager_llm", system_prompt="You coordinate."), + workers=[_agent("Sub Agent", "worker_llm", description="Writes poems")], + ) + + # Turn 1: three delegations to the same worker in one AIMessage. Turn 2: terminate. + manager_responses = [ + AIMessage( + content="", + tool_calls=[ + { + "name": "__delegate_to__sub_agent", + "args": {"task": "Spanish poem"}, + "id": "call_1", + }, + { + "name": "__delegate_to__sub_agent", + "args": {"task": "French poem"}, + "id": "call_2", + }, + { + "name": "__delegate_to__sub_agent", + "args": {"task": "German poem"}, + "id": "call_3", + }, + ], + ), + AIMessage(content="Here are your three poems."), + ] + worker_responses = [AIMessage(content=f"poem #{i}") for i in range(1, 6)] + + compiled = _load_with_fake_llms( + mw, + manager_llm=_fake_llm(*manager_responses), + worker_llm=_fake_llm(*worker_responses), + ) + + result = compiled.invoke( + {"messages": [HumanMessage(content="Write 3 poems via sub-agents.")]}, + {"configurable": {"thread_id": "mw-multi"}}, + ) + messages = result["messages"] + + tool_msgs = [m for m in messages if type(m).__name__ == "ToolMessage"] + answered = sorted(m.tool_call_id for m in tool_msgs) + assert answered == ["call_1", "call_2", "call_3"] + assert all(m.content.startswith("poem #") for m in tool_msgs) + + +def test_nested_manager_workers_compiles_recursively() -> None: + """A worker that is itself a ManagerWorkers compiles through the same dispatch and + is wired in as a subgraph node of the outer parent graph.""" + from langchain_core.messages import AIMessage + + inner_mw = ManagerWorkers( + name="Inner", + group_manager=_agent("InnerManager", "inner_llm", system_prompt="Manage leaves."), + workers=[_agent("Leaf", "leaf_llm", description="Leaf task")], + ) + outer_mw = ManagerWorkers( + name="Outer", + group_manager=_agent("OuterManager", "outer_llm", system_prompt="Manage subteams."), + workers=[inner_mw], + ) + + compiled = _load_with_fake_llms(outer_mw, default=_fake_llm(AIMessage(content="Done."))) + + assert "inner" in compiled.builder.nodes + + +def test_rejects_non_agent_group_manager() -> None: + """A nested ManagerWorkers as group_manager is valid per the pyagentspec + validators, but the adapter needs a chat-LLM emitting tool_calls to route on.""" + from langgraph.checkpoint.memory import MemorySaver + + from pyagentspec.adapters.langgraph import AgentSpecLoader + + inner_mw = ManagerWorkers( + name="Inner", + group_manager=_agent("Inner", "i"), + workers=[_agent("Leaf", "l")], + ) + outer_mw = ManagerWorkers( + name="Outer", + group_manager=inner_mw, + workers=[_agent("Other", "o")], + ) + + loader = AgentSpecLoader(tool_registry={}, checkpointer=MemorySaver()) + with pytest.raises(NotImplementedError, match="group_manager must be an Agent"): + loader.load_component(outer_mw) + + +def test_workers_with_name_slug_collision_are_rejected() -> None: + """Two workers whose names normalize to the same node identifier would silently + overwrite each other in the parent graph; raise at load time instead.""" + from langgraph.checkpoint.memory import MemorySaver + + from pyagentspec.adapters.langgraph import AgentSpecLoader + + # Both worker names normalize to "helper_a". + mw = ManagerWorkers( + name="T", + group_manager=_agent("M", "m"), + workers=[_agent("Helper A", "a"), _agent("helper-a", "b")], + ) + + loader = AgentSpecLoader(tool_registry={}, checkpointer=MemorySaver()) + with pytest.raises(ValueError, match="collide after normalization"): + loader.load_component(mw) + + +def test_worker_events_stream_natively_namespaced_under_worker_node() -> None: + """Regression: a worker's token events must stream under the worker node's + checkpoint namespace so a consumer can attribute them to the sub-agent. The + wrapper must inherit the ambient run config; a fresh thread_id would detach the + worker into an unattributable top-level ``agent:`` run.""" + import asyncio + + from langchain_core.language_models.fake_chat_models import GenericFakeChatModel + from langchain_core.messages import AIMessage, HumanMessage + from langchain_core.runnables import RunnableConfig + from langgraph.graph import END, START, MessagesState, StateGraph + + from pyagentspec.adapters.langgraph._managerworkers import _wrap_worker_for_subgraph + + # A minimal worker graph that streams some content. + wmodel = GenericFakeChatModel(messages=iter([AIMessage(content="Saturn has rings")] * 9)) + wb = StateGraph(MessagesState) + + async def _wagent(state: Any, config: RunnableConfig) -> Any: + return {"messages": [await wmodel.ainvoke(state["messages"], config=config)]} + + wb.add_node("agent", _wagent) + wb.add_edge(START, "agent") + wb.add_edge("agent", END) + worker_graph = wb.compile() + + # Parent: a plain manager node emits the delegate tool call, then routes to the + # wrapped worker node. + pb = StateGraph(MessagesState) + + def _manager(state: Any) -> Any: + return { + "messages": [ + AIMessage( + content="", + tool_calls=[ + { + "name": "__delegate_to__research_helper", + "args": {"task": "Saturn"}, + "id": "c1", + } + ], + ) + ] + } + + pb.add_node("__manager__", _manager) + pb.add_node("research_helper", _wrap_worker_for_subgraph(worker_graph, "research_helper")) + pb.add_edge(START, "__manager__") + pb.add_edge("__manager__", "research_helper") + pb.add_edge("research_helper", END) + parent = pb.compile() + + async def _collect() -> Any: + namespaces = [] + async for ev in parent.astream_events( + {"messages": [HumanMessage(content="hi")]}, + {"configurable": {"thread_id": "t"}}, + version="v2", + ): + if ev["event"] == "on_chat_model_stream": + ns = (ev.get("metadata") or {}).get("langgraph_checkpoint_ns", "") + namespaces.append(ns) + return namespaces + + namespaces = asyncio.run(_collect()) + assert namespaces, "expected the worker to emit token-stream events" + assert all(ns.startswith("research_helper:") for ns in namespaces), namespaces + + +def test_is_delegation_tool_name_matches_only_the_synthetic_prefix() -> None: + from pyagentspec.adapters.langgraph._managerworkers import ( + DELEGATE_TOOL_PREFIX, + is_delegation_tool_name, + ) + + assert DELEGATE_TOOL_PREFIX == "__delegate_to__" + assert is_delegation_tool_name("__delegate_to__research_helper") + assert not is_delegation_tool_name("get_weather") + # A real tool plausibly named delegate_to_ is not a delegation. + assert not is_delegation_tool_name("delegate_to_someone") + # Nor is one merely *containing* the prefix mid-name. + assert not is_delegation_tool_name("please__delegate_to__someone") + assert not is_delegation_tool_name(None) + assert not is_delegation_tool_name(123) + + +def test_manager_workers_leaves_astream_events_unwrapped() -> None: + """The delegation protocol is deliberately visible: only stream/astream are + patched (for the ManagerWorkersExecutionSpan); nothing wraps ``astream_events`` + to scrub the delegation tool calls.""" + mw = ManagerWorkers( + name="Team", + group_manager=_agent("Coordinator", "manager_llm"), + workers=[_agent("Research Helper", "worker_llm")], + ) + + compiled = _load_with_fake_llms(mw, default=_fake_llm()) + + assert "stream" in compiled.__dict__ and "astream" in compiled.__dict__ + assert "astream_events" not in compiled.__dict__ + + +def test_managerworkers_infers_inputs_from_group_manager_prompt() -> None: + """A ManagerWorkers exposes the group manager's prompt placeholders as inputs, so + a flow AgentNode wrapping it declares input ports a DataFlowEdge can resolve.""" + manager = _agent( + "manager", + "manager_llm", + system_prompt="Translate the following to Arabic:\n\n{{joke}}\n\nMake {{count}} variants.", + ) + worker = _agent("worker", "worker_llm", system_prompt="You translate.") + mw = ManagerWorkers(name="mw", group_manager=manager, workers=[worker]) + + assert sorted(p.title for p in (mw.inputs or [])) == ["count", "joke"] + + +def test_managerworkers_infers_outputs_from_group_manager() -> None: + """Symmetric with inputs: a ManagerWorkers exposes the group manager's outputs.""" + from pyagentspec.property import StringProperty + + manager = _agent( + "manager", + "manager_llm", + system_prompt="Answer the question.", + outputs=[StringProperty(title="answer")], + ) + worker = _agent("worker", "worker_llm", system_prompt="You help.") + mw = ManagerWorkers(name="mw", group_manager=manager, workers=[worker]) + + assert [p.title for p in (mw.outputs or [])] == ["answer"] + + +def _flow_with_manager_workers_step(outputs: List[Property]) -> Any: + """A start → AgentNode(ManagerWorkers) → end flow whose end node exposes + ``outputs``, with the data edges resolving the manager's ``joke`` input and + every output.""" + from pyagentspec.flows.edges import ControlFlowEdge, DataFlowEdge + from pyagentspec.flows.flow import Flow + from pyagentspec.flows.nodes import AgentNode, EndNode, StartNode + from pyagentspec.property import StringProperty + + joke = StringProperty(title="joke") + manager = _agent( + "manager", + "manager_llm", + system_prompt="Translate the following to Arabic:\n\n{{joke}}", + outputs=outputs, + ) + worker = _agent("worker", "worker_llm", system_prompt="You translate.") + mw = ManagerWorkers(name="translator", group_manager=manager, workers=[worker]) + assert [p.title for p in (mw.inputs or [])] == ["joke"] + + manager_node = AgentNode(name="manager_node", agent=mw) + start_node = StartNode(name="start", inputs=[joke]) + end_node = EndNode(name="end", outputs=outputs) + return Flow( + name="flow", + start_node=start_node, + nodes=[start_node, manager_node, end_node], + control_flow_connections=[ + ControlFlowEdge(name="start_to_node", from_node=start_node, to_node=manager_node), + ControlFlowEdge(name="node_to_end", from_node=manager_node, to_node=end_node), + ], + data_flow_connections=[ + DataFlowEdge( + name="joke_edge", + source_node=start_node, + source_output=joke.title, + destination_node=manager_node, + destination_input=joke.title, + ), + ] + + [ + DataFlowEdge( + name=f"{output.title}_edge", + source_node=manager_node, + source_output=output.title, + destination_node=end_node, + destination_input=output.title, + ) + for output in outputs + ], + outputs=outputs, + ) + + +def test_managerworkers_runs_as_a_flow_step_with_data_edge_inputs() -> None: + """A ManagerWorkers flow step loads with its data edge resolved and executes: + loading proves the node exposes the ``joke`` input the edge targets, running + proves the manager's answer comes back as the node's single string output.""" + from langchain_core.language_models.fake_chat_models import FakeMessagesListChatModel + from langchain_core.messages import AIMessage + from langgraph.checkpoint.memory import MemorySaver + + from pyagentspec.adapters.langgraph import AgentSpecLoader + from pyagentspec.adapters.langgraph._langgraphconverter import AgentSpecToLangGraphConverter + from pyagentspec.property import StringProperty + + # The final message has no tool_calls → the manager routes to END without delegating. + fake_llm = _fake_llm(AIMessage(content="لماذا...")) + flow = _flow_with_manager_workers_step(outputs=[StringProperty(title="translated")]) + + loader = AgentSpecLoader(tool_registry={}, checkpointer=MemorySaver()) + with patch.object( + AgentSpecToLangGraphConverter, + "_llm_convert_to_langgraph", + autospec=True, + side_effect=lambda self_obj, llm_config, *a, **k: fake_llm, + ), patch.object( + FakeMessagesListChatModel, + "bind_tools", + new=lambda self_obj, *a, **k: self_obj, + ): + compiled = loader.load_component(flow) + result = compiled.invoke( + { + "inputs": {"joke": "Why did the car..."}, + "messages": [{"role": "user", "content": ""}], + }, + {"configurable": {"thread_id": "managerworkers-node"}}, + ) + + assert result["outputs"]["translated"] == "لماذا..." + + +def test_managerworkers_flow_step_with_unsupported_outputs_fails_at_conversion() -> None: + """A ManagerWorkers flow step supports a single string output only; any other + shape must be rejected when the flow is converted, not once the step runs.""" + from langgraph.checkpoint.memory import MemorySaver + + from pyagentspec.adapters.langgraph import AgentSpecLoader + from pyagentspec.property import StringProperty + + flow = _flow_with_manager_workers_step( + outputs=[StringProperty(title="translated"), StringProperty(title="notes")] + ) + + loader = AgentSpecLoader(tool_registry={}, checkpointer=MemorySaver()) + with pytest.raises(NotImplementedError, match="single string output"): + loader.load_component(flow) diff --git a/pyagentspec/tests/serialization/test_managerworkers.py b/pyagentspec/tests/serialization/test_managerworkers.py index fefb4271..dba17837 100644 --- a/pyagentspec/tests/serialization/test_managerworkers.py +++ b/pyagentspec/tests/serialization/test_managerworkers.py @@ -9,6 +9,7 @@ from pyagentspec.agent import Agent from pyagentspec.llms import VllmConfig from pyagentspec.managerworkers import ManagerWorkers +from pyagentspec.property import StringProperty from pyagentspec.serialization import AgentSpecDeserializer, AgentSpecSerializer from pyagentspec.versioning import AgentSpecVersionEnum @@ -109,3 +110,60 @@ def test_deserializing_managerworkers_with_unsupported_version_raises_error( with pytest.raises(ValueError, match="Invalid agentspec_version"): AgentSpecDeserializer().from_yaml(serialized_managerworkers) + + +def test_managerworkers_infers_manager_ios_in_current_version() -> None: + llm_config = VllmConfig(name="model", model_id="model_id", url="https://example.com") + manager = Agent( + name="manager", + llm_config=llm_config, + system_prompt="Manage the team.", + outputs=[StringProperty(title="answer")], + ) + worker = Agent(name="worker", llm_config=llm_config, system_prompt="Help the manager.") + manager_workers = ManagerWorkers( + name="team", + group_manager=manager, + workers=[worker], + outputs=[StringProperty(title="answer")], + ) + + assert manager_workers.outputs == manager.outputs + assert manager_workers.min_agentspec_version == AgentSpecVersionEnum.current_version + with pytest.raises(ValueError, match="Invalid agentspec_version"): + AgentSpecSerializer().to_dict( + manager_workers, agentspec_version=AgentSpecVersionEnum.v26_3_0 + ) + + +def test_managerworkers_without_ios_with_manager_ios_is_legacy_compatible() -> None: + llm_config = VllmConfig(name="model", model_id="model_id", url="https://example.com") + manager = Agent( + name="manager", + llm_config=llm_config, + system_prompt="Manage {{question}}.", + inputs=[StringProperty(title="question")], + outputs=[StringProperty(title="answer")], + ) + worker = Agent(name="worker", llm_config=llm_config, system_prompt="Help the manager.") + manager_workers = ManagerWorkers( + name="team", + group_manager=manager, + workers=[worker], + inputs=[], + outputs=[], + ) + + assert manager_workers.inputs == [] + assert manager_workers.outputs == [] + assert manager_workers.min_agentspec_version == AgentSpecVersionEnum.v25_4_2 + assert manager_workers.max_agentspec_version == AgentSpecVersionEnum.v26_3_0 + + serialized = AgentSpecSerializer().to_dict( + manager_workers, agentspec_version=AgentSpecVersionEnum.v26_3_0 + ) + deserialized = AgentSpecDeserializer().from_dict(serialized) + + assert isinstance(deserialized, ManagerWorkers) + assert deserialized.inputs == [] + assert deserialized.outputs == [] diff --git a/pyagentspec/tests/validation/test_agentic_patterns_validation.py b/pyagentspec/tests/validation/test_agentic_patterns_validation.py index 81a13af8..2eeab694 100644 --- a/pyagentspec/tests/validation/test_agentic_patterns_validation.py +++ b/pyagentspec/tests/validation/test_agentic_patterns_validation.py @@ -14,6 +14,7 @@ from pyagentspec.flows.nodes.startnode import StartNode from pyagentspec.llms import OpenAiConfig from pyagentspec.managerworkers import ManagerWorkers +from pyagentspec.property import FloatProperty, ListProperty, StringProperty from pyagentspec.swarm import Swarm @@ -79,6 +80,85 @@ def test_managerworkers_with_different_agentic_components_can_be_validated() -> ) +def test_managerworkers_with_ios_matching_the_group_manager_can_be_validated() -> None: + manager_agent = Agent( + name="manager_agent", + system_prompt="Answer about {{topic}}.", + llm_config=OpenAiConfig(name="default", model_id="test_model"), + outputs=[StringProperty(title="answer")], + ) + worker_agent = Agent( + name="worker_agent", + system_prompt="You help.", + llm_config=OpenAiConfig(name="default", model_id="test_model"), + ) + + # The I/Os of a ManagerWorkers must be the I/Os of its group manager (same name + # and type); redeclaring them explicitly is valid. + _ = ManagerWorkers( + name="managerworkers", + group_manager=manager_agent, + workers=[worker_agent], + inputs=[StringProperty(title="topic")], + outputs=[StringProperty(title="answer")], + ) + + +def test_managerworkers_with_ios_not_matching_the_group_manager_raises_errors() -> None: + manager_agent = Agent( + name="manager_agent", + system_prompt="Answer about {{topic}}.", + llm_config=OpenAiConfig(name="default", model_id="test_model"), + outputs=[StringProperty(title="answer")], + ) + worker_agent = Agent( + name="worker_agent", + system_prompt="You help.", + llm_config=OpenAiConfig(name="default", model_id="test_model"), + ) + + # Same title as a group manager input, but a different type. + with pytest.raises(ValueError, match="must match the inputs of its group manager"): + ManagerWorkers( + name="managerworkers", + group_manager=manager_agent, + workers=[worker_agent], + inputs=[FloatProperty(title="topic")], + ) + + # Same title as a group manager output, but a different type. + with pytest.raises(ValueError, match="must match the outputs of its group manager"): + ManagerWorkers( + name="managerworkers", + group_manager=manager_agent, + workers=[worker_agent], + outputs=[FloatProperty(title="answer")], + ) + + # The comparison must also inspect nested schema types, rather than only the + # top-level ``array`` type. + list_output_manager = manager_agent.model_copy( + update={"outputs": [ListProperty(title="answer", item_type=StringProperty())]} + ) + with pytest.raises(ValueError, match="must match the outputs of its group manager"): + ManagerWorkers( + name="managerworkers", + group_manager=list_output_manager, + workers=[worker_agent], + outputs=[ListProperty(title="answer", item_type=FloatProperty())], + ) + + # A title the group manager does not declare is rejected by the base + # ComponentWithIO validation. + with pytest.raises(ValueError, match="expected only properties with the titles"): + ManagerWorkers( + name="managerworkers", + group_manager=manager_agent, + workers=[worker_agent], + outputs=[StringProperty(title="answer"), StringProperty(title="extra")], + ) + + def test_swarm_with_empty_relationships_raises_errors() -> None: first_agent = Agent( name="first_agent", diff --git a/tsagentspec/src/agents/manager-workers.ts b/tsagentspec/src/agents/manager-workers.ts index 5e32e4e1..03aaea42 100644 --- a/tsagentspec/src/agents/manager-workers.ts +++ b/tsagentspec/src/agents/manager-workers.ts @@ -3,7 +3,7 @@ */ import { z } from "zod"; import { ComponentWithIOSchema } from "../component.js"; -import type { Property } from "../property.js"; +import { propertiesHaveSameType, type Property } from "../property.js"; // z.record(z.unknown()) is used instead of AgenticComponentUnion to break a circular // dependency (ManagerWorkers -> AgenticComponentUnion -> ManagerWorkers). Validation of @@ -18,6 +18,43 @@ export const ManagerWorkersSchema = ComponentWithIOSchema.extend({ export type ManagerWorkers = z.infer; +function getComponentProperties( + component: Record, + field: "inputs" | "outputs", +): Property[] { + const properties = component[field]; + return Array.isArray(properties) ? (properties as Property[]) : []; +} + +function validatePropertiesMatchManager( + properties: Property[], + managerProperties: Property[], + kind: "inputs" | "outputs", +): void { + const managerPropertiesByTitle = new Map( + managerProperties.map((property) => [property.title, property]), + ); + const propertiesByTitle = new Map( + properties.map((property) => [property.title, property]), + ); + if ( + propertiesByTitle.size !== properties.length || + propertiesByTitle.size !== managerPropertiesByTitle.size + ) { + throw new Error( + `The ${kind} of a ManagerWorkers must match the ${kind} of its group manager.`, + ); + } + for (const property of properties) { + const managerProperty = managerPropertiesByTitle.get(property.title); + if (!managerProperty || !propertiesHaveSameType(property, managerProperty)) { + throw new Error( + `The ${kind} of a ManagerWorkers must match the ${kind} of its group manager.`, + ); + } + } +} + export function createManagerWorkers(opts: { name: string; groupManager: Record; @@ -31,9 +68,19 @@ export function createManagerWorkers(opts: { if (opts.workers.some(w => w === opts.groupManager)) { throw new Error("Group manager cannot be a worker."); } + const managerInputs = getComponentProperties(opts.groupManager, "inputs"); + const managerOutputs = getComponentProperties(opts.groupManager, "outputs"); + if (opts.inputs !== undefined) { + validatePropertiesMatchManager(opts.inputs, managerInputs, "inputs"); + } + if (opts.outputs !== undefined) { + validatePropertiesMatchManager(opts.outputs, managerOutputs, "outputs"); + } return Object.freeze( ManagerWorkersSchema.parse({ ...opts, + inputs: opts.inputs ?? managerInputs, + outputs: opts.outputs ?? managerOutputs, componentType: "ManagerWorkers" as const, }), ); diff --git a/tsagentspec/src/versioning.ts b/tsagentspec/src/versioning.ts index 94df5c04..fa566920 100644 --- a/tsagentspec/src/versioning.ts +++ b/tsagentspec/src/versioning.ts @@ -10,14 +10,17 @@ export const AgentSpecVersion = { V25_4_1: "25.4.1", V25_4_2: "25.4.2", V26_1_0: "26.1.0", + V26_1_2: "26.1.2", V26_2_0: "26.2.0", + V26_3_0: "26.3.0", + V26_4_0: "26.4.0", } as const; export type AgentSpecVersion = (typeof AgentSpecVersion)[keyof typeof AgentSpecVersion]; /** The current (latest) agent spec version */ -export const CURRENT_VERSION: AgentSpecVersion = AgentSpecVersion.V26_2_0; +export const CURRENT_VERSION: AgentSpecVersion = AgentSpecVersion.V26_4_0; /** Field name for the agentspec version in serialized JSON/YAML */ export const AGENTSPEC_VERSION_FIELD_NAME = "agentspec_version"; diff --git a/tsagentspec/tests/agents/manager-workers.test.ts b/tsagentspec/tests/agents/manager-workers.test.ts index 3f58b895..7daac563 100644 --- a/tsagentspec/tests/agents/manager-workers.test.ts +++ b/tsagentspec/tests/agents/manager-workers.test.ts @@ -3,6 +3,8 @@ import { createManagerWorkers, createAgent, createOpenAiCompatibleConfig, + numberProperty, + stringProperty, } from "../../src/index.js"; function makeLlmConfig() { @@ -118,4 +120,42 @@ describe("ManagerWorkers", () => { }), ).toThrow("Group manager cannot be a worker."); }); + + it("should infer the group manager's inputs and outputs", () => { + const topic = stringProperty({ title: "topic" }); + const answer = stringProperty({ title: "answer" }); + const manager = createAgent({ + name: "manager", + llmConfig: makeLlmConfig(), + systemPrompt: "Manage the team.", + inputs: [topic], + outputs: [answer], + }); + const mw = createManagerWorkers({ + name: "test-mw", + groupManager: manager, + workers: [makeAgent("worker")], + }); + + expect(mw.inputs).toEqual([topic]); + expect(mw.outputs).toEqual([answer]); + }); + + it("should reject explicit I/O that differs from the group manager", () => { + const manager = createAgent({ + name: "manager", + llmConfig: makeLlmConfig(), + systemPrompt: "Manage the team.", + outputs: [stringProperty({ title: "answer" })], + }); + + expect(() => + createManagerWorkers({ + name: "test-mw", + groupManager: manager, + workers: [makeAgent("worker")], + outputs: [numberProperty({ title: "answer" })], + }), + ).toThrow("outputs of a ManagerWorkers must match"); + }); }); diff --git a/tsagentspec/tests/versioning.test.ts b/tsagentspec/tests/versioning.test.ts index d8b63f9c..06d13be4 100644 --- a/tsagentspec/tests/versioning.test.ts +++ b/tsagentspec/tests/versioning.test.ts @@ -17,12 +17,15 @@ describe("AgentSpecVersion", () => { expect(AgentSpecVersion.V25_4_1).toBe("25.4.1"); expect(AgentSpecVersion.V25_4_2).toBe("25.4.2"); expect(AgentSpecVersion.V26_1_0).toBe("26.1.0"); + expect(AgentSpecVersion.V26_1_2).toBe("26.1.2"); expect(AgentSpecVersion.V26_2_0).toBe("26.2.0"); + expect(AgentSpecVersion.V26_3_0).toBe("26.3.0"); + expect(AgentSpecVersion.V26_4_0).toBe("26.4.0"); }); it("should set CURRENT_VERSION to the latest version", () => { - expect(CURRENT_VERSION).toBe("26.2.0"); - expect(CURRENT_VERSION).toBe(AgentSpecVersion.V26_2_0); + expect(CURRENT_VERSION).toBe("26.4.0"); + expect(CURRENT_VERSION).toBe(AgentSpecVersion.V26_4_0); }); it("should define the version field name", () => {