From da18bcd338154d9cc8bb53f5a191cbbc3e7f2070 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 23:40:12 -0700 Subject: [PATCH] refactor(tools): MCP run/discovery/registration log-call folding, _publish_error helper --- tools/mcp_tool_discovery.py | 10 ++--- tools/mcp_tool_handlers.py | 23 +++++----- tools/mcp_tool_registration.py | 25 +++++------ tools/mcp_tool_server_run.py | 76 ++++++++++++++-------------------- 4 files changed, 57 insertions(+), 77 deletions(-) diff --git a/tools/mcp_tool_discovery.py b/tools/mcp_tool_discovery.py index fd2c9cebfe..763150e1f7 100644 --- a/tools/mcp_tool_discovery.py +++ b/tools/mcp_tool_discovery.py @@ -160,10 +160,8 @@ def _ensure_lazy_server_connected(server_name: str) -> bool: for tool_name in phantom_names: registry.deregister(tool_name, scope=_core._server_registry_scope(server_name)) _core._forget_mcp_tool_server(tool_name) - logger.info( - "MCP server '%s': deregistered %d phantom cached tool(s) not served live (stale " - "schema-cache fingerprint %s): %s", - server_name, len(phantom_names), stale_fingerprint, ", ".join(phantom_names)) + logger.info("MCP server '%s': deregistered %d phantom cached tool(s) not served live (stale schema-cache " + "fingerprint %s): %s", server_name, len(phantom_names), stale_fingerprint, ", ".join(phantom_names)) return server is not None and server.session is not None @@ -304,8 +302,8 @@ def _run_discovery_pass(new_servers: Dict[str, dict]) -> None: with _core._lock: stale = [n for n in new_servers if n in _core._server_connecting] if stale: - logger.warning("MCP discovery %s while %d server(s) were still connecting; " - "clearing stale connecting set: %s", how, len(stale), ", ".join(stale)) + logger.warning("MCP discovery %s while %d server(s) were still connecting; clearing stale " + "connecting set: %s", how, len(stale), ", ".join(stale)) _core._server_connecting.difference_update(stale) for _sn in stale: _core._server_connect_errors.setdefault(_sn, f"Connection attempt {how} during discovery") diff --git a/tools/mcp_tool_handlers.py b/tools/mcp_tool_handlers.py index 1d4b7eed92..670e98a457 100644 --- a/tools/mcp_tool_handlers.py +++ b/tools/mcp_tool_handlers.py @@ -178,11 +178,11 @@ def _handle_session_expired_and_retry(server_name: str, exc: BaseException, retr srv = _lookup_reconnectable_server(server_name, require_loop=True) if srv is None: return None - logger.info("MCP server '%s': %s failed with session-expired error (%s); " - "signalling transport reconnect and retrying once.", server_name, op_description, exc) + logger.info("MCP server '%s': %s failed with session-expired error (%s); signalling transport reconnect " + "and retrying once.", server_name, op_description, exc) if not _core._signal_reconnect_and_wait(server_name, srv, op_description=op_description, timeout=15): - logger.warning("MCP server '%s': reconnect did not ready within 15s after " - "session-expired error; falling through to error response.", server_name) + logger.warning("MCP server '%s': reconnect did not ready within 15s after session-expired error; " + "falling through to error response.", server_name) return None return _retry_once(server_name, retry_call, op_description, "session reconnect") @@ -200,8 +200,8 @@ def _handle_stdio_child_exited_and_retry(server_name: str, exc: Exception, retry reconnected = False srv = _lookup_reconnectable_server(server_name) if srv is not None: - logger.info("MCP server '%s': %s found the stdio subprocess dead (%s); " - "respawning and retrying once.", server_name, op_description, exc) + logger.info("MCP server '%s': %s found the stdio subprocess dead (%s); respawning and retrying once.", + server_name, op_description, exc) if _mcp_loop_running(): reconnected = _core._signal_reconnect_and_wait( server_name, srv, op_description=op_description, timeout=_core._STDIO_RESPAWN_WAIT_SEC) @@ -219,8 +219,8 @@ def _handle_stdio_child_exited_and_retry(server_name: str, exc: Exception, retry return _record_call_outcome(server_name, retry_call()) except _StdioChildExited as retry_exc: # Died again right after respawn: broken server; run()'s budget takes it to the park. - logger.warning("MCP server '%s': %s stdio subprocess exited again right " - "after respawn (%s); not retrying further.", server_name, op_description, retry_exc) + logger.warning("MCP server '%s': %s stdio subprocess exited again right after respawn (%s); not retrying " + "further.", server_name, op_description, retry_exc) return _strike( server_name, f"MCP server '{server_name}' respawned its stdio subprocess and it exited again " @@ -389,8 +389,8 @@ def _make_tool_handler(server_name: str, tool_name: str, tool_timeout: float): op = f"tools/call {tool_name}" def _handler(args: dict, **kwargs) -> str: - # Security boundary: untrusted-server write tools need approval before ANY transport - # work, including the lazy first-use spawn below. + # Security boundary: untrusted-server write tools need approval before ANY transport work + # (including the lazy first-use spawn). error = _trust_gate_check(server_name, tool_name) or _check_circuit_breaker(server_name) if error is not None: return error @@ -401,8 +401,7 @@ def _make_tool_handler(server_name: str, tool_name: str, tool_timeout: float): async def _call(): _mark_server_call_started(server) async with server._rpc_lock, _track_inflight_rpc(server, server_name, op): - # Snapshot contextvars so an elicitation callback (fired on the MCP recv loop, - # which doesn't inherit them) can replay them for gateway platform / session routing. + # Snapshot contextvars for the elicitation callback (MCP recv loop doesn't inherit them). server._pending_call_context = contextvars.copy_context() try: result = await _call_tool_racing_stdio_death(server, server_name, tool_name, args) diff --git a/tools/mcp_tool_registration.py b/tools/mcp_tool_registration.py index 885602a63a..00449ad300 100644 --- a/tools/mcp_tool_registration.py +++ b/tools/mcp_tool_registration.py @@ -34,8 +34,8 @@ def _normalize_server_trust(value: Any) -> str: text = str(value).strip().lower() if text in (_core._TRUST_FULL, _core._TRUST_UNTRUSTED): return text - logger.warning( - "MCP trust: unrecognized trust value %r — treating as 'untrusted' (valid values: full, untrusted)", value) + logger.warning("MCP trust: unrecognized trust value %r — treating as 'untrusted' (valid values: full, untrusted)", + value) return _core._TRUST_UNTRUSTED @@ -226,17 +226,14 @@ def _resolve_name_collisions(name: str, candidates: List[_Candidate]) -> List[_C if len(native_origins) == 1 and utility_origins: shadowed.update((registry_name, o) for o in utility_origins) logger.info( - "MCP server '%s': generated utility %s normalizes onto server-native %s — keeping the " - "native tool and dropping the utility (the utility only applies when the server has no " - "such tool of its own)", + "MCP server '%s': generated utility %s normalizes onto server-native %s — keeping the native tool " + "and dropping the utility (the utility only applies when the server has no such tool of its own)", name, ", ".join(utility_origins), native_origins[0]) continue ambiguous[registry_name] = sorted(origins) for registry_name, origins in sorted(ambiguous.items()): - logger.error( - "MCP server '%s': name normalization collision for '%s' from %s; skipping every colliding " - "entry instead of choosing an arbitrary handler", - name, registry_name, ", ".join(origins)) + logger.error("MCP server '%s': name normalization collision for '%s' from %s; skipping every colliding " + "entry instead of choosing an arbitrary handler", name, registry_name, ", ".join(origins)) return [c for c in unique if c.registry_name not in ambiguous and (c.registry_name, c.origin) not in shadowed] @@ -248,13 +245,11 @@ def _log_foreign_owner(name: str, c: _Candidate, existing_toolset: str, lazy: bo name, c.registry_name, existing_toolset) return if existing_toolset.startswith("mcp-"): - log, fmt = logger.error, ( - "MCP server '%s': %s normalizes to '%s', already owned by MCP toolset '%s' " - "— skipping to preserve the existing owner") + logger.error("MCP server '%s': %s normalizes to '%s', already owned by MCP toolset '%s' — skipping to " + "preserve the existing owner", name, c.origin, c.registry_name, existing_toolset) else: - log, fmt = logger.warning, ( - "MCP server '%s': %s (→ '%s') collides with built-in tool in toolset '%s' — skipping to preserve built-in") - log(fmt, name, c.origin, c.registry_name, existing_toolset) + logger.warning("MCP server '%s': %s (→ '%s') collides with built-in tool in toolset '%s' — skipping to " + "preserve built-in", name, c.origin, c.registry_name, existing_toolset) def _register_candidates(name: str, candidates: List[_Candidate], *, check_fn: Callable, diff --git a/tools/mcp_tool_server_run.py b/tools/mcp_tool_server_run.py index 362ef77e73..6559be2ffe 100644 --- a/tools/mcp_tool_server_run.py +++ b/tools/mcp_tool_server_run.py @@ -81,10 +81,8 @@ class MCPServerRunMixin: await self._keepalive_probe() except Exception as exc: root = _core._unwrap_exception_group(exc) - logger.warning( - "MCP server '%s' keepalive failed, triggering " - "reconnect (state: connected → degraded): %s: %s", - self.name, type(root).__name__, root) + logger.warning("MCP server '%s' keepalive failed, triggering reconnect (state: connected → " + "degraded): %s: %s", self.name, type(root).__name__, root) self.mark_suspect(f"keepalive failed: {type(root).__name__}: {root}") self._reconnect_event.set() break @@ -124,8 +122,8 @@ class MCPServerRunMixin: self._reconnect_event.clear() if await self._wait_for_reconnect_or_shutdown(timeout=_core._PARKED_RETRY_INTERVAL) == "shutdown": return True - logger.debug("MCP server '%s': attempting revival %s (self-probe or explicit " - "reconnect request); rebuilding transport.", self.name, revival_reason) + logger.debug("MCP server '%s': attempting revival %s (self-probe or explicit reconnect request); " + "rebuilding transport.", self.name, revival_reason) return False async def _prepare_run(self, config: dict) -> bool: @@ -148,9 +146,8 @@ class MCPServerRunMixin: self._elicitation = (_core.ElicitationHandler(self.name, elicitation_config, owner=self) if elicitation_config.get("enabled", True) and _core._MCP_ELICITATION_TYPES else None) if "url" in config and "command" in config: - logger.warning("MCP server '%s' has both 'url' and 'command' in config. " - "Using HTTP transport ('url'). Remove 'command' to silence " - "this warning.", self.name) + logger.warning("MCP server '%s' has both 'url' and 'command' in config. Using HTTP transport " + "('url'). Remove 'command' to silence this warning.", self.name) if not self._is_http(): return True try: @@ -165,13 +162,16 @@ class MCPServerRunMixin: ssl_verify=config.get("ssl_verify", True), client_cert=_core._resolve_client_cert(self.name, config)) except (_core.InvalidMcpUrlError, _core.NonMcpEndpointError) as exc: - # Fail fast and non-retryably: publish the error to start(). logger.warning("%s", exc) - self._error = exc - self._ready.set() + self._publish_error(exc) # fail fast and non-retryably return False return True + def _publish_error(self, exc: BaseException) -> None: + """Hand *exc* to the waiting ``start()``.""" + self._error = exc + self._ready.set() + async def run(self, config: dict): """Long-lived: connecting -> connected -> (degraded -> parked -> revived)*. Unproven drops and transport errors charge a rapid-drop budget with jittered backoff; exhausting it (or @@ -205,8 +205,8 @@ class MCPServerRunMixin: if self._shutdown_event.is_set(): return False if lifecycle_reason == "recycle": - logger.info("MCP server '%s': stdio session recycled after %s; " - "waiting for lazy reconnect", self.name, self._recycled_reason) + logger.info("MCP server '%s': stdio session recycled after %s; waiting for lazy reconnect", + self.name, self._recycled_reason) self.session = None # Dormant until a lazy call wakes it (untimed: nothing to self-probe). return await self._wait_for_reconnect_or_shutdown() != "shutdown" @@ -215,9 +215,8 @@ class MCPServerRunMixin: # A clean return is NOT proof of health (a flapper handshakes fine, then drops). Only a # PROVEN session clears the budget; a teardown race is recovery, never a park charge. if self._teardown_race and not self._session_proven: - logger.info("MCP server '%s': reconnect after teardown race " - "(in-flight calls were failed); not charging the " - "rapid-drop budget", self.name) + logger.info("MCP server '%s': reconnect after teardown race (in-flight calls were failed); " + "not charging the rapid-drop budget", self.name) self._teardown_race, budget.backoff = False, 1.0 elif self._session_proven: self._reconnect_retries, budget.backoff = 0, 1.0 @@ -225,10 +224,8 @@ class MCPServerRunMixin: self._reconnect_retries += 1 if self._reconnect_retries > _core._MAX_RECONNECT_RETRIES: logger.warning( - "MCP server '%s': %d consecutive reconnects " - "without a healthy session (rapid-drop budget " - "exhausted), parking; will self-probe every %ds " - "until it recovers (state: degraded → parked)", + "MCP server '%s': %d consecutive reconnects without a healthy session (rapid-drop budget " + "exhausted), parking; will self-probe every %ds until it recovers (state: degraded → parked)", self.name, _core._MAX_RECONNECT_RETRIES, _core._PARKED_RETRY_INTERVAL) if not await self._park_and_rearm("from parked state", budget): return False @@ -247,8 +244,7 @@ class MCPServerRunMixin: async def _park_initial_failure(self, exc: Exception, revival_reason: str, budget: "_RetryBudget") -> bool: """Publish ``exc`` to ``start()``, park, and on revival reset every counter. False on shutdown.""" - self._error = exc - self._ready.set() + self._publish_error(exc) if await self._park(revival_reason): return False budget.initial_retries = self._reconnect_retries = 0 @@ -267,9 +263,8 @@ class MCPServerRunMixin: root = _core._unwrap_exception_group(exc) failure_class = _core._classify_mcp_failure(root) if self._is_recycled_stdio(): - logger.warning("MCP server '%s': lazy reconnect after stdio recycle " - "failed, marking unavailable while retrying: %s: %s", - self.name, type(root).__name__, root) + logger.warning("MCP server '%s': lazy reconnect after stdio recycle failed, marking unavailable " + "while retrying: %s: %s", self.name, type(root).__name__, root) self._recycled_reason = None # Initial-connect ladder (a startup blip must not kill the server); gated on # _ever_connected, not _ready (which clears every reconnect cycle). @@ -284,16 +279,13 @@ class MCPServerRunMixin: self._reconnect_retries += 1 if self._reconnect_retries > _core._MAX_RECONNECT_RETRIES: logger.warning( - "MCP server '%s' failed after %d reconnection attempts, " - "parking; will self-probe every %ds until it recovers " - "(state: degraded → parked): %s: %s", - self.name, _core._MAX_RECONNECT_RETRIES, _core._PARKED_RETRY_INTERVAL, - type(root).__name__, root) + "MCP server '%s' failed after %d reconnection attempts, parking; will self-probe every %ds " + "until it recovers (state: degraded → parked): %s: %s", + self.name, _core._MAX_RECONNECT_RETRIES, _core._PARKED_RETRY_INTERVAL, type(root).__name__, root) return await self._park_and_rearm("from parked state", budget) - logger.debug("MCP server '%s' connection lost (attempt %d/%d), " - "reconnecting in %.0fs: %s: %s", - self.name, self._reconnect_retries, _core._MAX_RECONNECT_RETRIES, - budget.backoff, type(root).__name__, root) + logger.debug("MCP server '%s' connection lost (attempt %d/%d), reconnecting in %.0fs: %s: %s", + self.name, self._reconnect_retries, _core._MAX_RECONNECT_RETRIES, budget.backoff, + type(root).__name__, root) await self._backoff_sleep(budget) return not self._shutdown_event.is_set() @@ -311,8 +303,7 @@ class MCPServerRunMixin: budget.initial_retries += 1 if budget.initial_retries > _core._MAX_INITIAL_CONNECT_RETRIES: logger.warning( - "MCP server '%s' failed initial connection after " - "%d attempts, parking until a reconnect is " + "MCP server '%s' failed initial connection after %d attempts, parking until a reconnect is " "requested (state: connecting → parked): %s: %s", self.name, _core._MAX_INITIAL_CONNECT_RETRIES, type(root).__name__, root) return await self._park_initial_failure(exc, "after initial connection failures", budget) @@ -322,8 +313,7 @@ class MCPServerRunMixin: type(root).__name__, root) await self._backoff_sleep(budget) if self._shutdown_event.is_set(): - self._error = exc - self._ready.set() + self._publish_error(exc) return not self._shutdown_event.is_set() async def _on_permanent_error(self, root: BaseException, budget: "_RetryBudget") -> bool: @@ -333,8 +323,7 @@ class MCPServerRunMixin: self._permanent_grace_used = True self.mark_suspect(f"auth error on proven session: {root}") logger.warning( - "MCP server '%s': auth error on a previously " - "healthy session — marking suspect and forcing " + "MCP server '%s': auth error on a previously healthy session — marking suspect and forcing " "one reconnect instead of parking (state: connected → suspect): %s: %s", self.name, type(root).__name__, root) self._reconnect_retries, budget.backoff = 0, 1.0 @@ -342,9 +331,8 @@ class MCPServerRunMixin: return not self._shutdown_event.is_set() # Deterministic failure on a working server: park now. logger.warning( - "MCP server '%s' hit a permanent error, parking " - "without retries; will self-probe every %ds (state: connected → parked): %s: %s", - self.name, _core._PARKED_RETRY_INTERVAL, type(root).__name__, root) + "MCP server '%s' hit a permanent error, parking without retries; will self-probe every %ds " + "(state: connected → parked): %s: %s", self.name, _core._PARKED_RETRY_INTERVAL, type(root).__name__, root) return await self._park_and_rearm("from parked state (permanent error)", budget) async def start(self, config: dict):