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:
rodrigo
2026-08-20 03:29:49 -03:00
committed by kshitij
parent f39931afd9
commit a1c83ef901
4 changed files with 526 additions and 31 deletions

View File

@@ -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

View File

@@ -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(

View File

@@ -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"
)

View File

@@ -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
)