From 8c78d1f41f726fe6cc5496c30548b6c4708ab4bc Mon Sep 17 00:00:00 2001 From: March7 <132431415+March-77@users.noreply.github.com> Date: Wed, 22 Jul 2026 14:59:33 +0800 Subject: [PATCH] fix(subagents): load user-scoped skills (#4356) --- README.md | 2 +- backend/AGENTS.md | 1 + .../harness/deerflow/subagents/executor.py | 9 +- backend/tests/test_subagent_executor.py | 84 +++++++++++++++---- 4 files changed, 76 insertions(+), 20 deletions(-) diff --git a/README.md b/README.md index 5d678a259..5ae6cd711 100644 --- a/README.md +++ b/README.md @@ -756,7 +756,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. 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. +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 configured skills are resolved from the same user-scoped catalog as the lead agent, so user-owned custom skills remain available without exposing another user's version. 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 890c3aede..fb2b85bcb 100644 --- a/backend/AGENTS.md +++ b/backend/AGENTS.md @@ -499,6 +499,7 @@ Proxied through nginx: `/api/langgraph/*` → Gateway LangGraph-compatible runti ### Subagent System (`packages/harness/deerflow/subagents/`) **Built-in Agents**: `general-purpose` (all tools except `task`) and `bash` (command specialist) +**User-scoped Skills**: Subagents resolve their configured skills through `get_or_new_user_skill_storage(user_id)` using the parent runtime identity, with `DEFAULT_USER_ID` only when no identity is available. This keeps custom-skill shadowing and visibility aligned with the lead agent instead of reading the global-only catalog. **Execution**: Dual thread pool - `_scheduler_pool` (3 workers) + `_execution_pool` (3 workers) **Concurrency and total delegation cap**: `MAX_CONCURRENT_SUBAGENTS = 3` is enforced by `SubagentLimitMiddleware` (truncates excess tool calls in `after_model`; runtime `max_concurrent_subagents` is clamped to 2-4). The same middleware also enforces `subagents.max_total_per_run` (default 6, config schema 1-50, runtime override `max_total_subagents` clamped to the same range) against current-run entries in the durable delegation ledger, so a long lead-agent run cannot bypass concurrency limits by launching repeated legal-sized batches at each planning checkpoint, but historical delegations from previous runs in the same thread do not consume the new run's budget. The lead-agent prompt uses the same clamped values, so model-visible limits match enforcement. Gateway `run_agent()` and embedded `DeerFlowClient.stream()` both provide a per-invocation `run_id` in runtime context; `DeerFlowClient.stream()` also tags its input `HumanMessage` with that same id so durable-context capture can identify the current request boundary. Gateway resume paths may not append a new `HumanMessage`, so the worker also exposes the pre-run checkpoint's message ids in runtime context; durable-context capture uses that as the current-run boundary and never re-tags older task calls as the resumed run. When no delegation slots remain, task calls are stripped, provider raw tool-call metadata is synced, `finish_reason` is forced to `stop`, and a visible "subagent delegation limit" note is appended so the agent can synthesize already-collected results. Default subagent timeout `subagents.timeout_seconds=1800` (30 min) and built-in `general-purpose` `max_turns=150` (raised from 100/15-min so deep-research subtasks stop hitting `GraphRecursionError` out of the box) **Flow**: `task()` tool → `SubagentExecutor` → background thread → poll 5s → SSE events → result. `task_started` carries the resolved effective model name. The per-subagent `SubagentTokenCollector` publishes a cumulative usage snapshot to the shared `SubagentResult` after every completed LLM response; the next `task_running` event carries that snapshot, so collapsed workspace cards can update without re-accounting parent-run totals. Terminal ToolMessage metadata (`subagent_model_name`, `subagent_token_usage`) and the persisted `subagent.end` event retain the model/usage after reload; absent provider usage stays absent rather than being estimated as zero. diff --git a/backend/packages/harness/deerflow/subagents/executor.py b/backend/packages/harness/deerflow/subagents/executor.py index 346a16529..f6e12a0d6 100644 --- a/backend/packages/harness/deerflow/subagents/executor.py +++ b/backend/packages/harness/deerflow/subagents/executor.py @@ -27,6 +27,7 @@ from deerflow.authz.principal import normalize_authz_attributes from deerflow.config import get_app_config from deerflow.config.app_config import AppConfig from deerflow.models import create_chat_model +from deerflow.runtime.user_context import DEFAULT_USER_ID from deerflow.skills.tool_policy import filter_tools_by_skill_allowed_tools from deerflow.skills.types import Skill from deerflow.subagents.config import SubagentConfig, resolve_subagent_model_name @@ -566,10 +567,14 @@ class SubagentExecutor: return [] try: - from deerflow.skills.storage import get_or_new_skill_storage + from deerflow.skills.storage import get_or_new_user_skill_storage storage_kwargs = {"app_config": self.app_config} if self.app_config is not None else {} - storage = await asyncio.to_thread(get_or_new_skill_storage, **storage_kwargs) + storage = await asyncio.to_thread( + get_or_new_user_skill_storage, + self.user_id or DEFAULT_USER_ID, + **storage_kwargs, + ) # Use asyncio.to_thread to avoid blocking the event loop (LangGraph ASGI requirement) all_skills = await asyncio.to_thread(storage.load_skills, enabled_only=True) logger.info(f"[trace={self.trace_id}] Subagent {self.config.name} loaded {len(all_skills)} enabled skills from disk") diff --git a/backend/tests/test_subagent_executor.py b/backend/tests/test_subagent_executor.py index af4e02001..8392d402c 100644 --- a/backend/tests/test_subagent_executor.py +++ b/backend/tests/test_subagent_executor.py @@ -356,11 +356,12 @@ class TestAgentConstruction: skill_file.write_text("Use demo skill", encoding="utf-8") captured: dict[str, object] = {} - def fake_get_or_new_skill_storage(*, app_config=None): + def fake_get_or_new_user_skill_storage(user_id, *, app_config=None): + captured["user_id"] = user_id captured["app_config"] = app_config return SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="demo-skill", skill_file=skill_file)]) - monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_skill_storage", fake_get_or_new_skill_storage) + monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_user_skill_storage", fake_get_or_new_user_skill_storage) executor = SubagentExecutor( config=base_config, @@ -372,10 +373,59 @@ class TestAgentConstruction: skills = await executor._load_skills() messages = await executor._load_skill_messages(skills) - assert captured["app_config"] is app_config + assert captured == {"user_id": "default", "app_config": app_config} assert len(messages) == 1 assert "Use demo skill" in messages[0].content + @pytest.mark.anyio + async def test_load_skills_uses_each_subagent_users_scoped_storage( + self, + classes, + base_config, + monkeypatch: pytest.MonkeyPatch, + ): + SubagentExecutor = classes["SubagentExecutor"] + app_config = SimpleNamespace(models=[SimpleNamespace(name="default-model")]) + storage_calls: list[tuple[str, object]] = [] + + def user_storage(user_id: str, *, app_config=None): + storage_calls.append((user_id, app_config)) + return SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="shared-skill", owner=user_id)]) + + global_storage = MagicMock(side_effect=AssertionError("subagents must not read the global-only skill catalog")) + monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_skill_storage", global_storage) + monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_user_skill_storage", user_storage) + + alice = SubagentExecutor(config=base_config, tools=[], app_config=app_config, thread_id="alice-thread", user_id="alice") + bob = SubagentExecutor(config=base_config, tools=[], app_config=app_config, thread_id="bob-thread", user_id="bob") + + alice_skills = await alice._load_skills() + bob_skills = await bob._load_skills() + + assert [skill.owner for skill in alice_skills] == ["alice"] + assert [skill.owner for skill in bob_skills] == ["bob"] + assert storage_calls == [("alice", app_config), ("bob", app_config)] + global_storage.assert_not_called() + + @pytest.mark.anyio + async def test_load_skills_defaults_missing_user_to_default_scope( + self, + classes, + base_config, + monkeypatch: pytest.MonkeyPatch, + ): + SubagentExecutor = classes["SubagentExecutor"] + user_storage = MagicMock(return_value=SimpleNamespace(load_skills=lambda *, enabled_only: [])) + global_storage = MagicMock(side_effect=AssertionError("subagents must not read the global-only skill catalog")) + monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_skill_storage", global_storage) + monkeypatch.setattr(sys.modules["deerflow.skills.storage"], "get_or_new_user_skill_storage", user_storage) + + executor = SubagentExecutor(config=base_config, tools=[], thread_id="test-thread", user_id=None) + + assert await executor._load_skills() == [] + user_storage.assert_called_once_with("default") + global_storage.assert_not_called() + @pytest.mark.anyio async def test_load_skill_messages_escapes_untrusted_name_and_content( self, @@ -427,8 +477,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="my-skill", skill_file=skill_file, allowed_tools=None)]), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="my-skill", skill_file=skill_file, allowed_tools=None)]), ) executor = SubagentExecutor( @@ -465,8 +515,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), ) executor = SubagentExecutor( @@ -510,8 +560,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="my-skill", skill_file=skill_file, allowed_tools=None)]), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="my-skill", skill_file=skill_file, allowed_tools=None)]), ) SubagentExecutor = classes["SubagentExecutor"] @@ -546,8 +596,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), ) monkeypatch.setattr(executor_module, "get_app_config", lambda: SimpleNamespace(tool_search=SimpleNamespace(enabled=True))) @@ -587,8 +637,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), ) monkeypatch.setattr(executor_module, "get_app_config", lambda: SimpleNamespace(tool_search=SimpleNamespace(enabled=False))) @@ -632,8 +682,8 @@ class TestAgentConstruction: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: []), ) monkeypatch.setattr(executor_module, "get_app_config", lambda: SimpleNamespace(tool_search=SimpleNamespace(enabled=True))) @@ -1418,8 +1468,8 @@ class TestAsyncExecutionPath: monkeypatch.setattr( sys.modules["deerflow.skills.storage"], - "get_or_new_skill_storage", - lambda *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="regression-skill", skill_file=skill_dir / "SKILL.md", allowed_tools=None)]), + "get_or_new_user_skill_storage", + lambda user_id, *, app_config=None: SimpleNamespace(load_skills=lambda *, enabled_only: [SimpleNamespace(name="regression-skill", skill_file=skill_dir / "SKILL.md", allowed_tools=None)]), ) captured_states: list[dict] = []