fix(mcp): unattended paths never open browser OAuth; gateway follows mcp_servers edits
A gateway with an OAuth MCP server whose refresh token expired opened a new authorize tab every 300s, all night (92 tabs). Four defects stacked: - The parked-server self-probe re-entered the SDK's authorization-code flow with interactive OAuth enabled. The timed wake is unattended by definition: `_wait_for_reconnect_or_shutdown` now distinguishes "self-probe" from an explicit "reconnect", and `_park` flips the task-local `_oauth_interactive_enabled` off before a self-probe revival. - Gateway MCP discovery (startup, `/reload-mcp`, hot-added multiplex profiles) ran interactive, unlike the CLI's background discovery. All three now run under `suppress_interactive_oauth()`; an expired token parks with the `hermes mcp login` hint instead of a browser. - `_is_interactive()` trusted `sys.stdin.isatty()`, which the Windows CRT reports True for a DEVNULL/detached stdin. `_stdin_is_console()` confirms with `GetConsoleMode` on Windows. - Removing an `mcp_servers` entry (or `enabled: false`) never reached a running gateway; the parked server probed forever. New `reconcile_mcp_servers_with_config()` tears down dropped/disabled servers (via `shutdown_mcp_servers(names=...)`) and connects new ones; a housekeeping chore runs it when config.yaml's (mtime, size) changes. `_select_new_servers` also stops nudging disabled parked servers. Fixes #81830. Fixes the browser-storm item of #96320.
This commit is contained in:
+61
-10
@@ -1808,18 +1808,25 @@ async def _discover_gateway_mcp_tools(config: object) -> None:
|
||||
Under multiplex, run it once per served profile inside that profile's ``_profile_runtime_scope`` and
|
||||
carry the scope into the executor thread with ``copy_context()`` (the same shape as
|
||||
``_run_in_executor_with_context``). See #95518.
|
||||
|
||||
No gateway run can complete a browser OAuth flow (nobody watches its stdout; on Windows its
|
||||
DEVNULL stdin even passes ``isatty``), so discovery runs with interactive OAuth suppressed — the
|
||||
same gate the CLI's background discovery uses. An expired token then parks the server with an
|
||||
actionable ``hermes mcp login`` warning instead of opening an authorize tab.
|
||||
"""
|
||||
from tools.mcp_oauth import suppress_interactive_oauth
|
||||
from tools.mcp_tool_discovery import discover_mcp_tools
|
||||
loop = asyncio.get_running_loop()
|
||||
if not getattr(config, "multiplex_profiles", False):
|
||||
await loop.run_in_executor(None, discover_mcp_tools)
|
||||
return
|
||||
for profile_name, profile_home in _multiplex_profile_homes(config):
|
||||
try:
|
||||
with _profile_runtime_scope(Path(profile_home)):
|
||||
await loop.run_in_executor(None, copy_context().run, discover_mcp_tools)
|
||||
except Exception:
|
||||
logger.warning("MCP tool discovery failed for profile '%s'", profile_name, exc_info=True)
|
||||
with suppress_interactive_oauth():
|
||||
if not getattr(config, "multiplex_profiles", False):
|
||||
await loop.run_in_executor(None, copy_context().run, discover_mcp_tools)
|
||||
return
|
||||
for profile_name, profile_home in _multiplex_profile_homes(config):
|
||||
try:
|
||||
with _profile_runtime_scope(Path(profile_home)):
|
||||
await loop.run_in_executor(None, copy_context().run, discover_mcp_tools)
|
||||
except Exception:
|
||||
logger.warning("MCP tool discovery failed for profile '%s'", profile_name, exc_info=True)
|
||||
|
||||
|
||||
def _platform_has_bot_credential(platform: "Platform", platform_config: "PlatformConfig") -> bool:
|
||||
@@ -4516,6 +4523,49 @@ def _housekeeping_memory_trim() -> None:
|
||||
trim_memory(reason="messaging gateway housekeeping")
|
||||
|
||||
|
||||
def _mcp_config_reconciler(runner=None):
|
||||
"""Chore keeping live MCP servers in step with ``mcp_servers`` on disk: an entry the user
|
||||
removed (or disabled) after boot must stop — a parked one otherwise self-probes every
|
||||
``_PARKED_RETRY_INTERVAL`` for the life of the process (and, before the OAuth gating in this
|
||||
same change, opened a browser tab each time). One ``stat`` per profile per tick; the reconcile
|
||||
runs only when ``config.yaml``'s (mtime, size) changed. Interactive OAuth is suppressed — this
|
||||
runs on a housekeeping thread nobody is watching."""
|
||||
from hermes_cli.config import get_config_path
|
||||
seen: dict = {}
|
||||
|
||||
def _sig(path) -> tuple:
|
||||
try:
|
||||
st = os.stat(path)
|
||||
return (st.st_mtime_ns, st.st_size)
|
||||
except OSError:
|
||||
return (None, None)
|
||||
|
||||
def _reconcile_current(label: str) -> None:
|
||||
from tools.mcp_oauth import suppress_interactive_oauth
|
||||
from tools.mcp_tool_discovery import reconcile_mcp_servers_with_config
|
||||
path = get_config_path()
|
||||
sig = _sig(path)
|
||||
prev = seen.get(label)
|
||||
seen[label] = sig
|
||||
if prev is None or prev == sig:
|
||||
return # first tick just records the baseline; startup discovery already ran
|
||||
with suppress_interactive_oauth():
|
||||
result = reconcile_mcp_servers_with_config()
|
||||
if result["removed"] or result["added"]:
|
||||
logger.info("MCP config changed (%s): removed=%s added=%s", label, result["removed"], result["added"])
|
||||
|
||||
def _tick() -> None:
|
||||
config = getattr(runner, "config", None)
|
||||
if not getattr(config, "multiplex_profiles", False):
|
||||
_reconcile_current("default")
|
||||
return
|
||||
for profile_name, profile_home in _multiplex_profile_homes(config):
|
||||
with _profile_runtime_scope(Path(profile_home)):
|
||||
_reconcile_current(str(profile_name))
|
||||
|
||||
return _tick
|
||||
|
||||
|
||||
def _drain_restart_safe_cron_deliveries(adapters, loop, runner=None) -> None:
|
||||
"""Drain each profile's worker queue through its matching live adapters. A credential-less satellite
|
||||
profile (empty adapter map) drains through the primary's adapters routed by its own profile routes."""
|
||||
@@ -4565,7 +4615,8 @@ def _start_gateway_housekeeping(
|
||||
(60, "Org sync pull tick", _housekeeping_org_skill_sync),
|
||||
(60, "Auto-archive tick", _housekeeping_auto_archive),
|
||||
(1, "Deferred FTS retry tick", _housekeeping_deferred_fts_retry),
|
||||
(1, "gateway housekeeping memory trim", _housekeeping_memory_trim)]
|
||||
(1, "gateway housekeeping memory trim", _housekeeping_memory_trim),
|
||||
(1, "MCP config reconcile", _mcp_config_reconciler(runner))]
|
||||
|
||||
logger.info("Gateway housekeeping started (interval=%ds)", interval)
|
||||
tick_count = 0
|
||||
|
||||
@@ -165,11 +165,12 @@ class GatewayProfileReconcileMixin:
|
||||
from contextvars import copy_context
|
||||
with _log_suppressed(logging.DEBUG, "log routing refresh failed", exc_info=True):
|
||||
_enable_multiplex_log_routing(self.config)
|
||||
from tools.mcp_oauth import suppress_interactive_oauth
|
||||
loop = asyncio.get_running_loop()
|
||||
for profile_name, profile_home in profile_homes:
|
||||
try:
|
||||
from tools.mcp_tool_discovery import discover_mcp_tools
|
||||
with _profile_runtime_scope(Path(profile_home)):
|
||||
with _profile_runtime_scope(Path(profile_home)), suppress_interactive_oauth():
|
||||
await loop.run_in_executor(None, copy_context().run, discover_mcp_tools)
|
||||
except Exception:
|
||||
logger.warning("MCP tool discovery failed for profile '%s'", profile_name, exc_info=True)
|
||||
|
||||
+5
-2
@@ -2353,8 +2353,11 @@ class GatewayTurnMixin:
|
||||
await self._run_in_executor_with_context(lambda: shutdown_mcp_servers(scope=reload_scope))
|
||||
# Explicit reload also re-probes tool availability (check_fn).
|
||||
reprobe_tool_availability()
|
||||
# Reconnect by discovering tools (reads config.yaml fresh).
|
||||
new_tools = await self._run_in_executor_with_context(discover_mcp_tools)
|
||||
# Reconnect by discovering tools (reads config.yaml fresh). A chat command cannot finish
|
||||
# a browser OAuth flow either: an expired token parks with a `hermes mcp login` hint.
|
||||
from tools.mcp_oauth import suppress_interactive_oauth
|
||||
with suppress_interactive_oauth():
|
||||
new_tools = await self._run_in_executor_with_context(discover_mcp_tools)
|
||||
|
||||
connected_servers = _scoped_server_names()
|
||||
if reload_scope is not None:
|
||||
|
||||
@@ -0,0 +1,50 @@
|
||||
"""No gateway MCP path may start a browser OAuth flow (nobody can complete it), and the gateway
|
||||
tracks ``mcp_servers`` edits made after boot."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from pathlib import Path
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.config import GatewayConfig
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_gateway_startup_discovery_suppresses_interactive_oauth(monkeypatch):
|
||||
import gateway.run as gateway_run
|
||||
from tools import mcp_tool_discovery as _mcp_discovery
|
||||
from tools.mcp_oauth import _is_interactive, force_interactive_oauth
|
||||
|
||||
seen: list = []
|
||||
monkeypatch.setattr(_mcp_discovery, "discover_mcp_tools", lambda: seen.append(_is_interactive()) or [])
|
||||
with force_interactive_oauth(): # even a "forced interactive" parent context is overridden
|
||||
await gateway_run._discover_gateway_mcp_tools(GatewayConfig(multiplex_profiles=False))
|
||||
assert seen == [False]
|
||||
|
||||
|
||||
def test_mcp_config_reconciler_runs_only_when_config_changes(monkeypatch, tmp_path: Path):
|
||||
import gateway.run as gateway_run
|
||||
from tools import mcp_tool_discovery as _mcp_discovery
|
||||
from tools.mcp_oauth import _is_interactive
|
||||
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
cfg = tmp_path / "config.yaml"
|
||||
cfg.write_text("mcp_servers:\n linear:\n url: https://x/mcp\n")
|
||||
calls: list = []
|
||||
|
||||
def fake_reconcile():
|
||||
calls.append(_is_interactive())
|
||||
return {"removed": ["linear"], "added": []}
|
||||
|
||||
monkeypatch.setattr(_mcp_discovery, "reconcile_mcp_servers_with_config", fake_reconcile)
|
||||
tick = gateway_run._mcp_config_reconciler(runner=None)
|
||||
|
||||
tick() # baseline only: startup discovery already reflects this file
|
||||
tick()
|
||||
assert calls == []
|
||||
cfg.write_text("model:\n default: x\n") # user removes the entry; size changes -> new signature
|
||||
tick()
|
||||
assert calls == [False], "reconcile must run once per change, with interactive OAuth suppressed"
|
||||
tick()
|
||||
assert calls == [False]
|
||||
@@ -0,0 +1,27 @@
|
||||
"""``_is_interactive()`` must not trust ``isatty()`` alone on Windows: the CRT reports a DEVNULL /
|
||||
detached stdin as a TTY, so the background gateway looked interactive and launched browser OAuth."""
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
@pytest.mark.windows_only
|
||||
def test_devnull_stdin_is_not_a_console_on_windows():
|
||||
code = ("import sys, os; sys.path.insert(0, os.getcwd()); "
|
||||
"from tools.mcp_oauth import _stdin_is_console; print(_stdin_is_console(), sys.stdin.isatty())")
|
||||
proc = subprocess.run([sys.executable, "-c", code], stdin=subprocess.DEVNULL,
|
||||
capture_output=True, text=True, cwd=os.getcwd(), timeout=60)
|
||||
console, isatty = proc.stdout.split()
|
||||
assert isatty == "True", "premise: the Windows CRT calls DEVNULL a tty (else this guard is moot)"
|
||||
assert console == "False"
|
||||
|
||||
|
||||
def test_non_tty_stdin_is_not_interactive(monkeypatch):
|
||||
import io
|
||||
from tools import mcp_oauth
|
||||
monkeypatch.setattr(mcp_oauth.sys, "stdin", io.StringIO())
|
||||
assert mcp_oauth._stdin_is_console() is False
|
||||
assert mcp_oauth._is_interactive() is False
|
||||
@@ -0,0 +1,74 @@
|
||||
"""A parked MCP server's timed self-probe is unattended: it must never start a browser OAuth
|
||||
flow. Left interactive, an expired refresh token opened a new authorize tab every
|
||||
``_PARKED_RETRY_INTERVAL`` for the life of the gateway (92 tabs in one night)."""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
@pytest.mark.no_isolate
|
||||
def test_self_probe_revival_runs_with_interactive_oauth_suppressed(monkeypatch, tmp_path):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
|
||||
from tools import mcp_tool
|
||||
from tools.mcp_oauth import _is_interactive, force_interactive_oauth
|
||||
from tools.mcp_tool import MCPServerTask
|
||||
|
||||
monkeypatch.setattr(mcp_tool, "_MAX_INITIAL_CONNECT_RETRIES", 0)
|
||||
monkeypatch.setattr(mcp_tool, "_PARKED_RETRY_INTERVAL", 0.05)
|
||||
_real_sleep = asyncio.sleep
|
||||
|
||||
async def _fast_sleep(_delay, *a, **kw):
|
||||
await _real_sleep(0)
|
||||
|
||||
monkeypatch.setattr(mcp_tool.asyncio, "sleep", _fast_sleep)
|
||||
attempts: list = []
|
||||
|
||||
class _Task(MCPServerTask):
|
||||
def _is_http(self):
|
||||
return False
|
||||
|
||||
def _deregister_tools(self):
|
||||
self._registered_tool_names = []
|
||||
|
||||
async def _run_stdio(self, config):
|
||||
attempts.append(_is_interactive())
|
||||
if len(attempts) == 1:
|
||||
raise RuntimeError("down") # parks at once (initial ladder = 0 retries)
|
||||
self.session = object()
|
||||
await self._wait_for_lifecycle_event()
|
||||
return "shutdown"
|
||||
|
||||
async def _scenario():
|
||||
# The gateway/TTY parent context IS interactive; the self-probe must still not be.
|
||||
with force_interactive_oauth():
|
||||
task = _Task("srv")
|
||||
run_task = asyncio.ensure_future(task.run({"command": "x"}))
|
||||
for _ in range(400):
|
||||
await _real_sleep(0.01)
|
||||
if len(attempts) >= 2:
|
||||
break
|
||||
await task.shutdown()
|
||||
await asyncio.wait_for(run_task, timeout=5)
|
||||
|
||||
asyncio.run(_scenario())
|
||||
assert attempts[0] is True, "the attended initial connect may still prompt"
|
||||
assert len(attempts) >= 2 and attempts[-1] is False, (
|
||||
"timed self-probe revival ran with interactive OAuth enabled -> browser tab per probe")
|
||||
|
||||
|
||||
def test_explicit_reconnect_request_keeps_the_callers_interactivity():
|
||||
"""Only the TIMED wake is unattended; ``_reconnect_event`` (OAuth recovery, /mcp refresh) is an
|
||||
explicit request and is reported as such."""
|
||||
from tools.mcp_tool import MCPServerTask
|
||||
|
||||
async def _scenario():
|
||||
task = MCPServerTask("srv")
|
||||
task._reconnect_event.set()
|
||||
assert await task._wait_for_reconnect_or_shutdown(timeout=1) == "reconnect"
|
||||
assert await task._wait_for_reconnect_or_shutdown(timeout=0.01) == "self-probe"
|
||||
task._shutdown_event.set()
|
||||
assert await task._wait_for_reconnect_or_shutdown(timeout=0.01) == "shutdown"
|
||||
|
||||
asyncio.run(_scenario())
|
||||
@@ -0,0 +1,72 @@
|
||||
"""Live MCP servers follow ``mcp_servers`` as it is on disk: an entry removed (or disabled) after
|
||||
boot is torn down instead of self-probing for the life of the process."""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
|
||||
@pytest.mark.no_isolate
|
||||
def test_reconcile_tears_down_server_dropped_from_config(monkeypatch, tmp_path):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
from tools import mcp_tool
|
||||
from tools import mcp_tool_config as _config
|
||||
from tools import mcp_tool_discovery as disc
|
||||
from tools import mcp_tool_loop as _loop
|
||||
from tools.mcp_tool import MCPServerTask
|
||||
|
||||
configured = {"linear": {"url": "https://mcp.example.test/mcp", "auth": "oauth"}}
|
||||
monkeypatch.setattr(_config, "_load_mcp_config", lambda: dict(configured))
|
||||
discovered: list = []
|
||||
monkeypatch.setattr(disc, "discover_mcp_tools", lambda *a, **k: discovered.append(1) or [])
|
||||
|
||||
_loop._ensure_mcp_loop()
|
||||
srv = MCPServerTask("linear")
|
||||
|
||||
async def _park():
|
||||
srv._task = asyncio.ensure_future(srv._wait_for_reconnect_or_shutdown())
|
||||
|
||||
asyncio.run_coroutine_threadsafe(_park(), mcp_tool._mcp_loop).result(5)
|
||||
with mcp_tool._lock:
|
||||
mcp_tool._servers["linear"] = srv
|
||||
mcp_tool._server_scope_keys["linear"] = None
|
||||
try:
|
||||
assert disc.reconcile_mcp_servers_with_config() == {"removed": [], "added": []}
|
||||
assert "linear" in mcp_tool._servers and not discovered
|
||||
|
||||
configured.clear() # user deletes the entry
|
||||
assert disc.reconcile_mcp_servers_with_config()["removed"] == ["linear"]
|
||||
assert "linear" not in mcp_tool._servers
|
||||
assert srv._shutdown_event.is_set()
|
||||
|
||||
configured["notion"] = {"url": "https://mcp.notion.test/mcp"}
|
||||
assert disc.reconcile_mcp_servers_with_config()["added"] == ["notion"]
|
||||
assert discovered, "a newly configured server must go through discovery"
|
||||
finally:
|
||||
with mcp_tool._lock:
|
||||
mcp_tool._servers.pop("linear", None)
|
||||
mcp_tool._server_scope_keys.pop("linear", None)
|
||||
_loop._stop_mcp_loop()
|
||||
|
||||
|
||||
def test_disabled_entry_counts_as_dropped(monkeypatch, tmp_path):
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
from tools import mcp_tool
|
||||
from tools import mcp_tool_config as _config
|
||||
from tools import mcp_tool_discovery as disc
|
||||
from tools import mcp_tool_lifecycle as _lifecycle
|
||||
|
||||
monkeypatch.setattr(_config, "_load_mcp_config", lambda: {"linear": {"url": "https://x/mcp", "enabled": False}})
|
||||
torn_down: list = []
|
||||
monkeypatch.setattr(_lifecycle, "shutdown_mcp_servers", lambda **kw: torn_down.append(kw))
|
||||
monkeypatch.setattr(disc, "discover_mcp_tools", lambda *a, **k: [])
|
||||
with mcp_tool._lock:
|
||||
mcp_tool._servers["linear"] = object()
|
||||
mcp_tool._server_scope_keys["linear"] = None
|
||||
try:
|
||||
assert disc.reconcile_mcp_servers_with_config()["removed"] == ["linear"]
|
||||
assert torn_down == [{"scope": None, "names": {"linear"}}]
|
||||
finally:
|
||||
with mcp_tool._lock:
|
||||
mcp_tool._servers.pop("linear", None)
|
||||
mcp_tool._server_scope_keys.pop("linear", None)
|
||||
+23
-4
@@ -169,16 +169,35 @@ def _cached_redirect(storage: "HermesTokenStorage | None") -> "tuple[str | None,
|
||||
return uri, port
|
||||
|
||||
|
||||
def _stdin_is_console() -> bool:
|
||||
"""A human can type on stdin. ``isatty()`` alone is wrong on Windows: the CRT reports True for
|
||||
a DEVNULL / detached / CREATE_NO_WINDOW stdin (the gateway's), so a background process looked
|
||||
interactive and launched browser OAuth flows nobody could finish. Confirm with the console API
|
||||
there: ``GetConsoleMode`` fails on anything that is not a real console handle."""
|
||||
try:
|
||||
if not sys.stdin.isatty():
|
||||
return False
|
||||
except (AttributeError, ValueError):
|
||||
return False
|
||||
if os.name != "nt":
|
||||
return True
|
||||
try:
|
||||
import ctypes
|
||||
import msvcrt
|
||||
handle = msvcrt.get_osfhandle(sys.stdin.fileno())
|
||||
mode = ctypes.c_ulong()
|
||||
return bool(ctypes.windll.kernel32.GetConsoleMode(ctypes.c_void_p(handle), ctypes.byref(mode)))
|
||||
except Exception:
|
||||
return False
|
||||
|
||||
|
||||
def _is_interactive() -> bool:
|
||||
"""True if we can reasonably expect to interact with a user."""
|
||||
if not _oauth_interactive_enabled.get():
|
||||
return False
|
||||
if _oauth_interactive_forced.get():
|
||||
return True
|
||||
try:
|
||||
return sys.stdin.isatty()
|
||||
except (AttributeError, ValueError):
|
||||
return False
|
||||
return _stdin_is_console()
|
||||
|
||||
|
||||
def _raise_if_non_interactive(lead: str) -> None:
|
||||
|
||||
@@ -255,8 +255,9 @@ def _select_new_servers(servers: Dict[str, dict]) -> Dict[str, dict]:
|
||||
if keys[k] not in _core._servers and keys[k] not in _core._server_connecting
|
||||
and keys[k] not in _core._lazy_server_configs
|
||||
and _enabled(v) and not _connect_cooldown_active(k)}
|
||||
stale_cached = [_core._servers[keys[k]] for k in servers
|
||||
if keys[k] in _core._servers and getattr(_core._servers[keys[k]], "session", None) is None]
|
||||
stale_cached = [_core._servers[keys[k]] for k, v in servers.items()
|
||||
if keys[k] in _core._servers and _enabled(v)
|
||||
and getattr(_core._servers[keys[k]], "session", None) is None]
|
||||
for srv_name in new_servers:
|
||||
_core._server_connecting.add(keys[srv_name])
|
||||
_core._server_scope_keys[keys[srv_name]] = current_scope
|
||||
@@ -473,6 +474,33 @@ def discover_mcp_tools(allowed_mcp_names: Optional[List[str]] = None) -> List[st
|
||||
cookie.release()
|
||||
|
||||
|
||||
def reconcile_mcp_servers_with_config() -> Dict[str, List[str]]:
|
||||
"""Bring the live server set in step with ``mcp_servers`` as it is on disk NOW: tear down
|
||||
servers that were removed from config or set ``enabled: false`` (a parked server keeps
|
||||
self-probing forever otherwise — for hours after the user deleted its entry), then connect
|
||||
anything newly configured via :func:`discover_mcp_tools`. Scoped to the current registry
|
||||
scope (one multiplexed profile's config prunes only its own connections). Returns
|
||||
``{"removed": [...], "added": [...]}``; a no-op when nothing changed."""
|
||||
servers = _config._load_mcp_config()
|
||||
wanted = {name for name, cfg in servers.items() if _enabled(cfg)}
|
||||
scope = _core._mcp_registry_scope()
|
||||
with _core._lock:
|
||||
owned = [key for key, owner in _core._server_scope_keys.items() if owner == scope]
|
||||
live = {_key_name(key) for key in owned if key in _core._servers}
|
||||
stale = sorted(live - wanted)
|
||||
if stale:
|
||||
logger.info("MCP server(s) %s no longer in config (or disabled); disconnecting", ", ".join(stale))
|
||||
_lifecycle.shutdown_mcp_servers(scope=scope, names=set(stale))
|
||||
with _core._lock:
|
||||
known = {_key_name(key) for key, owner in _core._server_scope_keys.items()
|
||||
if owner == scope and (key in _core._servers or key in _core._server_connecting)}
|
||||
known |= {_key_name(key) for key in _core._lazy_server_configs}
|
||||
added = sorted(wanted - known)
|
||||
if added:
|
||||
discover_mcp_tools()
|
||||
return {"removed": stale, "added": added}
|
||||
|
||||
|
||||
def is_mcp_tool_parallel_safe(tool_name: str) -> bool:
|
||||
"""True when the tool's server opted into ``supports_parallel_tool_calls`` (provenance
|
||||
captured at registration, never the ambiguous ``mcp__{server}__{tool}`` shape)."""
|
||||
|
||||
+18
-14
@@ -127,28 +127,32 @@ def _reregister_orphaned_adopters() -> None:
|
||||
reset_hermes_home_override(home_token)
|
||||
|
||||
|
||||
def shutdown_mcp_servers(*, scope: Optional[str] = None):
|
||||
def shutdown_mcp_servers(*, scope: Optional[str] = None, names: Optional[set] = None):
|
||||
"""Close MCP server connections (in parallel) and stop the background loop. Each server
|
||||
Task is signalled to exit its own ``async with`` so the anyio cancel-scope cleanup runs in
|
||||
the Task that opened it. ``scope`` restricts teardown to one multiplexed profile's servers
|
||||
(its ``/reload-mcp`` must not kill other profiles') and leaves the shared loop running if
|
||||
anything else is still connected."""
|
||||
anything else is still connected. ``names`` restricts it further to those server names
|
||||
(dropped-from-config pruning); other servers' bookkeeping is untouched."""
|
||||
from tools.mcp_tool_scope import _key_name
|
||||
with _core._lock:
|
||||
selected = [key for key in _core._servers if scope is None or _core._server_scope_keys.get(key) == scope]
|
||||
if names is not None:
|
||||
selected = [key for key in selected if _key_name(key) in names]
|
||||
servers_snapshot = [_core._servers[key] for key in selected]
|
||||
selected_status = (
|
||||
set(_core._servers) | set(_core._server_scope_keys)
|
||||
| set(_core._server_tool_scopes)
|
||||
| set(_core._server_connecting) | set(_core._server_connect_errors)
|
||||
if scope is None else {
|
||||
key for key, owner in _core._server_scope_keys.items() if owner == scope
|
||||
}
|
||||
)
|
||||
if names is not None:
|
||||
selected_status = set(selected)
|
||||
elif scope is None:
|
||||
selected_status = (
|
||||
set(_core._servers) | set(_core._server_scope_keys)
|
||||
| set(_core._server_tool_scopes)
|
||||
| set(_core._server_connecting) | set(_core._server_connect_errors))
|
||||
else:
|
||||
selected_status = {key for key, owner in _core._server_scope_keys.items() if owner == scope}
|
||||
# Adopters of the connections being torn down lose their overlays with the tasks' own
|
||||
# ``_deregister_tools``; remember them so the next discovery pass re-registers them
|
||||
# (``_reregister_orphaned_adopters``).
|
||||
if scope is not None:
|
||||
from tools.mcp_tool_scope import _key_name
|
||||
for key in selected:
|
||||
for adopter in _core._server_tool_scopes.get(key, ()):
|
||||
if adopter != scope:
|
||||
@@ -176,7 +180,7 @@ def shutdown_mcp_servers(*, scope: Optional[str] = None):
|
||||
_core._servers.pop(key, None)
|
||||
_core._server_scope_keys.pop(key, None)
|
||||
clear_selected_status()
|
||||
_clear_connect_cooldowns(None if scope is None else selected_status)
|
||||
_clear_connect_cooldowns(None if scope is None and names is None else selected_status)
|
||||
|
||||
with _core._lock:
|
||||
loop = _core._mcp_loop
|
||||
@@ -195,8 +199,8 @@ def shutdown_mcp_servers(*, scope: Optional[str] = None):
|
||||
with _core._lock:
|
||||
if not servers_snapshot:
|
||||
clear_selected_status()
|
||||
_clear_connect_cooldowns(None if scope is None else selected_status)
|
||||
_loop._stop_mcp_loop(only_if_idle=scope is not None)
|
||||
_clear_connect_cooldowns(None if scope is None and names is None else selected_status)
|
||||
_loop._stop_mcp_loop(only_if_idle=scope is not None or names is not None)
|
||||
|
||||
|
||||
def _take_reapable_pids(include_active: bool, server_name: Optional[str]) -> tuple[Dict[int, str], Dict[int, int]]:
|
||||
|
||||
@@ -110,8 +110,8 @@ class MCPServerRunMixin:
|
||||
return "reconnect"
|
||||
|
||||
async def _wait_for_reconnect_or_shutdown(self, timeout: Optional[float] = None) -> str:
|
||||
"""Parked wait: ``"shutdown"`` or ``"reconnect"`` (explicit, or the ``timeout`` self-probe;
|
||||
event cleared first). Shutdown wins a tie."""
|
||||
"""Parked wait: ``"shutdown"``, ``"reconnect"`` (explicit request; event cleared first) or
|
||||
``"self-probe"`` (``timeout`` elapsed with neither). Shutdown wins a tie."""
|
||||
shutdown_task, reconnect_task = self._event_waiters()
|
||||
try:
|
||||
await asyncio.wait({shutdown_task, reconnect_task}, return_when=asyncio.FIRST_COMPLETED, timeout=timeout)
|
||||
@@ -119,6 +119,8 @@ class MCPServerRunMixin:
|
||||
await self._cancel_waiters(shutdown_task, reconnect_task)
|
||||
if self._shutdown_event.is_set():
|
||||
return "shutdown"
|
||||
if not self._reconnect_event.is_set():
|
||||
return "self-probe"
|
||||
self._reconnect_event.clear()
|
||||
return "reconnect"
|
||||
|
||||
@@ -138,10 +140,19 @@ class MCPServerRunMixin:
|
||||
self._was_parked = True
|
||||
self._deregister_tools()
|
||||
self._reconnect_event.clear()
|
||||
if await self._wait_for_reconnect_or_shutdown(timeout=_core._PARKED_RETRY_INTERVAL) == "shutdown":
|
||||
outcome = await self._wait_for_reconnect_or_shutdown(timeout=_core._PARKED_RETRY_INTERVAL)
|
||||
if outcome == "shutdown":
|
||||
return True
|
||||
logger.debug("MCP server '%s': attempting revival %s (self-probe or explicit reconnect request); "
|
||||
"rebuilding transport.", self.name, revival_reason)
|
||||
# Nobody asked for this revival: a self-probe must never open a browser OAuth flow. The
|
||||
# OAuth provider runs inside THIS task (the SDK's auth flow sits in the transport), so a
|
||||
# task-local ContextVar reaches it; it stays set for the task's life — every later
|
||||
# revival of a once-parked server is unattended too. Left interactive, an expired
|
||||
# refresh token opened a new authorize tab every _PARKED_RETRY_INTERVAL, all night.
|
||||
if outcome == "self-probe":
|
||||
from tools.mcp_oauth import _oauth_interactive_enabled
|
||||
_oauth_interactive_enabled.set(False)
|
||||
logger.debug("MCP server '%s': attempting revival %s (%s); rebuilding transport.",
|
||||
self.name, revival_reason, outcome)
|
||||
return False
|
||||
|
||||
async def _prepare_run(self, config: dict) -> bool:
|
||||
|
||||
@@ -647,6 +647,10 @@ If you change MCP config, use:
|
||||
|
||||
This reloads MCP servers from config and refreshes the available tool list. It is also the explicit way to re-probe availability-gated tools (Docker, `HASS_TOKEN`, OAuth…): a session's tool set is otherwise frozen, so a credential or daemon that appears mid-session is only picked up on `/reload-mcp`, `/new`, or context compaction. For runtime tool changes pushed by the server itself, see [Dynamic Tool Discovery](#dynamic-tool-discovery) above.
|
||||
|
||||
A running messaging gateway (`hermes gateway run`) also watches `config.yaml` on its own: within about a minute of you removing an `mcp_servers` entry or setting `enabled: false`, that server's connection is torn down; a newly added entry is connected. No restart or `/reload-mcp` needed for the edit to take effect.
|
||||
|
||||
**Expired OAuth tokens in the background.** The gateway, `/reload-mcp`, and the periodic self-probe of a parked server never open a browser — nobody is there to complete the flow. When a refresh token dies, the server parks with a warning in `gateway.log` and you re-authorize once with `hermes mcp login <server>` (or the Desktop/dashboard *Authorize* button); the parked server picks the new token up on its next probe.
|
||||
|
||||
### Toolsets
|
||||
|
||||
Each configured MCP server also creates a runtime toolset when it contributes at least one registered tool:
|
||||
|
||||
Reference in New Issue
Block a user