Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions helpers/mcp_handler.py
Original file line number Diff line number Diff line change
Expand Up @@ -1276,6 +1276,20 @@ async def call_tool(
T = TypeVar("T")


def _stop_worker_loop(worker: DeferredTask) -> None:
"""Stop a timed-out worker's event loop without waiting on it.

Abandoning the worker leaves its loop spinning at 100% CPU holding the GIL
for the lifetime of the process. Stopping the loop ends the spin and lets
the thread exit. This deliberately does not drain tasks or join the thread:
both are unbounded on a wedged loop, which is precisely why the timeout
path could not terminate the worker before.
"""
loop = getattr(worker.event_loop_thread, "loop", None)
if loop is not None and loop.is_running():
loop.call_soon_threadsafe(loop.stop)


class MCPClientBase(ABC):
# server: Union[MCPServerLocal, MCPServerRemote] # Defined in __init__
# tools: List[dict[str, Any]] # Defined in __init__
Expand Down Expand Up @@ -1330,6 +1344,7 @@ async def _run_isolated_operation(
finally:
if timed_out:
worker.kill(terminate_thread=False)
_stop_worker_loop(worker)
else:
worker.kill(terminate_thread=True)

Expand Down
54 changes: 54 additions & 0 deletions tests/test_mcp_worker_termination.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,54 @@
"""A timed-out MCP worker must not be left spinning.

Before this fix the timeout path abandoned the worker. A worker whose event
loop keeps cycling then spins at 100% CPU holding the GIL for the lifetime of
the process; observed in the wild at 11h55m of CPU on a single leaked thread.
"""
import asyncio
import threading
import time

from helpers.defer import DeferredTask
from helpers.mcp_handler import _stop_worker_loop


async def _spinner():
# Keeps the loop RUNNING (and burning CPU) rather than blocking it --
# this is the state a wedged MCP worker is actually left in.
while True:
await asyncio.sleep(0)


def test_stop_worker_loop_ends_spin_without_blocking_caller():
task = DeferredTask(thread_name=f"test-mcp-spin-{threading.get_ident()}")
task.start_task(_spinner)
try:
deadline = time.monotonic() + 5.0
while time.monotonic() < deadline:
loop = task.event_loop_thread.loop
if loop is not None and loop.is_running():
break
time.sleep(0.05)
thread = task.event_loop_thread.thread
assert thread is not None and thread.is_alive()

started = time.monotonic()
task.kill(terminate_thread=False)
_stop_worker_loop(task)
elapsed = time.monotonic() - started

# Must not wait on the wedged loop: no drain, no join.
assert elapsed < 1.0, f"caller blocked for {elapsed:.2f}s"

thread.join(timeout=10)
assert not thread.is_alive(), "worker thread survived; it would spin forever"
finally:
loop = getattr(task.event_loop_thread, "loop", None)
if loop is not None and loop.is_running():
loop.call_soon_threadsafe(loop.stop)


def test_stop_worker_loop_is_safe_when_loop_already_gone():
task = DeferredTask(thread_name=f"test-mcp-gone-{threading.get_ident()}")
task.kill(terminate_thread=True)
_stop_worker_loop(task) # must not raise