diff --git a/docs/supported-integrations/deepagents.mdx b/docs/supported-integrations/deepagents.mdx index e7f3c5bb6..a89886687 100644 --- a/docs/supported-integrations/deepagents.mdx +++ b/docs/supported-integrations/deepagents.mdx @@ -94,7 +94,9 @@ when you need to capture Deep Agents skill or subagent marks. The integration works correctly when: - The Deep Agents run completes and prints a final response. -- The `deepagents-request` scope contains the top-level agent execution. +- The `deepagents-request` scope contains a semantic `main-agent` Agent scope. +- In-process delegated agents appear as nested Agent scopes with names derived + from callback metadata. Unnamed top-level runs use the `DeepAgent` fallback. - Skill, subagent, and human-in-the-loop marks appear when those features are exercised. ## Observability @@ -106,7 +108,9 @@ human-in-the-loop lifecycle events. It captures: - LangChain model and tool calls through NeMo Relay managed execution. -- LangGraph run scopes through callbacks. +- Semantic orchestrator and in-process subagent Agent scopes through the + Deep Agents callback handler. Internal LangGraph node runs do not create + additional Agent scopes in this specialized handler. - Human-in-the-loop interrupt and resume marks. - Configured skills and subagent summaries at agent-run start. - Automatic `skill.load` marks when a Deep Agents tool requests a complete @@ -114,8 +118,14 @@ It captures: - In-process dictionary-style subagents with the same NeMo Relay middleware, so their model and tool calls are captured when Deep Agents invokes them. -Remote graphs or processes still need NeMo Relay instrumentation in that graph -or process to capture their internal model and tool calls. +The automatically configured `general-purpose` subagent is identified from +Deep Agents callback metadata and appears as a nested Agent scope when invoked. +Deep Agents does not inherit arbitrary parent middleware into that generated +subagent, so its internal model and tool calls are not managed by Relay. +Precompiled local subagent runnables can expose a semantic run boundary through +the callback, but their internal model and tool calls require separate Relay +instrumentation. Remote graphs or processes likewise need Relay instrumentation +inside that graph or process to capture their internal calls. Refer to [Observability](/configure-plugins/observability/about) for details on exporting NeMo Relay observability data to third-party systems. diff --git a/python/nemo_relay/integrations/deepagents/callbacks.py b/python/nemo_relay/integrations/deepagents/callbacks.py index 28ae7fa8a..37344011e 100644 --- a/python/nemo_relay/integrations/deepagents/callbacks.py +++ b/python/nemo_relay/integrations/deepagents/callbacks.py @@ -7,6 +7,7 @@ from collections.abc import Mapping, Sequence from typing import Any +from uuid import UUID from nemo_relay.integrations.deepagents._events import emit_mark, event_base_name from nemo_relay.integrations.langgraph.callbacks import NemoRelayCallbackHandler as LangGraphNemoRelayCallbackHandler @@ -15,12 +16,70 @@ class NemoRelayDeepAgentsCallbackHandler(LangGraphNemoRelayCallbackHandler): - """Bridge Deep Agents LangGraph lifecycle events to NeMo Relay marks.""" + """Bridge semantic Deep Agents runs and LangGraph lifecycle events to NeMo Relay.""" def __init__(self, *args: Any, **kwargs: Any) -> None: super().__init__(*args, **kwargs) self._hitl_interrupts: set[_GraphEventKey] = set() + def on_chain_start( + self, + serialized: dict[str, Any], + inputs: dict[str, Any], + *, + run_id: UUID, + parent_run_id: UUID | None = None, + tags: list[str] | None = None, + metadata: dict[str, Any] | None = None, + **kwargs: Any, + ) -> Any: + """Push scopes for Deep Agents orchestrators and subagents, not graph nodes.""" + agent_name = self._semantic_agent_name(kwargs.get("name"), metadata) + if agent_name is None: + return None + + scope_metadata = dict(metadata or {}) + scope_metadata.update( + { + "integration": "deepagents", + "deepagents_agent_name": agent_name, + "deepagents_agent_role": "orchestrator" if parent_run_id is None else "subagent", + } + ) + scope_kwargs = dict(kwargs) + scope_kwargs["name"] = agent_name + return super().on_chain_start( + serialized, + inputs, + run_id=run_id, + parent_run_id=parent_run_id, + tags=tags, + metadata=scope_metadata, + **scope_kwargs, + ) + + @staticmethod + def _semantic_agent_name(name: Any, metadata: Mapping[str, Any] | None) -> str | None: + if metadata is None: + return None + + versions = metadata.get("lc_versions") + if not isinstance(versions, Mapping) or "deepagents" not in versions: + return None + + langgraph_node = metadata.get("langgraph_node") + if langgraph_node is not None and langgraph_node == name: + return None + + configured_name = metadata.get("lc_agent_name") + if isinstance(configured_name, str) and configured_name and name == configured_name: + return configured_name + + if metadata.get("ls_integration") == "deepagents": + return "DeepAgent" + + return None + def _emit_graph_mark(self, name: str, data: dict[str, Any]) -> None: key = self._graph_event_key(data) if name == "Graph Interrupt" and self._has_hitl_interrupt(data): diff --git a/python/tests/integrations/deepagents_tests/test_deepagents_integration.py b/python/tests/integrations/deepagents_tests/test_deepagents_integration.py index 324e0b533..80db5aaaa 100644 --- a/python/tests/integrations/deepagents_tests/test_deepagents_integration.py +++ b/python/tests/integrations/deepagents_tests/test_deepagents_integration.py @@ -459,6 +459,151 @@ def test_callback_handler_falls_back_for_non_hitl_interrupt( assert "deepagents_kind" not in _mark_metadata(marks[0]) +def test_callback_handler_creates_only_semantic_agent_scopes( + subscribed_events: list[nemo_relay.Event], + callback_handler: deepagents_integration.NemoRelayDeepAgentsCallbackHandler, +): + root_run_id = uuid4() + node_run_id = uuid4() + subagent_run_id = uuid4() + + with nemo_relay.scope.scope("request", nemo_relay.ScopeType.Agent): + callback_handler.on_chain_start( + {}, + {"messages": ["hello"]}, + run_id=root_run_id, + name="main-agent", + metadata={ + "ls_integration": "deepagents", + "lc_agent_name": "main-agent", + "lc_versions": {"deepagents": "0.7.4"}, + }, + ) + callback_handler.on_chain_start( + {}, + {"messages": ["hello"]}, + run_id=node_run_id, + parent_run_id=root_run_id, + name="main-agent", + metadata={ + "ls_integration": "deepagents", + "lc_agent_name": "main-agent", + "lc_versions": {"deepagents": "0.7.4"}, + "langgraph_node": "main-agent", + }, + ) + callback_handler.on_chain_start( + {}, + {"messages": ["review"]}, + run_id=subagent_run_id, + parent_run_id=node_run_id, + name="reviewer", + metadata={ + "ls_integration": "langchain_create_agent", + "lc_agent_name": "reviewer", + "lc_versions": {"deepagents": "0.7.4"}, + }, + ) + callback_handler.on_chain_end({"messages": ["done"]}, run_id=root_run_id) + callback_handler.on_chain_end({"messages": ["reviewed"]}, run_id=subagent_run_id) + callback_handler.on_chain_end({"messages": ["ignored"]}, run_id=node_run_id) + + nemo_relay.subscribers.flush() + scope_events = [event for event in subscribed_events if isinstance(event, nemo_relay.ScopeEvent)] + assert [(event.scope_category, event.name) for event in scope_events] == [ + ("start", "request"), + ("start", "main-agent"), + ("start", "reviewer"), + ("end", "reviewer"), + ("end", "main-agent"), + ("end", "request"), + ] + starts = {event.name: event for event in scope_events if event.scope_category == "start"} + assert starts["main-agent"].metadata == { + "ls_integration": "deepagents", + "lc_agent_name": "main-agent", + "lc_versions": {"deepagents": "0.7.4"}, + "integration": "deepagents", + "deepagents_agent_name": "main-agent", + "deepagents_agent_role": "orchestrator", + "langchain_run_id": str(root_run_id), + } + assert starts["reviewer"].metadata == { + "ls_integration": "langchain_create_agent", + "lc_agent_name": "reviewer", + "lc_versions": {"deepagents": "0.7.4"}, + "integration": "deepagents", + "deepagents_agent_name": "reviewer", + "deepagents_agent_role": "subagent", + "langchain_run_id": str(subagent_run_id), + } + assert starts["main-agent"].parent_uuid == starts["request"].uuid + assert starts["reviewer"].parent_uuid == starts["main-agent"].uuid + + +def test_callback_handler_uses_fallback_for_unnamed_agent( + subscribed_events: list[nemo_relay.Event], + callback_handler: deepagents_integration.NemoRelayDeepAgentsCallbackHandler, +): + run_id = uuid4() + + with nemo_relay.scope.scope("request", nemo_relay.ScopeType.Agent): + callback_handler.on_chain_start( + {}, + {"messages": ["hello"]}, + run_id=run_id, + name=None, + metadata={ + "ls_integration": "deepagents", + "lc_versions": {"deepagents": "0.7.4"}, + }, + ) + callback_handler.on_chain_end({"messages": ["done"]}, run_id=run_id) + + nemo_relay.subscribers.flush() + scope_events = [event for event in subscribed_events if isinstance(event, nemo_relay.ScopeEvent)] + assert [(event.scope_category, event.name) for event in scope_events] == [ + ("start", "request"), + ("start", "DeepAgent"), + ("end", "DeepAgent"), + ("end", "request"), + ] + + +def test_callback_handler_closes_semantic_agent_scope_on_error( + subscribed_events: list[nemo_relay.Event], + callback_handler: deepagents_integration.NemoRelayDeepAgentsCallbackHandler, +): + run_id = uuid4() + error = RuntimeError("agent failed") + + with nemo_relay.scope.scope("request", nemo_relay.ScopeType.Agent): + callback_handler.on_chain_start( + {}, + {"messages": ["hello"]}, + run_id=run_id, + name="main-agent", + metadata={ + "ls_integration": "deepagents", + "lc_agent_name": "main-agent", + "lc_versions": {"deepagents": "0.7.4"}, + }, + ) + callback_handler.on_chain_error(error, run_id=run_id) + + nemo_relay.subscribers.flush() + agent_end = next( + event + for event in subscribed_events + if isinstance(event, nemo_relay.ScopeEvent) and event.scope_category == "end" and event.name == "main-agent" + ) + assert isinstance(agent_end.metadata, dict) + agent_end_metadata = cast(dict[str, Any], agent_end.metadata) + assert agent_end_metadata["otel.status_code"] == "ERROR" + assert agent_end_metadata["otel.status_description"] == "agent failed" + assert agent_end_metadata["deepagents_agent_role"] == "orchestrator" + + def test_add_nemo_relay_integration_preserves_backend(deepagents_integration_module: types.ModuleType): mock_backend = MagicMock(name="mock_backend") mock_compiled_subagent = MagicMock(name="mock_compiled_subagent") @@ -547,13 +692,14 @@ def test_e2e_agent( ], ) agent = create_deep_agent(**kwargs) + callback = deepagents_integration_module.NemoRelayDeepAgentsCallbackHandler() with nemo_relay.scope.scope("deepagents-request", nemo_relay.ScopeType.Agent): input_payload = {"messages": [{"role": "user", "content": "Create a file named turtle."}]} if use_async: - result = asyncio.run(agent.ainvoke(input_payload)) + result = asyncio.run(agent.ainvoke(input_payload, config={"callbacks": [callback]})) else: - result = agent.invoke(input_payload) + result = agent.invoke(input_payload, config={"callbacks": [callback]}) nemo_relay.subscribers.flush() assert (tmp_path / "turtle").read_text() == "shell" @@ -575,6 +721,7 @@ def test_e2e_agent( expected_events = [ "scope.start.deepagents-request", + "scope.start.main-agent", "mark..DeepAgents Skills Configured", "scope.start.mock-model", "scope.end.mock-model", @@ -583,17 +730,101 @@ def test_e2e_agent( "scope.start.mock-model", "scope.end.mock-model", "scope.start.task", + "scope.start.reviewer", "mark..DeepAgents Skills Configured", "scope.start.mock-model", "scope.end.mock-model", + "scope.end.reviewer", "scope.end.task", "scope.start.mock-model", "scope.end.mock-model", + "scope.end.main-agent", "scope.end.deepagents-request", ] event_strings = [f"{event.kind}.{getattr(event, 'scope_category', '')}.{event.name}" for event in subscribed_events] assert event_strings == expected_events + scope_starts = { + event.name: event + for event in subscribed_events + if isinstance(event, nemo_relay.ScopeEvent) and event.scope_category == "start" and event.name != "mock-model" + } + model_starts = [ + event + for event in subscribed_events + if isinstance(event, nemo_relay.ScopeEvent) and event.scope_category == "start" and event.name == "mock-model" + ] + assert scope_starts["main-agent"].parent_uuid == scope_starts["deepagents-request"].uuid + assert scope_starts["write_file"].parent_uuid == scope_starts["main-agent"].uuid + assert scope_starts["task"].parent_uuid == scope_starts["main-agent"].uuid + assert scope_starts["reviewer"].parent_uuid == scope_starts["main-agent"].uuid + assert [event.parent_uuid for event in model_starts] == [ + scope_starts["main-agent"].uuid, + scope_starts["main-agent"].uuid, + scope_starts["reviewer"].uuid, + scope_starts["main-agent"].uuid, + ] + + +@pytest.mark.parametrize("use_async", [False, True]) +def test_e2e_general_purpose_subagent_scope( + use_async: bool, + subscribed_events: list[nemo_relay.Event], + deepagents_integration_module: types.ModuleType, +): + from deepagents import create_deep_agent + from langchain_core.messages import AIMessage + + model = _mock_deepagents_chat_model( + responses=[ + AIMessage( + content="", + tool_calls=[ + { + "name": "task", + "args": { + "description": "Return one concise result.", + "subagent_type": "general-purpose", + }, + "id": "call-1", + } + ], + ), + AIMessage(content="general-purpose result"), + AIMessage(content="done"), + ] + ) + agent = create_deep_agent( + **deepagents_integration_module.add_nemo_relay_integration( + model=model, + tools=[], + name="main-agent", + ) + ) + callback = deepagents_integration_module.NemoRelayDeepAgentsCallbackHandler() + + with nemo_relay.scope.scope("deepagents-request", nemo_relay.ScopeType.Agent): + input_payload = {"messages": [{"role": "user", "content": "Delegate this task."}]} + if use_async: + result = asyncio.run(agent.ainvoke(input_payload, config={"callbacks": [callback]})) + else: + result = agent.invoke(input_payload, config={"callbacks": [callback]}) + + nemo_relay.subscribers.flush() + assert result["messages"][-1].content == "done" + starts = [ + event + for event in subscribed_events + if isinstance(event, nemo_relay.ScopeEvent) and event.scope_category == "start" + ] + main_agent = next(event for event in starts if event.name == "main-agent") + general_purpose = next(event for event in starts if event.name == "general-purpose") + model_starts = [event for event in starts if event.name == "mock-model"] + assert general_purpose.parent_uuid == main_agent.uuid + assert isinstance(general_purpose.metadata, dict) + general_purpose_metadata = cast(dict[str, Any], general_purpose.metadata) + assert general_purpose_metadata["deepagents_agent_role"] == "subagent" + assert [event.parent_uuid for event in model_starts] == [main_agent.uuid, main_agent.uuid] def test_e2e_agent_exports_openinference_output_contract(