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

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 14 additions & 4 deletions docs/supported-integrations/deepagents.mdx
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -106,16 +108,24 @@ 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
`SKILL.md` read; this is distinct from the configured-skills summary.
- 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.
61 changes: 60 additions & 1 deletion python/nemo_relay/integrations/deepagents/callbacks.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand All @@ -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
Comment thread
bbednarski9 marked this conversation as resolved.

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"
Comment thread
bbednarski9 marked this conversation as resolved.

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):
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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")
Expand Down Expand Up @@ -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"
Expand All @@ -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",
Expand All @@ -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]
Comment thread
coderabbitai[bot] marked this conversation as resolved.


def test_e2e_agent_exports_openinference_output_contract(
Expand Down
Loading