diff --git a/gateway/run.py b/gateway/run.py index 2916822b4c..1c7a68c8f5 100644 --- a/gateway/run.py +++ b/gateway/run.py @@ -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 diff --git a/gateway/run_profile_reconcile.py b/gateway/run_profile_reconcile.py index 1b0ed3b69b..f433c126dc 100644 --- a/gateway/run_profile_reconcile.py +++ b/gateway/run_profile_reconcile.py @@ -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 diff --git a/tests/gateway/test_cron_delivery_housekeeping.py b/tests/gateway/test_cron_delivery_housekeeping.py index ce1dede5bf..c1c21c4966 100644 --- a/tests/gateway/test_cron_delivery_housekeeping.py +++ b/tests/gateway/test_cron_delivery_housekeeping.py @@ -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", diff --git a/tests/gateway/test_gateway_mcp_oauth_unattended.py b/tests/gateway/test_gateway_mcp_oauth_unattended.py index ea9c069d9d..52a69d34f0 100644 --- a/tests/gateway/test_gateway_mcp_oauth_unattended.py +++ b/tests/gateway/test_gateway_mcp_oauth_unattended.py @@ -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" diff --git a/tests/tools/test_mcp_reconcile_with_config.py b/tests/tools/test_mcp_reconcile_with_config.py index 755c40f8af..50b5bf334b 100644 --- a/tests/tools/test_mcp_reconcile_with_config.py +++ b/tests/tools/test_mcp_reconcile_with_config.py @@ -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) diff --git a/tools/mcp_tool_discovery.py b/tools/mcp_tool_discovery.py index f0d31f6a65..d424b9b366 100644 --- a/tools/mcp_tool_discovery.py +++ b/tools/mcp_tool_discovery.py @@ -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: