From de55982c5acd6edf5ccaa1a3a1d20a65185325ee Mon Sep 17 00:00:00 2001 From: Nan Gao Date: Thu, 16 Jul 2026 01:03:08 +0200 Subject: [PATCH] fix(subagents): preserve parent checkpoint namespace (#4215) * fix(subagents): preserve parent checkpoint namespace * test(subagents): align stream isolation coverage --- README.md | 2 +- backend/AGENTS.md | 1 + .../harness/deerflow/subagents/executor.py | 7 +- backend/tests/test_subagent_executor.py | 197 ++++++++++++++++++ 4 files changed, 204 insertions(+), 3 deletions(-) diff --git a/README.md b/README.md index 5303e70c3..487a8f110 100644 --- a/README.md +++ b/README.md @@ -735,7 +735,7 @@ Use `/compact` in the Web UI composer to summarize older context for the current Complex tasks rarely fit in a single pass. DeerFlow decomposes them. -The lead agent can spawn sub-agents on the fly — each with its own scoped context, tools, and termination conditions. Sub-agents run in parallel when possible, report back structured results, and the lead agent synthesizes everything into a coherent output. Long-running sub-agents compact older history when summarization is enabled and re-inject the summary as guarded, hidden durable context before continuing, so recent assistant/tool activity remains grounded in the task. Provider/model request failures are reported as failed sub-agent tasks rather than successful results, so the lead agent and Web UI can react to them correctly. Collapsed sub-agent cards show the effective model and, when the provider returns usage metadata, a cumulative token total that updates after each completed sub-agent LLM call and persists after a reload. When token usage tracking is enabled, completed sub-agent usage is also attributed back to the dispatching step. +The lead agent can spawn sub-agents on the fly — each with its own scoped context, tools, and termination conditions. Sub-agents run in parallel when possible, report back structured results, and the lead agent synthesizes everything into a coherent output. Their internal AI and tool messages stay scoped to the delegated graph instead of entering the parent chat stream. Long-running sub-agents compact older history when summarization is enabled and re-inject the summary as guarded, hidden durable context before continuing, so recent assistant/tool activity remains grounded in the task. Provider/model request failures are reported as failed sub-agent tasks rather than successful results, so the lead agent and Web UI can react to them correctly. Collapsed sub-agent cards show the effective model and, when the provider returns usage metadata, a cumulative token total that updates after each completed sub-agent LLM call and persists after a reload. When token usage tracking is enabled, completed sub-agent usage is also attributed back to the dispatching step. This is how DeerFlow handles tasks that take minutes to hours: a research task might fan out into a dozen sub-agents, each exploring a different angle, then converge into a single report — or a website — or a slide deck with generated visuals. One harness, many hands. diff --git a/backend/AGENTS.md b/backend/AGENTS.md index a811c8f6f..30e8d97b2 100644 --- a/backend/AGENTS.md +++ b/backend/AGENTS.md @@ -392,6 +392,7 @@ Proxied through nginx: `/api/langgraph/*` → Gateway LangGraph-compatible runti **Step capture & persistence (#3779)**: `executor.py` captures both assistant turns (`AIMessage`) **and** tool outputs (`ToolMessage`) via `subagents/step_events.py::capture_new_step_messages`, which walks the *newly-appended tail* of each `stream_mode="values"` chunk (not just `messages[-1]`) so a multi-tool-call turn — where LangGraph's `ToolNode` appends several `ToolMessage`s in one super-step — keeps every tool output instead of dropping all but the last. `runtime/runs/worker.py::_SubagentEventBuffer` additionally persists these `task_*` custom events to the `RunEventStore` as `subagent.start`/`subagent.step`/`subagent.end` (`category="subagent"`, `task_id` in `metadata`). It **batches** writes via `put_batch` (flushing on a terminal `subagent.end`, at `FLUSH_THRESHOLD` events, and in the worker's `finally`) rather than one `put()` per step, since `put()` is a documented low-frequency path (per-thread advisory lock per call) and a deep subagent (`max_turns=150`) emits hundreds of steps on the hot stream loop. `build_subagent_step` caps both the per-step `text` and each tool call's serialized `args` at `SUBAGENT_STEP_MAX_CHARS` (flagged `truncated` / `args_truncated`) so a large `write_file`/`bash` payload can't produce an unbounded row. The dedicated category keeps them out of `list_messages` (the thread feed) while `list_events` returns them for the frontend's fetch-on-expand backfill. `list_events` accepts `task_id` (filters on `metadata["task_id"]` — SQL-side in `DbRunEventStore` via `event_metadata["task_id"].as_string()`, in-memory in the JSONL/memory stores) plus an `after_seq` forward cursor, so the card pages through one subagent's steps without the run-wide `limit` truncating the tail (no schema migration: the filter rides the existing run-scoped index). `step_events.py` is a pure, unit-tested layer (`build_subagent_step` / `subagent_run_event`). **History contraction (#3875 Phase 3)**: `capture_new_step_messages` assumes append-only growth, but `DeerFlowSummarizationMiddleware` rewrites the messages channel via `RemoveMessage(id=REMOVE_ALL_MESSAGES)`, shrinking `len(messages)` below the cursor mid-run. On contraction (`total < processed_count`) the cursor resets to the new tail; `capture_step_message`'s id/content dedup prevents re-emitting pre-compaction steps, so steps appended after the compaction point are still captured instead of being dropped until `total` overtakes the stale cursor. **Deferred MCP tools** (if `tool_search.enabled`): `SubagentExecutor._build_initial_state` assembles deferral after policy filtering via the shared `assemble_deferred_tools` (fail-closed), appends the `tool_search` tool, injects the `` section into the subagent's `SystemMessage`, and threads the setup to `_create_agent`, which attaches `McpRoutingMiddleware` (when PR1 routing metadata matches deferred tools) before `DeferredToolFilterMiddleware` through `build_subagent_runtime_middlewares(...)`. Subagents thus withhold full MCP schemas until promotion, same as the lead agent; each task run gets a fresh `ThreadState` so promotion is isolated per run **Checkpointer isolation**: Subagent graphs are compiled with `checkpointer=False` to avoid inheriting the parent run's checkpointer, since subagents are one-shot and never resume. +**Checkpoint lineage / stream isolation**: `_aexecute` deliberately omits checkpoint-coordinate keys (`thread_id`, `checkpoint_ns`, `checkpoint_id`, `checkpoint_map`) from the child `RunnableConfig`. LangGraph must inherit those coordinates from the copied parent ContextVar so the delegated graph retains a non-root subgraph namespace; explicitly re-supplying even the same parent `thread_id` starts a new root lineage on LangGraph 1.2.6+ and can route child AI/tool frames into the parent `messages` stream. DeerFlow business components still receive the parent `thread_id` through `runtime.context`, which is the preferred lookup path for sandbox, middleware, and attribution code. Regression coverage in `tests/test_subagent_executor.py::TestSubagentCheckpointLineage` keeps the invocation-contract assertion active on every supported version and version-gates the production-shaped parent-stream test to LangGraph 1.2.6+, where the leak exists. ### Tool System (`packages/harness/deerflow/tools/`) diff --git a/backend/packages/harness/deerflow/subagents/executor.py b/backend/packages/harness/deerflow/subagents/executor.py index 8bfcc1648..7885970d9 100644 --- a/backend/packages/harness/deerflow/subagents/executor.py +++ b/backend/packages/harness/deerflow/subagents/executor.py @@ -728,7 +728,11 @@ class SubagentExecutor: collector_caller = f"subagent:{self.config.name}" collector = SubagentTokenCollector(caller=collector_caller) - # Build config with thread_id for sandbox access and recursion limit + # Do not put checkpoint coordinates (thread_id/checkpoint_ns/etc.) + # in the child config. LangGraph inherits those coordinates from + # the ambient parent run so this execution keeps its subgraph + # namespace. Business consumers receive thread_id via ``context`` + # below instead. run_config: RunnableConfig = { "recursion_limit": self.config.max_turns, "callbacks": [collector], @@ -767,7 +771,6 @@ class SubagentExecutor: context: dict[str, Any] = {} if self.thread_id: - run_config["configurable"] = {"thread_id": self.thread_id} context["thread_id"] = self.thread_id if self.app_config is not None: context["app_config"] = self.app_config diff --git a/backend/tests/test_subagent_executor.py b/backend/tests/test_subagent_executor.py index f01d9ebbf..7bf2df1ee 100644 --- a/backend/tests/test_subagent_executor.py +++ b/backend/tests/test_subagent_executor.py @@ -7,6 +7,7 @@ Covers: - Error handling in both sync and async paths - Async tool support (MCP tools) - Cooperative cancellation via cancel_event +- Parent/child checkpoint-lineage and message-stream isolation Note: Due to circular import issues in the main codebase, conftest.py mocks deerflow.subagents.executor. This test file uses delayed import via fixture to test @@ -18,11 +19,13 @@ import importlib import sys import threading from datetime import datetime +from importlib.metadata import version as package_version from pathlib import Path from types import ModuleType, SimpleNamespace from unittest.mock import MagicMock, patch import pytest +from packaging.version import Version from deerflow.skills.types import Skill @@ -39,6 +42,8 @@ _MOCKED_MODULE_NAMES = [ "deerflow.skills.storage", ] +_LANGGRAPH_HAS_ROOT_LINEAGE_STREAM_REGRESSION = Version(package_version("langgraph")) >= Version("1.2.6") + def _default_app_config(): return SimpleNamespace(tool_search=SimpleNamespace(enabled=False)) @@ -2540,6 +2545,198 @@ class _FakeStreamAgent: yield # pragma: no cover - make this an async generator +class TestSubagentCheckpointLineage: + """Keep delegated graphs on the parent run's checkpoint lineage.""" + + @pytest.mark.anyio + async def test_aexecute_leaves_checkpoint_coordinates_to_parent_context( + self, + classes, + monkeypatch, + ): + """A delegated graph must not declare a new root checkpoint lineage. + + LangGraph treats any explicitly supplied checkpoint coordinate as an + independent lineage. In particular, re-supplying the parent's own + ``thread_id`` clears the ambient ``checkpoint_ns`` on LangGraph 1.2.6+, + so the child is routed as a root graph instead of a subgraph. + """ + executor_module = importlib.import_module("deerflow.subagents.executor") + monkeypatch.setattr(executor_module, "build_tracing_callbacks", lambda: []) + + executor = classes["SubagentExecutor"]( + config=classes["SubagentConfig"]( + name="general-purpose", + description="Checkpoint lineage test agent", + system_prompt="You are a checkpoint lineage test agent.", + max_turns=5, + timeout_seconds=30, + ), + tools=[], + parent_model="test-model", + thread_id="parent-thread-1", + trace_id="trace-lineage-1", + ) + fake_agent = _FakeStreamAgent() + + async def build_initial_state(task): + return ({"messages": [classes["HumanMessage"](content=task)]}, [], None) + + monkeypatch.setattr(executor, "_build_initial_state", build_initial_state) + monkeypatch.setattr(executor, "_create_agent", lambda *args, **kwargs: fake_agent) + + await executor._aexecute("do something") + + assert fake_agent.captured_config is not None + configurable = fake_agent.captured_config.get("configurable") or {} + checkpoint_coordinates = { + "thread_id", + "checkpoint_ns", + "checkpoint_id", + "checkpoint_map", + } + assert checkpoint_coordinates.isdisjoint(configurable), f"subagent invocation must inherit checkpoint coordinates from the parent context, got {configurable!r}" + assert fake_agent.captured_context is not None + assert fake_agent.captured_context["thread_id"] == "parent-thread-1" + + @pytest.mark.anyio + @pytest.mark.skipif( + not _LANGGRAPH_HAS_ROOT_LINEAGE_STREAM_REGRESSION, + reason="root-lineage message leak only manifests on LangGraph >=1.2.6", + ) + async def test_parent_message_stream_excludes_delegated_graph_messages( + self, + classes, + monkeypatch, + ): + """Child AI/tool frames stay outside the parent's messages stream.""" + from langchain_core.messages import AIMessage, ToolMessage + from langgraph.checkpoint.memory import MemorySaver + from langgraph.graph import END, START, MessagesState, StateGraph + + executor_module = importlib.import_module("deerflow.subagents.executor") + monkeypatch.setattr(executor_module, "build_tracing_callbacks", lambda: []) + + child_builder = StateGraph(MessagesState) + child_builder.add_node( + "child_model", + lambda _state: { + "messages": [ + AIMessage( + content="", + id="child-ai-sentinel", + tool_calls=[ + { + "name": "child_tool", + "args": {}, + "id": "child-tool-call", + "type": "tool_call", + } + ], + ) + ] + }, + ) + child_builder.add_node( + "child_tool", + lambda _state: { + "messages": [ + ToolMessage( + content="CHILD_TOOL_SENTINEL", + name="child_tool", + tool_call_id="child-tool-call", + id="child-tool-sentinel", + ) + ] + }, + ) + child_builder.add_node( + "child_final", + lambda _state: { + "messages": [ + AIMessage( + content="CHILD_FINAL_SENTINEL", + id="child-final-sentinel", + ) + ] + }, + ) + child_builder.add_edge(START, "child_model") + child_builder.add_edge("child_model", "child_tool") + child_builder.add_edge("child_tool", "child_final") + child_builder.add_edge("child_final", END) + child_graph = child_builder.compile(checkpointer=False) + + executor = classes["SubagentExecutor"]( + config=classes["SubagentConfig"]( + name="general-purpose", + description="Stream isolation test agent", + system_prompt="You are a stream isolation test agent.", + max_turns=5, + timeout_seconds=30, + ), + tools=[], + parent_model="test-model", + thread_id="parent-thread-1", + trace_id="trace-stream-isolation-1", + ) + + async def build_initial_state(task): + return ({"messages": [classes["HumanMessage"](content=task)]}, [], None) + + monkeypatch.setattr(executor, "_build_initial_state", build_initial_state) + monkeypatch.setattr(executor, "_create_agent", lambda *args, **kwargs: child_graph) + + async def delegate(_state): + task_id = executor.execute_async("run the child graph") + try: + deadline = asyncio.get_running_loop().time() + 5 + while True: + result = executor_module.get_background_task_result(task_id) + if result is not None and result.status.is_terminal: + break + if asyncio.get_running_loop().time() >= deadline: + pytest.fail("background subagent did not complete") + await asyncio.sleep(0.001) + assert result.status.value == "completed" + finally: + executor_module.cleanup_background_task(task_id) + return { + "messages": [ + AIMessage( + content="PARENT_FINAL_SENTINEL", + id="parent-final-sentinel", + ) + ] + } + + parent_builder = StateGraph(MessagesState) + parent_builder.add_node("delegate", delegate) + parent_builder.add_edge(START, "delegate") + parent_builder.add_edge("delegate", END) + parent_graph = parent_builder.compile(checkpointer=MemorySaver()) + + streamed_messages = [ + message + async for message, _metadata in parent_graph.astream( + {"messages": [classes["HumanMessage"](content="delegate")]}, + config={"configurable": {"thread_id": "parent-thread-1"}}, + stream_mode="messages", + ) + ] + + streamed_ids = {message.id for message in streamed_messages} + assert "parent-final-sentinel" in streamed_ids + assert ( + not { + "child-ai-sentinel", + "child-tool-sentinel", + "child-final-sentinel", + } + & streamed_ids + ) + + class TestSubagentTracingWiring: """Verify the subagent graph-root tracing wiring matches the lead agent."""