fix(mcp): cap discovery pass timeout and derive lock waiter budget from it
A bounded discovery pass runs ceil(N/cap) 120s waves, so a fleet of slow servers now legitimately occupies the cross-process discovery lock past the 120s the unbounded scatter used to need. The lock waiter (240 x 0.5s) would give up mid-pass, fail over, and discover unguarded beside the still-connecting holder - respawning the duplicate stdio trees this PR removes (#117373 review). Cap the pass at _MCP_DISCOVERY_PASS_MAX_SEC = 300 so one stuck connect cannot pin the calling thread (and the lock) for the tens of minutes 120*waves allows, and derive _MCP_DISCOVERY_LOCK_MAX_RETRIES from the same ceiling (320s waiter budget > 300s pass) so the waiter always outlasts a legitimate holder.
This commit is contained in:
@@ -2837,6 +2837,39 @@ class TestDiscoveryConnectConcurrency:
|
||||
for name in server_names:
|
||||
_servers.pop(name, None)
|
||||
|
||||
def test_pass_timeout_capped_by_waiter_budget(self):
|
||||
"""The discovery pass timeout is capped so a slow pass cannot outlive the
|
||||
cross-process lock waiter's fail-over budget: with the connect cap,
|
||||
ceil(N/cap) waves make a pass legitimately long, and an uncapped
|
||||
120s-per-wave timeout would let a lock loser run unguarded discovery
|
||||
beside a still-connecting holder (#117373 review)."""
|
||||
from tools import mcp_tool_discovery as _discovery
|
||||
from tools.mcp_tool import _MCP_DISCOVERY_LOCK_MAX_RETRIES, _MCP_DISCOVERY_LOCK_RETRY_DELAY_S
|
||||
from tools.mcp_tool import _MCP_DISCOVERY_PASS_MAX_SEC
|
||||
|
||||
waiter_budget = _MCP_DISCOVERY_LOCK_MAX_RETRIES * _MCP_DISCOVERY_LOCK_RETRY_DELAY_S
|
||||
assert waiter_budget > _MCP_DISCOVERY_PASS_MAX_SEC, (
|
||||
f"waiter budget {waiter_budget}s must outlast the pass ceiling "
|
||||
f"{_MCP_DISCOVERY_PASS_MAX_SEC}s or a slow pass re-opens unguarded discovery")
|
||||
|
||||
captured = {}
|
||||
|
||||
def fake_run_on_mcp_loop(factory, timeout=None):
|
||||
captured["timeout"] = timeout
|
||||
return None
|
||||
|
||||
# 40 servers = 14 waves at cap 3: uncapped, the pass would block 28 min
|
||||
# and outlive the waiter budget by 26+ minutes.
|
||||
server_names = {f"srv{i}": {} for i in range(40)}
|
||||
with patch("tools.mcp_tool_discovery._loop._run_on_mcp_loop",
|
||||
side_effect=fake_run_on_mcp_loop):
|
||||
_discovery._run_discovery_pass(server_names)
|
||||
|
||||
assert captured["timeout"] == _MCP_DISCOVERY_PASS_MAX_SEC, (
|
||||
f"pass timeout must be capped at {_MCP_DISCOVERY_PASS_MAX_SEC}s, "
|
||||
f"got {captured['timeout']}s (uncapped: {120 * 14}s vs waiter "
|
||||
f"budget {waiter_budget}s)")
|
||||
|
||||
|
||||
class TestMCPSelectiveToolLoading:
|
||||
"""Tests for per-server MCP filtering and utility tool policies."""
|
||||
|
||||
@@ -681,8 +681,16 @@ def _server_visible_in_scope(key, scope: Optional[str]) -> bool:
|
||||
# See issue #62771.
|
||||
_LOCK_UNAVAILABLE: Any = object() # sentinel: locking broken/unavailable
|
||||
_MCP_DISCOVERY_LOCK_PATH: Optional[str] = None # resolved lazily
|
||||
# Bounded wait when another process holds the lock.
|
||||
_MCP_DISCOVERY_LOCK_MAX_RETRIES, _MCP_DISCOVERY_LOCK_RETRY_DELAY_S = 240, 0.5
|
||||
# A discovery pass (bounded gathers, one 120 s budget per wave) may legitimately
|
||||
# run past the 120 s a scatter completes in. The waiter must outlast the worst
|
||||
# legitimate holder (pass ceiling + slack), or it fails over at 120 s, discovers
|
||||
# unguarded beside a still-connecting holder, and spawns duplicate stdio trees.
|
||||
# See #117373: the concurrency cap made the old 120 s waiter budget stale.
|
||||
_MCP_DISCOVERY_PASS_MAX_SEC = 300 # overall pass ceiling
|
||||
_MCP_DISCOVERY_LOCK_RETRY_DELAY_S = 0.5
|
||||
# Waiter budget (max_retries * delay) must outlast the pass ceiling: 320 s > 300 s.
|
||||
_MCP_DISCOVERY_LOCK_MAX_RETRIES = int(
|
||||
_MCP_DISCOVERY_PASS_MAX_SEC / _MCP_DISCOVERY_LOCK_RETRY_DELAY_S) + 20
|
||||
|
||||
|
||||
# ---- BEGIN PLUGIN-COMPAT (revert-scheduled; see COMPAT_MANIFEST.md) ----
|
||||
|
||||
@@ -414,9 +414,13 @@ def _run_discovery_pass(new_servers: Dict[str, dict]) -> None:
|
||||
# Budget scales with the concurrency cap: a bounded gather finishes in
|
||||
# ceil(N/cap) waves, so N > cap multiplies the wall clock the base
|
||||
# (unbounded) gather never needed. 120s per wave keeps the original
|
||||
# per-wave ceiling; a slow fleet aborts later, not never.
|
||||
# per-wave ceiling; a slow fleet aborts later, not never. Capped by
|
||||
# _MCP_DISCOVERY_PASS_MAX_SEC so one stuck connect cannot pin the
|
||||
# calling thread (and the cross-process discovery lock) for tens of
|
||||
# minutes; the lock waiter's budget is derived from the same ceiling.
|
||||
waves = max(1, -(-len(new_servers) // _DISCOVERY_CONNECT_CONCURRENCY))
|
||||
_loop._run_on_mcp_loop(lambda: _discover_all(new_servers), timeout=120 * waves)
|
||||
timeout = min(120 * waves, _core._MCP_DISCOVERY_PASS_MAX_SEC)
|
||||
_loop._run_on_mcp_loop(lambda: _discover_all(new_servers), timeout=timeout)
|
||||
except (TimeoutError, InterruptedError) as _e:
|
||||
# Stranded _server_connecting entries would block future reconnects.
|
||||
how = "timed out" if isinstance(_e, TimeoutError) else "interrupted"
|
||||
|
||||
Reference in New Issue
Block a user