Skip to content
62 changes: 61 additions & 1 deletion gateway/run.py
Original file line number Diff line number Diff line change
Expand Up @@ -5908,6 +5908,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
_restart_command_source: Optional[SessionSource] = None
_stop_task: Optional[asyncio.Task] = None
_restart_task: Optional[asyncio.Task] = None
_shutdown_task: Optional[asyncio.Task] = None
_profile_failed_platforms: Optional[Dict[str, Dict[Platform, asyncio.Task]]] = None
_systemd_watchdog: Optional[Any] = None
_startup_restore_in_progress: bool = False
Expand Down Expand Up @@ -6113,6 +6114,10 @@ def __init__(self, config: Optional[GatewayConfig] = None):
self._booted_from_restart: bool = False
self._stop_task: Optional[asyncio.Task] = None
self._restart_task: Optional[asyncio.Task] = None
# Strong reference to the task the SIGINT/SIGTERM handler creates for
# stop(). Same role as _restart_task on the SIGUSR1 path — see the
# comment in request_restart() and shutdown_signal_handler().
self._shutdown_task: Optional[asyncio.Task] = None
self._executor_lock = threading.Lock()
self._executor: Optional[concurrent.futures.ThreadPoolExecutor] = None
# Set on gateway stop so the recreate-on-shutdown path can't resurrect
Expand Down Expand Up @@ -13025,6 +13030,40 @@ async def _stop_systemd_watchdog(self) -> None:
self._systemd_watchdog = None
await watchdog.stop()

def _schedule_shutdown_task(self) -> "asyncio.Task":
"""Schedule ``stop()`` from a signal handler, keeping a strong reference.

Signal handlers cannot await, so the SIGINT/SIGTERM handler has to
schedule ``stop()`` as a task. The event loop keeps only a *weak*
reference to a task, so a bare ``asyncio.create_task(self.stop())``
can be garbage-collected while still pending — the same hazard
``request_restart()`` keeps ``self._restart_task`` to avoid on the
SIGUSR1 path.

``stop()`` only becomes self-anchoring once it reaches
``self._stop_task = asyncio.create_task(_stop_impl())``. A collection
before that point means ``_stop_impl`` is never created at all, so the
gateway does not tear down: it keeps running until the service
manager's stop timeout escalates to SIGKILL, with in-flight turns
undrained and sessions never finalized.

A repeat signal must not re-point the anchor. The handler re-enters in
ordinary operation — a second Ctrl+C, a SIGINT followed by the service
manager's SIGTERM, or ``_run_planned_stop_watcher`` driving the same
callable from its polling thread — and overwriting the attribute would
drop the only strong reference to the task performing the live
teardown. Returning the in-flight task instead loses nothing:
``stop()`` short-circuits on ``self._stop_task``, so a second task
would only have awaited the first one's work.

Returns the task that owns the shutdown, new or already running.
"""
existing = self._shutdown_task
if existing is not None and not existing.done():
return existing
self._shutdown_task = asyncio.create_task(self.stop())
return self._shutdown_task

async def stop(
self,
*,
Expand Down Expand Up @@ -13386,6 +13425,18 @@ def _phase_elapsed() -> float:
# into this _stop_impl and skip _shutdown_event.set() /
# _exit_code = 75 (#12875). It self-terminates anyway.
continue
if _task is self._shutdown_task:
# Same shape as _restart_task: the task the SIGINT/SIGTERM
# handler created is sitting in `await self._stop_task`
# right now, i.e. it is awaiting *this* coroutine.
# Cancelling it would propagate CancelledError back into
# _stop_impl and skip the tail of teardown. Like
# _restart_task it is deliberately not added to
# _background_tasks, but the exemption belongs here with
# the other two: this loop is the single place that
# enforces the invariant, so any future path that does
# track the handle cannot silently cancel the teardown.
continue
_task.cancel()
self._background_tasks.clear()

Expand Down Expand Up @@ -27765,7 +27816,16 @@ def shutdown_signal_handler(received_signal=None):
)
except Exception as _e:
logger.debug("spawn_async_diagnostic failed: %s", _e)
asyncio.create_task(runner.stop())
# Schedule teardown while keeping a strong reference to the task: a
# bare asyncio.create_task() here leaves the event loop holding only a
# weak reference, so a still-pending shutdown can be garbage-collected
# and the gateway then never tears down at all. This is the path every
# `systemctl stop`, `hermes gateway stop` and interactive Ctrl+C takes,
# and it re-enters (second Ctrl+C, SIGINT then SIGTERM, or the
# planned-stop watcher thread racing a real signal), so the anchor must
# not be re-pointed while a shutdown is still in flight. See
# GatewayRunner._schedule_shutdown_task.
runner._schedule_shutdown_task()

def restart_signal_handler():
runner.request_restart(detached=False, via_service=True)
Expand Down
177 changes: 177 additions & 0 deletions tests/gateway/test_shutdown_task_anchor.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,177 @@
"""Regression: the SIGINT/SIGTERM shutdown task must stay strongly referenced.

The gateway's signal handler cannot await, so it schedules ``stop()`` as a
task. For a long time it did that with a bare ``asyncio.create_task(...)`` and
threw the handle away. The event loop keeps only a weak reference to a task,
so a still-pending task can be garbage-collected mid-flight — which is exactly
why the *other* signal path (``request_restart``, SIGUSR1) holds
``self._restart_task`` and says so in a comment.

If the shutdown task is collected before ``stop()`` reaches
``self._stop_task = asyncio.create_task(_stop_impl())``, no teardown coroutine
is ever created: the gateway simply keeps running until the service manager's
stop timeout escalates to SIGKILL, so in-flight turns are never drained and
sessions are never finalized.

``GatewayRunner._schedule_shutdown_task`` owns that behaviour and is what the
signal handler calls.
"""

from __future__ import annotations

import asyncio
import contextlib
from unittest.mock import AsyncMock, patch

import pytest

from gateway.run import GatewayRunner
from tests.gateway.restart_test_helpers import make_restart_runner


def _make_runner_with_blocking_stop():
"""A bare runner whose ``stop()`` parks until the returned event is set."""
runner, adapter = make_restart_runner()
release = asyncio.Event()
calls: list[int] = []

async def _blocking_stop() -> None:
calls.append(1)
await release.wait()

runner.stop = _blocking_stop
runner._schedule_shutdown_task = GatewayRunner._schedule_shutdown_task.__get__(
runner, GatewayRunner
)
return runner, adapter, release, calls


async def _drain(*tasks: asyncio.Task) -> None:
for task in tasks:
task.cancel()
with contextlib.suppress(asyncio.CancelledError):
await task


def test_gateway_runner_declares_a_shutdown_task_slot():
"""The anchor needs a class-level default so bare runners inherit ``None``.

Shutdown-path tests build runners via ``object.__new__`` (``stop()`` even
carries a getattr-guard for them), so the attribute has to exist on the
class the way ``_stop_task`` and ``_restart_task`` do.
"""
assert GatewayRunner._stop_task is None
assert GatewayRunner._restart_task is None
assert GatewayRunner._shutdown_task is None


@pytest.mark.asyncio
async def test_scheduling_shutdown_anchors_the_task_on_the_runner():
"""The scheduled task is reachable from the runner, not just from the loop.

The event loop holds only a weak reference to a task, so without this the
shutdown can be collected while still pending.
"""
runner, _adapter, release, calls = _make_runner_with_blocking_stop()

task = runner._schedule_shutdown_task()
await asyncio.sleep(0)

assert runner._shutdown_task is task
assert task.done() is False
assert calls == [1]

release.set()
await task


@pytest.mark.asyncio
async def test_a_repeat_signal_does_not_replace_the_in_flight_shutdown_task():
"""A second signal must reuse the task that owns the live teardown.

The handler re-enters in ordinary operation: a second Ctrl+C, a SIGINT
followed by the service manager's SIGTERM, or ``_run_planned_stop_watcher``
driving the same callable from its polling thread. Re-pointing the anchor
would drop the only strong reference to the running shutdown.
"""
runner, _adapter, release, calls = _make_runner_with_blocking_stop()

first = runner._schedule_shutdown_task()
await asyncio.sleep(0)
second = runner._schedule_shutdown_task()
await asyncio.sleep(0)

assert second is first
assert runner._shutdown_task is first
assert calls == [1], "the repeat signal started a second stop()"

release.set()
await first


@pytest.mark.asyncio
async def test_a_later_signal_schedules_again_once_the_previous_one_finished():
"""The guard must not wedge the gateway if an earlier shutdown completed."""
runner, _adapter, release, calls = _make_runner_with_blocking_stop()

first = runner._schedule_shutdown_task()
await asyncio.sleep(0)
release.set()
await first

second = runner._schedule_shutdown_task()
await asyncio.sleep(0)

assert second is not first
assert runner._shutdown_task is second
assert calls == [1, 1]

await second


@pytest.mark.asyncio
async def test_stop_does_not_cancel_the_anchored_shutdown_task():
"""``_stop_impl``'s cancel sweep must skip ``_shutdown_task``.

The anchored task is parked in ``await self._stop_task`` while the sweep
runs, so cancelling it would push ``CancelledError`` into the very
``_stop_impl`` doing the cancelling. This is why the shutdown task cannot
simply be parked in ``_background_tasks`` — the same reason ``_stop_task``
and ``_restart_task`` are already exempt.
"""
runner, adapter = make_restart_runner()
runner._restart_drain_timeout = 0.0
adapter.disconnect = AsyncMock()

async def _park() -> None:
await asyncio.Event().wait()

shutdown_task = asyncio.create_task(_park())
control_task = asyncio.create_task(_park())
await asyncio.sleep(0)

runner._shutdown_task = shutdown_task
runner._background_tasks = {shutdown_task, control_task}

try:
with (
patch("gateway.status.remove_pid_file"),
patch("gateway.status.write_runtime_status"),
patch("agent.auxiliary_client.shutdown_cached_clients"),
):
await runner.stop()

for _ in range(10):
if control_task.done():
break
await asyncio.sleep(0)

assert control_task.cancelled() is True, (
"an ordinary background task should still be swept by _stop_impl"
)
assert shutdown_task.done() is False, (
"_stop_impl cancelled the shutdown task, which is awaiting the very "
"_stop_task that runs this sweep"
)
finally:
await _drain(shutdown_task, control_task)
Loading