Skip to content
Merged
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
52 changes: 49 additions & 3 deletions tests/tools/test_mcp_tool.py
Original file line number Diff line number Diff line change
Expand Up @@ -466,10 +466,56 @@ async def _sample():
]
assert runtime_warnings == []

def test_activity_callback_fires_during_long_call(self):
"""A blocking MCP call ticks the activity heartbeat while it waits.

Regression for the NetSuite-hive stuck-kill class (#cc4a1b4f): a single
synchronous ns_runreport call blocked _run_on_mcp_loop for >900s with
zero activity touches, starving the kanban worker heartbeat so the
dispatcher killed the worker as stale. The poll loop must fire the
thread-local activity callback on each iteration so liveness is
reported throughout the wait.
"""
import tools.mcp_tool as mcp

# Real background event loop, same shape as the production _mcp_loop.
loop = asyncio.new_event_loop()
t = threading.Thread(target=loop.run_forever, daemon=True)
t.start()

# Coroutine that blocks long enough for several poll iterations
# (the loop polls every 0.1s).
async def _slow():
await asyncio.sleep(0.35)
return "done"

recorded = []

def _recording_touch(state, label):
recorded.append((dict(state), label))

try:
with patch.object(mcp, "_mcp_loop", loop):
with patch(
"tools.environments.base.touch_activity_if_due",
side_effect=_recording_touch,
):
result = mcp._run_on_mcp_loop(_slow, timeout=10)
finally:
loop.call_soon_threadsafe(loop.stop)
t.join(timeout=2)
loop.close()

assert result == "done"
# The callback must have been invoked at least once while waiting.
assert recorded, "activity callback never fired during the MCP wait"
# State carries the heartbeat cadence and start bookkeeping.
first_state, first_label = recorded[0]
assert first_label == "waiting for MCP tool call"
assert first_state["interval"] == 30.0
assert "start" in first_state and "last_touch" in first_state


# ---------------------------------------------------------------------------
# Tool handler
# ---------------------------------------------------------------------------

class TestToolHandler:
"""Tool handlers are sync functions that schedule work on the MCP loop."""
Expand Down
29 changes: 29 additions & 0 deletions tools/mcp_tool.py
Original file line number Diff line number Diff line change
Expand Up @@ -2465,10 +2465,27 @@ def _run_on_mcp_loop(coro_or_factory, timeout: float = 30):

Poll in short intervals so the calling agent thread can honor user
interrupts while the MCP work is still running on the background loop.

While polling, fire the thread-local activity callback (set by the agent
before dispatching the tool — see ``agent/tool_executor.py``) at a steady
cadence. A single synchronous MCP/hive tool call (e.g. ``ns_runreport``
against the NetSuite hive) can block this loop for many minutes; without
the periodic touch, the agent's last-activity timestamp goes stale, the
gateway inactivity monitor (or the kanban dispatcher's stuck-worker
watchdog, which is bridged through ``_touch_activity``) sees no liveness,
and an actively-running worker gets killed mid-call. The terminal tool's
``_wait_for_process`` already heartbeats this way; this mirrors it for the
MCP transport so long hive calls are survivable (#cc4a1b4f).
"""
from tools.interrupt import is_interrupted
from agent.async_utils import safe_schedule_threadsafe

try:
from tools.environments.base import touch_activity_if_due
except Exception: # pragma: no cover - defensive import guard
def touch_activity_if_due(state: dict, label: str) -> None:
return None

with _lock:
loop = _mcp_loop
if loop is None or not loop.is_running():
Expand All @@ -2486,12 +2503,24 @@ def _run_on_mcp_loop(coro_or_factory, timeout: float = 30):
raise RuntimeError("MCP event loop unavailable (failed to schedule)")
start_time = time.monotonic()
deadline = None if timeout is None else start_time + timeout
# Activity-heartbeat state for the long-call liveness touch. Cadence of
# 30s is well inside the gateway inactivity window and the kanban
# stuck-after-seconds default (15 min), and the underlying bridge is
# itself rate-limited to one write per 60s so an over-eager cadence is
# harmless.
_activity_state = {
"start": start_time,
"last_touch": start_time,
"interval": 30.0,
}

while True:
if is_interrupted():
future.cancel()
raise InterruptedError("User sent a new message")

touch_activity_if_due(_activity_state, "waiting for MCP tool call")

wait_timeout = 0.1
if deadline is not None:
remaining = deadline - time.monotonic()
Expand Down
Loading