diff --git a/EvoScientist/EvoScientist.py b/EvoScientist/EvoScientist.py index 30ce5c8..9908ae9 100644 --- a/EvoScientist/EvoScientist.py +++ b/EvoScientist/EvoScientist.py @@ -549,6 +549,105 @@ def _maybe_swap_async_subagents( return out +def _route_async_specs_through_evo_middleware( + subs: list, base_middleware: list, *, cfg=None +) -> list: + """Move ``AsyncSubAgent`` specs from ``subs`` into ``EvoAsyncSubAgentMiddleware``. + + Deepagents' ``create_deep_agent`` auto-composes the vanilla + ``AsyncSubAgentMiddleware`` when it sees ``graph_id``-carrying entries + in ``subagents=``. We need our payload-aware subclass to handle those + (see ``EvoScientist/middleware/expert_async_subagent.py`` for the + upstream-workaround rationale). To prevent the auto-composition and + route all async dispatch through our subclass, we strip AsyncSubAgent + specs from ``subs`` here and hand them to our middleware. + + Also folds in ``AsyncSubAgent`` specs for installed + ``default_dispatch: async`` expert skills — all pointing at the shared + ``expert-container-async`` graph, marked ``is_expert=True`` so the + middleware requires a payload with ``skill_name``. + + Returns: + ``subs`` with ``graph_id``-carrying entries removed. Safe to pass + as ``create_deep_agent(subagents=...)`` — the async-auto-compose + branch is skipped for empty async lists. + """ + from .middleware.async_watcher import AsyncWatcherMiddleware + from .middleware.expert_async_subagent import EvoAsyncSubAgentMiddleware + from .subagents.expert_container_async import build_expert_async_subagent_specs + + cfg = cfg if cfg is not None else _ensure_config() + + async_specs = [s for s in subs if "graph_id" in s] + sync_subs = [s for s in subs if "graph_id" not in s] + expert_specs = build_expert_async_subagent_specs(cfg=cfg) + async_specs.extend(expert_specs) + + if async_specs: + # ``_maybe_swap_async_subagents`` installs the model-passthrough patch + # only when the yaml-async spec list is non-empty. An expert-only setup + # (no ``writing-agent`` / ``data-analysis-agent`` / ``scheduler`` in + # yaml) would otherwise miss the patch entirely, so we install it here + # too. Idempotent — the shared ``_model_passthrough_patched`` flag + # guards against double-patching. + from .llm.patches import _patch_deepagents_model_passthrough + + _patch_deepagents_model_passthrough() + + # Prepend rather than append so the ``## Async subagents`` prompt + # section stays in the stable prefix. Appending pushes it past the + # volatile memory tail, invalidating the cached prefix on every + # memory change. + base_middleware.insert( + 0, EvoAsyncSubAgentMiddleware(async_subagents=async_specs) + ) + + # Extend AsyncWatcherMiddleware's client cache with expert specs so + # start_async_task launches for experts spawn a completion watcher — + # otherwise the watcher's ``get_async(agent_name)`` KeyErrors on the + # expert name, no notification is enqueued, and the main agent never + # learns the task finished. ``_maybe_swap_async_subagents`` above only + # populates the watcher with YAML-defined async subagents (writing-agent, + # data-analysis-agent, scheduler); this hook folds in the experts too. + if expert_specs: + watcher = next( + (m for m in base_middleware if isinstance(m, AsyncWatcherMiddleware)), + None, + ) + if watcher is not None: + # The mutation reaches through two layers of private state: + # ``AsyncWatcherMiddleware._clients`` (our own) and + # ``_ClientCache._agents`` (upstream deepagents). If upstream ever + # renames ``_agents`` or wraps it in an immutable snapshot, the + # ``.update(...)`` below silently lands on nothing — expert + # completion nudges then stop firing without a diagnostic surface. + # Convert that silent-drop into a grep-able error line and bail + # out of the extension path; expert dispatches still work, just + # without completion notifications until upstream drift is fixed. + if not hasattr(watcher._clients, "_agents"): + logging.getLogger(__name__).error( + "AsyncWatcherMiddleware._clients has no `_agents` slot — " + "deepagents internal renamed; expert completion " + "notifications will not fire until the extension hook is " + "updated to the new attribute name." + ) + return sync_subs + watcher._clients._agents.update({s["name"]: s for s in expert_specs}) + else: + # No YAML async subagents were registered, so ``_maybe_swap`` did + # not install the watcher. Install it now so experts still get + # completion notifications. + from .cli import async_notifier + + base_middleware.append( + AsyncWatcherMiddleware( + {s["name"]: s for s in expert_specs}, + notifier=async_notifier, + ) + ) + return sync_subs + + def _build_base_kwargs( base_backend, base_middleware, *, cfg=None, chat_model=None, workspace_dir=None ): @@ -557,12 +656,7 @@ def _build_base_kwargs( from .utils import load_subagents cfg = cfg if cfg is not None else _ensure_config() - # `skill_manager` is registered here in addition to `base_tools` because - # expert subagents resolve their default toolset from `tool_registry` (see - # `_DEFAULT_EXPERT_TOOLS` in expert_container.py). Without this entry the - # tool silently misses from every expert sub-agent — e.g. idea-brainstorm - # can't run its `paper-navigator` precondition check. - tool_registry = {"think_tool": think_tool, "skill_manager": skill_manager} + tool_registry = {"think_tool": think_tool} if os.environ.get("TAVILY_API_KEY"): tool_registry["tavily_search"] = tavily_search base_tools = [think_tool, skill_manager] @@ -582,6 +676,10 @@ def _build_base_kwargs( subs, workspace_dir=workspace_dir, cfg=cfg, chat_model=chat_model ) subs = _maybe_swap_async_subagents(subs, base_middleware, cfg=cfg) + # Route AsyncSubAgent specs (both standard and expert) through + # EvoAsyncSubAgentMiddleware so the payload-aware start_async_task tool + # replaces upstream's non-parameterisable one. + subs = _route_async_specs_through_evo_middleware(subs, base_middleware, cfg=cfg) return { "name": "EvoScientist", "model": chat_model if chat_model is not None else _ensure_chat_model(), @@ -634,10 +732,7 @@ def load_mcp_and_build_kwargs( workspace_dir=workspace_dir, ) - # Match `_build_base_kwargs`: register `skill_manager` in the registry so - # expert subagents (which resolve tools via `_DEFAULT_EXPERT_TOOLS` from - # `expert_container.py`) actually get it. - tool_registry = {"think_tool": think_tool, "skill_manager": skill_manager} + tool_registry = {"think_tool": think_tool} if os.environ.get("TAVILY_API_KEY"): tool_registry["tavily_search"] = tavily_search base_tools = [think_tool, skill_manager] @@ -672,6 +767,10 @@ def load_mcp_and_build_kwargs( # Swap selected sub-agents to AsyncSubAgent (must happen AFTER MCP injection # since async sub-agents are remote graphs that load their own tools). subs = _maybe_swap_async_subagents(subs, base_middleware, cfg=cfg) + # Mirror the base path: route AsyncSubAgent specs through + # EvoAsyncSubAgentMiddleware so the payload-aware start_async_task tool + # is the one composed into the main agent. + subs = _route_async_specs_through_evo_middleware(subs, base_middleware, cfg=cfg) return { "name": "EvoScientist", diff --git a/EvoScientist/commands/implementation/experts.py b/EvoScientist/commands/implementation/experts.py index 7193a87..27c237a 100644 --- a/EvoScientist/commands/implementation/experts.py +++ b/EvoScientist/commands/implementation/experts.py @@ -203,19 +203,30 @@ class ExpertCommand(Command): dispatchable = {s.name for s in _dispatchable_experts()} if target not in dispatchable: + from ...subagents.expert_container import is_async_dispatch_available from ...tools.skills_manager import list_expert_skills - installed = {s.name for s in list_expert_skills(include_system=True)} - if target in installed: + installed = {s.name: s for s in list_expert_skills(include_system=True)} + match = installed.get(target) + if match is None: ctx.ui.append_system( - f"Expert '{target}' can't be dispatched (empty SKILL.md body " - "or name collision with a built-in sub-agent).", + f"No expert skill named '{target}'. `/experts` lists " + "installed ones.", + style="red", + ) + elif ( + match.default_dispatch == "async" and not is_async_dispatch_available() + ): + ctx.ui.append_system( + f"Expert '{target}' declares async dispatch but async " + "dispatch is unavailable (enable_async_subagents is off " + "or langgraph dev is not reachable).", style="red", ) else: ctx.ui.append_system( - f"No expert skill named '{target}'. `/experts` lists " - "installed ones.", + f"Expert '{target}' can't be dispatched (empty SKILL.md body " + "or name collision with a built-in sub-agent).", style="red", ) return diff --git a/EvoScientist/langgraph_dev/graphs.py b/EvoScientist/langgraph_dev/graphs.py index 839da43..a9a4ead 100644 --- a/EvoScientist/langgraph_dev/graphs.py +++ b/EvoScientist/langgraph_dev/graphs.py @@ -29,10 +29,19 @@ from EvoScientist.memory.agents import ( ) from EvoScientist.memory.types import MemorySourceType from EvoScientist.subagents._factory import build_async_subagent_graph +from EvoScientist.subagents.expert_container_async import ( + build_expert_container_async_graph, +) writing_agent = build_async_subagent_graph("writing-agent") data_analysis_agent = build_async_subagent_graph("data-analysis-agent") scheduler = build_async_subagent_graph("scheduler") +# Generic async container for expert-skill dispatch. One graph, parameterised +# per invocation by the ``skill_name`` payload the main agent passes through +# ``EvoAsyncSubAgentMiddleware.start_async_task``. Any installed expert skill +# dispatches through this graph; the loader middleware resolves the skill +# body at model-call time. +expert_container_async = build_expert_container_async_graph() evomemory_subagent_worker = build_memory_worker_graph(MemorySourceType.SUBAGENT) evomemory_turn_worker = build_memory_worker_graph(MemorySourceType.TURN) evomemory_observation_linker = build_observation_linker_graph() diff --git a/EvoScientist/langgraph_dev/langgraph.json b/EvoScientist/langgraph_dev/langgraph.json index c834540..56a3720 100644 --- a/EvoScientist/langgraph_dev/langgraph.json +++ b/EvoScientist/langgraph_dev/langgraph.json @@ -5,6 +5,7 @@ "writing-agent": "EvoScientist.langgraph_dev.graphs:writing_agent", "data-analysis-agent": "EvoScientist.langgraph_dev.graphs:data_analysis_agent", "scheduler": "EvoScientist.langgraph_dev.graphs:scheduler", + "expert-container-async": "EvoScientist.langgraph_dev.graphs:expert_container_async", "evomemory-subagent-worker": "EvoScientist.langgraph_dev.graphs:evomemory_subagent_worker", "evomemory-turn-worker": "EvoScientist.langgraph_dev.graphs:evomemory_turn_worker", "evomemory-observation-linker": "EvoScientist.langgraph_dev.graphs:evomemory_observation_linker", diff --git a/EvoScientist/middleware/active_team.py b/EvoScientist/middleware/active_team.py index bf12596..58d7667 100644 --- a/EvoScientist/middleware/active_team.py +++ b/EvoScientist/middleware/active_team.py @@ -2,7 +2,21 @@ Reads ``configurable.active_teams: list[str]`` on every model call and appends a system-prompt cue biasing the main agent to consult the -user-invited expert(s) via ``task({subagent_type: ...})``. +user-invited expert(s). + +The cue is dispatch-aware: for each active expert the middleware looks up +its ``default_dispatch`` (via ``list_expert_skills()``) and emits the +tool-shape cue that matches how the expert actually runs: + +- ``sync`` / ``panel`` / unset -> ``task({subagent_type: 'X', ...})``. +- ``async`` -> ``start_async_task(subagent_type: 'X', payload: {skill_name: + 'X', output_path: '...'})``, plus a reminder that ``check_async_task`` + returns the status/result later. + +Without the dispatch-aware branch, an ``async`` expert like +``literature-review`` gets told to use ``task()``, which routes it back +through the sync ``SubAgentMiddleware`` rather than the async graph the +container registers for it. Backend-stateless team binding: WebUI sends ``active_teams`` on every ``stream.submit()`` for as long as the invited expert is active; this @@ -40,25 +54,46 @@ from langchain.agents.middleware.types import ( ModelResponse, ) +# Per-expert cue shapes. Composed inside ``_TEMPLATE_SINGLE`` / +# ``_TEMPLATE_MULTI`` at render time so the wrapping tags stay in sync +# with the count of active experts. +_SYNC_CUE = ( + "Consult it via `task({{subagent_type: '{expert}', description: '...'}})`. " + "It runs synchronously and returns its result to the same turn." +) + +_ASYNC_CUE = ( + "Consult it via `start_async_task(description: '.md`` verbatim>', " + "subagent_type: '{expert}')`. It runs in the background and returns a " + "task_id immediately; the result artifact is written to the path you " + "named in the description. Use ``check_async_task`` to poll status " + "when the user asks. On ``status: 'success'`` the ``result`` field " + "contains a JSON envelope with ``output_path``, a one-paragraph " + "``summary``, and a skill-defined ``metadata`` block (fields vary by " + "expert — e.g. ``word_count`` / ``section_count`` / ``citations_used`` " + "for surveys). Render ``summary`` and ``metadata`` directly to the " + "user — do not re-read the artifact to build a synopsis." +) + _TEMPLATE_SINGLE = ( "\n" "The user has invited the expert `{expert}` to this thread. " - "Consult it via `task({{subagent_type: '{expert}', ...}})` for " - "requests within its scope. It stays available for the whole session " - "until the user dismisses it.\n" + "{cue} " + "It stays available for the whole session until the user dismisses it.\n" "" ) -_TEMPLATE_MULTI = ( +_TEMPLATE_MULTI_HEADER = ( "\n" - "The user has invited the following experts to this thread: " - "{experts}. Consult any of them via " - "`task({{subagent_type: '', ...}})` based on which fits " - "the current request. Do not consult an expert if the request is " - "clearly outside its scope.\n" - "" + "The user has invited the following experts to this thread: {experts}. " + "Consult the right one for the current request; do not consult an expert " + "if the request is clearly outside its scope. Per-expert dispatch:\n" ) +_TEMPLATE_MULTI_FOOTER = "" + def _read_active_teams() -> list[str]: """Read ``configurable.active_teams`` from the current RunnableConfig. @@ -85,28 +120,84 @@ def _read_active_teams() -> list[str]: return [t for t in raw if isinstance(t, str) and t] +def _dispatch_by_name() -> dict[str, str]: + """Return ``{skill_name: default_dispatch}`` for currently dispatchable experts. + + Fresh filesystem read every call so a ``skill_manager install `` + is visible on the next turn without an agent rebuild. Cheap at current + scale (a handful of skills, cached bodies). + + Sourced from ``list_dispatchable_experts`` — which filters empty-body + experts, name collisions with reserved sub-agents, AND async-declared + experts when async dispatch is unavailable (``enable_async_subagents`` + off or langgraph dev unreachable). Keeps the cue honest: any expert + named here can actually be reached by the tool shape the cue advertises. + + On import failure returns an empty dict — the middleware then emits no + cue, matching the outside-runnable-context no-op path. + """ + try: + from ..subagents.expert_container import list_dispatchable_experts + except Exception: + return {} + try: + return {s.name: s.default_dispatch for s in list_dispatchable_experts()} + except Exception: + return {} + + +def _cue_shape_for(dispatch: str, expert: str) -> str: + """Return the ``task()`` / ``start_async_task(...)`` cue for one expert.""" + if dispatch == "async": + return _ASYNC_CUE.format(expert=expert) + return _SYNC_CUE.format(expert=expert) + + class ActiveTeamMiddleware(AgentMiddleware): """Bias delegation toward the user's active expert(s) on every turn.""" name = "active_team" def _cue_for(self, experts: list[str]) -> str: + """Render the cue over the dispatchable subset of ``experts``. + + Invited experts that aren't currently dispatchable (uninstalled, + empty body, name collision, or async-declared while async + dispatch is unavailable) are dropped from the cue — pointing the + model at a tool shape that will fail is worse than saying nothing. + Returns the empty string when nothing survives the filter; caller + skips the system-prompt append in that case. + """ + dispatch_map = _dispatch_by_name() + experts = [e for e in experts if e in dispatch_map] + if not experts: + return "" if len(experts) == 1: - return _TEMPLATE_SINGLE.format(expert=experts[0]) + expert = experts[0] + cue = _cue_shape_for(dispatch_map[expert], expert) + return _TEMPLATE_SINGLE.format(expert=expert, cue=cue) experts_str = ", ".join(f"`{e}`" for e in experts) - return _TEMPLATE_MULTI.format(experts=experts_str) + per_expert_lines = "\n".join( + f"- `{e}`: {_cue_shape_for(dispatch_map[e], e)}" for e in experts + ) + return ( + _TEMPLATE_MULTI_HEADER.format(experts=experts_str) + + per_expert_lines + + "\n" + + _TEMPLATE_MULTI_FOOTER + ) def modify_request(self, request: ModelRequest) -> ModelRequest: """Append the active-expert cue to the request's system message.""" experts = _read_active_teams() if not experts: return request + cue = self._cue_for(experts) + if not cue: + return request from .utils import append_to_system_message - new_system = append_to_system_message( - request.system_message, - self._cue_for(experts), - ) + new_system = append_to_system_message(request.system_message, cue) return request.override(system_message=new_system) def wrap_model_call( diff --git a/EvoScientist/middleware/expert_async_subagent.py b/EvoScientist/middleware/expert_async_subagent.py new file mode 100644 index 0000000..0d2920b --- /dev/null +++ b/EvoScientist/middleware/expert_async_subagent.py @@ -0,0 +1,272 @@ +"""Skill-name-injecting AsyncSubAgentMiddleware for expert dispatch. + +Upstream ``deepagents.AsyncSubAgentMiddleware`` hardcodes the invocation +input to ``{"messages": [{"role": "user", "content": description}]}`` — no +way for ``start_async_task`` to pass per-run state to the target graph. That +blocks the generic-container async pattern we need for agent-teams' expert +dispatch (one container graph, parameterised by which skill is active via +``skill_name`` in the initial state). + +Multiple community issues on the deepagents tracker target this gap +(``#2440``, ``#3838``, ``#4668``, ``#606``, ``#2512``) and the maintainers +have been closing implementation PRs (``#2617``, ``#3839``, ``#4669``) with +process-gate comments, none assigned. Upstream fix is not expected on any +predictable timeline; this subclass gives us the mechanism locally. + +Design +------ +- Subclass ``AsyncSubAgentMiddleware``; call ``super().__init__()`` for spec + validation + default 5-tool build, then swap in a start tool that injects + ``skill_name=subagent_type`` by construction (keeping check / update / + cancel / list unchanged). +- The tool signature matches upstream exactly: ``(description, subagent_type, + runtime)``. No LLM-visible ``payload`` field: every value the middleware + can derive itself (the skill name) is injected inside the middleware, not + entrusted to a channel the model can get wrong. Any run-specific + information the model uniquely holds (e.g. the desired ``output_path``) + belongs in the description string. +- Extend the ``AsyncSubAgent`` typed dict with an optional ``is_expert`` + marker so the middleware knows when to add ``skill_name`` to the run + input. Standard specs (``writing-agent`` / ``data-analysis-agent`` / + ``scheduler``) reach ``client.runs.create`` with the upstream shape. + +If deepagents ever lands a skill-name-passthrough of its own, delete this +file and rebind ``EvoAsyncSubAgentMiddleware`` → ``AsyncSubAgentMiddleware`` +in one commit; the state-schema shape on the container graph doesn't change. + +Do NOT add ``from __future__ import annotations`` to this module. langchain's +``StructuredTool._injected_args_keys`` uses ``inspect.signature(fn)`` (raw +annotations, not ``get_type_hints``) to decide which parameters are injected +runtime args. With PEP 563 in effect ``runtime: ToolRuntime`` becomes the +string ``"ToolRuntime"``, fails the ``issubclass(type_, _DirectlyInjectedToolArg)`` +check, and gets stripped from tool_input at parse time — the coroutine is +then called without ``runtime`` and raises ``TypeError``. +""" + +import logging +from datetime import UTC, datetime +from typing import Any, NotRequired + +from deepagents.middleware.async_subagents import ( + ASYNC_TASK_TOOL_DESCRIPTION, + AsyncSubAgent, + AsyncSubAgentMiddleware, + AsyncTask, + StartAsyncTaskSchema, + _build_cancel_tool, + _build_check_tool, + _build_list_tasks_tool, + _build_update_tool, + _ClientCache, + _validate_agent_type, +) +from langchain.tools import ToolRuntime +from langchain_core.messages import ToolMessage +from langchain_core.tools import StructuredTool +from langgraph.types import Command + +_logger = logging.getLogger(__name__) + + +class ExpertAsyncSubAgent(AsyncSubAgent): + """AsyncSubAgent spec extended with the expert-dispatch marker. + + Same wire fields as upstream ``AsyncSubAgent`` plus an internal + ``is_expert`` marker. Expert specs get ``skill_name`` injected into + the run input by construction so the shared container graph knows + which persona to load; standard specs reach ``runs.create`` with the + upstream shape. + """ + + is_expert: NotRequired[bool] + + +def _build_run_input( + spec: AsyncSubAgent, subagent_type: str, description: str +) -> dict[str, Any]: + """Build the ``input`` dict for ``client.runs.create``. + + ``skill_name`` is injected by construction for expert specs — never + accepted from the LLM, because the value is derivable from + ``subagent_type`` and every LLM-authored field is a field the LLM can + get wrong (silently overwriting ``messages`` was the pre-fix bug). + Standard specs (``writing-agent`` / ``data-analysis-agent`` / + ``scheduler``) reach ``runs.create`` with the upstream single-key shape. + """ + input_dict: dict[str, Any] = { + "messages": [{"role": "user", "content": description}] + } + if spec.get("is_expert"): + input_dict["skill_name"] = subagent_type + return input_dict + + +def _build_task_envelope( + subagent_type: str, thread_id: str, run_id: str, tool_call_id: str +) -> Command: + """Wrap a successful launch in the ``Command`` shape the router expects.""" + now = datetime.now(UTC).strftime("%Y-%m-%dT%H:%M:%SZ") + task: AsyncTask = { + "task_id": thread_id, + "agent_name": subagent_type, + "thread_id": thread_id, + "run_id": run_id, + "status": "running", + "created_at": now, + "last_checked_at": now, + "last_updated_at": now, + } + msg = f"Launched async subagent. task_id: {thread_id}" + return Command( + update={ + "messages": [ToolMessage(msg, tool_call_id=tool_call_id)], + "async_tasks": {thread_id: task}, + } + ) + + +def _build_expert_start_tool( + agent_map: dict[str, AsyncSubAgent], + clients: _ClientCache, + tool_description: str, +) -> StructuredTool: + """Build the skill-name-injecting ``start_async_task`` tool. + + Tool signature is upstream's exact shape (``description``, + ``subagent_type``, ``runtime``). For expert specs the middleware + injects ``skill_name=subagent_type`` into the run input before + dispatch, so the container graph resolves the right persona without + the model contributing (or being able to corrupt) that value. + """ + + def start_async_task( + description: str, + subagent_type: str, + runtime: ToolRuntime, + ) -> str | Command: + error = _validate_agent_type(agent_map, subagent_type) + if error: + return error + spec = agent_map[subagent_type] + input_dict = _build_run_input(spec, subagent_type, description) + try: + client = clients.get_sync(subagent_type) + thread = client.threads.create() + run = client.runs.create( + thread_id=thread["thread_id"], + assistant_id=spec["graph_id"], + input=input_dict, + ) + except Exception as e: + _logger.warning( + "Failed to launch async subagent '%s': %s", subagent_type, e + ) + return f"Failed to launch async subagent '{subagent_type}': {e}" + return _build_task_envelope( + subagent_type, thread["thread_id"], run["run_id"], runtime.tool_call_id + ) + + async def astart_async_task( + description: str, + subagent_type: str, + runtime: ToolRuntime, + ) -> str | Command: + error = _validate_agent_type(agent_map, subagent_type) + if error: + return error + spec = agent_map[subagent_type] + input_dict = _build_run_input(spec, subagent_type, description) + try: + client = clients.get_async(subagent_type) + thread = await client.threads.create() + run = await client.runs.create( + thread_id=thread["thread_id"], + assistant_id=spec["graph_id"], + input=input_dict, + ) + except Exception as e: + _logger.warning( + "Failed to launch async subagent '%s': %s", subagent_type, e + ) + return f"Failed to launch async subagent '{subagent_type}': {e}" + return _build_task_envelope( + subagent_type, thread["thread_id"], run["run_id"], runtime.tool_call_id + ) + + return StructuredTool.from_function( + name="start_async_task", + func=start_async_task, + coroutine=astart_async_task, + description=tool_description, + infer_schema=False, + args_schema=StartAsyncTaskSchema, + ) + + +class EvoAsyncSubAgentMiddleware(AsyncSubAgentMiddleware): + """AsyncSubAgentMiddleware with skill-name-injecting ``start_async_task``. + + Composes exactly like upstream — same constructor kwargs, same + ``system_prompt`` handling, same ``wrap_model_call`` / ``awrap_model_call``, + same tool signature (``description``, ``subagent_type``, ``runtime``). + Only difference: for expert specs (``is_expert=True``) the middleware + injects ``skill_name=subagent_type`` into ``client.runs.create(input=...)`` + so the shared container graph resolves the right persona. + + Existing async subagents (``writing-agent``, ``data-analysis-agent``, + ``scheduler``) work unchanged — they are declared without ``is_expert`` + and reach ``runs.create`` with the upstream single-key shape. + """ + + def __init__( + self, + *, + async_subagents: list[AsyncSubAgent], + system_prompt: str | None = None, + ) -> None: + # Install the model-passthrough patch BEFORE ``super().__init__(...)`` + # so upstream's ``_build_async_subagent_tools`` sees the patched + # ``_build_start_tool`` / ``_build_update_tool`` module attributes. + # Idempotent (guarded by ``_model_passthrough_patched`` in + # ``llm/patches.py``), so re-invocation on repeated middleware + # construction is a no-op. Without this, super()'s vanilla tools + # would still ignore ``cfg.model`` — including ``update_async_task``, + # which we inherit unchanged below. + from ..llm.patches import ( + _ClientCacheProxy, + _patch_deepagents_model_passthrough, + ) + + _patch_deepagents_model_passthrough() + + # Upstream's __init__ validates spec shape, builds the default 5-tool + # list, and composes the system_prompt. Delegate to it, then swap in + # the skill-name-injecting start tool. This wastes one tool-build cycle + # (~microseconds at construction) but avoids duplicating upstream's + # validation and system-prompt-composition logic. Pass ``system_prompt`` + # through unchanged — deepagents 0.7.0 dropped its ``ASYNC_TASK_SYSTEM_PROMPT`` + # default text; callers that want extra guidance in the async-task + # section of the prompt now supply it explicitly. + super().__init__( + async_subagents=async_subagents, + system_prompt=system_prompt, + ) + agent_map: dict[str, AsyncSubAgent] = {a["name"]: a for a in async_subagents} + # Wrap the client cache in ``_ClientCacheProxy`` so ``client.runs.create`` + # in our replacement start tool (and in the rebuilt check / update / + # cancel / list tools below) injects ``configurable.model`` / + # ``configurable.model_provider`` per run. ``_ClientCacheProxy`` exposes + # the same ``get_sync`` / ``get_async`` surface as ``_ClientCache``, so + # the upstream tool builders accept it without a type change. + clients = _ClientCacheProxy(_ClientCache(agent_map)) + agents_desc = "\n".join( + f"- {a['name']}: {a['description']}" for a in async_subagents + ) + launch_desc = ASYNC_TASK_TOOL_DESCRIPTION.format(available_agents=agents_desc) + self.tools = [ + _build_expert_start_tool(agent_map, clients, launch_desc), + _build_check_tool(clients), + _build_update_tool(agent_map, clients), + _build_cancel_tool(clients), + _build_list_tasks_tool(clients), + ] diff --git a/EvoScientist/middleware/utils.py b/EvoScientist/middleware/utils.py index 6d967c0..b8db047 100644 --- a/EvoScientist/middleware/utils.py +++ b/EvoScientist/middleware/utils.py @@ -104,3 +104,35 @@ def append_to_system_message( if system_message is None: return SystemMessage(content=new_blocks) return system_message.model_copy(update={"content": new_blocks}) + + +def replace_block_by_sentinel( + system_message: SystemMessage | None, + sentinel: str, + replacement: str, +) -> SystemMessage | None: + """Swap the block containing ``sentinel`` for ``replacement`` text. + + Iterates ``system_message.content_blocks`` and returns a new + ``SystemMessage`` whose block-list has the first block containing + ``sentinel`` replaced by ``{"type": "text", "text": replacement}``. + Other blocks and metadata (``additional_kwargs``, ``id``, ``name``, + ``response_metadata``) are preserved via ``model_copy``. + + Returns ``None`` when no block carries the sentinel — the caller + decides fallback policy (typically log-and-append rather than + hard-fail, so a deepagents base-stack refactor degrades gracefully + instead of killing the graph). + """ + if system_message is None: + return None + blocks = list(system_message.content_blocks) + for i, block in enumerate(blocks): + if isinstance(block, dict) and sentinel in block.get("text", ""): + new_blocks = [ + *blocks[:i], + {"type": "text", "text": replacement}, + *blocks[i + 1 :], + ] + return system_message.model_copy(update={"content": new_blocks}) + return None diff --git a/EvoScientist/stream/display.py b/EvoScientist/stream/display.py index 01c4b11..be4b2bf 100644 --- a/EvoScientist/stream/display.py +++ b/EvoScientist/stream/display.py @@ -1492,6 +1492,7 @@ def _run_streaming( on_stream_event=on_stream_event, status_footer_builder=status_footer_builder, metadata=metadata, + configurable_extra=configurable_extra, hitl_prompt_fn=hitl_prompt_fn, ask_user_prompt_fn=ask_user_prompt_fn, cancel_scope=cancel_scope, @@ -1824,6 +1825,7 @@ def _run_streaming( on_stream_event=on_stream_event, status_footer_builder=status_footer_builder, metadata=metadata, + configurable_extra=configurable_extra, hitl_prompt_fn=hitl_prompt_fn, ask_user_prompt_fn=ask_user_prompt_fn, cancel_scope=cancel_scope, diff --git a/EvoScientist/subagents/expert_container.py b/EvoScientist/subagents/expert_container.py index c7ef1b4..8ab96eb 100644 --- a/EvoScientist/subagents/expert_container.py +++ b/EvoScientist/subagents/expert_container.py @@ -166,17 +166,54 @@ def _reserved_subagent_names() -> frozenset[str]: return _reserved_subagent_names_cache -def list_dispatchable_experts(*, include_system: bool = True) -> list[SkillInfo]: - """Experts that will actually be dispatchable via ``task()``. +def is_async_dispatch_available(cfg: Any | None = None) -> bool: + """Return True when async-declared experts can actually be dispatched. - Combines ``list_expert_skills`` with the same filters - ``build_expert_subagent_specs`` (empty body) and - ``_fold_expert_subagents`` (name collision with yaml sub-agents or - ``general-purpose``) apply at construction time. Callers surfacing - experts to the user (e.g. the ``/expert`` slash command) should use - this instead of ``list_expert_skills`` directly, otherwise they can - accept a name that will silently misroute or per-turn error at - dispatch time. + Both gates must hold: ``enable_async_subagents`` opt-in AND langgraph + dev subprocess reachable. Same predicate + ``build_expert_async_subagent_specs`` uses at spec-build time, factored + out so ``list_dispatchable_experts`` (invite whitelist) and + ``ActiveTeamMiddleware`` (system-prompt cue) can honour it without + each re-deriving. + """ + if cfg is None: + from ..config import get_effective_config + + cfg = get_effective_config() + if not getattr(cfg, "enable_async_subagents", False): + return False + from ..langgraph_dev.manager import is_async_subagents_available + + return is_async_subagents_available() + + +def list_dispatchable_experts( + *, include_system: bool = True, cfg: Any | None = None +) -> list[SkillInfo]: + """Experts eligible for the ``/expert`` invite whitelist. + + Covers **both** sync (``task()``) and async (``start_async_task``) + dispatch shapes — an async-declared expert must remain invitable when + async dispatch is registered, or the active-team cue that instructs + ``start_async_task(...)`` never fires. Do NOT add a ``default_dispatch + == "async"`` exclusion here; the sync-specific exclusion lives in + ``build_expert_subagent_specs``, which is a different surface. + + Async-declared experts are dropped when async dispatch is NOT + registered (``enable_async_subagents=false`` or langgraph dev + unreachable). Advertising an expert that resolves to a + ``start_async_task`` tool that either doesn't exist or doesn't list + the expert is worse than an honest refusal — see reviewer thread on + PR #391. + + Combines ``list_expert_skills`` with the two filters that + construction-time paths already apply (empty body via + ``build_expert_subagent_specs``; name collision with yaml sub-agents + or ``general-purpose`` via ``_fold_expert_subagents``). Callers + surfacing experts to the user (e.g. the ``/expert`` slash command) + should use this instead of ``list_expert_skills`` directly, otherwise + they can accept a name that will silently misroute or per-turn error + at dispatch time. Read-only filter — construction-time warnings for empty-body / colliding experts are emitted by ``build_expert_subagent_specs`` and @@ -185,12 +222,15 @@ def list_dispatchable_experts(*, include_system: bool = True) -> list[SkillInfo] from ..tools.skills_manager import list_expert_skills reserved = _reserved_subagent_names() + async_available = is_async_dispatch_available(cfg=cfg) dispatchable: list[SkillInfo] = [] for info in list_expert_skills(include_system=include_system): if not _body_of(info).strip(): continue if info.name in reserved: continue + if info.default_dispatch == "async" and not async_available: + continue dispatchable.append(info) return dispatchable @@ -200,7 +240,7 @@ def build_expert_subagent_specs( *, include_system: bool = True, ) -> list[dict[str, Any]]: - """Build spec dicts for every installed expert skill. + """Build spec dicts for every installed sync-dispatched expert skill. Thin wrapper over ``list_expert_skills()`` + ``build_expert_subagent_spec``. Called by the main-agent construction path (``_build_base_kwargs``) to @@ -210,11 +250,20 @@ def build_expert_subagent_specs( personaless expert advertised in the ``task`` tool schema would let the orchestrator dispatch to a blank system prompt, a worse failure mode than the expert being absent. + + Experts declared with ``default_dispatch: async`` are excluded — those + are folded via ``build_expert_async_subagent_specs`` in + ``expert_container_async.py`` and reached through + ``EvoAsyncSubAgentMiddleware.start_async_task`` instead of the sync + ``task`` tool. A skill that appears in both lists would produce two + competing tool schemas for the same subagent name. """ from ..tools.skills_manager import list_expert_skills specs: list[dict[str, Any]] = [] for info in list_expert_skills(include_system=include_system): + if info.default_dispatch == "async": + continue if not _body_of(info).strip(): _logger.warning( "Expert skill %r: SKILL.md body is empty; skipping registration.", diff --git a/EvoScientist/subagents/expert_container_async.py b/EvoScientist/subagents/expert_container_async.py new file mode 100644 index 0000000..e0b75fa --- /dev/null +++ b/EvoScientist/subagents/expert_container_async.py @@ -0,0 +1,348 @@ +"""Async container graph for expert-skill dispatch. + +One generic graph that reads ``skill_name`` from initial state and loads +that expert skill's ``SKILL.md`` body as the sub-agent's system prompt at +invocation time. Registered once in ``langgraph.json``; parameterised per +run via the payload the main agent's ``start_async_task`` passes through +:class:`EvoScientist.middleware.expert_async_subagent.EvoAsyncSubAgentMiddleware`. + +Why one shared graph rather than one graph per expert: the per-expert +alternative (each expert registered statically in ``langgraph.json``) +doesn't scale to ``skill_manager install `` at runtime, because a +new expert would need a repo edit + a langgraph dev restart. The +generic-container approach preserves the "installable async expert" +story. + +State schema +------------ +Extends ``DeepAgentState`` with ``skill_name`` and ``output_path`` as +optional keys. Both arrive via the ``payload`` the main agent passes; the +middleware validates presence and halts with a clear error message when +either is missing (rather than falling back to an ambient default that +would silently produce the wrong survey). +""" + +from __future__ import annotations + +import logging +from collections.abc import Awaitable, Callable +from typing import Any, NotRequired + +from deepagents.graph import DeepAgentState +from langchain.agents.middleware.types import ( + AgentMiddleware, + ModelRequest, + ModelResponse, +) +from langchain_core.messages import SystemMessage + +_logger = logging.getLogger(__name__) + +# Sentinel token embedded in ``_FALLBACK_SYSTEM_PROMPT`` so +# :class:`ExpertSkillLoaderMiddleware` can locate that block inside the +# base-stack-composed ``system_message`` and swap it for the persona in +# place — preserving the base-stack sections (todo guidance, ``## `task```, +# ``## Filesystem Tools``, ``## Skills System``, ``## Async subagents``, +# ...) that would otherwise be dropped by an ``override(system_message=...)``. +# Double-underscore ASCII cannot occur incidentally in skill bodies or +# base-stack prose, so a substring scan is precise. HTML-comment tokens +# were rejected because pretty-printers can strip them and Claude has +# been observed to echo them back into prose. +_PERSONA_SENTINEL = "__EVOEXPERT_PERSONA_SLOT__" + +_FALLBACK_SYSTEM_PROMPT = ( + f"{_PERSONA_SENTINEL}\n\n" + "Fallback: the expert loader middleware failed to resolve ``skill_name``. " + "Return an error envelope naming the failure and halt." +) + + +class ExpertContainerState(DeepAgentState): + """State schema for the async expert container graph. + + Adds ``skill_name`` as a ``NotRequired`` key — + ``EvoAsyncSubAgentMiddleware`` injects it by construction from + ``subagent_type`` so the loader middleware knows which persona to + load. ``output_path`` is NOT in state; the main agent embeds the + desired artifact path in the task description (natural language) and + the expert's SKILL.md contract instructs the LLM to pin it in + ``write_todos`` on turn 1 — surviving summarization via langgraph's + todo composition. + """ + + skill_name: NotRequired[str] + + +class ExpertSkillLoaderMiddleware(AgentMiddleware[Any, Any, Any]): + """Load the expert skill's SKILL.md body as system prompt on every model call. + + Reads ``state.skill_name``, resolves the corresponding installed expert + skill via ``list_expert_skills()``, composes the system message from + ``role`` + SKILL.md body (mirrors the sync path in + ``EvoScientist.subagents.expert_container._compose_system_prompt``), + and overrides ``request.system_message`` before the handler runs. + + The container graph's static ``system_prompt`` at construction time is a + minimal fallback; this middleware is the load-bearing component. If the + skill_name is missing or resolves to no installed skill, appends an + explicit error to the system message so the LLM immediately halts and + returns an error envelope (rather than answering as an ambient + generalist). + """ + + name = "expert_skill_loader" + + def _compose_prompt(self, state: dict[str, Any]) -> str: + """Look up the skill and compose its system prompt. + + Returns the composed prompt string, or an error-cue string when the + skill can't be loaded. Never raises — errors are surfaced through + the LLM's system prompt so it can return a well-formed error + envelope rather than crash the graph mid-turn. + + Emits a "Runtime context" tail re-asserting ``skill_name`` on every + model call so the expert knows its own persona name after + summarization. Any run-specific values the LLM needs (output path, + user goal, ...) travel via the initial user message — the middleware + does not touch them. + """ + skill_name = state.get("skill_name") + if not skill_name: + return ( + "ERROR: The async expert container was invoked without a " + "``skill_name`` in state. This is a wiring bug in whichever " + "middleware invoked ``start_async_task``. Return an error " + "envelope naming the missing field and halt." + ) + + # Lazy import — the loader is a per-turn call, so this stays cheap. + from ..tools.skills_manager import list_expert_skills + + experts = list_expert_skills(include_system=True) + match = next((s for s in experts if s.name == skill_name), None) + if match is None: + installed = ", ".join(sorted(s.name for s in experts)) or "(none)" + return ( + f"ERROR: Expert skill '{skill_name}' is not installed. " + f"Installed experts: {installed}. Return an error envelope " + "with status='error' explaining the skill is missing." + ) + + # Mirrors the sync fold-in's empty-body skip in + # ``expert_container.py::build_expert_subagent_specs``. A skill with + # only a role line (or nothing) would otherwise run against a + # persona-less system prompt — a worse failure mode than the expert + # being absent. Prefer a well-formed error envelope over silent + # nonsense. + if not (match.body or "").strip(): + return ( + f"ERROR: Expert skill '{skill_name}' has an empty SKILL.md " + "body — the persona / pipeline the sub-agent needs is missing. " + "This is a skill-authoring bug; the sub-agent cannot proceed. " + "Return an error envelope with status='error' naming the " + "empty skill." + ) + + # Compose: role prepend (if present) + body + runtime-context tail. + body = match.body or "" + head = f"You are {match.role}.\n\n" if match.role else "" + runtime_block = ( + "\n---\n\n" + "## Runtime context (injected by the container)\n\n" + f"- ``skill_name``: ``{skill_name}``\n" + ) + return (head + body).rstrip() + runtime_block + + def _compose_system_message(self, request: ModelRequest[Any]) -> SystemMessage: + """Return the base-stack ``system_message`` with the persona swapped in. + + Locates the fallback block by ``_PERSONA_SENTINEL`` and replaces its + text with the composed persona, preserving every other block + (``## `task```, ``## Filesystem Tools``, ``## Skills System``, …). + When the sentinel isn't found (e.g. deepagents renamed the + ``system_prompt=`` handling), degrade to appending the persona so + the base-stack sections stay live and the expert stays operational; + surface the drift in logs so it can't rot silently. + """ + from ..middleware.utils import ( + append_to_system_message, + replace_block_by_sentinel, + ) + + composed = self._compose_prompt(request.state) + new_system = replace_block_by_sentinel( + request.system_message, _PERSONA_SENTINEL, composed + ) + if new_system is None: + _logger.warning( + "Expert persona sentinel %r not found in system_message; " + "appending persona instead of replacing the fallback block. " + "Deepagents' base-stack composition may have changed.", + _PERSONA_SENTINEL, + ) + new_system = append_to_system_message(request.system_message, composed) + return new_system + + def wrap_model_call( + self, + request: ModelRequest[Any], + handler: Callable[[ModelRequest[Any]], ModelResponse[Any]], + ) -> ModelResponse[Any]: + return handler( + request.override(system_message=self._compose_system_message(request)) + ) + + async def awrap_model_call( + self, + request: ModelRequest[Any], + handler: Callable[[ModelRequest[Any]], Awaitable[ModelResponse[Any]]], + ) -> ModelResponse[Any]: + return await handler( + request.override(system_message=self._compose_system_message(request)) + ) + + +def build_expert_async_subagent_specs(cfg: Any | None = None) -> list[dict[str, Any]]: + """Build ``AsyncSubAgent``-shaped specs for every ``default_dispatch: async`` expert. + + Each spec is a dict pointing at the shared ``expert-container-async`` graph + with ``is_expert=True``. The main agent's + ``EvoAsyncSubAgentMiddleware.start_async_task`` uses ``is_expert`` to + require a ``payload`` including ``skill_name``. + + Returns an empty list when ``cfg.enable_async_subagents`` is not set, or + when the langgraph dev subprocess isn't reachable. Both are the same + conditions used by ``_maybe_swap_async_subagents`` to gate the existing + ``writing-agent`` / ``data-analysis-agent`` / ``scheduler`` async + subagents — keeps behaviour consistent across sync-fallback situations. + """ + from ..config import get_effective_config + + cfg = cfg if cfg is not None else get_effective_config() + if not getattr(cfg, "enable_async_subagents", False): + return [] + # Same reachability guard used for standard async subagents in + # ``_maybe_swap_async_subagents``. + from ..langgraph_dev.manager import is_async_subagents_available + + if not is_async_subagents_available(): + return [] + + from ..tools.skills_manager import list_expert_skills + from .expert_container import _reserved_subagent_names + + port = int(getattr(cfg, "langgraph_dev_port", 6174)) + # Mirror the sync fold-in's name guard in ``_fold_expert_subagents``: + # yaml async sub-agents (``writing-agent``, ``data-analysis-agent``, + # ``scheduler``, …) and ``general-purpose`` share the async spec pool + # by name, so an expert skill named after any of them would raise + # ``ValueError: Duplicate async subagent names`` inside + # ``AsyncSubAgentMiddleware.__init__`` and kill CLI startup. The local + # ``seen`` set catches the second failure trigger: two workspace-tier + # skills that share the same frontmatter ``name`` (workspace listing + # uses ``check_seen=False``, so both survive to this point). + taken = set(_reserved_subagent_names()) + specs: list[dict[str, Any]] = [] + for skill in list_expert_skills(include_system=True): + if skill.default_dispatch != "async": + continue + # Same empty-body skip the sync fold-in enforces in + # ``expert_container.py::build_expert_subagent_specs``. Advertising + # a body-less expert in ``start_async_task``'s tool schema, then + # rejecting it at loader time, wastes a launch round-trip; filter + # upstream so ``start_async_task`` never sees the broken skill. + if not (skill.body or "").strip(): + _logger.warning( + "Expert skill %r: SKILL.md body is empty; skipping " + "async-dispatch registration.", + skill.name, + ) + continue + if skill.name in taken: + _logger.warning( + "Expert skill %r collides with an existing async sub-agent " + "name; skipping async-dispatch registration.", + skill.name, + ) + continue + taken.add(skill.name) + specs.append( + { + "name": skill.name, + "description": skill.description, + "graph_id": "expert-container-async", + "url": f"http://localhost:{port}", + "is_expert": True, + } + ) + return specs + + +def build_expert_container_async_graph() -> Any: + """Build the async expert container graph. + + Called once at langgraph dev startup. The returned graph accepts + ``{messages, skill_name}`` as initial state; the + :class:`ExpertSkillLoaderMiddleware` resolves ``skill_name`` on every + model call and injects the matching SKILL.md body as system prompt. + + Tool set is intentionally minimal (``think_tool`` only) — matches the + sync ``expert_container`` factory. Once the per-skill ``allowed-tools`` + follow-up ships, the tool list will union with the skill's declared + tools. + """ + from deepagents import create_deep_agent + + from ..config import apply_config_to_env, get_effective_config + from ..EvoScientist import ( + _ensure_chat_model, + _ensure_general_purpose_subagent, + _get_default_backend, + _get_default_middleware, + _inject_subagent_middleware, + ) + from ..tools import think_tool + + cfg = get_effective_config() + apply_config_to_env(cfg) + + subagents: list[dict[str, Any]] = [] + # The async expert runs as its own graph inside langgraph-dev, so it + # does NOT inherit the main agent's subagent stack. Without at least + # one dispatchable target, the expert's ``task()`` tool has no target + # to call — long-horizon experts (e.g. literature-review's Phase 2 + # ``paper-navigator`` fan-out) rely on this. ``general-purpose`` is + # deepagents' default sub-agent and provides the fan-out capability + # symmetric with what the main agent has by default. + # + # Trade-off: an expert LLM that never uses ``task()`` still pays the + # schema-description tokens on every model call. Acceptable at v1; if + # experts routinely abuse the fan-out or ignore it entirely, a + # per-skill ``include_general_purpose: false`` frontmatter opt-out + # is a natural follow-up. + _ensure_general_purpose_subagent(subagents) + _inject_subagent_middleware(subagents) + + middleware = [ + # Loader runs FIRST so downstream middleware sees the composed + # system_message. Ordering matters — put ExpertSkillLoaderMiddleware + # before context editing / error normalisation so they operate on + # the already-composed prompt. + ExpertSkillLoaderMiddleware(), + *_get_default_middleware( + for_async_subagent=True, + memory_source_agent="expert-container-async", + ), + ] + + return create_deep_agent( + name="expert-container-async", + model=_ensure_chat_model(), + system_prompt=_FALLBACK_SYSTEM_PROMPT, + tools=[think_tool], + skills=["/skills/"], + backend=_get_default_backend(), + middleware=middleware, + subagents=subagents, + state_schema=ExpertContainerState, + ).with_config({"recursion_limit": cfg.recursion_limit}) diff --git a/EvoScientist/tools/skill_manager.py b/EvoScientist/tools/skill_manager.py index 37c58cd..bb2b62b 100644 --- a/EvoScientist/tools/skill_manager.py +++ b/EvoScientist/tools/skill_manager.py @@ -194,16 +194,16 @@ def skill_manager( f"Skill not found: {name}. " f"Use action='list' with include_system=True to see all available skills." ) - lines = [ - f"Name: {info.name}", - f"Description: {info.description}", - f"Source: {info.source}", - ] # ``Path`` reports the sandbox-visible virtual mount segment (the skill's # directory name under ``/skills/``), not the host filesystem path. # Surfacing the host path (e.g. ``/home/.../EvoScientist/skills/``) # invited ``cd && …`` chains that fail in the sandbox. - lines.append(f"Path: /skills/{info.path.name}") + lines = [ + f"Name: {info.name}", + f"Description: {info.description}", + f"Source: {info.source}", + f"Path: /skills/{info.path.name}", + ] if info.tags: lines.append(f"Tags: {', '.join(info.tags)}") # Expert-skill surface (agent-teams v1): only shown when the diff --git a/EvoScientist/tools/skills_manager.py b/EvoScientist/tools/skills_manager.py index 34931c9..a995b99 100644 --- a/EvoScientist/tools/skills_manager.py +++ b/EvoScientist/tools/skills_manager.py @@ -118,7 +118,7 @@ class SkillInfo: byline: str = "" # WebUI gallery byline capability_tags: list[str] = field(default_factory=list) # WebUI chips avatar_hint: str = "" # WebUI icon hint - default_dispatch: str = "" # "sync" | "panel" (expert skills only) + default_dispatch: str = "" # "sync" | "panel" | "async" (expert skills only) # SKILL.md body (post-frontmatter). Populated by ``_parse_skill_md`` so # the expert-container factory doesn't have to re-read the file on every # main-agent construction. Empty for skills built by hand or when the @@ -415,7 +415,20 @@ def _parse_skill_md(skill_md_path: Path, *, source: str = "") -> SkillInfo: capability_tags = _normalize_tags(frontmatter.get("capability_tags")) avatar_hint = frontmatter.get("avatar_hint") or "" raw_dispatch = frontmatter.get("default_dispatch") - default_dispatch = raw_dispatch if raw_dispatch in ("sync", "panel") else "" + default_dispatch = ( + raw_dispatch if raw_dispatch in ("sync", "panel", "async") else "" + ) + # Mirror the ``raw_type`` handling above. A typo like + # ``default_dispatch: asnyc`` would otherwise be indistinguishable + # from an unset field and would register the skill under sync + # dispatch — the opposite of the author's intent. + if raw_dispatch is not None and raw_dispatch != default_dispatch: + _logger.warning( + "Skill %r: unrecognized default_dispatch %r in frontmatter; " + "treating as unset (skill will register under sync dispatch).", + frontmatter.get("name") or parent.name, + raw_dispatch, + ) # ``.get("name", parent.name)`` only defaults on missing key; a # present-but-empty ``name:`` yields None, which would flow into # ``SkillInfo.name`` and slip past the ``_fold_expert_subagents`` diff --git a/tests/test_active_team_middleware.py b/tests/test_active_team_middleware.py index 7666460..3ad140e 100644 --- a/tests/test_active_team_middleware.py +++ b/tests/test_active_team_middleware.py @@ -120,11 +120,24 @@ def test_middleware_no_op_when_active_teams_empty_list(mock_get_config): assert modified is request +def _mock_expert(name: str, dispatch: str) -> MagicMock: + """Build a MagicMock ``SkillInfo`` with the given dispatch shape. + + ``name`` on ``MagicMock`` must be set via attribute assignment; passing + ``name=`` to the constructor names the mock instance itself. + """ + info = MagicMock(default_dispatch=dispatch) + info.name = name + return info + + +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") @patch("langgraph.config.get_config") -def test_middleware_appends_single_expert_cue(mock_get_config): +def test_middleware_appends_single_expert_cue(mock_get_config, mock_dispatchable): mock_get_config.return_value = { "configurable": {"active_teams": ["idea-brainstorm"]}, } + mock_dispatchable.return_value = [_mock_expert("idea-brainstorm", "sync")] middleware = ActiveTeamMiddleware() modified = middleware.modify_request(_request()) text = _system_text(modified) @@ -134,31 +147,148 @@ def test_middleware_appends_single_expert_cue(mock_get_config): assert "base system" in text # original preserved +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") @patch("langgraph.config.get_config") -def test_middleware_appends_multi_expert_cue(mock_get_config): +def test_middleware_appends_multi_expert_cue(mock_get_config, mock_dispatchable): mock_get_config.return_value = { "configurable": {"active_teams": ["idea-brainstorm", "literature-review"]}, } + mock_dispatchable.return_value = [ + _mock_expert("idea-brainstorm", "sync"), + _mock_expert("literature-review", "sync"), + ] middleware = ActiveTeamMiddleware() modified = middleware.modify_request(_request()) text = _system_text(modified) assert "" in text assert "`idea-brainstorm`" in text assert "`literature-review`" in text - assert "Consult any of them" in text + # Multi-cue: header names both experts, then per-expert dispatch lines follow. + assert "The user has invited the following experts" in text + assert "Per-expert dispatch" in text assert "base system" in text +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") @patch("langgraph.config.get_config") -def test_middleware_appends_cue_for_unknown_expert_names(mock_get_config): - """Middleware doesn't validate names against the registry; main decides.""" +def test_middleware_omits_cue_for_undispatchable_names( + mock_get_config, mock_dispatchable +): + """Names not in ``list_dispatchable_experts`` are dropped from the cue. + + Covers uninstalled experts, empty-body experts, name collisions, and + async-declared experts when async dispatch is unavailable — anything + the model would find missing at dispatch time. + """ mock_get_config.return_value = { "configurable": {"active_teams": ["nonexistent-expert"]}, } + mock_dispatchable.return_value = [] # nothing dispatchable + request = _request() + middleware = ActiveTeamMiddleware() + modified = middleware.modify_request(request) + # No cue appended — modify_request returns the original request untouched. + assert modified is request + + +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") +@patch("langgraph.config.get_config") +def test_middleware_uses_start_async_task_cue_for_async_dispatch( + mock_get_config, mock_dispatchable +): + """An expert declared ``default_dispatch: async`` gets the async cue. + + Only reaches the cue when async dispatch is actually registered — the + honest-surface filter in ``list_dispatchable_experts`` drops + async-declared experts otherwise. + """ + mock_get_config.return_value = { + "configurable": {"active_teams": ["literature-review"]}, + } + mock_dispatchable.return_value = [_mock_expert("literature-review", "async")] + middleware = ActiveTeamMiddleware() modified = middleware.modify_request(_request()) text = _system_text(modified) - assert "`nonexistent-expert`" in text + assert "" in text + assert "start_async_task(" in text + assert "subagent_type: 'literature-review'" in text + # Post-X-4: no payload dict. The cue instructs the main agent to embed + # the desired output path directly in the description string. + assert "payload" not in text + assert "output path" in text.lower() or "output_path" in text + assert "check_async_task" in text + # Sync cue must NOT be advertised for async experts. + assert "Consult it via `task(" not in text + + +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") +@patch("langgraph.config.get_config") +def test_middleware_uses_task_cue_for_sync_dispatch(mock_get_config, mock_dispatchable): + """Sync-dispatched experts get the ``task()`` cue.""" + mock_get_config.return_value = { + "configurable": {"active_teams": ["idea-brainstorm"]}, + } + mock_dispatchable.return_value = [_mock_expert("idea-brainstorm", "sync")] + + middleware = ActiveTeamMiddleware() + modified = middleware.modify_request(_request()) + text = _system_text(modified) + assert "Consult it via `task(" in text + assert "runs synchronously" in text + # No async-specific fragments for a sync expert. + assert "start_async_task(" not in text + assert "output_path" not in text + + +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") +@patch("langgraph.config.get_config") +def test_middleware_multi_mixed_dispatch(mock_get_config, mock_dispatchable): + """When both sync and async experts are active, each gets its own cue.""" + mock_get_config.return_value = { + "configurable": {"active_teams": ["idea-brainstorm", "literature-review"]}, + } + mock_dispatchable.return_value = [ + _mock_expert("idea-brainstorm", "sync"), + _mock_expert("literature-review", "async"), + ] + + middleware = ActiveTeamMiddleware() + modified = middleware.modify_request(_request()) + text = _system_text(modified) + # Both cue shapes appear once each in the per-expert block. + assert text.count("`task(") == 1 + assert text.count("start_async_task(") == 1 + assert "`idea-brainstorm`:" in text + assert "`literature-review`:" in text + + +@patch("EvoScientist.subagents.expert_container.list_dispatchable_experts") +@patch("langgraph.config.get_config") +def test_middleware_drops_invited_expert_that_is_not_dispatchable( + mock_get_config, mock_dispatchable +): + """An async-declared expert stays invited across a config change, but + when async dispatch turns unavailable it drops out of + ``list_dispatchable_experts``. The cue must not mention it — otherwise + the model is told to reach for a tool that either doesn't exist or + doesn't list the expert.""" + mock_get_config.return_value = { + "configurable": { + "active_teams": ["idea-brainstorm", "literature-review"], + }, + } + # literature-review invited but not dispatchable this turn. + mock_dispatchable.return_value = [_mock_expert("idea-brainstorm", "sync")] + + middleware = ActiveTeamMiddleware() + modified = middleware.modify_request(_request()) + text = _system_text(modified) + # Single-cue shape (only one expert survived the filter). + assert "" in text + assert "`idea-brainstorm`" in text + assert "literature-review" not in text + assert "start_async_task(" not in text @patch("langgraph.config.get_config", side_effect=RuntimeError("outside context")) diff --git a/tests/test_expert_async_subagent.py b/tests/test_expert_async_subagent.py new file mode 100644 index 0000000..c0a7535 --- /dev/null +++ b/tests/test_expert_async_subagent.py @@ -0,0 +1,340 @@ +"""Tests for the skill-name-injecting AsyncSubAgentMiddleware subclass.""" + +from __future__ import annotations + +from types import SimpleNamespace +from unittest.mock import AsyncMock, MagicMock, patch + +import pytest + +from EvoScientist.middleware.expert_async_subagent import ( + EvoAsyncSubAgentMiddleware, + _build_run_input, +) + + +class _TestPayloadValidationRemoved: + """Placeholder — the ``_payload_validation_error`` helper was deleted + when ``payload`` was dropped from the tool schema (PR #391 review, X-4). + The seven tests that lived here (``TestPayloadValidation``) no longer + apply: subagent_type is validated by ``_validate_agent_type``, + ``skill_name`` is injected by construction, and no other user-supplied + fields reach ``client.runs.create(input=...)``. See + ``TestBuildRunInput`` below and ``TestStartToolInvocation`` for the + replacement coverage. + """ + + +# ============================================================================= +# _build_run_input — the shared input-dict factory +# ============================================================================= + + +class TestBuildRunInput: + """``skill_name`` is injected for expert specs, absent for standard specs. + The description always lands in ``messages`` verbatim — no LLM-authored + key can overwrite it (was the pre-fix bug when ``payload`` was in scope). + """ + + def test_expert_spec_injects_skill_name(self): + spec = {"name": "e", "graph_id": "g", "is_expert": True} + result = _build_run_input(spec, "literature-review", "write a survey") + assert result == { + "messages": [{"role": "user", "content": "write a survey"}], + "skill_name": "literature-review", + } + + def test_standard_spec_matches_upstream_shape(self): + """Standard specs (writing-agent, scheduler, ...) reach ``runs.create`` + with the upstream single-key shape — no ``skill_name`` injected.""" + spec = {"name": "writing-agent", "graph_id": "writing_agent"} + result = _build_run_input(spec, "writing-agent", "hi") + assert result == {"messages": [{"role": "user", "content": "hi"}]} + + def test_is_expert_false_treated_as_standard(self): + """Explicit ``is_expert=False`` matches the default (absent) behaviour.""" + spec = {"name": "std", "graph_id": "writing_agent", "is_expert": False} + result = _build_run_input(spec, "std", "hi") + assert result == {"messages": [{"role": "user", "content": "hi"}]} + + def test_description_lands_verbatim(self): + """Regression guard against the pre-fix bug where an LLM-authored + ``payload`` could overwrite ``messages`` — description now travels + through a channel the LLM cannot corrupt.""" + spec = {"name": "e", "graph_id": "g", "is_expert": True} + result = _build_run_input( + spec, "e", "write to ./artifacts/e/foo.md a summary of X" + ) + assert result["messages"][0]["content"] == ( + "write to ./artifacts/e/foo.md a summary of X" + ) + + +# ============================================================================= +# EvoAsyncSubAgentMiddleware — end-to-end tool invocation +# ============================================================================= + + +def _standard_spec(): + return { + "name": "writing-agent", + "description": "std writer", + "graph_id": "writing_agent", + } + + +def _expert_spec(): + return { + "name": "literature-review", + "description": "expert lit review", + "graph_id": "expert_container", + "is_expert": True, + } + + +class TestMiddlewareConstruction: + def test_middleware_has_five_tools(self): + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + names = [t.name for t in mw.tools] + assert set(names) == { + "start_async_task", + "check_async_task", + "update_async_task", + "cancel_async_task", + "list_async_tasks", + } + + def test_start_tool_schema_matches_upstream(self): + """The tool signature returned to upstream's exact shape when + ``payload`` was dropped — schema is now ``deepagents``'s + ``StartAsyncTaskSchema``.""" + from deepagents.middleware.async_subagents import StartAsyncTaskSchema + + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + assert start.args_schema is StartAsyncTaskSchema + + def test_construction_rejects_empty_subagents(self): + with pytest.raises(ValueError, match="At least one async subagent"): + EvoAsyncSubAgentMiddleware(async_subagents=[]) + + def test_construction_rejects_duplicate_names(self): + with pytest.raises(ValueError, match="Duplicate"): + EvoAsyncSubAgentMiddleware( + async_subagents=[_standard_spec(), _standard_spec()] + ) + + +def _fake_sync_client(): + client = MagicMock() + client.threads.create.return_value = {"thread_id": "task-abc"} + client.runs.create.return_value = {"run_id": "run-xyz"} + return client + + +def _fake_async_client(): + client = MagicMock() + client.threads.create = AsyncMock(return_value={"thread_id": "task-abc"}) + client.runs.create = AsyncMock(return_value={"run_id": "run-xyz"}) + return client + + +class TestStartToolInvocation: + """Direct invocation of the start tool's sync function. + + Mocks ``_ClientCache.get_sync`` so we can assert on the ``input`` dict + handed to ``runs.create`` without any real network round-trip. + """ + + def test_start_injects_skill_name_for_expert_spec(self): + """The middleware sets ``input_dict['skill_name'] = subagent_type`` + by construction — the shared container graph resolves the right + persona without a payload dict crossing the LLM channel.""" + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_expert_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_sync_client() + with patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync", + return_value=client, + ): + result = start.func( + description="write to ./artifacts/literature-review/attn.md a survey on X", + subagent_type="literature-review", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + + client.runs.create.assert_called_once() + kwargs = client.runs.create.call_args.kwargs + assert kwargs["assistant_id"] == "expert_container" + assert kwargs["input"]["messages"] == [ + { + "role": "user", + "content": ( + "write to ./artifacts/literature-review/attn.md a survey on X" + ), + } + ] + assert kwargs["input"]["skill_name"] == "literature-review" + assert "payload" not in kwargs["input"] + assert "output_path" not in kwargs["input"] + # Return value stamps the task into async_tasks state. + assert "async_tasks" in result.update + assert "task-abc" in result.update["async_tasks"] + + def test_start_injects_cfg_model_into_configurable(self): + """cfg.model / cfg.provider land in ``config.configurable`` on every + ``runs.create`` so the deployed graph re-resolves its chat model per + run instead of using whatever was baked at container-build time. + Without this the ``/model`` CLI switch silently doesn't propagate to + expert launches. + """ + from EvoScientist.config.settings import EvoScientistConfig + + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_expert_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_sync_client() + fake_cfg = EvoScientistConfig(model="test-model-abc", provider="test-provider") + with ( + patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync", + return_value=client, + ), + patch("EvoScientist.EvoScientist._ensure_config", return_value=fake_cfg), + ): + start.func( + description="w", + subagent_type="literature-review", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + + kwargs = client.runs.create.call_args.kwargs + assert "config" in kwargs + configurable = kwargs["config"]["configurable"] + assert configurable["model"] == "test-model-abc" + assert configurable["model_provider"] == "test-provider" + + def test_start_standard_spec_matches_upstream_input_shape(self): + """Standard subagents (writing-agent, scheduler, ...) reach + ``runs.create`` with the upstream single-key ``messages`` shape.""" + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_sync_client() + with patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync", + return_value=client, + ): + start.func( + description="hi", + subagent_type="writing-agent", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + kwargs = client.runs.create.call_args.kwargs + assert kwargs["input"] == {"messages": [{"role": "user", "content": "hi"}]} + + def test_start_unknown_subagent_returns_error(self): + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + result = start.func( + description="hi", + subagent_type="does-not-exist", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + assert isinstance(result, str) + assert "Unknown async subagent type" in result + + +class TestAstartToolInvocation: + """Mirror ``TestStartToolInvocation`` against ``astart_async_task`` — the + coroutine langgraph_api actually runs in production. Pre-fix zero + coverage: X-iZhang flagged that a fix applied only to the sync body + would leave tests green and production broken.""" + + @pytest.mark.asyncio + async def test_astart_injects_skill_name_for_expert_spec(self): + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_expert_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_async_client() + with patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_async", + return_value=client, + ): + result = await start.coroutine( + description="write to ./artifacts/literature-review/attn.md a survey on X", + subagent_type="literature-review", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + + client.runs.create.assert_awaited_once() + kwargs = client.runs.create.await_args.kwargs + assert kwargs["assistant_id"] == "expert_container" + assert kwargs["input"]["skill_name"] == "literature-review" + assert kwargs["input"]["messages"][0]["content"].startswith( + "write to ./artifacts/literature-review/attn.md" + ) + assert "payload" not in kwargs["input"] + assert "async_tasks" in result.update + assert "task-abc" in result.update["async_tasks"] + + @pytest.mark.asyncio + async def test_astart_injects_cfg_model_into_configurable(self): + from EvoScientist.config.settings import EvoScientistConfig + + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_expert_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_async_client() + fake_cfg = EvoScientistConfig(model="test-model-abc", provider="test-provider") + with ( + patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_async", + return_value=client, + ), + patch("EvoScientist.EvoScientist._ensure_config", return_value=fake_cfg), + ): + await start.coroutine( + description="w", + subagent_type="literature-review", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + + kwargs = client.runs.create.await_args.kwargs + assert "config" in kwargs + configurable = kwargs["config"]["configurable"] + assert configurable["model"] == "test-model-abc" + assert configurable["model_provider"] == "test-provider" + + @pytest.mark.asyncio + async def test_astart_standard_spec_matches_upstream_input_shape(self): + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + client = _fake_async_client() + with patch( + "EvoScientist.middleware.expert_async_subagent._ClientCache.get_async", + return_value=client, + ): + await start.coroutine( + description="hi", + subagent_type="writing-agent", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + kwargs = client.runs.create.await_args.kwargs + assert kwargs["input"] == {"messages": [{"role": "user", "content": "hi"}]} + + @pytest.mark.asyncio + async def test_astart_unknown_subagent_returns_error(self): + mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()]) + start = next(t for t in mw.tools if t.name == "start_async_task") + + result = await start.coroutine( + description="hi", + subagent_type="does-not-exist", + runtime=SimpleNamespace(tool_call_id="tc1"), + ) + assert isinstance(result, str) + assert "Unknown async subagent type" in result diff --git a/tests/test_expert_container.py b/tests/test_expert_container.py index 107976f..6299951 100644 --- a/tests/test_expert_container.py +++ b/tests/test_expert_container.py @@ -3,6 +3,7 @@ from __future__ import annotations from pathlib import Path +from types import SimpleNamespace from unittest.mock import patch from EvoScientist.subagents.expert_container import ( @@ -10,6 +11,8 @@ from EvoScientist.subagents.expert_container import ( _compose_system_prompt, build_expert_subagent_spec, build_expert_subagent_specs, + is_async_dispatch_available, + list_dispatchable_experts, ) from EvoScientist.tools.skills_manager import SkillInfo @@ -444,102 +447,90 @@ class TestFoldExpertSubagents: # ============================================================================= -# End-to-end wiring — both kwargs builders register experts with skill_manager +# is_async_dispatch_available / list_dispatchable_experts honest surface # ============================================================================= -class TestExpertWiringInBuildKwargs: - """Regression pin on the two-line wiring in both construction paths: - (a) ``tool_registry["skill_manager"]`` is populated so - ``_DEFAULT_EXPERT_TOOLS`` resolves and (b) ``_fold_expert_subagents`` is - invoked with that registry so the expert spec's ``tools`` carry the real - callable. Drift on either half silently drops ``skill_manager`` from every - expert — the exact regression fix commit ``34586c7`` addressed. A single - installed expert exercises both halves in one call.""" +class TestIsAsyncDispatchAvailable: + """The gate ``list_dispatchable_experts`` and ``ActiveTeamMiddleware`` + both consult to decide whether async-declared experts can be surfaced.""" - @staticmethod - def _install_one_expert(root: Path) -> Path: - skill_dir = root / "expert-a" - skill_dir.mkdir(parents=True, exist_ok=True) - (skill_dir / "SKILL.md").write_text( - """--- -name: expert-a -description: A test expert -type: expert -role: test expert ---- + def test_false_when_flag_disabled(self): + cfg = SimpleNamespace(enable_async_subagents=False) + assert is_async_dispatch_available(cfg=cfg) is False -Second-person persona body content. -""" + def test_false_when_dev_unreachable(self): + cfg = SimpleNamespace(enable_async_subagents=True) + with patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=False, + ): + assert is_async_dispatch_available(cfg=cfg) is False + + def test_true_when_both_gates_pass(self): + cfg = SimpleNamespace(enable_async_subagents=True) + with patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ): + assert is_async_dispatch_available(cfg=cfg) is True + + +class TestListDispatchableExpertsAsyncFilter: + """``list_dispatchable_experts`` drops async-declared experts when async + dispatch isn't registered — sync-declared experts pass through, mirroring + honest advertising per the reviewer's ask on PR #391.""" + + def _skill(self, name: str, dispatch: str) -> SkillInfo: + return SkillInfo( + name=name, + description=f"{name} description", + path=Path("/tmp/nope"), + source="builtin", + type="expert", + role=f"{name} role", + default_dispatch=dispatch, + body="persona body\n", ) - empty = root / "empty" - empty.mkdir(exist_ok=True) - return empty - def test_build_base_kwargs_registers_expert_with_skill_manager(self, tmp_path): - from EvoScientist import EvoScientist as evo - from EvoScientist.tools import skill_manager - - empty = self._install_one_expert(tmp_path) + def test_async_expert_dropped_when_flag_disabled(self): + cfg = SimpleNamespace(enable_async_subagents=False) + skills = [self._skill("idea-brainstorm", "sync"), self._skill("lit", "async")] + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ): + result = list_dispatchable_experts(cfg=cfg) + assert [s.name for s in result] == ["idea-brainstorm"] + def test_async_expert_dropped_when_dev_unreachable(self): + cfg = SimpleNamespace(enable_async_subagents=True) + skills = [self._skill("idea-brainstorm", "sync"), self._skill("lit", "async")] with ( - patch("EvoScientist.paths.USER_SKILLS_DIR", tmp_path), - patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", empty), - patch("EvoScientist.EvoScientist.SKILLS_DIR", str(empty)), - patch.object(evo, "_inject_subagent_middleware", lambda subs, **k: None), - patch.object( - evo, - "_maybe_swap_async_subagents", - lambda subs, mw, cfg=None: subs, + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=False, ), ): - kwargs = evo._build_base_kwargs( - base_backend=None, - base_middleware=[], - chat_model=object(), - ) + result = list_dispatchable_experts(cfg=cfg) + assert [s.name for s in result] == ["idea-brainstorm"] - expert = next( - (s for s in kwargs["subagents"] if s.get("name") == "expert-a"), None - ) - assert expert is not None, "expert-a not folded into subagents" - assert skill_manager in expert["tools"], ( - "skill_manager missing from expert tools — tool_registry wiring drifted" - ) - - def test_load_mcp_and_build_kwargs_registers_expert_with_skill_manager( - self, tmp_path - ): - from EvoScientist import EvoScientist as evo - from EvoScientist.tools import skill_manager - - empty = self._install_one_expert(tmp_path) - - # Non-empty mcp_by_agent forces the MCP branch. Without it, - # load_mcp_and_build_kwargs delegates to _build_base_kwargs and we - # re-test the first path only. + def test_async_expert_included_when_registered(self): + cfg = SimpleNamespace(enable_async_subagents=True) + skills = [self._skill("idea-brainstorm", "sync"), self._skill("lit", "async")] with ( - patch("EvoScientist.paths.USER_SKILLS_DIR", tmp_path), - patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", empty), - patch("EvoScientist.EvoScientist.SKILLS_DIR", str(empty)), - patch.object(evo, "_load_mcp_tools_cached", return_value={"main": []}), - patch.object(evo, "_inject_subagent_middleware", lambda subs, **k: None), - patch.object( - evo, - "_maybe_swap_async_subagents", - lambda subs, mw, cfg=None: subs, + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, ), ): - kwargs = evo.load_mcp_and_build_kwargs( - base_backend=None, - base_middleware=[], - chat_model=object(), - ) - - expert = next( - (s for s in kwargs["subagents"] if s.get("name") == "expert-a"), None - ) - assert expert is not None, "expert-a not folded into subagents (MCP path)" - assert skill_manager in expert["tools"], ( - "skill_manager missing from expert tools on MCP path" - ) + result = list_dispatchable_experts(cfg=cfg) + assert {s.name for s in result} == {"idea-brainstorm", "lit"} diff --git a/tests/test_expert_container_async.py b/tests/test_expert_container_async.py new file mode 100644 index 0000000..e9006be --- /dev/null +++ b/tests/test_expert_container_async.py @@ -0,0 +1,294 @@ +"""Tests for the async expert container graph builder + loader middleware. + +The full ``build_expert_container_async_graph()`` factory is exercised end- +to-end at langgraph dev startup; here we cover the load-bearing piece — +``ExpertSkillLoaderMiddleware._compose_prompt`` — in isolation so a +regression on skill resolution surfaces without needing a live langgraph +subprocess. +""" + +from __future__ import annotations + +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import patch + +from langchain_core.messages import SystemMessage + +from EvoScientist.subagents.expert_container_async import ( + _PERSONA_SENTINEL, + ExpertContainerState, + ExpertSkillLoaderMiddleware, +) +from EvoScientist.tools.skills_manager import SkillInfo + +# Two base-stack witness blocks used across the wrap_model_call tests. Their +# contents mirror the section headers deepagents emits per-turn — regressing +# the compose logic would drop these from the composed system_message. +_TASK_WITNESS = "## `task` (subagent spawner)\n\nUse ``task`` to delegate ..." +_SKILLS_WITNESS = "## Skills System\n\nInstalled skills are mounted under ..." + +# ============================================================================= +# _compose_prompt — the load-bearing logic +# ============================================================================= + + +def _skill_info( + *, + name: str = "literature-review", + role: str = "literature-review strategist", + body: str = "You produce manuscript-quality surveys.\n\nPipeline: ...\n", + description: str = "d", +) -> SkillInfo: + return SkillInfo( + name=name, + description=description, + path=Path("/tmp/does-not-matter"), + source="builtin", + type="expert", + role=role, + body=body, + ) + + +class TestComposePrompt: + def test_returns_role_and_body_for_known_skill(self): + mw = ExpertSkillLoaderMiddleware() + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info()], + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + # Role prepended, body preserved, trailing newline guaranteed. + assert composed.startswith("You are literature-review strategist.") + assert "You produce manuscript-quality surveys." in composed + assert composed.endswith("\n") + + def test_omits_role_line_when_absent(self): + mw = ExpertSkillLoaderMiddleware() + info = _skill_info(role="", body="Second-person persona body.\n") + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[info], + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + assert not composed.startswith("You are ") + assert "Second-person persona body." in composed + + def test_missing_skill_name_returns_error_cue(self): + mw = ExpertSkillLoaderMiddleware() + composed = mw._compose_prompt({}) + assert composed.startswith("ERROR:") + assert "skill_name" in composed + assert "wiring bug" in composed + + def test_unknown_skill_returns_error_cue_with_installed_list(self): + mw = ExpertSkillLoaderMiddleware() + installed = [_skill_info(name="literature-review"), _skill_info(name="other")] + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=installed, + ): + composed = mw._compose_prompt({"skill_name": "not-installed"}) + assert composed.startswith("ERROR:") + assert "'not-installed' is not installed" in composed + # Names of the installed experts are listed so the LLM's error + # envelope can suggest the correct spelling. + assert "literature-review" in composed + assert "other" in composed + + def test_no_installed_experts_reports_none(self): + mw = ExpertSkillLoaderMiddleware() + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", return_value=[] + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + assert composed.startswith("ERROR:") + assert "(none)" in composed + + def test_empty_body_returns_error_cue(self): + """A skill with an empty SKILL.md body would otherwise run against a + persona-less system prompt (just the role line). Mirror the sync + fold-in's policy: refuse to compose a prompt at all and surface the + skill-authoring bug through the LLM's error envelope.""" + mw = ExpertSkillLoaderMiddleware() + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info(body="")], + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + assert composed.startswith("ERROR:") + assert "empty SKILL.md body" in composed + assert "literature-review" in composed # names the offending skill + + def test_whitespace_only_body_returns_error_cue(self): + """A body that's just whitespace (` \\n\\n`) is still empty in the + sense that matters — no persona, no pipeline. Same error cue.""" + mw = ExpertSkillLoaderMiddleware() + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info(body=" \n\n \n")], + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + assert composed.startswith("ERROR:") + assert "empty SKILL.md body" in composed + + def test_runtime_context_tail_surfaces_skill_name(self): + """The tail block re-asserts ``skill_name`` on every model call so the + expert knows its own persona name after summarization. Since + ``output_path`` moved to the task description (payload dropped in + PR #391 review X-4), the tail carries no path — the LLM pins it into + its own todo list per SKILL.md contract.""" + mw = ExpertSkillLoaderMiddleware() + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info()], + ): + composed = mw._compose_prompt({"skill_name": "literature-review"}) + assert "## Runtime context" in composed + assert "``skill_name``: ``literature-review``" in composed + # Path retention is no longer a middleware responsibility. + assert "``output_path``" not in composed + assert "verbatim" not in composed + + +# ============================================================================= +# ExpertContainerState — state schema smoke check +# ============================================================================= + + +class TestExpertContainerState: + """The state schema carries ``skill_name`` only. ``output_path`` was + dropped in PR #391 review X-4 — the main agent now embeds the desired + path in the task description (natural language) and the expert's + SKILL.md contract pins it via ``write_todos`` on turn 1.""" + + def test_state_shape(self): + # TypedDicts don't runtime-validate — assert the field is declared + # so downstream code can rely on ``state.get("skill_name")``. + annotations = ExpertContainerState.__annotations__ + assert "skill_name" in annotations + assert "output_path" not in annotations + + +# ============================================================================= +# wrap_model_call — override via ModelRequest.override +# ============================================================================= + + +def _system_message_with_sentinel_and_witnesses() -> SystemMessage: + """Base-stack-shaped ``SystemMessage``: the fallback (sentinel-bearing) + block, then two witness blocks that represent deepagents' composed + sections. This is the exact shape our middleware sees at model-call + time when the container graph was built with + ``system_prompt=_FALLBACK_SYSTEM_PROMPT`` and the base stack has + appended its sections on top.""" + from EvoScientist.subagents.expert_container_async import _FALLBACK_SYSTEM_PROMPT + + return SystemMessage( + content=[ + {"type": "text", "text": _FALLBACK_SYSTEM_PROMPT}, + {"type": "text", "text": _TASK_WITNESS}, + {"type": "text", "text": _SKILLS_WITNESS}, + ] + ) + + +def _mock_request(system_message: SystemMessage): + """Stub ``ModelRequest`` supporting ``state`` and ``override``. Returns + ``(request, seen, handler)`` — ``seen`` is a list the handler pushes the + post-override ``system_message`` into for post-call assertions.""" + seen: list[SystemMessage] = [] + overridden = SimpleNamespace() + + def override(*, system_message): + overridden.system_message = system_message + return overridden + + def handler(new_request): + seen.append(new_request.system_message) + return SimpleNamespace() + + request = SimpleNamespace( + state={"skill_name": "literature-review"}, + system_message=system_message, + override=override, + ) + return request, seen, handler + + +class TestWrapModelCall: + def test_wrap_composes_persona_into_base_stack_system_message(self): + """Persona swaps for the sentinel block; base-stack witness blocks + stay in place. The whole point of the fix — replacing the whole + system_message (the pre-fix behaviour) dropped every base-stack + section (measured live: 9,608 → 382 chars) and broke ``task()`` + for async experts.""" + mw = ExpertSkillLoaderMiddleware() + request, seen, handler = _mock_request( + _system_message_with_sentinel_and_witnesses() + ) + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info()], + ): + mw.wrap_model_call(request, handler) + + composed = seen[0] + block_texts = [b.get("text", "") for b in composed.content_blocks] + # Persona landed — role prepend visible. + assert any( + t.startswith("You are literature-review strategist.") for t in block_texts + ) + # Sentinel gone (block was replaced, not appended). + assert not any(_PERSONA_SENTINEL in t for t in block_texts) + # Witnesses preserved verbatim — the base-stack sections stay live. + assert _TASK_WITNESS in block_texts + assert _SKILLS_WITNESS in block_texts + # Block count unchanged — replace, not append. + assert len(block_texts) == 3 + + def test_wrap_appends_persona_when_sentinel_missing(self, caplog): + """When the sentinel block isn't found (e.g. deepagents refactors + how ``system_prompt=`` reaches ``content_blocks``), the persona is + appended instead of silently dropped, and the drift is logged.""" + import logging + + mw = ExpertSkillLoaderMiddleware() + # No sentinel block — only witnesses. + request, seen, handler = _mock_request( + SystemMessage( + content=[ + {"type": "text", "text": _TASK_WITNESS}, + {"type": "text", "text": _SKILLS_WITNESS}, + ] + ) + ) + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill_info()], + ), + caplog.at_level( + logging.WARNING, + logger="EvoScientist.subagents.expert_container_async", + ), + ): + mw.wrap_model_call(request, handler) + + composed = seen[0] + block_texts = [b.get("text", "") for b in composed.content_blocks] + # Persona appended as a new block. + assert any( + t.startswith("You are literature-review strategist.") for t in block_texts + ) + # Witnesses still present. + assert _TASK_WITNESS in block_texts + assert _SKILLS_WITNESS in block_texts + # Original two blocks + persona = 3. + assert len(block_texts) == 3 + # Drift-detected warning surfaced. + assert any( + _PERSONA_SENTINEL in r.message and "not found" in r.message + for r in caplog.records + ) diff --git a/tests/test_experts_command.py b/tests/test_experts_command.py index d910314..799aa0e 100644 --- a/tests/test_experts_command.py +++ b/tests/test_experts_command.py @@ -136,6 +136,32 @@ class TestExpertToggle: ) assert ctx.channel_runtime.active_teams == [] + async def test_async_expert_refused_with_reason_when_async_unavailable(self): + """When an installed expert declares ``default_dispatch: async`` but + async dispatch is unavailable, ``/expert`` must refuse with the specific + reason — not the empty-body / name-collision default — so the user + knows to enable ``enable_async_subagents`` or start langgraph dev. + Reviewer thread on PR #391.""" + ctx, ui = _make_ctx() + async_expert = _FakeSkillInfo( + name="literature-review", default_dispatch="async" + ) + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[async_expert], + ), + patch( + "EvoScientist.subagents.expert_container.is_async_dispatch_available", + return_value=False, + ), + ): + await ExpertCommand().execute(ctx, args=["literature-review"]) + assert any("async dispatch is unavailable" in text for text, _ in ui.lines), ( + f"expected honest async-unavailable message, got: {ui.lines}" + ) + assert ctx.channel_runtime.active_teams == [] + async def test_invite_adds_to_active_teams(self): ctx, ui = _make_ctx() with patch( diff --git a/tests/test_middleware_utils.py b/tests/test_middleware_utils.py new file mode 100644 index 0000000..5d6e9b5 --- /dev/null +++ b/tests/test_middleware_utils.py @@ -0,0 +1,65 @@ +"""Unit tests for EvoScientist.middleware.utils helpers. + +Focuses on the ``system_message`` composition primitives shared across +middleware modules. Model-side helpers (``disable_thinking``, +``disable_streaming``) are covered by the middleware suites that use them. +""" + +from __future__ import annotations + +from langchain_core.messages import SystemMessage + +from EvoScientist.middleware.utils import ( + replace_block_by_sentinel, +) + + +class TestReplaceBlockBySentinel: + """``replace_block_by_sentinel`` swaps the block containing the sentinel + for a replacement text block, preserving every other block. Used by + ``ExpertSkillLoaderMiddleware`` to inject the persona in place of the + graph's fallback block while keeping the base-stack sections intact.""" + + _SENTINEL = "__TEST_PERSONA_SLOT__" + + def test_swaps_matching_block(self): + original = SystemMessage( + content=[ + {"type": "text", "text": f"{self._SENTINEL}\n\nfallback"}, + {"type": "text", "text": "## `task` (subagent spawner)"}, + {"type": "text", "text": "## Skills System"}, + ] + ) + result = replace_block_by_sentinel(original, self._SENTINEL, "persona body") + assert result is not None + block_texts = [b.get("text", "") for b in result.content_blocks] + assert block_texts == [ + "persona body", + "## `task` (subagent spawner)", + "## Skills System", + ] + + def test_preserves_block_count(self): + original = SystemMessage( + content=[ + {"type": "text", "text": self._SENTINEL}, + {"type": "text", "text": "witness"}, + ] + ) + result = replace_block_by_sentinel(original, self._SENTINEL, "persona") + assert result is not None + assert len(list(result.content_blocks)) == len(list(original.content_blocks)) + + def test_returns_none_when_sentinel_missing(self): + """Signal path: caller decides fallback policy (typically log + append) + so a deepagents refactor degrades gracefully instead of hard-failing.""" + original = SystemMessage( + content=[ + {"type": "text", "text": "## `task` (subagent spawner)"}, + {"type": "text", "text": "## Skills System"}, + ] + ) + assert replace_block_by_sentinel(original, self._SENTINEL, "persona") is None + + def test_returns_none_when_message_is_none(self): + assert replace_block_by_sentinel(None, self._SENTINEL, "persona") is None diff --git a/tests/test_route_async_specs.py b/tests/test_route_async_specs.py new file mode 100644 index 0000000..bf6f109 --- /dev/null +++ b/tests/test_route_async_specs.py @@ -0,0 +1,401 @@ +"""Tests for the AsyncSubAgent → EvoAsyncSubAgentMiddleware routing helper. + +Covers: +- ``_route_async_specs_through_evo_middleware`` splits AsyncSubAgent specs + out of the ``subs`` list and folds them into the base middleware. +- ``build_expert_async_subagent_specs`` filters by + ``default_dispatch == "async"`` and respects the async-enable flag + + langgraph dev reachability. +- ``build_expert_subagent_specs`` (sync fold-in) excludes async experts so + a single skill never surfaces twice in the main agent's tool schema. +""" + +from __future__ import annotations + +from pathlib import Path +from types import SimpleNamespace +from unittest.mock import patch + +from EvoScientist.subagents.expert_container import build_expert_subagent_specs +from EvoScientist.subagents.expert_container_async import ( + build_expert_async_subagent_specs, +) +from EvoScientist.tools.skills_manager import SkillInfo + + +def _skill(name: str, dispatch: str) -> SkillInfo: + return SkillInfo( + name=name, + description=f"{name} description", + path=Path("/tmp/does-not-matter"), + source="builtin", + type="expert", + role=f"{name} role", + default_dispatch=dispatch, + body="body\n", + ) + + +# ============================================================================= +# build_expert_async_subagent_specs +# ============================================================================= + + +class TestBuildExpertAsyncSubagentSpecs: + def test_empty_when_async_disabled(self): + cfg = SimpleNamespace(enable_async_subagents=False) + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill("literature-review", "async")], + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + assert specs == [] + + def test_empty_when_langgraph_dev_unreachable(self): + cfg = SimpleNamespace(enable_async_subagents=True, langgraph_dev_port=6174) + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill("literature-review", "async")], + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=False, + ), + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + assert specs == [] + + def test_filters_by_default_dispatch(self): + """Only ``default_dispatch: async`` skills become AsyncSubAgent specs.""" + cfg = SimpleNamespace(enable_async_subagents=True, langgraph_dev_port=6174) + skills = [ + _skill("idea-brainstorm", "sync"), + _skill("literature-review", "async"), + _skill("panel-expert", "panel"), + ] + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + assert len(specs) == 1 + assert specs[0]["name"] == "literature-review" + assert specs[0]["graph_id"] == "expert-container-async" + assert specs[0]["is_expert"] is True + assert "http://localhost:6174" in specs[0]["url"] + + def test_empty_body_experts_skipped(self): + """Empty-body async experts are filtered out at spec-build time so + ``start_async_task``'s tool schema never advertises a broken skill. + Mirrors the sync fold-in in + ``expert_container.py::build_expert_subagent_specs``.""" + cfg = SimpleNamespace(enable_async_subagents=True, langgraph_dev_port=6174) + skills = [ + _skill("literature-review", "async"), # normal body from _skill() + _skill("empty-persona", "async"), + ] + # Second skill has no body — dataclass field default is ``""``, but + # helper sets it to "body\n" — override to empty. + skills[1].body = "" + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + assert [s["name"] for s in specs] == ["literature-review"] + + def test_reserved_name_collision_skipped(self, caplog): + """A skill named after a yaml async sub-agent (or ``general-purpose``) + must skip async-dispatch registration with a warning, not raise. Without + this guard ``AsyncSubAgentMiddleware.__init__`` would ``ValueError: + Duplicate async subagent names`` on the merged spec list and kill CLI + startup — see reviewer thread on PR #391.""" + import logging + + cfg = SimpleNamespace(enable_async_subagents=True, langgraph_dev_port=6174) + skills = [ + _skill("writing-agent", "async"), # collides with yaml async agent + _skill("literature-review", "async"), + ] + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + patch( + "EvoScientist.subagents.expert_container._reserved_subagent_names", + return_value=frozenset({"writing-agent", "general-purpose"}), + ), + caplog.at_level( + logging.WARNING, + logger="EvoScientist.subagents.expert_container_async", + ), + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + assert [s["name"] for s in specs] == ["literature-review"] + assert any( + "writing-agent" in r.message and "collides" in r.message + for r in caplog.records + ) + + def test_workspace_duplicate_name_skipped(self, caplog): + """Two workspace-tier expert skills sharing a frontmatter ``name`` must + register only the first — the workspace listing uses + ``check_seen=False`` so both survive to this point. Without a local + seen-set the second would collide inside + ``AsyncSubAgentMiddleware.__init__``.""" + import logging + + cfg = SimpleNamespace(enable_async_subagents=True, langgraph_dev_port=6174) + skills = [ + _skill("literature-review", "async"), + _skill("literature-review", "async"), # duplicate name + ] + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + patch( + "EvoScientist.subagents.expert_container._reserved_subagent_names", + return_value=frozenset({"general-purpose"}), + ), + caplog.at_level( + logging.WARNING, + logger="EvoScientist.subagents.expert_container_async", + ), + ): + specs = build_expert_async_subagent_specs(cfg=cfg) + # Only the first `literature-review` survives. + assert [s["name"] for s in specs] == ["literature-review"] + assert any( + "literature-review" in r.message and "collides" in r.message + for r in caplog.records + ) + + +# ============================================================================= +# build_expert_subagent_specs (sync side) — must exclude async experts +# ============================================================================= + + +class TestBuildExpertSubagentSpecsExcludesAsync: + """The sync fold-in must not emit specs for ``default_dispatch: async`` skills. + + A skill in both lists would produce two competing tool schemas — one + under ``task(subagent_type='')`` and one under + ``start_async_task(subagent_type='')`` — from the main agent's + perspective. Ambiguous. The partition is: async goes async, everything + else goes sync. + """ + + def test_async_expert_skipped_by_sync_fold_in(self): + skills = [ + _skill("idea-brainstorm", "sync"), + _skill("literature-review", "async"), + _skill("panel-expert", "panel"), + ] + with patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=skills, + ): + specs = build_expert_subagent_specs(tool_registry={}) + names = [s["name"] for s in specs] + assert "idea-brainstorm" in names + assert "panel-expert" in names + assert "literature-review" not in names + + +# ============================================================================= +# _route_async_specs_through_evo_middleware +# ============================================================================= + + +class TestRouteAsyncSpecs: + """The routing helper splits AsyncSubAgent specs from ``subs`` and + hands them to ``EvoAsyncSubAgentMiddleware``. Verifies: + - Sync subagents pass through untouched. + - AsyncSubAgent specs are stripped from the returned ``subs``. + - Expert async specs (from ``build_expert_async_subagent_specs``) are + merged in. + - The middleware is appended to ``base_middleware`` only when there + are async specs (either standard or expert). + """ + + def _cfg(self, *, enable_async: bool = True, port: int = 6174): + return SimpleNamespace( + enable_async_subagents=enable_async, langgraph_dev_port=port + ) + + def test_sync_subagents_pass_through(self): + from EvoScientist.EvoScientist import _route_async_specs_through_evo_middleware + + subs = [{"name": "sync-a", "system_prompt": ""}] + middleware: list = [] + # Disable async path via cfg + patched reachability. + with patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=False, + ): + result = _route_async_specs_through_evo_middleware( + subs, middleware, cfg=self._cfg(enable_async=False) + ) + assert result == [{"name": "sync-a", "system_prompt": ""}] + assert middleware == [] # no async → no middleware added + + def test_async_specs_moved_to_middleware(self): + from EvoScientist.EvoScientist import _route_async_specs_through_evo_middleware + from EvoScientist.middleware.expert_async_subagent import ( + EvoAsyncSubAgentMiddleware, + ) + + subs = [ + {"name": "sync-a", "system_prompt": ""}, + { + "name": "writing-agent", + "description": "std", + "graph_id": "writing_agent", + "url": "http://localhost:6174", + }, + ] + middleware: list = [] + # Disable expert-async fold-in to isolate the standard-spec routing. + with patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=False, + ): + result = _route_async_specs_through_evo_middleware( + subs, middleware, cfg=self._cfg(enable_async=False) + ) + # `writing-agent` stripped from subs (it has graph_id). + assert [s["name"] for s in result] == ["sync-a"] + # Middleware appended. + assert len(middleware) == 1 + assert isinstance(middleware[0], EvoAsyncSubAgentMiddleware) + + def test_expert_async_specs_merged_in(self): + from EvoScientist.EvoScientist import _route_async_specs_through_evo_middleware + from EvoScientist.middleware.expert_async_subagent import ( + EvoAsyncSubAgentMiddleware, + ) + + subs = [{"name": "sync-a", "system_prompt": ""}] + middleware: list = [] + cfg = self._cfg(enable_async=True) + # Enable expert-async by patching skills list + reachability. + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill("literature-review", "async")], + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + ): + result = _route_async_specs_through_evo_middleware( + subs, middleware, cfg=cfg + ) + # sync-a stays; middleware got the expert spec, and an + # AsyncWatcherMiddleware was installed so expert launches spawn + # completion watchers (previously the watcher's client cache had no + # entry for the expert name, KeyErrored on `get_async`, and silently + # dropped the notification). + from EvoScientist.middleware.async_watcher import AsyncWatcherMiddleware + + assert [s["name"] for s in result] == ["sync-a"] + assert len(middleware) == 2 + evo_mw = next( + m for m in middleware if isinstance(m, EvoAsyncSubAgentMiddleware) + ) + watcher_mw = next( + m for m in middleware if isinstance(m, AsyncWatcherMiddleware) + ) + # The middleware's start tool schema advertises literature-review. + start = next(t for t in evo_mw.tools if t.name == "start_async_task") + assert "literature-review" in start.description + # The watcher's client cache knows how to construct a client for the + # expert so the completion nudge can spawn. + assert "literature-review" in watcher_mw._clients._agents + + def test_watcher_cache_extends_when_yaml_watcher_preinstalled(self): + """Default deployed shape — ``_maybe_swap_async_subagents`` installed + ``AsyncWatcherMiddleware`` for a yaml async agent, then the routing + helper extends the cache with expert specs. Without the extension + branch, an expert completion nudge would KeyError on the watcher's + ``get_async()`` and silently drop the notification.""" + from EvoScientist.cli import async_notifier + from EvoScientist.EvoScientist import _route_async_specs_through_evo_middleware + from EvoScientist.middleware.async_watcher import AsyncWatcherMiddleware + from EvoScientist.middleware.expert_async_subagent import ( + EvoAsyncSubAgentMiddleware, + ) + + yaml_async_spec = { + "name": "writing-agent", + "description": "std", + "graph_id": "writing_agent", + "url": "http://localhost:6174", + } + subs = [{"name": "sync-a", "system_prompt": ""}, yaml_async_spec] + # Simulate the state after ``_maybe_swap_async_subagents``: watcher is + # already installed and carries the yaml async agent. + middleware: list = [ + AsyncWatcherMiddleware( + {"writing-agent": yaml_async_spec}, notifier=async_notifier + ) + ] + cfg = self._cfg(enable_async=True) + with ( + patch( + "EvoScientist.tools.skills_manager.list_expert_skills", + return_value=[_skill("literature-review", "async")], + ), + patch( + "EvoScientist.langgraph_dev.manager.is_async_subagents_available", + return_value=True, + ), + patch( + "EvoScientist.subagents.expert_container._reserved_subagent_names", + return_value=frozenset({"general-purpose"}), + ), + ): + result = _route_async_specs_through_evo_middleware( + subs, middleware, cfg=cfg + ) + # `writing-agent` (graph_id-carrying) stripped from subs; sync-a stays. + assert [s["name"] for s in result] == ["sync-a"] + evo_mw = next( + m for m in middleware if isinstance(m, EvoAsyncSubAgentMiddleware) + ) + watcher_mw = next( + m for m in middleware if isinstance(m, AsyncWatcherMiddleware) + ) + # Both the yaml async agent and the expert reach the start-task schema. + start = next(t for t in evo_mw.tools if t.name == "start_async_task") + assert "writing-agent" in start.description + assert "literature-review" in start.description + # The pre-existing watcher was extended in place — both names route. + assert "writing-agent" in watcher_mw._clients._agents + assert "literature-review" in watcher_mw._clients._agents diff --git a/tests/test_skills_manager.py b/tests/test_skills_manager.py index 175ff7f..8ebeb95 100644 --- a/tests/test_skills_manager.py +++ b/tests/test_skills_manager.py @@ -1063,7 +1063,26 @@ role: This should be ignored for r in caplog.records ) - def test_invalid_default_dispatch_falls_back_to_empty(self, tmp_path): + def test_valid_async_default_dispatch(self, tmp_path): + """``default_dispatch: async`` is a recognized value (agent-teams v2 dispatch mode).""" + skill_dir = tmp_path / "async-skill" + skill_dir.mkdir() + (skill_dir / "SKILL.md").write_text( + """--- +name: async-skill +description: uses async dispatch +type: expert +role: Some role +default_dispatch: async +--- + +# Body +""" + ) + result = _parse_skill_md(skill_dir / "SKILL.md") + assert result.default_dispatch == "async" + + def test_invalid_default_dispatch_falls_back_to_empty(self, tmp_path, caplog): skill_dir = tmp_path / "bad-dispatch" skill_dir.mkdir() (skill_dir / "SKILL.md").write_text( @@ -1078,9 +1097,22 @@ default_dispatch: asynchronous # Body """ ) - result = _parse_skill_md(skill_dir / "SKILL.md") + import logging + + with caplog.at_level( + logging.WARNING, logger="EvoScientist.tools.skills_manager" + ): + result = _parse_skill_md(skill_dir / "SKILL.md") assert result.type == "expert" assert result.default_dispatch == "" # rejected, not passed through + # A typo like ``asynchronous`` is silently indistinguishable from + # unset without the warning; the log line must name the offending + # value so authors can see why their expert didn't register async. + assert any( + "unrecognized default_dispatch" in rec.message + and "asynchronous" in rec.message + for rec in caplog.records + ) def test_capability_tags_accepts_comma_string(self, tmp_path): """capability_tags falls back to comma-separated string parsing (like `tags`).""" @@ -1656,3 +1688,250 @@ class TestSkillsChangedCallback: result = install_skill("/nonexistent/path", str(temp_skills_dir)) assert result["success"] is False assert good == [True] + + +class TestSkillManagerInfo: + """Tests for the skill_manager() tool's action='info' output. + + Guards the sandbox-visible ``Path: /skills/`` shape and the absence + of any host filesystem path in the agent-visible response. Agents burn + turns on ``cd && …`` chains whenever the host path leaks. + """ + + def _make_skill(self, parent, name, description="A skill"): + skill_dir = parent / name + skill_dir.mkdir() + (skill_dir / "SKILL.md").write_text( + f"---\nname: {name}\ndescription: {description}\n---\n" + ) + return skill_dir + + def test_info_reports_virtual_mount_path(self, tmp_path): + """``Path:`` is the sandbox-visible ``/skills/``, not the host path.""" + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + self._make_skill(tmp_path, "info-skill") + install_skill(str(tmp_path / "info-skill"), str(workspace_dir)) + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + ): + result = skill_manager.invoke({"action": "info", "name": "info-skill"}) + + assert "Path: /skills/info-skill" in result + + def test_info_omits_host_path(self, tmp_path): + """No host filesystem path leaks into the response. + + Stronger than a label-only check: catches any future refactor that + keeps the path visible under a different label (``Local:``, + ``Installed at:``, embedded in ``Source: …``). + """ + from EvoScientist.tools.skill_manager import skill_manager + from EvoScientist.tools.skills_manager import get_skill_info + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + self._make_skill(tmp_path, "host-leak-guard") + install_skill(str(tmp_path / "host-leak-guard"), str(workspace_dir)) + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + ): + info = get_skill_info("host-leak-guard") + result = skill_manager.invoke({"action": "info", "name": "host-leak-guard"}) + + assert str(info.path) not in result + + +class TestSkillManagerInstall: + """Tests for the skill_manager() tool's action='install' output shape. + + Covers both single-install and batch-install returns: + - Single: ``{"success": True, "name": ..., "path": ..., "description": ...}``. + - Batch: ``{"success": ..., "batch": True, "installed": [...], "failed": [...]}`` + with no top-level ``name``, ``path``, ``description``, or ``error``. + """ + + def _make_skill(self, parent, name, description="A skill"): + skill_dir = parent / name + skill_dir.mkdir() + (skill_dir / "SKILL.md").write_text( + f"---\nname: {name}\ndescription: {description}\n---\n" + ) + return skill_dir + + def test_install_single_reports_virtual_mount_path(self, tmp_path): + """Single install: ``Path: /skills/``, no host path.""" + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + self._make_skill(tmp_path, "solo-skill") + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + ): + result = skill_manager.invoke( + {"action": "install", "source": str(tmp_path / "solo-skill")} + ) + + assert "Successfully installed skill: solo-skill" in result + assert "Path: /skills/solo-skill" in result + + def test_install_single_omits_host_path(self, tmp_path): + """Single install: no host filesystem path leaks into the response. + + ``install_skill(source)`` defaults to ``global_install=True``, so the + skill lands under ``GLOBAL_SKILLS_DIR`` rather than ``USER_SKILLS_DIR``. + Checking against a narrower directory (e.g. workspace_dir) would pass + even without the scrub - the check has to cover every path the tool + might resolve to. ``tmp_path`` covers both patched dirs and the source + path used by ``install_skill``. + """ + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + self._make_skill(tmp_path, "leak-guard-install") + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + ): + result = skill_manager.invoke( + {"action": "install", "source": str(tmp_path / "leak-guard-install")} + ) + + assert str(tmp_path) not in result + + def test_install_batch_lists_each_skill_with_virtual_path(self, tmp_path): + """Batch install: one block per installed skill, each with ``Path: /skills/``.""" + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + pack = tmp_path / "pack" + pack.mkdir() + self._make_skill(pack, "alpha", description="first") + self._make_skill(pack, "beta", description="second") + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + ): + result = skill_manager.invoke({"action": "install", "source": str(pack)}) + + assert "Successfully installed skill: alpha" in result + assert "Successfully installed skill: beta" in result + assert "Path: /skills/alpha" in result + assert "Path: /skills/beta" in result + # Same leak guard as ``test_install_single_omits_host_path``: batch + # returns must not surface any host path either. Cover every dir the + # install might resolve to. + assert str(tmp_path) not in result + + def test_install_batch_all_fail_returns_error_list(self, tmp_path): + """Batch install where every skill fails must not KeyError on the + missing top-level ``error`` field. + + Pre-fix behavior: ``result['error']`` crashed because + ``_batch_install_local`` returns ``{"success": False, "batch": True, + "installed": [], "failed": [{"name": ..., "error": ...}]}`` with no + top-level ``error`` key. This test pins the guard. + """ + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + + batch_result = { + "success": False, + "batch": True, + "installed": [], + "failed": [ + {"name": "broken-a", "error": "corrupt frontmatter"}, + {"name": "broken-b", "error": "missing SKILL.md"}, + ], + } + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + patch( + "EvoScientist.tools.skills_manager.install_skill", + return_value=batch_result, + ), + ): + result = skill_manager.invoke( + {"action": "install", "source": "some/source"} + ) + + # No KeyError, and every failure surfaced. + assert "broken-a" in result + assert "corrupt frontmatter" in result + assert "broken-b" in result + assert "missing SKILL.md" in result + + def test_install_batch_partial_fail_surfaces_both(self, tmp_path): + """Batch install with partial failure lists successes AND failures. + + Pre-fix behavior: partial failures were silently dropped; only the + success blocks reached the agent. + """ + from EvoScientist.tools.skill_manager import skill_manager + + workspace_dir = tmp_path / "workspace" + workspace_dir.mkdir() + global_dir = tmp_path / "global-empty" + global_dir.mkdir() + + partial_result = { + "success": True, + "batch": True, + "installed": [ + { + "name": "worked", + "path": str(workspace_dir / "worked"), + "description": "installed cleanly", + }, + ], + "failed": [ + {"name": "broken", "error": "corrupt frontmatter"}, + ], + } + + with ( + patch("EvoScientist.paths.USER_SKILLS_DIR", workspace_dir), + patch("EvoScientist.paths.GLOBAL_SKILLS_DIR", global_dir), + patch( + "EvoScientist.tools.skills_manager.install_skill", + return_value=partial_result, + ), + ): + result = skill_manager.invoke( + {"action": "install", "source": "some/source"} + ) + + assert "Successfully installed skill: worked" in result + + assert "Path: /skills/worked" in result + assert "broken" in result + assert "corrupt frontmatter" in result