fix(gateway): run shutdown tool-subprocess kill off the event loop
_stop_kill_tool_subprocesses fans out into blocking per-target kills (checkpoint I/O, subprocess.run, sandbox exec) but ran inline on the asyncio loop from _stop_interrupt_remaining_work and _stop_release_runtime_state. Offload the sweep via asyncio.to_thread (loop default executor, clear of the gateway-owned executor quiesce), awaited in phase order. _stop_release_runtime_state is now async; _stop_impl awaits it. Fixes #116327.
This commit is contained in:
@@ -1681,6 +1681,24 @@ class GatewayShutdownMixin:
|
||||
_step("cleanup_all_browsers", _cleanup_browsers)
|
||||
return _marked_cron_jobs
|
||||
|
||||
@staticmethod
|
||||
async def _stop_kill_tool_subprocesses_off_loop(phase: str) -> list:
|
||||
"""Run _stop_kill_tool_subprocesses in a worker thread; returns cron job IDs marked interrupted.
|
||||
|
||||
``kill_all`` fans out into per-target ``kill_process`` calls that do blocking work
|
||||
(registry checkpoint disk I/O, ``subprocess.run`` for systemd scopes, sandbox exec),
|
||||
so running the sweep inline would monopolize the gateway event loop (#116327).
|
||||
Offloaded with ``asyncio.to_thread`` — the loop's default executor, deliberately NOT
|
||||
the gateway-owned ``self._executor``, which ``_stop_quiesce_and_close_session_dbs``
|
||||
drains right after this phase. Phase order is preserved: callers await this before
|
||||
cron notices / adapter teardown. If the surrounding stop task is cancelled while the
|
||||
worker runs, the thread is left to finish on its own; the thread-based shutdown
|
||||
watchdog remains the hard backstop.
|
||||
"""
|
||||
return await asyncio.to_thread(
|
||||
GatewayShutdownMixin._stop_kill_tool_subprocesses, phase
|
||||
)
|
||||
|
||||
async def _stop_begin_teardown(self, ctx: "GatewayShutdownMixin._StopContext") -> None:
|
||||
"""Flag teardown, stop room worker/watchdog, notify sessions."""
|
||||
logger.info("Stopping gateway%s...", " for restart" if self._restart_requested else "")
|
||||
@@ -1785,7 +1803,8 @@ class GatewayShutdownMixin:
|
||||
self._interrupt_running_agents(reason)
|
||||
logger.debug("Re-signaled interrupt for work still live at settle-window exit")
|
||||
# Kill tool subprocesses NOW: deferring past adapter/DB teardown risks the systemd cgroup SIGKILL.
|
||||
_interrupted_cron_jobs = GatewayRunner._stop_kill_tool_subprocesses("post-interrupt")
|
||||
# Off-loop: the sweep does blocking kills that must not monopolize the event loop (#116327).
|
||||
_interrupted_cron_jobs = await GatewayRunner._stop_kill_tool_subprocesses_off_loop("post-interrupt")
|
||||
logger.info("Shutdown phase: post-interrupt tool kill done at +%.2fs", ctx.elapsed())
|
||||
# Last window with the transport up (the cron worker's own notice arrives after teardown).
|
||||
with _log_suppressed(logging.DEBUG, "Cron interrupt notification failed: %s"):
|
||||
@@ -1828,7 +1847,7 @@ class GatewayShutdownMixin:
|
||||
_profile_adapters.clear()
|
||||
logger.info("Shutdown phase: all adapters disconnected at +%.2fs", ctx.elapsed())
|
||||
|
||||
def _stop_release_runtime_state(self, ctx: "GatewayShutdownMixin._StopContext") -> None:
|
||||
async def _stop_release_runtime_state(self, ctx: "GatewayShutdownMixin._StopContext") -> None:
|
||||
"""Cancel background tasks, flush pending messages, clear per-session state, final tool kill."""
|
||||
from gateway.run import GatewayRunner
|
||||
for _task in list(self._background_tasks):
|
||||
@@ -1865,7 +1884,8 @@ class GatewayShutdownMixin:
|
||||
getattr(self, _attr).clear()
|
||||
self._shutdown_event.set()
|
||||
# Global catch-all subprocess kill (safe to repeat) for the graceful path and late respawns.
|
||||
GatewayRunner._stop_kill_tool_subprocesses("final-cleanup")
|
||||
# Off-loop: same blocking sweep as the post-interrupt kill (#116327).
|
||||
await GatewayRunner._stop_kill_tool_subprocesses_off_loop("final-cleanup")
|
||||
logger.info("Shutdown phase: final-cleanup tool kill done at +%.2fs", ctx.elapsed())
|
||||
# Reap the auxiliary-client cache: clients bound to dead worker-thread loops leak httpx transports.
|
||||
def _reap_aux_clients() -> None:
|
||||
@@ -2043,7 +2063,7 @@ class GatewayShutdownMixin:
|
||||
if ctx.timed_out:
|
||||
await GatewayRunner._stop_interrupt_remaining_work(self, ctx)
|
||||
await GatewayRunner._stop_finalize_agents_and_adapters(self, ctx)
|
||||
GatewayRunner._stop_release_runtime_state(self, ctx)
|
||||
await GatewayRunner._stop_release_runtime_state(self, ctx)
|
||||
GatewayRunner._stop_quiesce_and_close_session_dbs(self, timeout, ctx)
|
||||
GatewayRunner._stop_persist_exit_state(self, ctx)
|
||||
finally:
|
||||
|
||||
87
tests/gateway/test_shutdown_kill_off_loop.py
Normal file
87
tests/gateway/test_shutdown_kill_off_loop.py
Normal file
@@ -0,0 +1,87 @@
|
||||
"""Gateway shutdown must not run process teardown on the event loop (#116327).
|
||||
|
||||
``ProcessRegistry.kill_all()`` is synchronous and blocking (per-target
|
||||
``kill_process`` does disk I/O via ``_write_checkpoint`` and may spawn
|
||||
``subprocess.run`` up to 15 s via ``_stop_systemd_unit``). Driving it inline
|
||||
from ``_stop_interrupt_remaining_work`` monopolizes the asyncio loop. These
|
||||
tests pin the invariant: the kill sweep runs off-loop, in phase order.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import threading
|
||||
import time
|
||||
|
||||
import pytest
|
||||
|
||||
from gateway.run_shutdown import GatewayShutdownMixin
|
||||
from tests.gateway.restart_test_helpers import make_restart_runner
|
||||
|
||||
|
||||
def _make_phase_runner(monkeypatch, events):
|
||||
runner, _adapter = make_restart_runner()
|
||||
runner._restart_drain_timeout = 0.01
|
||||
|
||||
loop_thread = threading.current_thread()
|
||||
|
||||
def _fake_kill_all(task_id=None):
|
||||
events.append(("kill_all", threading.current_thread()))
|
||||
return 2
|
||||
|
||||
import tools.process_registry as _pr
|
||||
|
||||
monkeypatch.setattr(_pr.process_registry, "kill_all", _fake_kill_all)
|
||||
monkeypatch.setattr(
|
||||
"cron.scheduler.mark_running_jobs_interrupted", lambda *a, **k: []
|
||||
)
|
||||
monkeypatch.setattr("tools.async_delegation.interrupt_all", lambda *a, **k: 0)
|
||||
monkeypatch.setattr(
|
||||
"tools.terminal_tool_lifecycle.cleanup_all_environments",
|
||||
lambda: events.append(("cleanup_envs", threading.current_thread())),
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
"tools.browser_tool_lifecycle.cleanup_all_browsers",
|
||||
lambda: events.append(("cleanup_browsers", threading.current_thread())),
|
||||
)
|
||||
return runner, loop_thread
|
||||
|
||||
|
||||
def _make_ctx():
|
||||
ctx = GatewayShutdownMixin._StopContext(deferred_count=lambda: 0)
|
||||
ctx.started_at = time.monotonic()
|
||||
return ctx
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_post_interrupt_kill_runs_off_event_loop(monkeypatch):
|
||||
"""kill_all must execute on a worker thread, not the loop thread (#116327)."""
|
||||
events: list = []
|
||||
runner, loop_thread = _make_phase_runner(monkeypatch, events)
|
||||
|
||||
await runner._stop_interrupt_remaining_work(_make_ctx())
|
||||
|
||||
kill_threads = [t for name, t in events if name == "kill_all"]
|
||||
assert kill_threads, f"expected kill_all to run, got events: {events}"
|
||||
for t in kill_threads:
|
||||
assert t is not loop_thread, (
|
||||
"kill_all ran on the event-loop thread; it must be offloaded "
|
||||
"via asyncio.to_thread (#116327)"
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_post_interrupt_kill_preserves_phase_order(monkeypatch):
|
||||
"""Offloading must not reorder the teardown sequence within the phase."""
|
||||
events: list = []
|
||||
runner, _loop_thread = _make_phase_runner(monkeypatch, events)
|
||||
|
||||
await runner._stop_interrupt_remaining_work(_make_ctx())
|
||||
|
||||
names = [name for name, _t in events]
|
||||
assert "kill_all" in names
|
||||
assert names.index("kill_all") < names.index("cleanup_envs")
|
||||
assert names.index("cleanup_envs") < names.index("cleanup_browsers")
|
||||
# The loop must still be responsive while the sweep runs: a callback
|
||||
# scheduled during teardown fires without waiting for the phase.
|
||||
probe = asyncio.Event()
|
||||
asyncio.get_running_loop().call_soon(probe.set)
|
||||
await asyncio.wait_for(probe.wait(), timeout=5)
|
||||
Reference in New Issue
Block a user