fix(gateway): tag the loop-liveness and heartbeat-poll tasks as permanent supervised watchers (#84558)
#84327 excluded _spawn_supervised's permanent watchers (session-expiry, kanban, reconnect, the scale-to-zero watcher itself, ...) from _scale_to_zero_has_live_background_work() via a _hermes_supervised_watcher tag, because counting them made an armed gateway consider itself busy forever and never go dormant. Two more permanent, infinite-loop tasks are added to _background_tasks OUTSIDE _spawn_supervised and were untagged: - _loop_heartbeat_task (loop_heartbeat_forever, #66892): a `while True` loop started unconditionally in start() on every gateway boot. Extracted the inline spawn block into _start_loop_heartbeat_task() so it's independently testable, matching the existing _start_heartbeat_poller() pattern. - _heartbeat_poll_task (_poll_loop in _start_heartbeat_poller): also a `while True` loop, started the first time a session registers a heartbeat watch, and then permanent for the rest of the process. Because _loop_heartbeat_task starts on every boot, it alone made _scale_to_zero_has_live_background_work() return True forever on every armed instance, regardless of the #84327 fix -- confirmed empirically against the real method with the exact untagged-task shape this task has. Two new regression tests spawn each task through its real production entry point and assert the busy check returns False; both fail against the unfixed code (missing method / real assertion failure). Co-authored-by: pierrenode <298902573+pierrenode@users.noreply.github.com> Co-authored-by: Ben Barclay <ben@nousresearch.com>
This commit is contained in:
+36
-19
@@ -12222,6 +12222,36 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
logger.warning("Legacy session recovery on startup failed: %s", exc)
|
||||
return exact, fallback
|
||||
|
||||
def _start_loop_heartbeat_task(self) -> None:
|
||||
"""Start the loop-liveness heartbeat task (#66892), idempotent.
|
||||
|
||||
An asyncio task so a frozen loop stops refreshing
|
||||
``state/gateway.heartbeat``. Cancelled with the other background
|
||||
tasks during stop(). Best-effort — a liveness probe must never be
|
||||
able to abort startup.
|
||||
"""
|
||||
try:
|
||||
_existing_hb = getattr(self, "_loop_heartbeat_task", None)
|
||||
if _existing_hb is not None and not _existing_hb.done():
|
||||
return
|
||||
self._loop_heartbeat_task = asyncio.create_task(
|
||||
loop_heartbeat_forever(
|
||||
interval_s=DEFAULT_HEARTBEAT_INTERVAL_S,
|
||||
start_time=getattr(self, "_gateway_started_at", 0.0),
|
||||
)
|
||||
)
|
||||
# PERMANENT for the process lifetime, same as a
|
||||
# _spawn_supervised watcher — tag it so
|
||||
# _scale_to_zero_has_live_background_work() doesn't treat an
|
||||
# armed, otherwise-idle gateway as busy forever.
|
||||
self._loop_heartbeat_task._hermes_supervised_watcher = True # type: ignore[attr-defined]
|
||||
_bg = getattr(self, "_background_tasks", None)
|
||||
if _bg is not None:
|
||||
_bg.add(self._loop_heartbeat_task)
|
||||
self._loop_heartbeat_task.add_done_callback(_bg.discard)
|
||||
except Exception:
|
||||
logger.debug("Failed to start gateway loop heartbeat", exc_info=True)
|
||||
|
||||
async def start(self) -> bool:
|
||||
"""
|
||||
Start the gateway and all configured platform adapters.
|
||||
@@ -12991,25 +13021,7 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
self._install_plugin_message_injector()
|
||||
self._update_runtime_status("running")
|
||||
|
||||
# Loop-liveness heartbeat (#66892): an asyncio task so a frozen loop
|
||||
# stops refreshing ``state/gateway.heartbeat``. Cancelled with the
|
||||
# other background tasks during stop(). Best-effort — a liveness probe
|
||||
# must never be able to abort startup.
|
||||
try:
|
||||
_existing_hb = getattr(self, "_loop_heartbeat_task", None)
|
||||
if _existing_hb is None or _existing_hb.done():
|
||||
self._loop_heartbeat_task = asyncio.create_task(
|
||||
loop_heartbeat_forever(
|
||||
interval_s=DEFAULT_HEARTBEAT_INTERVAL_S,
|
||||
start_time=getattr(self, "_gateway_started_at", 0.0),
|
||||
)
|
||||
)
|
||||
_bg = getattr(self, "_background_tasks", None)
|
||||
if _bg is not None:
|
||||
_bg.add(self._loop_heartbeat_task)
|
||||
self._loop_heartbeat_task.add_done_callback(_bg.discard)
|
||||
except Exception:
|
||||
logger.debug("Failed to start gateway loop heartbeat", exc_info=True)
|
||||
self._start_loop_heartbeat_task()
|
||||
|
||||
# Emit gateway:startup hook
|
||||
hook_count = len(self.hooks.loaded_hooks)
|
||||
@@ -21249,6 +21261,11 @@ class GatewayRunner(GatewayAuthorizationMixin, GatewayKanbanWatchersMixin, Gatew
|
||||
try:
|
||||
task = asyncio.create_task(_poll_loop())
|
||||
self._heartbeat_poll_task = task
|
||||
# PERMANENT once started (an infinite while-True loop, no exit
|
||||
# condition) — same as a _spawn_supervised watcher. Tag it so
|
||||
# _scale_to_zero_has_live_background_work() doesn't treat a
|
||||
# gateway with an active heartbeat watch as busy forever.
|
||||
task._hermes_supervised_watcher = True # type: ignore[attr-defined]
|
||||
_bg = getattr(self, "_background_tasks", None)
|
||||
if _bg is not None:
|
||||
_bg.add(task)
|
||||
|
||||
@@ -370,3 +370,48 @@ async def test_done_supervised_watcher_is_ignored_either_way():
|
||||
await t
|
||||
r._background_tasks = {t}
|
||||
assert r._scale_to_zero_has_live_background_work() is False
|
||||
|
||||
|
||||
# ── permanent tasks spawned OUTSIDE _spawn_supervised must also be tagged ──
|
||||
#
|
||||
# _loop_heartbeat_task and _heartbeat_poll_task are both infinite while-True
|
||||
# loops added to _background_tasks via plain asyncio.create_task() + manual
|
||||
# add(), NOT through _spawn_supervised — so they were untagged and defeated
|
||||
# the fix above: _loop_heartbeat_task starts unconditionally on every
|
||||
# gateway boot (start()), which would make the busy check return True
|
||||
# forever regardless of the _spawn_supervised fix, on every armed instance.
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_loop_heartbeat_task_does_not_block_idle():
|
||||
r = GatewayRunner.__new__(GatewayRunner)
|
||||
r._running = True
|
||||
r._background_tasks = set()
|
||||
r._loop_heartbeat_task = None
|
||||
r._gateway_started_at = time.time()
|
||||
|
||||
r._start_loop_heartbeat_task()
|
||||
await asyncio.sleep(0) # let the task start
|
||||
try:
|
||||
assert r._scale_to_zero_has_live_background_work() is False
|
||||
finally:
|
||||
r._loop_heartbeat_task.cancel()
|
||||
await asyncio.gather(r._loop_heartbeat_task, return_exceptions=True)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_heartbeat_poll_task_does_not_block_idle():
|
||||
r = GatewayRunner.__new__(GatewayRunner)
|
||||
r._running = True
|
||||
r._background_tasks = set()
|
||||
r._heartbeat_poll_task = None
|
||||
r._heartbeat_watch = {}
|
||||
r._running_agents = {}
|
||||
|
||||
r._start_heartbeat_poller()
|
||||
await asyncio.sleep(0) # let the task start
|
||||
try:
|
||||
assert r._scale_to_zero_has_live_background_work() is False
|
||||
finally:
|
||||
r._heartbeat_poll_task.cancel()
|
||||
await asyncio.gather(r._heartbeat_poll_task, return_exceptions=True)
|
||||
|
||||
Reference in New Issue
Block a user