0410b40f57
* feat: inherit the caller's model for async sub-agent launch and update * docs: tighten middleware related docstrings
984 lines
39 KiB
Python
984 lines
39 KiB
Python
"""Tests for the skill-name-injecting AsyncSubAgentMiddleware subclass."""
|
|
|
|
from __future__ import annotations
|
|
|
|
import asyncio
|
|
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")
|
|
|
|
# Patch the resolve-on-miss walk so the negative-miss path stays
|
|
# hermetic — an unpatched call would read the real skills tree.
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[],
|
|
):
|
|
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")
|
|
|
|
# Patch the resolve-on-miss walk — see the sync twin.
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[],
|
|
):
|
|
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
|
|
|
|
|
|
def _newly_installed_expert_spec():
|
|
"""An expert spec as ``build_expert_async_subagent_specs`` would return
|
|
it for a skill installed after the agent was built."""
|
|
return {
|
|
"name": "brand-new-expert",
|
|
"description": "freshly installed expert",
|
|
"graph_id": "expert-container-async",
|
|
"is_expert": True,
|
|
}
|
|
|
|
|
|
class TestResolveOnMiss:
|
|
"""Resolve-on-miss: an unknown ``subagent_type`` that names a real,
|
|
newly installed expert becomes dispatchable on the first launch —
|
|
no agent rebuild, no restart. A name that is still unknown after one
|
|
resolution walk gets upstream's error with the refreshed type list."""
|
|
|
|
def test_unknown_expert_resolves_and_dispatches(self):
|
|
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.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
),
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync",
|
|
return_value=client,
|
|
),
|
|
):
|
|
result = start.func(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
# Dispatch succeeded rather than returning the unknown-type error.
|
|
assert "async_tasks" in result.update
|
|
kwargs = client.runs.create.call_args.kwargs
|
|
assert kwargs["input"]["skill_name"] == "brand-new-expert"
|
|
|
|
def test_resolution_updates_the_watcher_dict(self):
|
|
"""The watcher holds a SEPARATE agent dict from ``agent_map``; the
|
|
resolution must land in both or the completion notification for the
|
|
newly resolved expert silently never fires (the watcher's
|
|
``get_async`` KeyError is swallowed by its ``try/except``)."""
|
|
watcher_agents: dict = {}
|
|
mw = EvoAsyncSubAgentMiddleware(
|
|
async_subagents=[_standard_spec()], watcher_agents=watcher_agents
|
|
)
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
client = _fake_sync_client()
|
|
with (
|
|
patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
),
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync",
|
|
return_value=client,
|
|
),
|
|
):
|
|
start.func(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
assert "brand-new-expert" in watcher_agents
|
|
|
|
def test_resolution_never_overwrites_existing_entries(self):
|
|
"""``setdefault`` semantics: a spec already in ``agent_map`` keeps its
|
|
identity — an overwrite could smuggle in a spec the running agent
|
|
was not validated against (the constructor already raised on
|
|
duplicate names at build time)."""
|
|
incumbent = {
|
|
"name": "literature-review",
|
|
"description": "original description",
|
|
"graph_id": "incumbent-graph",
|
|
"is_expert": True,
|
|
}
|
|
challenger = {
|
|
"name": "literature-review",
|
|
"description": "different description",
|
|
"graph_id": "challenger-graph",
|
|
"is_expert": True,
|
|
}
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[incumbent])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
# The miss-walk returns BOTH a new expert and a same-name challenger
|
|
# for the incumbent; the dispatch goes to the new name so the walk
|
|
# runs, then to the incumbent to observe which spec survived.
|
|
client = _fake_sync_client()
|
|
with (
|
|
patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[challenger, _newly_installed_expert_spec()],
|
|
),
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync",
|
|
return_value=client,
|
|
),
|
|
):
|
|
start.func(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
start.func(
|
|
description="hi",
|
|
subagent_type="literature-review",
|
|
runtime=SimpleNamespace(tool_call_id="tc2"),
|
|
)
|
|
|
|
# The incumbent's graph_id served both the survivor check and the
|
|
# dispatch: had the challenger overwritten it, this would be
|
|
# "challenger-graph".
|
|
assistant_ids = [
|
|
call.kwargs["assistant_id"] for call in client.runs.create.call_args_list
|
|
]
|
|
assert "incumbent-graph" in assistant_ids
|
|
assert "challenger-graph" not in assistant_ids
|
|
|
|
def test_negative_miss_returns_error_with_refreshed_list(self):
|
|
"""A hallucinated name is still an error after the one resolution
|
|
walk — and the message's allowed-type list now includes names the
|
|
walk just added (the second ``_validate_agent_type`` call reads the
|
|
mutated map)."""
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
):
|
|
result = start.func(
|
|
description="hi",
|
|
subagent_type="still-does-not-exist",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
assert isinstance(result, str)
|
|
assert "Unknown async subagent type" in result
|
|
assert "brand-new-expert" in result
|
|
|
|
def test_resolution_uses_the_construction_cfg(self):
|
|
"""The miss-walk must spec against the cfg the agent was constructed
|
|
with, not a fresh ``get_effective_config()`` read. Re-deriving config
|
|
at dispatch time would let a mid-session ``langgraph_dev_port`` change
|
|
spec a newly resolved expert onto a port the running dev subprocess
|
|
is not on — dispatch accepts the name, only ``runs.create`` fails."""
|
|
construction_cfg = SimpleNamespace(enable_async_subagents=True)
|
|
mw = EvoAsyncSubAgentMiddleware(
|
|
async_subagents=[_standard_spec()], cfg=construction_cfg
|
|
)
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
captured: dict = {}
|
|
|
|
def capture_cfg(cfg=None, **kwargs):
|
|
captured["cfg"] = cfg
|
|
return [_newly_installed_expert_spec()]
|
|
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
side_effect=capture_cfg,
|
|
):
|
|
start.func(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
assert captured["cfg"] is construction_cfg
|
|
|
|
|
|
class TestAstartResolveOnMiss:
|
|
"""Async twins of ``TestResolveOnMiss`` — the coroutine langgraph_api
|
|
actually runs in production."""
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_unknown_expert_resolves_and_dispatches(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()
|
|
to_thread_calls = []
|
|
|
|
async def _fake_to_thread(fn, *args):
|
|
to_thread_calls.append(fn.__name__)
|
|
return fn(*args)
|
|
|
|
with (
|
|
patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
),
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_async",
|
|
return_value=client,
|
|
),
|
|
patch("asyncio.to_thread", new=_fake_to_thread),
|
|
):
|
|
result = await start.coroutine(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
assert "async_tasks" in result.update
|
|
kwargs = client.runs.create.await_args.kwargs
|
|
assert kwargs["input"]["skill_name"] == "brand-new-expert"
|
|
# The resolution ran off the event loop — langgraph-dev's blockbuster
|
|
# guard turns a skills-tree walk on the loop into a BlockingError.
|
|
assert to_thread_calls == ["_resolve_merge_validate"]
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_resolution_updates_the_watcher_dict(self):
|
|
watcher_agents: dict = {}
|
|
mw = EvoAsyncSubAgentMiddleware(
|
|
async_subagents=[_standard_spec()], watcher_agents=watcher_agents
|
|
)
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
client = _fake_async_client()
|
|
with (
|
|
patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
),
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_async",
|
|
return_value=client,
|
|
),
|
|
):
|
|
await start.coroutine(
|
|
description="hi",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
assert "brand-new-expert" in watcher_agents
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_negative_miss_returns_error(self):
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[],
|
|
):
|
|
result = await start.coroutine(
|
|
description="hi",
|
|
subagent_type="still-does-not-exist",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
assert isinstance(result, str)
|
|
assert "Unknown async subagent type" in result
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_negative_miss_returns_refreshed_error(self):
|
|
"""The async miss path must honor the threaded call's return value:
|
|
the error comes from the worker's merge-and-validate under the
|
|
lock, so its allowed-type list already includes the names the walk
|
|
just merged. A caller that dropped the ``to_thread`` result and
|
|
re-derived the error from a stale message would lose the new
|
|
names."""
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
):
|
|
result = await start.coroutine(
|
|
description="hi",
|
|
subagent_type="still-does-not-exist",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
assert isinstance(result, str)
|
|
assert "Unknown async subagent type" in result
|
|
assert "brand-new-expert" in result
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_known_name_dispatch_skips_the_lock(self):
|
|
"""A known-name dispatch on the event loop must never touch
|
|
``_resolve_lock``: the miss check is a keyed lookup, and all lock
|
|
work lives on the ``to_thread`` worker. Holding the lock from this
|
|
coroutine pins the property — the dispatch completes while the
|
|
lock is unavailable. The pre-reshape shape ran its validation
|
|
under the lock on the loop and hung here until the timeout."""
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
client = _fake_async_client()
|
|
acquired = mw._resolve_lock.acquire()
|
|
assert acquired
|
|
try:
|
|
with patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_async",
|
|
return_value=client,
|
|
):
|
|
result = await asyncio.wait_for(
|
|
start.coroutine(
|
|
description="hi",
|
|
subagent_type="writing-agent",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
),
|
|
timeout=2.0,
|
|
)
|
|
finally:
|
|
mw._resolve_lock.release()
|
|
assert "async_tasks" in result.update
|
|
|
|
|
|
class TestResolveOnMissLocking:
|
|
"""The resolver's merge and the start tool's map iteration serialize on
|
|
one lock. Deterministic, blocking-based — no timing lottery: each test
|
|
blocks a participant on an event we control and asserts the other side
|
|
genuinely waits for the lock."""
|
|
|
|
def test_resolver_merge_waits_for_the_lock(self):
|
|
"""With the lock held by an unrelated holder, the resolver's merge
|
|
must not insert into ``agent_map`` until the lock is released.
|
|
Without the lock parameter (or without locking in the resolver),
|
|
the ``setdefault`` lands immediately and the mid-hold assertion
|
|
fails. The return value is the refreshed validation: ``None`` once
|
|
the merged name resolves."""
|
|
import threading
|
|
import time
|
|
|
|
from EvoScientist.middleware.expert_async_subagent import (
|
|
_resolve_merge_validate,
|
|
)
|
|
|
|
agent_map: dict = {"writing-agent": _standard_spec()}
|
|
watcher_agents: dict = {}
|
|
lock = threading.Lock()
|
|
done = threading.Event()
|
|
|
|
def resolver():
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
):
|
|
result = _resolve_merge_validate(
|
|
agent_map, watcher_agents, None, "brand-new-expert", lock
|
|
)
|
|
assert result is None
|
|
done.set()
|
|
|
|
with lock:
|
|
thread = threading.Thread(target=resolver)
|
|
thread.start()
|
|
time.sleep(0.05)
|
|
# The merge is locked out while we hold the lock.
|
|
assert "brand-new-expert" not in agent_map
|
|
|
|
thread.join(timeout=5)
|
|
assert done.is_set()
|
|
assert "brand-new-expert" in agent_map
|
|
assert "brand-new-expert" in watcher_agents
|
|
|
|
def test_validate_blocks_while_resolver_holds_the_lock(self):
|
|
"""End to end through the middleware's own lock, on the SYNC tool
|
|
path (a blocked coroutine would freeze the event loop, making the
|
|
blocking unobservable from the same loop; the sync variant shares
|
|
the identical locked-validation closure). A resolver whose merge
|
|
blocks on an event we control holds the lock; a concurrent
|
|
``start_async_task`` at a KNOWN name (validation only, no
|
|
resolution) must not complete while the lock is held — its
|
|
``_validate_agent_type`` joins over ``agent_map`` under the same
|
|
lock. Without the lock, the known-name dispatch completes during
|
|
the resolver's block and the ``thread.is_alive()`` assertion
|
|
fails."""
|
|
import threading
|
|
import time
|
|
|
|
from EvoScientist.middleware import expert_async_subagent as mod
|
|
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
start = next(t for t in mw.tools if t.name == "start_async_task")
|
|
|
|
resolver_entered = threading.Event()
|
|
resolver_release = threading.Event()
|
|
orig_merge = mod._merge_expert_specs
|
|
|
|
def blocking_merge(agent_map, watcher_agents, specs):
|
|
resolver_entered.set()
|
|
assert resolver_release.wait(timeout=10)
|
|
orig_merge(agent_map, watcher_agents, specs)
|
|
|
|
def miss_dispatch():
|
|
with patch(
|
|
"EvoScientist.subagents.expert_container_async"
|
|
".build_expert_async_subagent_specs",
|
|
return_value=[_newly_installed_expert_spec()],
|
|
):
|
|
return start.func(
|
|
description="one",
|
|
subagent_type="brand-new-expert",
|
|
runtime=SimpleNamespace(tool_call_id="tc1"),
|
|
)
|
|
|
|
def known_name_dispatch(done_event):
|
|
client = _fake_sync_client()
|
|
with patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync",
|
|
return_value=client,
|
|
):
|
|
start.func(
|
|
description="two",
|
|
subagent_type="writing-agent",
|
|
runtime=SimpleNamespace(tool_call_id="tc2"),
|
|
)
|
|
done_event.set()
|
|
|
|
with patch.object(mod, "_merge_expert_specs", blocking_merge):
|
|
# Thread A: a miss -> resolver enters the merge, acquires the
|
|
# lock, and blocks on our event.
|
|
t1 = threading.Thread(target=miss_dispatch)
|
|
t1.start()
|
|
assert resolver_entered.wait(timeout=10)
|
|
|
|
# Thread B: a KNOWN name -> validation only. Must block on the
|
|
# lock the resolver holds.
|
|
b_done = threading.Event()
|
|
t2 = threading.Thread(target=known_name_dispatch, args=(b_done,))
|
|
t2.start()
|
|
time.sleep(0.1)
|
|
assert t2.is_alive()
|
|
assert not b_done.is_set()
|
|
|
|
resolver_release.set()
|
|
t1.join(timeout=10)
|
|
t2.join(timeout=10)
|
|
assert b_done.is_set()
|
|
assert not t1.is_alive()
|
|
assert not t2.is_alive()
|
|
|
|
|
|
class TestCallerModelInheritance:
|
|
"""start / update forward the *caller's* per-run model into ``runs.create``,
|
|
beating the config-default.
|
|
|
|
This is the bill-the-config-default bug on the ``langgraph_server`` backend:
|
|
the model-passthrough proxy runs inside the dev-server process, where
|
|
``_ensure_config()`` reports the server's config-default (e.g. a billed
|
|
``gemini-3-flash-preview``) rather than the CLI's per-run choice. The
|
|
launching run's real model reaches the tool as
|
|
``runtime.config.configurable.model``, so it must win — otherwise a
|
|
sub-agent launched (or continued) while the caller is on a free model
|
|
silently bills the config-default.
|
|
"""
|
|
|
|
def _runtime(self, *, model="free", provider="openrouter", state=None):
|
|
ns = SimpleNamespace(
|
|
tool_call_id="tc1",
|
|
config={"configurable": {"model": model, "model_provider": provider}},
|
|
)
|
|
if state is not None:
|
|
ns.state = state
|
|
return ns
|
|
|
|
def _cfg_default(self):
|
|
from EvoScientist.config.settings import EvoScientistConfig
|
|
|
|
return EvoScientistConfig(model="gemini-3-flash-preview", provider="openrouter")
|
|
|
|
def _tracked_task(self, agent_name="writing-agent"):
|
|
return {
|
|
"task_id": "task-abc",
|
|
"agent_name": agent_name,
|
|
"thread_id": "task-abc",
|
|
"run_id": "old-run",
|
|
"status": "running",
|
|
"created_at": "2026-05-07T00:00:00Z",
|
|
"last_checked_at": "2026-05-07T00:00:00Z",
|
|
"last_updated_at": "2026-05-07T00:00:00Z",
|
|
}
|
|
|
|
def test_start_forwards_caller_model_over_cfg(self):
|
|
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,
|
|
),
|
|
patch(
|
|
"EvoScientist.EvoScientist._ensure_config",
|
|
return_value=self._cfg_default(),
|
|
),
|
|
):
|
|
start.func(
|
|
description="w",
|
|
subagent_type="literature-review",
|
|
runtime=self._runtime(),
|
|
)
|
|
|
|
configurable = client.runs.create.call_args.kwargs["config"]["configurable"]
|
|
assert configurable["model"] == "free"
|
|
assert configurable["model_provider"] == "openrouter"
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_astart_forwards_caller_model_over_cfg(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,
|
|
),
|
|
patch(
|
|
"EvoScientist.EvoScientist._ensure_config",
|
|
return_value=self._cfg_default(),
|
|
),
|
|
):
|
|
await start.coroutine(
|
|
description="w",
|
|
subagent_type="literature-review",
|
|
runtime=self._runtime(),
|
|
)
|
|
|
|
configurable = client.runs.create.await_args.kwargs["config"]["configurable"]
|
|
assert configurable["model"] == "free"
|
|
assert configurable["model_provider"] == "openrouter"
|
|
|
|
def test_update_forwards_caller_model_over_cfg(self):
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
update = next(t for t in mw.tools if t.name == "update_async_task")
|
|
|
|
client = _fake_sync_client()
|
|
state = {"async_tasks": {"task-abc": self._tracked_task()}}
|
|
with (
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_sync",
|
|
return_value=client,
|
|
),
|
|
patch(
|
|
"EvoScientist.EvoScientist._ensure_config",
|
|
return_value=self._cfg_default(),
|
|
),
|
|
):
|
|
update.func(
|
|
task_id="task-abc",
|
|
message="keep going",
|
|
runtime=self._runtime(state=state),
|
|
)
|
|
|
|
kwargs = client.runs.create.call_args.kwargs
|
|
assert kwargs["config"]["configurable"]["model"] == "free"
|
|
# Upstream update semantics preserved by delegation.
|
|
assert kwargs["multitask_strategy"] == "interrupt"
|
|
|
|
@pytest.mark.asyncio
|
|
async def test_aupdate_forwards_caller_model_over_cfg(self):
|
|
mw = EvoAsyncSubAgentMiddleware(async_subagents=[_standard_spec()])
|
|
update = next(t for t in mw.tools if t.name == "update_async_task")
|
|
|
|
client = _fake_async_client()
|
|
state = {"async_tasks": {"task-abc": self._tracked_task()}}
|
|
with (
|
|
patch(
|
|
"EvoScientist.middleware.expert_async_subagent._ClientCache.get_async",
|
|
return_value=client,
|
|
),
|
|
patch(
|
|
"EvoScientist.EvoScientist._ensure_config",
|
|
return_value=self._cfg_default(),
|
|
),
|
|
):
|
|
await update.coroutine(
|
|
task_id="task-abc",
|
|
message="keep going",
|
|
runtime=self._runtime(state=state),
|
|
)
|
|
|
|
kwargs = client.runs.create.await_args.kwargs
|
|
assert kwargs["config"]["configurable"]["model"] == "free"
|
|
assert kwargs["multitask_strategy"] == "interrupt"
|
|
|
|
def test_caller_scope_reset_after_start(self):
|
|
"""The contextvar must not leak past the tool call — a later launch
|
|
with no override falls back to the config-default, not the prior
|
|
caller's model."""
|
|
from EvoScientist.llm import patches as patches_mod
|
|
|
|
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,
|
|
),
|
|
patch(
|
|
"EvoScientist.EvoScientist._ensure_config",
|
|
return_value=self._cfg_default(),
|
|
),
|
|
):
|
|
start.func(
|
|
description="w",
|
|
subagent_type="literature-review",
|
|
runtime=self._runtime(),
|
|
)
|
|
# Reset restores the default (None) — nothing leaks to the next launch.
|
|
assert not patches_mod._caller_configurable.get()
|