Files
EvoScientist/EvoScientist/EvoScientist.py
T
Wiktor Cupiał 9e51ec6fdd feat(cmd): add /model-fallback command (#196)
* feat(cmd): add /model-fallback command

* fix: apply feedback

* fix: lock usage with _fallback_chain

* fix: apply feedback

* fix: apply feedback

* feat: add tests

* fix: tests

* Update EvoScientist/middleware/model_fallback.py

Co-authored-by: dinos <dinospk1999@gmail.com>

---------

Co-authored-by: dinos <dinospk1999@gmail.com>
2026-05-07 15:35:26 +02:00

636 lines
23 KiB
Python

"""EvoScientist Agent graph construction.
This module defines the agent graph and its factory functions. All heavy
initialization (deepagents, backends, LLM, middleware) is deferred to first
use so that importing this module is fast and non-agent CLI commands
(``EvoSci config list``, ``EvoSci onboard``) never pay the cost.
Usage:
from EvoScientist import EvoScientist_agent
# Notebook / programmatic usage
for state in EvoScientist_agent.stream(
{"messages": [HumanMessage(content="your question")]},
config={"configurable": {"thread_id": "1"}},
stream_mode="values",
):
...
"""
import json
import logging
import os
from datetime import datetime
from pathlib import Path
from langchain.agents.middleware import AgentMiddleware, HumanInTheLoopMiddleware
from . import paths as _paths_mod
from .config import apply_config_to_env, get_effective_config
from .paths import set_active_workspace, set_workspace_root
from .prompts import RESEARCHER_INSTRUCTIONS, get_system_prompt
# Suppress noisy warnings from deepagents skill loader (non-string frontmatter fields, etc.)
logging.getLogger("deepagents.middleware.skills").setLevel(logging.ERROR)
# =============================================================================
# Constants
# =============================================================================
SUBAGENTS_CONFIG = Path(__file__).parent / "subagents"
SKILLS_DIR = str(Path(__file__).parent / "skills")
# =============================================================================
# Lazy state — initialized on first use, not at import time
# =============================================================================
_config = None
_chat_model = None
# Track the (model, provider) binding of _chat_model so cache invalidates
# when config.model/provider change (e.g. via /model). Without this,
# _ensure_chat_model() returns the stale cached instance even after
# _ensure_config(new_cfg) has overwritten the active config — causing
# /model switch to lag one step (see issue #179).
_chat_model_key: tuple[str | None, str | None] | None = None
# Cache MCP tools by the effective config signature to avoid reconnecting
# to MCP servers on every `/new` when config is unchanged.
_MCP_TOOLS_CACHE_KEY: str | None = None
_MCP_TOOLS_CACHE_VALUE: dict[str, list] | None = None
# Default agent (no checkpointer) — used by langgraph dev / LangSmith / notebooks.
# Lazily constructed on first access so MCP tools are included without
# spawning subprocesses at import time.
_EvoScientist_agent = None
# =============================================================================
# Lazy initialization helpers
# =============================================================================
def _ensure_config(config=None):
"""Return cached config. If *config* is passed, cache and use it."""
global _config
if config is not None:
_config = config
apply_config_to_env(_config)
if _config is None:
_config = get_effective_config()
apply_config_to_env(_config)
return _config
def _replace_chat_model(instance, key: tuple[str | None, str | None]) -> None:
"""Install a new chat model and propagate the related invariants.
Single write point for ``_chat_model`` / ``_chat_model_key`` /
``_EvoScientist_agent``: both ``_ensure_chat_model`` (cache-miss
rebuild) and ``set_chat_model`` (explicit switch via ``/model``)
funnel through here so the three globals can never drift.
"""
global _chat_model, _chat_model_key, _EvoScientist_agent
_chat_model = instance
_chat_model_key = key
# The lazy default agent captured a reference to the previous
# ``_chat_model`` at build time, so it must be rebuilt on next access.
_EvoScientist_agent = None
def _ensure_chat_model():
"""Return cached chat model, rebuilding if cfg.model/provider changed.
The cache key is the current config's ``(model, provider)``. If it
differs from the key that built ``_chat_model``, rebuild — this makes
``create_cli_agent(config=temp_cfg)`` bind the freshly requested model
into the new agent without requiring callers to interleave
``set_chat_model()`` calls in any particular order.
"""
from .llm import get_chat_model
cfg = _ensure_config()
key = (cfg.model, cfg.provider)
if _chat_model is None or _chat_model_key != key:
_replace_chat_model(
get_chat_model(model=cfg.model, provider=cfg.provider),
key,
)
return _chat_model
def set_chat_model(model: str, provider: str | None = None):
"""Replace the cached chat model with a new one.
Called by ``/model`` to switch the LLM mid-session. No-op when the
cache already holds the requested ``(model, provider)`` — avoids
spawning a second ``get_chat_model`` instance (and its HTTP client)
under the ``/model`` flow where ``_ensure_chat_model`` has already
rebuilt ``_chat_model`` during the preceding ``_load_agent`` call.
Returns the current chat model instance.
"""
from .llm import get_chat_model
key = (model, provider)
if _chat_model is None or _chat_model_key != key:
_replace_chat_model(get_chat_model(model=model, provider=provider), key)
return _chat_model
# =============================================================================
# MCP caching
# =============================================================================
def _load_mcp_config_once() -> tuple[str, dict]:
"""Load MCP config and return ``(signature, config)``."""
from .mcp.client import load_mcp_config
cfg = load_mcp_config()
if not cfg:
return "", {}
try:
sig = json.dumps(cfg, sort_keys=True, ensure_ascii=True)
except TypeError:
sig = repr(cfg)
return sig, cfg
def _load_mcp_tools_cached(on_progress=None) -> dict[str, list]:
"""Load MCP tools with config-aware caching.
Args:
on_progress: Optional per-server progress callback forwarded to
:func:`EvoScientist.mcp.load_mcp_tools`. Only invoked on a
cache miss — cached replays don't re-emit progress events.
"""
global _MCP_TOOLS_CACHE_KEY, _MCP_TOOLS_CACHE_VALUE
from .mcp import load_mcp_tools
cfg_key, cfg = _load_mcp_config_once()
if not cfg_key:
_MCP_TOOLS_CACHE_KEY = ""
_MCP_TOOLS_CACHE_VALUE = {}
return {}
if _MCP_TOOLS_CACHE_KEY == cfg_key and _MCP_TOOLS_CACHE_VALUE is not None:
return {k: list(v) for k, v in _MCP_TOOLS_CACHE_VALUE.items()}
loaded = load_mcp_tools(config=cfg, on_progress=on_progress)
_MCP_TOOLS_CACHE_KEY = cfg_key
_MCP_TOOLS_CACHE_VALUE = {k: list(v) for k, v in loaded.items()}
return {k: list(v) for k, v in loaded.items()}
# =============================================================================
# Agent construction helpers
# =============================================================================
def _inject_subagent_middleware(subs: list[dict]) -> None:
"""Ensure every subagent gets error handling and context management middleware.
Without this, subagent tool errors are caught by LangGraph's default
ToolNode handler which produces terse messages without tracebacks or
retry guidance — reducing the subagent's ability to self-recover.
"""
from .middleware import (
ContextOverflowMapperMiddleware,
ToolErrorHandlerMiddleware,
create_context_editing_middleware,
)
for sa in subs:
sa.setdefault("middleware", []).extend(
[
# No ``model=`` — subagents share the main agent's model,
# so defer to the factory's ``_ensure_chat_model()`` fallback.
create_context_editing_middleware(),
ToolErrorHandlerMiddleware(),
ContextOverflowMapperMiddleware(),
]
)
def _build_prompt_refs() -> dict:
"""Build prompt references with the current date (not frozen at import)."""
return {
"RESEARCHER_INSTRUCTIONS": RESEARCHER_INSTRUCTIONS.format(
date=datetime.now().strftime("%Y-%m-%d"),
),
}
def _maybe_swap_async_subagents(subs: list) -> list:
"""Replace ``_async``-flagged sub-agents with ``AsyncSubAgent`` specs when enabled.
Reads the ``_async`` field carried through by ``utils.load_subagents._build_one``
(sourced from each yaml's ``async: true`` flag). When
``config.enable_async_subagents`` is also set, those sub-agents are
swapped from synchronous in-process dicts to ``AsyncSubAgent`` references
pointing at the langgraph dev graph of the same name.
The deployed graphs live in ``EvoScientist.langgraph_dev.graphs`` and
are registered in ``EvoScientist/langgraph_dev/langgraph.json``.
Adding a new async sub-agent requires no change here — flip
``async: true`` in its yaml and create the matching deployment graph.
All return paths strip the internal ``_async`` field from sub-agent dicts
before handoff, since deepagents may schema-validate the kwarg.
"""
cfg = _ensure_config()
if not getattr(cfg, "enable_async_subagents", False):
# Async fully disabled — strip the internal flag before handoff.
for s in subs:
s.pop("_async", None)
return subs
# Guard: if the langgraph dev subprocess never came up (port conflict,
# binary missing, etc.), routing sub-agents to a dead URL produces hangs
# and confusing tool errors. Fall back to in-process sync delegation.
from .langgraph_dev.manager import is_async_subagents_available
if not is_async_subagents_available():
logging.getLogger(__name__).warning(
"enable_async_subagents=true but langgraph dev is not reachable; "
"falling back to in-process sync delegation for all sub-agents."
)
# Strip the internal ``_async`` flag (carried from ``load_subagents``)
# before sub-agents reach deepagents — it's never a deepagents key.
for s in subs:
s.pop("_async", None)
return subs
# The ``_async`` flag was set by ``utils.load_subagents._build_one`` from
# each yaml's ``async:`` field. No need to re-parse the yaml files here.
async_specs: dict[str, str] = {
s["name"]: s.get("description", "") for s in subs if s.get("_async")
}
if not async_specs:
for s in subs:
s.pop("_async", None)
return subs
from deepagents import AsyncSubAgent
port = int(getattr(cfg, "langgraph_dev_port", 6174))
out = []
# MCP tools routed to async sub-agents (via ``expose_to: <name>`` in
# mcp.yaml) ARE delivered — the deployed factory
# ``subagents/_factory.py:build_async_subagent_graph`` loads its own MCP
# connection per server (cost: one extra MCP server subprocess per
# exposed server, since stdio transports can't share across processes).
for s in subs:
name = s.get("name")
if name in async_specs:
out.append(
AsyncSubAgent(
name=name,
description=async_specs[name],
graph_id=name,
url=f"http://localhost:{port}",
)
)
else:
# Strip the internal flag before handoff to deepagents.
s.pop("_async", None)
out.append(s)
return out
def _build_base_kwargs(base_backend, base_middleware):
"""Build agent kwargs *without* MCP (fast, no subprocess spawning)."""
from .tools import skill_manager, tavily_search, think_tool
from .utils import load_subagents
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]
subs = load_subagents(
SUBAGENTS_CONFIG,
tool_registry=tool_registry,
prompt_refs=_build_prompt_refs(),
)
_inject_subagent_middleware(subs)
subs = _maybe_swap_async_subagents(subs)
return {
"name": "EvoScientist",
"model": _ensure_chat_model(),
"tools": list(base_tools),
"backend": base_backend,
"subagents": subs,
"middleware": base_middleware,
"system_prompt": get_system_prompt(),
"skills": ["/skills/"],
}
def load_mcp_and_build_kwargs(base_backend, base_middleware, *, on_mcp_progress=None):
"""Load MCP tools (cached by config) and build agent kwargs.
Re-connects to MCP servers only when the effective MCP config changes.
Falls back to base kwargs if no MCP configured.
Args:
on_mcp_progress: Optional per-server progress callback. Forwarded
to the MCP loader so UIs can render live status.
"""
from .tools import skill_manager, tavily_search, think_tool
from .utils import load_subagents
mcp_by_agent = _load_mcp_tools_cached(on_progress=on_mcp_progress)
if not mcp_by_agent:
return _build_base_kwargs(base_backend, base_middleware)
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]
# Fresh tool registry — start from base tools + MCP tools
registry = dict(tool_registry)
for tools in mcp_by_agent.values():
for t in tools:
registry[t.name] = t
mcp_main = mcp_by_agent.pop("main", [])
subs = load_subagents(
SUBAGENTS_CONFIG,
tool_registry=registry,
prompt_refs=_build_prompt_refs(),
)
_inject_subagent_middleware(subs)
# Inject MCP tools into subagents by name
for sa in subs:
if sa_tools := mcp_by_agent.get(sa["name"], []):
sa.setdefault("tools", []).extend(sa_tools)
# 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)
return {
"name": "EvoScientist",
"model": _ensure_chat_model(),
"tools": base_tools + mcp_main,
"backend": base_backend,
"subagents": subs,
"middleware": base_middleware,
"system_prompt": get_system_prompt(),
"skills": ["/skills/"],
}
# =============================================================================
# Default agent (langgraph dev / notebooks)
# =============================================================================
def _get_default_backend():
"""Build the default composite backend from current paths."""
from deepagents.backends import CompositeBackend, FilesystemBackend
from .backends import CustomSandboxBackend, MergedSkillsBackend
workspace_dir = str(_paths_mod.WORKSPACE_ROOT)
set_active_workspace(workspace_dir)
memory_dir = str(_paths_mod.MEMORIES_DIR)
user_skills_dir = str(_paths_mod.USER_SKILLS_DIR)
global_skills_dir = str(_paths_mod.GLOBAL_SKILLS_DIR)
ws_backend = CustomSandboxBackend(
root_dir=workspace_dir,
virtual_mode=True,
timeout=300,
)
sk_backend = MergedSkillsBackend(
primary_dir=user_skills_dir,
global_dir=global_skills_dir,
secondary_dir=SKILLS_DIR,
)
mem_backend = FilesystemBackend(
root_dir=memory_dir,
virtual_mode=True,
)
return CompositeBackend(
default=ws_backend,
routes={
"/skills/": sk_backend,
"/memories/": mem_backend,
},
)
def _get_default_middleware():
"""Build the default middleware list."""
from .middleware import (
ContextOverflowMapperMiddleware,
ModelFallbackMiddleware,
ToolErrorHandlerMiddleware,
create_context_editing_middleware,
create_memory_middleware,
create_tool_selector_middleware,
load_fallback_chain,
)
cfg = _ensure_config()
if cfg.model_fallbacks:
load_fallback_chain(cfg.model_fallbacks)
model = _ensure_chat_model()
memory_dir = str(_paths_mod.MEMORIES_DIR)
mw = [
create_context_editing_middleware(model),
ModelFallbackMiddleware(),
ContextOverflowMapperMiddleware(),
ToolErrorHandlerMiddleware(),
*create_tool_selector_middleware(model=model),
create_memory_middleware(memory_dir, extraction_model=model),
]
if cfg.enable_ask_user and not cfg.auto_mode:
from .middleware.ask_user import AskUserMiddleware
mw.insert(0, AskUserMiddleware())
return mw
def _get_default_agent():
"""Build the default agent (with MCP, no checkpointer) on first access.
When invoked from the langgraph dev subprocess (env var
``EVOSCIENTIST_DEPLOYED_NO_MCP=true``, set by
``langgraph_dev.manager.start_langgraph_dev``), MCP loading is skipped to
avoid duplicating the CLI's MCP server pool — the deployed main agent
is currently only reachable via HTTP for Web UI / SDK clients (none in
use yet), so paying for a second copy of the same MCP servers is pure
waste. Re-enable later by removing the env var when MCP-needing remote
callers are introduced.
"""
global _EvoScientist_agent
if _EvoScientist_agent is None:
from deepagents import create_deep_agent
cfg = _ensure_config()
be = _get_default_backend()
mw = _get_default_middleware()
# HITL on main agent only (mirrors create_cli_agent). Use middleware,
# not interrupt_on= kwarg — the kwarg propagates to every subagent and
# breaks parallel execute calls (multi-pending-interrupt LangGraph
# error). See PR #202.
if not cfg.auto_approve:
mw.append(HumanInTheLoopMiddleware(interrupt_on={"execute": True}))
if os.environ.get("EVOSCIENTIST_DEPLOYED_NO_MCP", "").lower() == "true":
kwargs = _build_base_kwargs(be, mw)
else:
kwargs = load_mcp_and_build_kwargs(be, mw)
_EvoScientist_agent = create_deep_agent(
**kwargs,
).with_config({"recursion_limit": cfg.recursion_limit})
return _EvoScientist_agent
def __getattr__(name: str):
if name == "EvoScientist_agent":
return _get_default_agent()
# Backward compat for module-level names
if name == "chat_model":
return _ensure_chat_model()
if name == "SYSTEM_PROMPT":
return get_system_prompt()
if name == "backend":
return _get_default_backend()
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
# =============================================================================
# CLI agent factory
# =============================================================================
def create_cli_agent(
workspace_dir: str | None = None,
checkpointer=None,
config=None,
*,
on_mcp_progress=None,
):
"""Create agent with checkpointer for CLI multi-turn support.
A fresh backend is constructed on every call using the current
``paths.WORKSPACE_ROOT`` (or the explicit *workspace_dir*), so
runtime ``set_workspace_root()`` changes are always respected.
Args:
workspace_dir: Per-session workspace directory. If ``None``,
defaults to the current ``paths.WORKSPACE_ROOT``.
checkpointer: Optional LangGraph checkpointer. If ``None``,
falls back to ``InMemorySaver`` (non-persistent).
config: Optional pre-loaded ``EvoScientistConfig``. If ``None``,
loads from file/env/defaults. Passing this avoids double
loading when the CLI has already loaded config.
"""
import os as _os
from deepagents import create_deep_agent
from deepagents.backends import CompositeBackend, FilesystemBackend
from . import paths as _paths
from .backends import CustomSandboxBackend, MergedSkillsBackend
from .middleware import (
ContextOverflowMapperMiddleware,
ModelFallbackMiddleware,
ToolErrorHandlerMiddleware,
create_context_editing_middleware,
create_memory_middleware,
create_tool_selector_middleware,
load_fallback_chain,
)
cfg = _ensure_config(config)
if cfg.model_fallbacks:
load_fallback_chain(cfg.model_fallbacks)
if checkpointer is None:
from langgraph.checkpoint.memory import InMemorySaver
checkpointer = InMemorySaver()
# When no explicit workspace_dir is provided, apply config.default_workdir
# as a fallback. This covers direct callers (notebooks, iMessage server)
# that never call set_workspace_root() themselves. CLI callers always
# pass workspace_dir explicitly, so their --workdir is never overwritten.
if workspace_dir is None:
if cfg.default_workdir:
set_workspace_root(
_os.path.abspath(_os.path.expanduser(cfg.default_workdir))
)
workspace_dir = str(_paths.WORKSPACE_ROOT)
# Read paths dynamically so runtime set_workspace_root() changes are picked up
_mem_dir = str(_paths.MEMORIES_DIR)
_usr_skills_dir = str(_paths.USER_SKILLS_DIR)
_global_skills_dir = str(_paths.GLOBAL_SKILLS_DIR)
# Always construct fresh backends from current paths (avoids stale
# module-level backend when workspace root changed at runtime).
set_active_workspace(workspace_dir)
ws_backend = CustomSandboxBackend(
root_dir=workspace_dir,
virtual_mode=True,
timeout=300,
)
sk_backend = MergedSkillsBackend(
primary_dir=_usr_skills_dir,
global_dir=_global_skills_dir,
secondary_dir=SKILLS_DIR,
)
mem_backend = FilesystemBackend(
root_dir=_mem_dir,
virtual_mode=True,
)
be = CompositeBackend(
default=ws_backend,
routes={
"/skills/": sk_backend,
"/memories/": mem_backend,
},
)
model = _ensure_chat_model()
mw: list[AgentMiddleware] = [
create_context_editing_middleware(model),
ModelFallbackMiddleware(),
ContextOverflowMapperMiddleware(),
ToolErrorHandlerMiddleware(),
*create_tool_selector_middleware(model=model),
create_memory_middleware(_mem_dir, extraction_model=model),
]
if cfg.enable_ask_user and not cfg.auto_mode:
from .middleware.ask_user import AskUserMiddleware
mw.insert(0, AskUserMiddleware())
# HITL on main agent only — passing `interrupt_on=` to create_deep_agent
# would propagate it to every subagent, breaking parallel execute calls
# (multi-pending-interrupt LangGraph error).
if not cfg.auto_approve:
mw.append(HumanInTheLoopMiddleware(interrupt_on={"execute": True}))
# Re-load MCP tools from current config (picks up /mcp add changes)
kwargs = load_mcp_and_build_kwargs(be, mw, on_mcp_progress=on_mcp_progress)
return create_deep_agent(
**kwargs,
checkpointer=checkpointer,
).with_config({"recursion_limit": cfg.recursion_limit})