fix(cron): a server parked on a permanent error blocks the job again; warn once per outage
The reconnecting exemption keyed only on `_ever_connected`, so a server that
connected once and then parked on a PERMANENT error (revoked credentials,
endpoint gone: `_park_and_rearm('from parked state (permanent error)')`) looked
identical to a router-reboot park. Its self-probe fails the same way every
interval, so the job ran tool-less on every tick forever with only a gateway
WARNING and never received the one-shot blocked_config alert it used to get.
Record the revival reason `_park` was given on MCPServerTask (`_park_reason`,
cleared when a session proves healthy) and have `mcp_server_reconnecting`
return False for a permanent-error park, restoring the block and its alert.
Transient parks (network blip, rapid-drop budget exhausted) still return True.
Also dedupe the exemption's WARNING per job+server for the length of one
outage, like the one-shot blocked_config alert, instead of logging every tick;
the entry drops once the server resolves tools again so the next outage warns.
This commit is contained in:
@@ -312,6 +312,10 @@ def _preflight_check_skills(job: dict) -> Optional[str]:
|
||||
return None
|
||||
|
||||
|
||||
# (job id, server name) pairs already warned about as reconnecting; see _empty_requested_mcp_toolsets.
|
||||
_RECONNECTING_WARNED: set = set()
|
||||
|
||||
|
||||
def _empty_requested_mcp_toolsets(job: dict, cfg: dict) -> Optional[str]:
|
||||
"""Reason when an MCP server the job's own ``enabled_toolsets`` names resolves to zero tools.
|
||||
|
||||
@@ -333,10 +337,19 @@ def _empty_requested_mcp_toolsets(job: dict, cfg: dict) -> Optional[str]:
|
||||
# that did resolve rather than losing a whole tick to a minute of downtime (#112871). Only a
|
||||
# server that never connected for this profile is judged below.
|
||||
reconnecting = sorted(name for name in missing if mcp_server_reconnecting(name))
|
||||
if reconnecting:
|
||||
job_id = str(job.get("id", "?"))
|
||||
# One WARNING per job+server per outage (like the one-shot blocked_config alert), not one per
|
||||
# tick; the entry drops once the server is back so the next outage warns again.
|
||||
_RECONNECTING_WARNED.difference_update(
|
||||
key for key in list(_RECONNECTING_WARNED) if key[0] == job_id and key[1] not in reconnecting)
|
||||
unwarned = [name for name in reconnecting if (job_id, name) not in _RECONNECTING_WARNED]
|
||||
if unwarned:
|
||||
_RECONNECTING_WARNED.update((job_id, name) for name in unwarned)
|
||||
logger.warning(
|
||||
"Job '%s': MCP server(s) %s named in enabled_toolsets are reconnecting — running "
|
||||
"without their tools this tick", job.get("id", "?"), ", ".join(reconnecting))
|
||||
"without their tools until they recover (a server parked on a permanent error blocks "
|
||||
"the job instead)", job_id, ", ".join(unwarned))
|
||||
if reconnecting:
|
||||
missing = [name for name in missing if name not in reconnecting]
|
||||
if not missing:
|
||||
return None
|
||||
|
||||
@@ -92,15 +92,17 @@ def test_requested_mcp_server_with_tools_runs(tmp_path):
|
||||
assert success is True and error is None
|
||||
|
||||
|
||||
def _park_notion(*, ever_connected: bool):
|
||||
def _park_notion(*, ever_connected: bool, park_reason=None):
|
||||
"""Install a sessionless ``notion`` run task (tools deregistered, alias still global) the way
|
||||
the MCP layer leaves a degraded/parked server; ``ever_connected`` separates a server that
|
||||
worked in this process and lost its network from one that never came up here."""
|
||||
worked in this process and lost its network from one that never came up here, and
|
||||
``park_reason`` is what ``_park`` recorded (permanent-error parks are not recovering)."""
|
||||
import tools.mcp_tool as core
|
||||
from tools.registry import registry
|
||||
|
||||
server = core.MCPServerTask("notion")
|
||||
server._ever_connected = ever_connected
|
||||
server._park_reason = park_reason
|
||||
registry.register_toolset_alias("notion", "mcp-notion")
|
||||
core._servers["notion"] = server
|
||||
return lambda: core._servers.pop("notion", None)
|
||||
@@ -132,3 +134,50 @@ def test_requested_mcp_server_never_connected_still_blocks(tmp_path):
|
||||
assert agent_built is False
|
||||
assert success is False
|
||||
assert error is not None and "[blocked_config]" in error and "notion" in error
|
||||
|
||||
|
||||
def test_requested_mcp_server_parked_on_permanent_error_blocks(tmp_path):
|
||||
"""A server that connected once and then parked on a PERMANENT error (revoked credentials,
|
||||
endpoint gone) is not recovering: its self-probe fails identically every time, so the job must
|
||||
take the one-shot blocked_config path instead of silently running tool-less forever."""
|
||||
undo = _park_notion(ever_connected=True, park_reason="from parked state (permanent error)")
|
||||
try:
|
||||
(success, _output, _final, error), agent_built = _run(
|
||||
_job(enabled_toolsets=["terminal", "notion"]), tmp_path)
|
||||
finally:
|
||||
undo()
|
||||
|
||||
assert agent_built is False
|
||||
assert success is False
|
||||
assert error is not None and "[blocked_config]" in error and "notion" in error
|
||||
|
||||
|
||||
def test_reconnecting_warning_is_once_per_job_and_server_per_outage(tmp_path, caplog):
|
||||
"""The reconnecting exemption warns once per job+server while the outage lasts (not every
|
||||
tick) and warns again once the server recovered and dropped out a second time."""
|
||||
import logging
|
||||
|
||||
import cron.scheduler_preflight as preflight
|
||||
|
||||
preflight._RECONNECTING_WARNED.clear()
|
||||
job = _job(enabled_toolsets=["terminal", "notion"])
|
||||
undo = _park_notion(ever_connected=True)
|
||||
try:
|
||||
with caplog.at_level(logging.WARNING, logger="cron.scheduler_preflight"):
|
||||
_run(job, tmp_path)
|
||||
_run(job, tmp_path)
|
||||
warned = [r for r in caplog.records if "are reconnecting" in r.getMessage()]
|
||||
assert len(warned) == 1
|
||||
# Server back: the dedupe entry drops, so the next outage warns again.
|
||||
_register = _register_notion_in_scope(None)
|
||||
try:
|
||||
_run(job, tmp_path)
|
||||
finally:
|
||||
_register()
|
||||
assert ("mcpjob", "notion") not in preflight._RECONNECTING_WARNED
|
||||
_run(job, tmp_path)
|
||||
warned = [r for r in caplog.records if "are reconnecting" in r.getMessage()]
|
||||
assert len(warned) == 2
|
||||
finally:
|
||||
undo()
|
||||
preflight._RECONNECTING_WARNED.clear()
|
||||
|
||||
47
tests/tools/test_mcp_server_reconnecting_predicate.py
Normal file
47
tests/tools/test_mcp_server_reconnecting_predicate.py
Normal file
@@ -0,0 +1,47 @@
|
||||
"""``mcp_server_reconnecting`` (cron preflight's "is this outage healing on its own?" question)
|
||||
must say True only for a transient park. A server parked on a PERMANENT error (revoked
|
||||
credentials, endpoint gone) fails its self-probe identically every time, so treating it as
|
||||
reconnecting would run the job tool-less forever with no blocked_config alert.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
|
||||
import pytest
|
||||
|
||||
import tools.mcp_tool as core
|
||||
from tools.mcp_tool_discovery import mcp_server_reconnecting
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def parked_server(monkeypatch):
|
||||
server = core.MCPServerTask("x")
|
||||
server._ever_connected = True
|
||||
server.session = None
|
||||
server._task = None
|
||||
monkeypatch.setitem(core._servers, "x", server)
|
||||
return server
|
||||
|
||||
|
||||
def test_transient_park_is_reconnecting(parked_server):
|
||||
assert mcp_server_reconnecting("x") is True
|
||||
parked_server._park_reason = "from parked state" # rapid-drop budget exhausted, still probing
|
||||
assert mcp_server_reconnecting("x") is True
|
||||
|
||||
|
||||
def test_permanent_error_park_is_not_reconnecting(parked_server):
|
||||
parked_server._park_reason = "from parked state (permanent error)"
|
||||
assert mcp_server_reconnecting("x") is False
|
||||
|
||||
|
||||
def test_park_records_reason_and_proven_session_clears_it(parked_server, monkeypatch):
|
||||
"""The run loop hands ``_park`` its revival reason; a proven healthy session forgets it."""
|
||||
async def _shutdown_at_once(timeout=None):
|
||||
return "shutdown"
|
||||
monkeypatch.setattr(parked_server, "_wait_for_reconnect_or_shutdown", _shutdown_at_once)
|
||||
|
||||
asyncio.run(parked_server._park("from parked state (permanent error)"))
|
||||
assert parked_server._park_reason == "from parked state (permanent error)"
|
||||
assert mcp_server_reconnecting("x") is False
|
||||
|
||||
parked_server._mark_session_proven()
|
||||
assert parked_server._park_reason is None
|
||||
@@ -317,7 +317,7 @@ class MCPServerTask(MCPServerRunMixin, MCPServerTransportMixin, MCPServerHealthM
|
||||
"_recycled_reason", "initialize_result", "_ping_unsupported", "_list_cache_meta",
|
||||
"_reconnect_retries", "_session_proven", "_was_parked", "_inflight_tasks", "_reconnecting",
|
||||
"_suspect_reason", "_teardown_race", "_permanent_grace_used", "_stdio_child_pids",
|
||||
"_ever_connected", "_sse_fallback")
|
||||
"_ever_connected", "_sse_fallback", "_park_reason")
|
||||
|
||||
def __init__(self, name: str):
|
||||
self.name = name
|
||||
@@ -350,6 +350,10 @@ class MCPServerTask(MCPServerRunMixin, MCPServerTransportMixin, MCPServerHealthM
|
||||
self._sse_fallback: bool = False
|
||||
# True from park until proven healthy again; logs the revival once.
|
||||
self._was_parked: bool = False
|
||||
# Why the server is parked (the revival_reason handed to _park), None once healthy again.
|
||||
# Lets cron preflight tell a network-blip park (recovering) from a permanent-error park
|
||||
# (revoked credentials, dead endpoint) that must not run the job tool-less forever.
|
||||
self._park_reason: Optional[str] = None
|
||||
# In-flight RPC tasks so a deliberate teardown fails them fast; _reconnecting is True
|
||||
# during that teardown so _track_inflight_rpc turns the cancel into a retryable error.
|
||||
# In-flight RPC bookkeeping (#48069 salvage): user-visible requests registered while running so a
|
||||
|
||||
@@ -613,13 +613,18 @@ def get_mcp_status(configured: Optional[Dict[str, dict]] = None, *, include_runt
|
||||
|
||||
def mcp_server_reconnecting(name: str) -> bool:
|
||||
"""True when this profile's connection to *name* connected once in this process and is now
|
||||
between sessions (degraded/parked): the run task is alive and self-probing, so the outage is
|
||||
environmental and heals on its own. A server that never connected here (wrong URL, other
|
||||
profile's credentials, permanent error) is not reconnecting. Reads cached state; never connects."""
|
||||
between sessions (degraded/parked) after a transient failure: the run task is alive and
|
||||
self-probing, so the outage is environmental and heals on its own. A server that never
|
||||
connected here (wrong URL, other profile's credentials) is not reconnecting, and neither is one
|
||||
parked on a PERMANENT error (revoked credentials, endpoint gone): its self-probe fails the same
|
||||
way every time, so callers must treat it as blocked rather than wait forever. Reads cached
|
||||
state; never connects."""
|
||||
with _core._lock:
|
||||
server = _core._servers.get(_resolve_server_key(name))
|
||||
if server is None or server.session is not None or not server._ever_connected:
|
||||
return False
|
||||
if server._park_reason and "permanent" in server._park_reason:
|
||||
return False
|
||||
return server._task is None or not server._task.done()
|
||||
|
||||
|
||||
|
||||
@@ -202,6 +202,7 @@ class MCPServerHealthMixin:
|
||||
return
|
||||
self._session_proven = True
|
||||
self._reconnect_retries = 0
|
||||
self._park_reason = None
|
||||
if self._was_parked:
|
||||
self._was_parked = False
|
||||
logger.warning("MCP server '%s': revived — session healthy again after "
|
||||
|
||||
@@ -175,6 +175,7 @@ class MCPServerRunMixin:
|
||||
# (#57129). An explicit _reconnect_event.set() (OAuth recovery, manual /mcp refresh) still wakes us
|
||||
# immediately.
|
||||
self._was_parked = True
|
||||
self._park_reason = revival_reason
|
||||
self._deregister_tools()
|
||||
self._reconnect_event.clear()
|
||||
outcome = await self._wait_for_reconnect_or_shutdown(timeout=_core._PARKED_RETRY_INTERVAL)
|
||||
|
||||
Reference in New Issue
Block a user