fix(mcp): reconcile chore lives in run_profile_reconcile; prune lazy + mid-connect servers
The MCP config reconciler was appended to the gateway/run.py facade; it moves to gateway/run_profile_reconcile.py, which already owns post-boot MCP discovery, and run.py keeps only the chore-table entry. reconcile_mcp_servers_with_config() also drops a schema-cache (lazy) registration whose entry is gone (its cached tools would otherwise stay callable and spawn the server on first use) and reports a dropped server still mid-connect as "pending"; the chore retries on the next tick without waiting for another config edit. test_cron_delivery_housekeeping neutralizes the chore: it pins the exact scope/drain sequence of the housekeeping loop and the new chore enters each profile's scope once per tick.
This commit is contained in:
@@ -4523,49 +4523,6 @@ 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."""
|
||||
@@ -4597,6 +4554,7 @@ def _start_gateway_housekeeping(
|
||||
"""Background thread for gateway-only periodic chores (NOT cron). Separate from the cron trigger
|
||||
so chores run under any ``CronScheduler`` provider (external scale-to-zero has no 60s loop).
|
||||
Cadences are ticks of ``interval``; inner gates own the real cadence."""
|
||||
from gateway.run_profile_reconcile import _mcp_config_reconciler
|
||||
chores: list[tuple[int, str, Any]] = []
|
||||
if adapters is not None or runner is not None:
|
||||
# Restart-safe cron workers run outside the gateway cgroup and queue their final send for
|
||||
|
||||
@@ -210,3 +210,50 @@ class GatewayProfileReconcileMixin:
|
||||
from plugins.memory.holographic.store import MemoryStore
|
||||
MemoryStore.release_all_under(home)
|
||||
logger.info("[MULTIPLEX] Profile '%s' deleted — %d adapter(s) stopped and unrouted", name, len(adapters))
|
||||
|
||||
|
||||
def _mcp_config_reconciler(runner=None):
|
||||
"""Housekeeping 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. One ``stat`` per profile per tick; the
|
||||
reconcile runs when ``config.yaml``'s (mtime, size) changed, and again on the next tick while a
|
||||
dropped server was still mid-connect (``pending``) and could not be torn down yet. Interactive
|
||||
OAuth is suppressed — this runs on a housekeeping thread nobody is watching."""
|
||||
from hermes_cli.config import get_config_path
|
||||
seen: dict = {}
|
||||
retry: set = set()
|
||||
|
||||
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
|
||||
sig = _sig(get_config_path())
|
||||
prev = seen.get(label)
|
||||
seen[label] = sig
|
||||
if label not in retry and (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()
|
||||
retry.discard(label)
|
||||
if result["pending"]:
|
||||
retry.add(label)
|
||||
if result["removed"] or result["added"]:
|
||||
logger.info("MCP config changed (%s): removed=%s added=%s", label, result["removed"], result["added"])
|
||||
|
||||
def _tick() -> None:
|
||||
from gateway.run import _multiplex_profile_homes, _profile_runtime_scope
|
||||
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
|
||||
|
||||
@@ -84,6 +84,8 @@ def test_multiplex_housekeeping_scopes_primary_and_drains_each_profile(
|
||||
yield
|
||||
|
||||
monkeypatch.setattr(gateway_run, "_profile_runtime_scope", fake_scope)
|
||||
from gateway import run_profile_reconcile
|
||||
monkeypatch.setattr(run_profile_reconcile, "_mcp_config_reconciler", lambda runner: lambda: None)
|
||||
monkeypatch.setattr(
|
||||
scheduler,
|
||||
"drain_delivery_queue",
|
||||
|
||||
@@ -24,7 +24,7 @@ async def test_gateway_startup_discovery_suppresses_interactive_oauth(monkeypatc
|
||||
|
||||
|
||||
def test_mcp_config_reconciler_runs_only_when_config_changes(monkeypatch, tmp_path: Path):
|
||||
import gateway.run as gateway_run
|
||||
from gateway.run_profile_reconcile import _mcp_config_reconciler
|
||||
from tools import mcp_tool_discovery as _mcp_discovery
|
||||
from tools.mcp_oauth import _is_interactive
|
||||
|
||||
@@ -35,10 +35,11 @@ def test_mcp_config_reconciler_runs_only_when_config_changes(monkeypatch, tmp_pa
|
||||
|
||||
def fake_reconcile():
|
||||
calls.append(_is_interactive())
|
||||
return {"removed": ["linear"], "added": []}
|
||||
return {"removed": ["linear"], "added": [], "pending": pending.copy()}
|
||||
|
||||
pending: list = []
|
||||
monkeypatch.setattr(_mcp_discovery, "reconcile_mcp_servers_with_config", fake_reconcile)
|
||||
tick = gateway_run._mcp_config_reconciler(runner=None)
|
||||
tick = _mcp_config_reconciler(runner=None)
|
||||
|
||||
tick() # baseline only: startup discovery already reflects this file
|
||||
tick()
|
||||
@@ -48,3 +49,10 @@ def test_mcp_config_reconciler_runs_only_when_config_changes(monkeypatch, tmp_pa
|
||||
assert calls == [False], "reconcile must run once per change, with interactive OAuth suppressed"
|
||||
tick()
|
||||
assert calls == [False]
|
||||
cfg.write_text("model:\n default: y\n")
|
||||
pending.append("linear") # dropped server was still mid-connect: retry next tick, unchanged file
|
||||
tick()
|
||||
pending.clear()
|
||||
tick()
|
||||
tick()
|
||||
assert calls == [False, False, False], "one retry after a pending teardown, then quiet again"
|
||||
|
||||
@@ -31,7 +31,7 @@ def test_reconcile_tears_down_server_dropped_from_config(monkeypatch, tmp_path):
|
||||
mcp_tool._servers["linear"] = srv
|
||||
mcp_tool._server_scope_keys["linear"] = None
|
||||
try:
|
||||
assert disc.reconcile_mcp_servers_with_config() == {"removed": [], "added": []}
|
||||
assert disc.reconcile_mcp_servers_with_config() == {"removed": [], "added": [], "pending": []}
|
||||
assert "linear" in mcp_tool._servers and not discovered
|
||||
|
||||
configured.clear() # user deletes the entry
|
||||
@@ -49,24 +49,39 @@ def test_reconcile_tears_down_server_dropped_from_config(monkeypatch, tmp_path):
|
||||
_loop._stop_mcp_loop()
|
||||
|
||||
|
||||
def test_disabled_entry_counts_as_dropped(monkeypatch, tmp_path):
|
||||
def test_disabled_lazy_and_connecting_entries(monkeypatch, tmp_path):
|
||||
"""``enabled: false`` counts as dropped; a schema-cache (lazy) registration loses its cached
|
||||
tools; a server still mid-connect is reported ``pending`` (torn down on a later pass)."""
|
||||
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
|
||||
from tools import mcp_tool_registration as _registration
|
||||
|
||||
monkeypatch.setattr(_config, "_load_mcp_config", lambda: {"linear": {"url": "https://x/mcp", "enabled": False}})
|
||||
torn_down: list = []
|
||||
deregistered: list = []
|
||||
monkeypatch.setattr(_lifecycle, "shutdown_mcp_servers", lambda **kw: torn_down.append(kw))
|
||||
monkeypatch.setattr(_registration, "_deregister_mcp_tool_all_scopes", lambda key, name: deregistered.append(name))
|
||||
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
|
||||
mcp_tool._server_scope_keys["notion"] = None
|
||||
mcp_tool._server_connecting.add("notion")
|
||||
mcp_tool._lazy_server_configs["asana"] = {"url": "https://a/mcp"}
|
||||
mcp_tool._lazy_server_tool_names["asana"] = ["mcp__asana__list"]
|
||||
try:
|
||||
assert disc.reconcile_mcp_servers_with_config()["removed"] == ["linear"]
|
||||
result = disc.reconcile_mcp_servers_with_config()
|
||||
assert result == {"removed": ["linear", "asana"], "added": [], "pending": ["notion"]}
|
||||
assert torn_down == [{"scope": None, "names": {"linear"}}]
|
||||
assert deregistered == ["mcp__asana__list"] and "asana" not in mcp_tool._lazy_server_configs
|
||||
finally:
|
||||
with mcp_tool._lock:
|
||||
mcp_tool._servers.pop("linear", None)
|
||||
mcp_tool._server_scope_keys.pop("linear", None)
|
||||
for name in ("linear", "notion"):
|
||||
mcp_tool._servers.pop(name, None)
|
||||
mcp_tool._server_scope_keys.pop(name, None)
|
||||
mcp_tool._server_connecting.discard("notion")
|
||||
mcp_tool._lazy_server_configs.pop("asana", None)
|
||||
mcp_tool._lazy_server_tool_names.pop("asana", None)
|
||||
|
||||
@@ -16,7 +16,7 @@ from tools import mcp_tool_lifecycle as _lifecycle
|
||||
from tools import mcp_tool_loop as _loop
|
||||
from tools import mcp_tool_registration as _registration
|
||||
from tools.mcp_tool_schema import MCP_TOOL_NAME_PREFIX
|
||||
from tools.mcp_tool_scope import _key_name, _resolve_server_key, _server_key
|
||||
from tools.mcp_tool_scope import _key_name, _key_scope, _resolve_server_key, _server_key
|
||||
|
||||
logger = logging.getLogger("tools.mcp_tool")
|
||||
|
||||
@@ -479,18 +479,25 @@ def reconcile_mcp_servers_with_config() -> Dict[str, List[str]]:
|
||||
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."""
|
||||
scope (one multiplexed profile's config prunes only its own connections). A lazily registered
|
||||
(schema-cache) server loses its cached tools; one still mid-connect cannot be torn down yet and
|
||||
is reported under ``"pending"`` so the caller retries. Returns
|
||||
``{"removed": [...], "added": [...], "pending": [...]}``; 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}
|
||||
connecting = {_key_name(key) for key in owned if key in _core._server_connecting}
|
||||
lazy = {key for key in _core._lazy_server_configs
|
||||
if _key_scope(key) == scope and _key_name(key) not in wanted}
|
||||
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))
|
||||
for key in lazy:
|
||||
_forget_lazy_server(key)
|
||||
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)}
|
||||
@@ -498,7 +505,19 @@ def reconcile_mcp_servers_with_config() -> Dict[str, List[str]]:
|
||||
added = sorted(wanted - known)
|
||||
if added:
|
||||
discover_mcp_tools()
|
||||
return {"removed": stale, "added": added}
|
||||
return {"removed": stale + sorted(_key_name(k) for k in lazy), "added": added,
|
||||
"pending": sorted(connecting - wanted)}
|
||||
|
||||
|
||||
def _forget_lazy_server(key) -> None:
|
||||
"""Drop a schema-cache (lazy) registration whose config entry is gone: its cached tools would
|
||||
otherwise stay callable and spawn the server on first use."""
|
||||
with _core._lock:
|
||||
_core._lazy_server_configs.pop(key, None)
|
||||
_core._lazy_server_fingerprints.pop(key, None)
|
||||
cached_names = _core._lazy_server_tool_names.pop(key, None) or []
|
||||
for tool_name in cached_names:
|
||||
_registration._deregister_mcp_tool_all_scopes(key, tool_name)
|
||||
|
||||
|
||||
def is_mcp_tool_parallel_safe(tool_name: str) -> bool:
|
||||
|
||||
Reference in New Issue
Block a user