fix(gateway): interlock the stale-heartbeat wedge verdict with a loop-scheduling witness
The off-loop heartbeat write broke the producer->consumer invariant #86860 depends on: file freshness no longer equals loop schedulability, yet the probe still classified a stale file as WEDGED — and WEDGED is destructive authority (SIGTERM -> SIGKILL, bypassing the #86684 cron drain floor). The measured motivating stall (112.6s max) exceeds the 90s stale budget, so a healthy loop blocked inside the watchdog's own write could be killed, and executor saturation produces the same false positive. The inverse edge also existed: an off-loop write landing after the loop froze refreshes the file mtime, manufacturing a false-fresh liveness proof. The gateway loop now also arms a loop-scheduling witness: a UNIX socket (state/gateway.loop-tick.<pid>.sock) answered by the loop itself via await asyncio.start_unix_server — socket-buffer writes, no fsync, no disk I/O, so it keeps working on the filesystem that stalls the heartbeat write. The heartbeat payload records whether the witness is armed (loop_tick_socket). The classifier is now two-witness: - socket answers -> ALIVE (file age irrelevant: a stalled write or saturated executor can no longer produce a wedge verdict) - file fresh, socket silent -> UNKNOWN (a late off-loop write can no longer manufacture a liveness proof) - file stale, socket silent, producer armed -> WEDGED (both witnesses agree the loop stopped scheduling) - legacy payload (no flag) -> unchanged single-witness contract: the legacy producer wrote on-loop, so staleness is still proof - any conflict/ambiguity -> UNKNOWN, never escalate Tests are a producer->consumer composition: a real heartbeat loop with a stalled write probes ALIVE while the file is past the stale budget, and launchd_restart fed by the real probe drains instead of escalating; a silent socket with a fresh file denies ALIVE; WEDGED requires the armed socket to agree; a bind-failed producer disables stale escalation; legacy payloads keep the old contract; a source-inspection test pins that the witness is awaited on the loop. Mutation-checked: reverting either source file fails the new tests. 45 tests pass across the watchdog suites; ruff clean.
This commit is contained in:
@@ -229,6 +229,24 @@ def get_loop_heartbeat_path(home: Optional[Path] = None) -> Path:
|
||||
return base.joinpath(*_HEARTBEAT_RELATIVE)
|
||||
|
||||
|
||||
def get_loop_tick_socket_path(
|
||||
home: Optional[Path] = None, pid: Optional[int] = None
|
||||
) -> Path:
|
||||
"""Return the loop-scheduling witness socket for ``pid``.
|
||||
|
||||
``<HERMES_HOME>/state/gateway.loop-tick.<pid>.sock`` — PID-suffixed so a
|
||||
leftover node from a previous process can never be mistaken for this
|
||||
gateway's witness. Served by the gateway loop itself (see
|
||||
``_tick_socket_handler``): an answer is direct proof that the loop is
|
||||
dispatching, which is exactly the property the heartbeat file lost when
|
||||
its write moved off-loop (#90502).
|
||||
"""
|
||||
base = home if home is not None else _process_hermes_home()
|
||||
return base.joinpath(
|
||||
"state", f"gateway.loop-tick.{int(pid if pid is not None else os.getpid())}.sock"
|
||||
)
|
||||
|
||||
|
||||
def get_shutdown_watchdog_dump_path(home: Optional[Path] = None) -> Path:
|
||||
"""Return the faulthandler / metadata dump path for a fired watchdog."""
|
||||
base = home if home is not None else _process_hermes_home()
|
||||
@@ -434,6 +452,30 @@ def arm_shutdown_watchdog(
|
||||
return done
|
||||
|
||||
|
||||
async def _tick_socket_handler(
|
||||
reader: asyncio.StreamReader, writer: asyncio.StreamWriter
|
||||
) -> None:
|
||||
"""Answer a liveness ping with one byte.
|
||||
|
||||
Runs on the gateway loop: the reply is produced only while the loop is
|
||||
actually dispatching, so a successful read is a witness of loop
|
||||
schedulability that no executor thread and no filesystem stall can
|
||||
refresh. A UNIX-socket write is a socket-buffer copy — no fsync, no
|
||||
disk I/O — so the witness keeps working on the exact filesystem that
|
||||
stalls the heartbeat write. Best-effort; never raises.
|
||||
"""
|
||||
try:
|
||||
writer.write(b"1")
|
||||
await writer.drain()
|
||||
except Exception:
|
||||
pass
|
||||
finally:
|
||||
try:
|
||||
writer.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
async def loop_heartbeat_forever(
|
||||
*,
|
||||
interval_s: float = DEFAULT_HEARTBEAT_INTERVAL_S,
|
||||
@@ -462,31 +504,78 @@ async def loop_heartbeat_forever(
|
||||
while the loop was wedged, which is exactly the signal the docstring above
|
||||
promises. And a single in-flight write at a time, so a 112s stall cannot pile
|
||||
up one queued thread per interval behind it.
|
||||
|
||||
Because the write is now off-loop, file freshness is no longer *proof* of
|
||||
loop schedulability: a stalled write or a saturated executor can age the file
|
||||
while the loop runs, and a write that lands after the loop froze can keep it
|
||||
fresh. The file therefore stops being sufficient authority on its own. This
|
||||
task also arms a loop-scheduling witness — a UNIX socket answered by the
|
||||
loop itself (``_tick_socket_handler``) — and records whether it is armed in
|
||||
the heartbeat payload (``loop_tick_socket``). External probes must require
|
||||
the witness to agree with file staleness before classifying a loop as
|
||||
wedged; see ``hermes_cli.gateway.probe_gateway_loop_liveness`` for the
|
||||
two-witness contract.
|
||||
"""
|
||||
try:
|
||||
interval = max(float(interval_s), 1.0)
|
||||
except (TypeError, ValueError):
|
||||
interval = DEFAULT_HEARTBEAT_INTERVAL_S
|
||||
|
||||
# Arm the loop-scheduling witness. Best-effort: a failed bind (permissions,
|
||||
# path length) must not abort the gateway or the file heartbeat — it only
|
||||
# disables the witness, and the payload flag tells probes that staleness is
|
||||
# no longer sufficient authority to escalate.
|
||||
tick_server = None
|
||||
tick_socket_path = None
|
||||
try:
|
||||
tick_socket_path = get_loop_tick_socket_path(home)
|
||||
tick_socket_path.parent.mkdir(parents=True, exist_ok=True)
|
||||
tick_server = await asyncio.start_unix_server(
|
||||
_tick_socket_handler, path=str(tick_socket_path)
|
||||
)
|
||||
except Exception:
|
||||
tick_server = None
|
||||
logger.warning(
|
||||
"Loop tick socket unavailable — liveness probes will have no "
|
||||
"loop-scheduling witness and will not escalate on a stale heartbeat",
|
||||
exc_info=True,
|
||||
)
|
||||
|
||||
async def _write_off_loop() -> None:
|
||||
# write_loop_heartbeat never raises, so a failure here is an executor
|
||||
# problem (shutdown, saturation) and must not kill the heartbeat task.
|
||||
try:
|
||||
await asyncio.to_thread(
|
||||
write_loop_heartbeat, start_time=start_time, home=home
|
||||
write_loop_heartbeat,
|
||||
start_time=start_time,
|
||||
home=home,
|
||||
extra={"loop_tick_socket": tick_server is not None},
|
||||
)
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception:
|
||||
logger.debug("Loop heartbeat write failed off-loop", exc_info=True)
|
||||
|
||||
# Immediate first write so monitors see a fresh file as soon as the
|
||||
# gateway is running, not after the first interval.
|
||||
await _write_off_loop()
|
||||
while True:
|
||||
if should_continue is not None and not should_continue():
|
||||
return
|
||||
await asyncio.sleep(interval)
|
||||
if should_continue is not None and not should_continue():
|
||||
return
|
||||
try:
|
||||
# Immediate first write so monitors see a fresh file as soon as the
|
||||
# gateway is running, not after the first interval.
|
||||
await _write_off_loop()
|
||||
while True:
|
||||
if should_continue is not None and not should_continue():
|
||||
return
|
||||
await asyncio.sleep(interval)
|
||||
if should_continue is not None and not should_continue():
|
||||
return
|
||||
await _write_off_loop()
|
||||
finally:
|
||||
if tick_server is not None:
|
||||
tick_server.close()
|
||||
try:
|
||||
await tick_server.wait_closed()
|
||||
except Exception:
|
||||
pass
|
||||
if tick_socket_path is not None:
|
||||
try:
|
||||
tick_socket_path.unlink(missing_ok=True)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
@@ -12,6 +12,7 @@ import os
|
||||
import shlex
|
||||
import shutil
|
||||
import signal
|
||||
import socket
|
||||
import subprocess
|
||||
import sys
|
||||
import textwrap
|
||||
@@ -379,24 +380,42 @@ def _wait_for_pid_exit(pid: int, timeout: float) -> bool:
|
||||
# rewrites ``state/gateway.heartbeat`` every 30s (#66892), so a frozen loop
|
||||
# stops refreshing the file while a busy-but-alive loop keeps refreshing it.
|
||||
#
|
||||
# ``probe_gateway_loop_liveness`` reads that signal (a local stat + JSON read,
|
||||
# instant — far inside the 10s query tier of the subprocess timeout doc) and
|
||||
# classifies the gateway BEFORE any drain wait begins:
|
||||
# Since #90502 the heartbeat write runs on a thread (a stalling filesystem
|
||||
# must not be able to block the loop the watchdog watches), which costs the
|
||||
# file its status as *proof*: a stalled write or a saturated executor can age
|
||||
# the file while the loop runs, and an off-loop write can land after the loop
|
||||
# froze, keeping the file fresh for a dead loop. The loop therefore also arms
|
||||
# a second witness — ``state/gateway.loop-tick.<pid>.sock``, a UNIX socket
|
||||
# answered by the loop itself — and records whether it is armed in the
|
||||
# heartbeat payload (``loop_tick_socket``).
|
||||
#
|
||||
# - ``alive`` — heartbeat is fresh: the loop is dispatching (possibly busy).
|
||||
# Callers must take the normal graceful-drain path, which
|
||||
# honours the in-flight cron drain floor (#86684).
|
||||
# - ``wedged`` — heartbeat belongs to this PID but has gone stale well past
|
||||
# several missed beats: the loop is provably dead. Draining
|
||||
# is pointless (nothing can run the drain), so callers may
|
||||
# escalate immediately via ``_escalate_wedged_gateway``.
|
||||
# - ``unknown`` — no heartbeat / unreadable / PID mismatch (older gateway,
|
||||
# still starting up, stale file from a previous process).
|
||||
# ``probe_gateway_loop_liveness`` reads both signals (a local stat + JSON
|
||||
# read + a bounded socket ping — instant, far inside the 10s query tier of
|
||||
# the subprocess timeout doc) and classifies the gateway BEFORE any drain
|
||||
# wait begins:
|
||||
#
|
||||
# - ``alive`` — the loop answered the tick socket, or the file is fresh and
|
||||
# the loop is not contradicted by the socket. Callers must
|
||||
# take the normal graceful-drain path, which honours the
|
||||
# in-flight cron drain floor (#86684).
|
||||
# - ``wedged`` — the heartbeat belongs to this PID, is stale well past
|
||||
# several missed beats, AND the tick socket is armed but does
|
||||
# not answer: both witnesses agree the loop is provably dead.
|
||||
# Draining is pointless (nothing can run the drain), so
|
||||
# callers may escalate immediately via
|
||||
# ``_escalate_wedged_gateway``.
|
||||
# - ``unknown`` — no heartbeat / unreadable / PID mismatch / witness conflict
|
||||
# (fresh file with a silent loop, armed socket unreachable).
|
||||
# Treated like ``alive``: never escalate on ambiguity.
|
||||
#
|
||||
# The distinction matters: only a *provably dead* loop may bypass the cron
|
||||
# drain floor. A merely busy gateway still answers the probe (fresh file)
|
||||
# and keeps its full drain budget.
|
||||
# drain floor. A merely busy gateway still answers the probe (socket ping)
|
||||
# and keeps its full drain budget — even when the filesystem is stalling the
|
||||
# heartbeat write (the incident that motivated #90502).
|
||||
#
|
||||
# Legacy gateways (no ``loop_tick_socket`` flag in the payload) wrote the
|
||||
# file on-loop, so their staleness remains proof and the old single-witness
|
||||
# contract is unchanged.
|
||||
|
||||
GATEWAY_LOOP_ALIVE = "alive"
|
||||
GATEWAY_LOOP_WEDGED = "wedged"
|
||||
@@ -406,19 +425,82 @@ GATEWAY_LOOP_UNKNOWN = "unknown"
|
||||
# Three missed beats is decisive without false-positiving on one slow write.
|
||||
DEFAULT_LOOP_LIVENESS_STALE_AFTER_S = 90.0
|
||||
|
||||
# Sentinel for "the producer never wrote the witness flag" (legacy payload).
|
||||
_LOOP_TICK_ABSENT = object()
|
||||
|
||||
|
||||
def _probe_loop_tick_socket(
|
||||
pid: int,
|
||||
home: Path | None,
|
||||
timeout: float = 1.0,
|
||||
) -> bool | None:
|
||||
"""Ping the loop-scheduling witness socket for ``pid``.
|
||||
|
||||
Returns:
|
||||
True — the loop answered: it is dispatching right now.
|
||||
False — a socket node exists for this PID but did not answer (the loop
|
||||
is not scheduling, or the node is a leftover from a dead
|
||||
listener).
|
||||
None — no socket node for this PID (legacy producer), or the path
|
||||
could not be resolved. Not evidence either way.
|
||||
"""
|
||||
try:
|
||||
from gateway.shutdown_watchdog import get_loop_tick_socket_path
|
||||
|
||||
path = get_loop_tick_socket_path(home, pid)
|
||||
if not path.is_socket():
|
||||
return None
|
||||
except Exception:
|
||||
return None
|
||||
sock = None
|
||||
try:
|
||||
sock = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
sock.settimeout(max(float(timeout), 0.0))
|
||||
sock.connect(str(path))
|
||||
return sock.recv(1) == b"1"
|
||||
except Exception:
|
||||
# ECONNREFUSED (node with no listener), timeout (loop not answering),
|
||||
# transient errors: the witness exists but is silent.
|
||||
return False
|
||||
finally:
|
||||
if sock is not None:
|
||||
try:
|
||||
sock.close()
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
def probe_gateway_loop_liveness(
|
||||
pid: int,
|
||||
*,
|
||||
stale_after: float = DEFAULT_LOOP_LIVENESS_STALE_AFTER_S,
|
||||
home: Path | None = None,
|
||||
tick_timeout: float = 1.0,
|
||||
) -> str:
|
||||
"""Classify a gateway PID's event loop as alive / wedged / unknown.
|
||||
|
||||
Reads the loop-liveness heartbeat file the gateway rewrites every 30s
|
||||
while its loop is dispatching. Never raises; any ambiguity (missing
|
||||
file, unreadable JSON, PID mismatch) returns ``GATEWAY_LOOP_UNKNOWN``
|
||||
so callers default to the safe graceful-drain path.
|
||||
Two witnesses:
|
||||
|
||||
- the loop-tick socket (``state/gateway.loop-tick.<pid>.sock``): answered
|
||||
by the gateway loop itself, so a reply is direct proof that the loop is
|
||||
dispatching. It is never refreshed by the heartbeat executor thread and
|
||||
never stalled by a filesystem that is slow to fsync.
|
||||
- the heartbeat file (``state/gateway.heartbeat``): rewritten every 30s
|
||||
on a thread since #90502, so freshness alone is no longer proof of loop
|
||||
schedulability — a stalled write (measured at 112.6s max on the
|
||||
incident box) or a saturated executor can age the file while the loop
|
||||
runs, and a write can land after the loop froze.
|
||||
|
||||
A stale file classifies as ``wedged`` only when the producer declared the
|
||||
tick socket armed (``loop_tick_socket: true`` in the payload) AND the
|
||||
socket does not answer — both witnesses agree the loop stopped
|
||||
scheduling. Any conflict or ambiguity returns ``unknown`` so callers keep
|
||||
the safe graceful-drain path. Legacy producers (payload without the
|
||||
flag) wrote the file on-loop, so their staleness remains proof and the
|
||||
old contract is unchanged.
|
||||
|
||||
Never raises; any ambiguity (missing file, unreadable JSON, PID mismatch)
|
||||
returns ``GATEWAY_LOOP_UNKNOWN``.
|
||||
"""
|
||||
try:
|
||||
stale_budget = max(float(stale_after), 0.0)
|
||||
@@ -437,10 +519,41 @@ def probe_gateway_loop_liveness(
|
||||
# No heartbeat for THIS process — old gateway version, still starting
|
||||
# up, or a stale file from a previous PID. Not evidence of a wedge.
|
||||
return GATEWAY_LOOP_UNKNOWN
|
||||
|
||||
witness = _probe_loop_tick_socket(pid, home, timeout=tick_timeout)
|
||||
if witness is True:
|
||||
# The loop answered a ping — it is dispatching right now. A stale
|
||||
# heartbeat file is a stalled write or a saturated executor, not a
|
||||
# wedge (#90502).
|
||||
return GATEWAY_LOOP_ALIVE
|
||||
|
||||
tick_armed = payload.get("loop_tick_socket", _LOOP_TICK_ABSENT)
|
||||
age = time.time() - mtime
|
||||
if age > stale_budget:
|
||||
if age <= stale_budget:
|
||||
if witness is False:
|
||||
# File fresh but the loop did not answer: an off-loop write can
|
||||
# land after the loop froze, so a fresh file is not a liveness
|
||||
# proof while the loop itself is silent.
|
||||
return GATEWAY_LOOP_UNKNOWN
|
||||
return GATEWAY_LOOP_ALIVE
|
||||
|
||||
# File is stale past the budget. The verdict now depends on what the
|
||||
# producer promised about its witness:
|
||||
if tick_armed is _LOOP_TICK_ABSENT:
|
||||
# Legacy producer: the write ran on-loop, so staleness really does
|
||||
# prove the loop stopped scheduling — old contract, unchanged.
|
||||
return GATEWAY_LOOP_WEDGED
|
||||
return GATEWAY_LOOP_ALIVE
|
||||
if tick_armed is not True:
|
||||
# New producer whose witness could not be armed (bind failed): the
|
||||
# write is off-loop, so staleness is NOT proof. Never escalate
|
||||
# without a witness.
|
||||
return GATEWAY_LOOP_UNKNOWN
|
||||
if witness is False:
|
||||
# Both witnesses agree the loop is not scheduling.
|
||||
return GATEWAY_LOOP_WEDGED
|
||||
# Armed producer but the socket is unreachable: ambiguity — never kill on
|
||||
# it. The graceful drain path remains the backstop.
|
||||
return GATEWAY_LOOP_UNKNOWN
|
||||
|
||||
|
||||
def _escalate_wedged_gateway(
|
||||
|
||||
@@ -382,3 +382,24 @@ def test_heartbeat_write_is_awaited_so_a_frozen_loop_still_goes_stale():
|
||||
"the heartbeat write is fire-and-forget; a frozen loop would keep the "
|
||||
"file fresh and the staleness signal would be lost"
|
||||
)
|
||||
|
||||
|
||||
def test_loop_scheduling_witness_is_served_by_the_loop_itself():
|
||||
"""The tick socket must be armed on the loop, never in a thread.
|
||||
|
||||
The two-witness contract in ``probe_gateway_loop_liveness`` rests on the
|
||||
socket being answered only while the loop is actually dispatching. If the
|
||||
server ever moved into the heartbeat's executor thread, a wedged loop
|
||||
could keep answering pings (same class of lie as a fire-and-forget file
|
||||
write) and the interlock would be void.
|
||||
"""
|
||||
src = pathlib.Path(
|
||||
inspect.getsourcefile(loop_heartbeat_forever) or ""
|
||||
).read_text()
|
||||
body = src[src.index("async def loop_heartbeat_forever("):]
|
||||
body = body[: body.index("\ndef ") if "\ndef " in body else len(body)]
|
||||
# Awaited directly on the loop task: a coroutine cannot run inside a
|
||||
# thread, so an awaited start_unix_server is structurally loop-owned.
|
||||
assert "await asyncio.start_unix_server(" in body, (
|
||||
"the loop-scheduling witness socket is not armed by the loop task"
|
||||
)
|
||||
|
||||
@@ -9,14 +9,23 @@ seconds. A busy-but-alive gateway (fresh heartbeat) must keep the full drain
|
||||
path — including the in-flight cron drain floor from #86684.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import json
|
||||
import os
|
||||
import socket
|
||||
import threading
|
||||
import time
|
||||
from unittest.mock import patch
|
||||
|
||||
import pytest
|
||||
|
||||
import hermes_cli.gateway as gateway_cli
|
||||
from gateway.shutdown_watchdog import get_loop_heartbeat_path, write_loop_heartbeat
|
||||
from gateway.shutdown_watchdog import (
|
||||
get_loop_heartbeat_path,
|
||||
get_loop_tick_socket_path,
|
||||
loop_heartbeat_forever,
|
||||
write_loop_heartbeat,
|
||||
)
|
||||
|
||||
|
||||
def _write_heartbeat(home, pid, age_s=0.0):
|
||||
@@ -29,6 +38,34 @@ def _write_heartbeat(home, pid, age_s=0.0):
|
||||
return path
|
||||
|
||||
|
||||
def _mark_witness_flag(home, armed, age_s=0.0):
|
||||
"""Set ``loop_tick_socket`` on the heartbeat payload; re-stamp mtime."""
|
||||
path = get_loop_heartbeat_path(home)
|
||||
payload = json.loads(path.read_text(encoding="utf-8"))
|
||||
payload["loop_tick_socket"] = armed
|
||||
path.write_text(json.dumps(payload), encoding="utf-8")
|
||||
if age_s:
|
||||
stamp = time.time() - age_s
|
||||
os.utime(path, (stamp, stamp))
|
||||
return path
|
||||
|
||||
|
||||
def _silent_socket_node(path):
|
||||
"""Create a socket node at ``path`` that never answers.
|
||||
|
||||
Bind + listen, then close the listener WITHOUT unlinking: the node stays,
|
||||
so a probe's connect() gets ECONNREFUSED — a witness that exists but is
|
||||
silent, exactly like a dead listener (or a loop that stopped scheduling).
|
||||
"""
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
srv = socket.socket(socket.AF_UNIX, socket.SOCK_STREAM)
|
||||
try:
|
||||
srv.bind(str(path))
|
||||
srv.listen(1)
|
||||
finally:
|
||||
srv.close()
|
||||
|
||||
|
||||
class TestProbeGatewayLoopLiveness:
|
||||
def test_fresh_heartbeat_is_alive(self, tmp_path):
|
||||
"""A gateway that refreshed its heartbeat recently is busy, not wedged."""
|
||||
@@ -261,3 +298,238 @@ class TestLaunchdRestartWedgedIntegration:
|
||||
gateway_cli.launchd_restart()
|
||||
assert "escalate" not in events
|
||||
assert ("drain", 4242, 195.0) in events
|
||||
|
||||
|
||||
class TestLoopTickWitness:
|
||||
"""Two-witness liveness (#90502 review).
|
||||
|
||||
The heartbeat write moved off-loop, so a stale file no longer proves a
|
||||
wedged loop and a fresh file no longer proves an alive one. The loop
|
||||
answers a UNIX socket instead; the probe only escalates when BOTH
|
||||
witnesses agree the loop stopped scheduling.
|
||||
"""
|
||||
|
||||
def test_stalled_heartbeat_write_never_escalates_a_running_loop(
|
||||
self, tmp_path, monkeypatch
|
||||
):
|
||||
"""Producer + consumer composition.
|
||||
|
||||
While the heartbeat write is stalled longer than the stale budget, a
|
||||
loop that demonstrably keeps dispatching must probe ALIVE — and a
|
||||
restart path fed by the real probe must take the graceful drain,
|
||||
never the bounded escalation. This is the exact false-positive the
|
||||
review called out: the measured fsync stall (112.6s max) exceeds the
|
||||
90s destructive-classifier threshold.
|
||||
"""
|
||||
pid = os.getpid()
|
||||
block_s = 1.5
|
||||
stale_after = 1.0
|
||||
|
||||
# First write lands immediately; every later write stalls like an
|
||||
# fsync on the incident filesystem, so the file ages past the budget
|
||||
# while the loop keeps running.
|
||||
def stalling_write(**_kwargs):
|
||||
if not get_loop_heartbeat_path(tmp_path).exists():
|
||||
return write_loop_heartbeat(**_kwargs)
|
||||
time.sleep(block_s)
|
||||
return write_loop_heartbeat(**_kwargs)
|
||||
|
||||
errors = []
|
||||
|
||||
async def producer() -> None:
|
||||
with patch(
|
||||
"gateway.shutdown_watchdog.write_loop_heartbeat", stalling_write
|
||||
):
|
||||
task = asyncio.create_task(
|
||||
loop_heartbeat_forever(interval_s=1.0, home=tmp_path)
|
||||
)
|
||||
try:
|
||||
# One interval (1.0s) elapses, the second write starts and
|
||||
# stalls; 1.35s in the file is stale but the loop ticks.
|
||||
await asyncio.sleep(1.35)
|
||||
finally:
|
||||
task.cancel()
|
||||
try:
|
||||
await task
|
||||
except asyncio.CancelledError:
|
||||
pass
|
||||
|
||||
def run_producer() -> None:
|
||||
try:
|
||||
asyncio.run(producer())
|
||||
except Exception as exc: # surfaced via errors after join
|
||||
errors.append(exc)
|
||||
|
||||
thread = threading.Thread(target=run_producer, daemon=True)
|
||||
thread.start()
|
||||
try:
|
||||
sock_path = get_loop_tick_socket_path(tmp_path, pid)
|
||||
deadline = time.monotonic() + 5.0
|
||||
while not sock_path.exists() and time.monotonic() < deadline:
|
||||
time.sleep(0.02)
|
||||
assert sock_path.exists(), "producer never armed the tick socket"
|
||||
|
||||
hb_path = get_loop_heartbeat_path(tmp_path)
|
||||
deadline = time.monotonic() + 5.0
|
||||
while True:
|
||||
try:
|
||||
age = time.time() - hb_path.stat().st_mtime
|
||||
except FileNotFoundError:
|
||||
# The socket is armed before the first write lands; the
|
||||
# file appears a tick later.
|
||||
age = 0.0
|
||||
if age > stale_after:
|
||||
break
|
||||
assert time.monotonic() < deadline, "heartbeat never went stale"
|
||||
time.sleep(0.02)
|
||||
|
||||
# The file is stale but the loop answers: ALIVE, not WEDGED.
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, stale_after=stale_after, tick_timeout=0.25
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_ALIVE
|
||||
)
|
||||
|
||||
# The restart path fed by the REAL probe must drain, not escalate.
|
||||
events = []
|
||||
monkeypatch.setattr(
|
||||
gateway_cli, "get_launchd_label", lambda: "ai.hermes.gateway"
|
||||
)
|
||||
monkeypatch.setattr(gateway_cli, "_launchd_domain", lambda: "gui/501")
|
||||
monkeypatch.setattr(gateway_cli, "_get_restart_drain_timeout", lambda: 180.0)
|
||||
monkeypatch.setattr(
|
||||
"gateway.status.get_running_pid", lambda *a, **k: pid
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli, "_request_gateway_self_restart", lambda pid: False
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli,
|
||||
"_escalate_wedged_gateway",
|
||||
lambda pid, **kw: events.append("escalate") or True,
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli,
|
||||
"terminate_pid",
|
||||
lambda pid, force=False: events.append(("kill" if force else "term", pid)),
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli,
|
||||
"_wait_for_gateway_exit",
|
||||
lambda timeout, force_after=None: events.append(("drain", timeout))
|
||||
or True,
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli.subprocess,
|
||||
"run",
|
||||
lambda *a, **k: events.append("kickstart")
|
||||
or __import__("types").SimpleNamespace(
|
||||
returncode=0, stdout="", stderr=""
|
||||
),
|
||||
)
|
||||
monkeypatch.setattr(
|
||||
gateway_cli, "_clear_launchd_unsupported_marker", lambda: None
|
||||
)
|
||||
# The real probe resolves the state dir through this hook.
|
||||
monkeypatch.setattr(
|
||||
"gateway.shutdown_watchdog._process_hermes_home", lambda: tmp_path
|
||||
)
|
||||
gateway_cli.launchd_restart()
|
||||
assert "escalate" not in events
|
||||
assert ("drain", 180.0) in events
|
||||
finally:
|
||||
thread.join(timeout=5.0)
|
||||
assert not errors, errors
|
||||
|
||||
def test_off_loop_completion_cannot_manufacture_fresh_liveness(self, tmp_path):
|
||||
"""A write landing after the loop froze must not look alive.
|
||||
|
||||
File fresh (a late off-loop write completed) + loop silent: the probe
|
||||
must NOT return ALIVE — the loop itself never answered, so freshness
|
||||
is not a liveness proof. UNKNOWN keeps the safe drain path while
|
||||
denying the false-fresh window the review described.
|
||||
"""
|
||||
pid = 4242
|
||||
_silent_socket_node(get_loop_tick_socket_path(tmp_path, pid))
|
||||
_write_heartbeat(tmp_path, pid, age_s=5.0)
|
||||
_mark_witness_flag(tmp_path, armed=True)
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, tick_timeout=0.2
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_UNKNOWN
|
||||
)
|
||||
|
||||
def test_true_wedge_requires_both_witnesses_to_agree(self, tmp_path):
|
||||
"""Stale file + armed but silent socket: both agree, so WEDGED.
|
||||
|
||||
The destructive path must still exist for genuinely dead loops — a
|
||||
frozen loop stops answering the socket, and the file goes stale.
|
||||
"""
|
||||
pid = 4242
|
||||
_silent_socket_node(get_loop_tick_socket_path(tmp_path, pid))
|
||||
_write_heartbeat(tmp_path, pid, age_s=600.0)
|
||||
_mark_witness_flag(tmp_path, armed=True, age_s=600.0)
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, tick_timeout=0.2
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_WEDGED
|
||||
)
|
||||
|
||||
def test_armed_witness_unreachable_is_unknown(self, tmp_path):
|
||||
"""Producer claims the witness is armed but no node exists: ambiguity.
|
||||
|
||||
Never kill on it — the graceful drain remains the backstop.
|
||||
"""
|
||||
pid = 4242
|
||||
_write_heartbeat(tmp_path, pid, age_s=600.0)
|
||||
_mark_witness_flag(tmp_path, armed=True, age_s=600.0)
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, tick_timeout=0.2
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_UNKNOWN
|
||||
)
|
||||
|
||||
def test_unarmed_witness_disables_stale_escalation(self, tmp_path):
|
||||
"""New producer whose bind failed: staleness is NOT proof.
|
||||
|
||||
The write is off-loop, so the file can age while the loop runs; with
|
||||
no witness, a stale file must never escalate.
|
||||
"""
|
||||
pid = 4242
|
||||
_write_heartbeat(tmp_path, pid, age_s=600.0)
|
||||
_mark_witness_flag(tmp_path, armed=False, age_s=600.0)
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, tick_timeout=0.2
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_UNKNOWN
|
||||
)
|
||||
|
||||
def test_legacy_payload_keeps_single_witness_contract(self, tmp_path):
|
||||
"""No witness flag = on-loop writer: staleness stays proof.
|
||||
|
||||
Old gateways never moved the write off-loop, so their stale file
|
||||
still means a dead loop — the legacy WEDGED verdict is unchanged.
|
||||
"""
|
||||
pid = 4242
|
||||
_write_heartbeat(tmp_path, pid, age_s=600.0)
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(pid, home=tmp_path)
|
||||
== gateway_cli.GATEWAY_LOOP_WEDGED
|
||||
)
|
||||
# And a fresh legacy file stays safe even if a dead-listener node
|
||||
# exists for the PID (leftover from a newer process): the silent
|
||||
# socket denies ALIVE, and UNKNOWN never escalates — the drain path
|
||||
# keeps the full budget either way.
|
||||
_write_heartbeat(tmp_path, pid, age_s=5.0)
|
||||
_silent_socket_node(get_loop_tick_socket_path(tmp_path, pid))
|
||||
assert (
|
||||
gateway_cli.probe_gateway_loop_liveness(
|
||||
pid, home=tmp_path, tick_timeout=0.2
|
||||
)
|
||||
== gateway_cli.GATEWAY_LOOP_UNKNOWN
|
||||
)
|
||||
|
||||
Reference in New Issue
Block a user