refactor(tools): MCP run/discovery/registration log-call folding, _publish_error helper
This commit is contained in:
@@ -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")
|
||||
|
||||
@@ -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)
|
||||
|
||||
@@ -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,
|
||||
|
||||
@@ -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):
|
||||
|
||||
Reference in New Issue
Block a user