fix(discord): liveness probe gains a dispatch-side dimension (#109521)
The Gateway WS health probe sampled only transport state — ready, open, heartbeat-ACK age, latency. A socket that stays ESTAB and keeps ACKing while zero gateway frames arrive (the #109521 "connected-but-deaf" incident) read healthy indefinitely, and the adapter went silent for hours with no log line and no watchdog firing. Two defects fixed: 1. Dispatch-side dimension. `on_socket_raw_receive` now stamps `_last_gateway_frame_at` for every inbound raw gateway frame — heartbeats and ACKs included, so a legitimately quiet server is not flagged. The health check gains an `event_silence` reason with its own bound, `websocket_event_max_silence_seconds` (default 300s, 0 disables the dimension alone). The stamp resets on `on_ready` so a reconnect never inherits pre-restart silence. Trip path is unchanged: consecutive failures -> retryable `discord_websocket_health_stale` -> the existing reconnect watcher builds a fresh adapter. 2. Silent probe disable. `_finite_positive_config_float` / `_config_int` mapped anything `float()` rejects ("15s", "nan", "true") to 0.0 with no log line, permanently disabling the watchdog invisibly. Unparsable and non-positive values now log one WARNING naming the knob and raw value. The new key rides the existing `_YAML_WEBSOCKET_LIVENESS_KEYS` seeding and is documented in the Discord guide's liveness section.
This commit is contained in:
@@ -1068,6 +1068,16 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
self._max_latency_seconds = self._finite_positive_config_float(
|
||||
"websocket_max_latency_seconds", 30.0,
|
||||
)
|
||||
# Dispatch-side dimension (#109521): ready/open/ACK/latency prove the transport, not
|
||||
# that events are still being DISPATCHED. Last raw gateway frame (any frame —
|
||||
# heartbeats, ACKs, presence — so a legitimately quiet server is not "dead") gives
|
||||
# the probe an event-age bound; a socket that stays ESTAB and keeps ACKing while
|
||||
# zero frames arrive is the connected-but-deaf fingerprint the old probe read healthy.
|
||||
self._event_max_silence_seconds = self._finite_positive_config_float(
|
||||
"websocket_event_max_silence_seconds", 300.0,
|
||||
)
|
||||
# perf_counter clock (same as _read_websocket_health's ack math); 0 = never saw a frame.
|
||||
self._last_gateway_frame_at: float = 0.0
|
||||
self._liveness_task: Optional[asyncio.Task] = None
|
||||
self._liveness_notification_task: Optional[asyncio.Task] = None
|
||||
# True while disconnect() intentionally closes discord.py (done callback: shutdown vs crash).
|
||||
@@ -1102,24 +1112,45 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
value = _scoped_gate_env(env_key) or None
|
||||
return default if value is None or value == "" else value
|
||||
|
||||
def _warn_liveness_config_disabled(self, key: str, raw: Any) -> None:
|
||||
"""One-shot warning when a liveness knob resolves to a disabling value (#109521).
|
||||
|
||||
Unparsable config (`"15s"`, `nan`, `true`) silently mapped to 0 and turned the whole
|
||||
watchdog off with no log line — indistinguishable from "the watchdog missed it".
|
||||
Loud beats silent: the operator's mitigation (cron-restart on no-`[Discord]`-lines)
|
||||
exists only because nothing in-process ever told them the probe was off.
|
||||
"""
|
||||
logger.warning(
|
||||
"[%s] Discord liveness knob %s=%r is not a usable positive number; "
|
||||
"the websocket liveness probe may be disabled by this value",
|
||||
self.name, key, raw,
|
||||
)
|
||||
|
||||
def _finite_positive_config_float(
|
||||
self, key: str, default: float, *, env_key: Optional[str] = None
|
||||
) -> float:
|
||||
"""Resolve a finite positive liveness duration; invalid values disable it."""
|
||||
"""Resolve a finite positive liveness duration; invalid values disable it (with a warning)."""
|
||||
raw = self._config_value(key, default, env_key=env_key)
|
||||
try:
|
||||
value = float(self._config_value(key, default, env_key=env_key))
|
||||
value = float(raw)
|
||||
except (TypeError, ValueError):
|
||||
self._warn_liveness_config_disabled(key, raw)
|
||||
return 0.0
|
||||
return value if math.isfinite(value) and value > 0 else 0.0
|
||||
if math.isfinite(value) and value > 0:
|
||||
return value
|
||||
self._warn_liveness_config_disabled(key, raw)
|
||||
return 0.0
|
||||
|
||||
def _config_int(self, key: str, default: int, *, env_key: Optional[str] = None) -> int:
|
||||
"""Resolve a positive liveness count; invalid values disable it."""
|
||||
value = self._config_value(key, default, env_key=env_key)
|
||||
if isinstance(value, bool):
|
||||
"""Resolve a positive liveness count; invalid values disable it (with a warning)."""
|
||||
raw = self._config_value(key, default, env_key=env_key)
|
||||
if isinstance(raw, bool):
|
||||
self._warn_liveness_config_disabled(key, raw)
|
||||
return 0
|
||||
try:
|
||||
return int(value)
|
||||
return int(raw)
|
||||
except (TypeError, ValueError):
|
||||
self._warn_liveness_config_disabled(key, raw)
|
||||
return 0
|
||||
|
||||
def _handle_bot_task_done(self, task: asyncio.Task) -> None:
|
||||
@@ -1227,6 +1258,9 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
logger.info("[%s] Connected as %s", adapter_self.name, adapter_self._client.user)
|
||||
await adapter_self._resolve_allowed_usernames()
|
||||
adapter_self._ready_event.set()
|
||||
# Fresh connection => no silence history: reset the dispatch-side stamp so a
|
||||
# reconnect never inherits the pre-restart silence (#109521).
|
||||
adapter_self._last_gateway_frame_at = time.perf_counter()
|
||||
if adapter_self._post_connect_task and not adapter_self._post_connect_task.done():
|
||||
adapter_self._post_connect_task.cancel()
|
||||
adapter_self._post_connect_task = asyncio.create_task(
|
||||
@@ -1235,6 +1269,17 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
if adapter_self._missed_message_backfill_enabled():
|
||||
adapter_self._ensure_missed_message_backfill_task()
|
||||
|
||||
@self._client.event
|
||||
async def on_socket_raw_receive(msg: str):
|
||||
"""Stamp every inbound gateway frame for the liveness probe's event-age check.
|
||||
|
||||
Fires for ALL frames — heartbeats, ACKs, presence — not just messages, so a
|
||||
legitimately quiet server is not "dead", but a socket that stays ESTAB while
|
||||
zero frames arrive (connected-but-deaf, #109521 incident 2) becomes visible.
|
||||
Intentionally bare: parsing happens in discord.py; this hook only clocks.
|
||||
"""
|
||||
adapter_self._last_gateway_frame_at = time.perf_counter()
|
||||
|
||||
@self._client.event
|
||||
async def on_message(message: DiscordMessage):
|
||||
await adapter_self._dispatch_discord_message(message)
|
||||
@@ -1590,6 +1635,7 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
or self._liveness_failure_threshold <= 0
|
||||
or self._heartbeat_ack_max_age_seconds <= 0
|
||||
or self._max_latency_seconds <= 0
|
||||
or self._event_max_silence_seconds <= 0
|
||||
):
|
||||
return
|
||||
if self._liveness_task and not self._liveness_task.done():
|
||||
@@ -1629,6 +1675,15 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
|
||||
return False, "latency_non_finite"
|
||||
if latency > self._max_latency_seconds:
|
||||
return False, "latency_exceeded"
|
||||
# Dispatch-side dimension (#109521): every check above inspects the transport, and a
|
||||
# socket that is ESTAB and still ACKing can deliver ZERO events for hours while
|
||||
# reading healthy. The last raw gateway frame (any frame, including heartbeats) is
|
||||
# the one signal that events are actually arriving. Never-seen-frame counts as
|
||||
# silent — on_ready fires long before any gap could be legitimate.
|
||||
if self._event_max_silence_seconds > 0:
|
||||
frame_age = time.perf_counter() - self._last_gateway_frame_at
|
||||
if self._last_gateway_frame_at <= 0 or frame_age > self._event_max_silence_seconds:
|
||||
return False, "event_silence"
|
||||
return True, "healthy"
|
||||
|
||||
async def _liveness_loop(self) -> None:
|
||||
@@ -6923,6 +6978,7 @@ _YAML_WEBSOCKET_LIVENESS_KEYS = (
|
||||
("websocket_liveness_failure_threshold", "liveness_failure_threshold", "HERMES_DISCORD_LIVENESS_FAILURE_THRESHOLD"),
|
||||
("websocket_heartbeat_ack_max_age_seconds", None, None),
|
||||
("websocket_max_latency_seconds", None, None),
|
||||
("websocket_event_max_silence_seconds", None, None),
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -95,6 +95,7 @@ def _make_adapter(
|
||||
threshold=1,
|
||||
max_ack_age=1.0,
|
||||
max_latency=1.0,
|
||||
max_silence=300.0,
|
||||
) -> DiscordAdapter:
|
||||
monkeypatch.setenv("HERMES_DISCORD_LIVENESS_INTERVAL_SECONDS", str(interval))
|
||||
monkeypatch.setenv("HERMES_DISCORD_LIVENESS_FAILURE_THRESHOLD", str(threshold))
|
||||
@@ -105,6 +106,7 @@ def _make_adapter(
|
||||
extra={
|
||||
"websocket_heartbeat_ack_max_age_seconds": max_ack_age,
|
||||
"websocket_max_latency_seconds": max_latency,
|
||||
"websocket_event_max_silence_seconds": max_silence,
|
||||
},
|
||||
)
|
||||
)
|
||||
@@ -122,6 +124,7 @@ class _BrokenWebSocket:
|
||||
("websocket_liveness_interval_seconds", "_liveness_interval_seconds", "nan"),
|
||||
("websocket_heartbeat_ack_max_age_seconds", "_heartbeat_ack_max_age_seconds", "inf"),
|
||||
("websocket_max_latency_seconds", "_max_latency_seconds", "-inf"),
|
||||
("websocket_event_max_silence_seconds", "_event_max_silence_seconds", "15s"),
|
||||
],
|
||||
)
|
||||
def test_nonfinite_liveness_config_disables_that_probe_dimension(monkeypatch, key, attribute, raw):
|
||||
@@ -132,6 +135,26 @@ def test_nonfinite_liveness_config_disables_that_probe_dimension(monkeypatch, ke
|
||||
assert getattr(adapter, attribute) == 0.0
|
||||
|
||||
|
||||
def test_unparsable_liveness_config_warns_instead_of_disabling_silently(monkeypatch, caplog):
|
||||
"""A knob value that can't parse must not disable the probe without a trace (#109521).
|
||||
|
||||
Pre-fix, ``websocket_liveness_interval_seconds: 15s`` mapped to 0.0 with no log line —
|
||||
the watchdog was off and the only visible symptom was hours of Discord silence.
|
||||
"""
|
||||
with caplog.at_level("WARNING", logger="plugins.platforms.discord.adapter"):
|
||||
adapter = DiscordAdapter(
|
||||
PlatformConfig(
|
||||
enabled=True,
|
||||
token="test-token",
|
||||
extra={"websocket_liveness_interval_seconds": "15s"},
|
||||
)
|
||||
)
|
||||
|
||||
assert adapter._liveness_interval_seconds == 0.0
|
||||
assert "websocket_liveness_interval_seconds" in caplog.text
|
||||
assert "15s" in caplog.text
|
||||
|
||||
|
||||
def test_default_liveness_bounds_trigger_timed_recovery(monkeypatch):
|
||||
for key in (
|
||||
"HERMES_DISCORD_LIVENESS_INTERVAL_SECONDS",
|
||||
@@ -145,6 +168,7 @@ def test_default_liveness_bounds_trigger_timed_recovery(monkeypatch):
|
||||
assert adapter._liveness_failure_threshold == 2
|
||||
assert adapter._heartbeat_ack_max_age_seconds == 60.0
|
||||
assert adapter._max_latency_seconds == 30.0
|
||||
assert adapter._event_max_silence_seconds == 300.0
|
||||
|
||||
|
||||
def test_platform_config_extra_overrides_process_liveness_bridge(monkeypatch):
|
||||
@@ -299,3 +323,170 @@ async def test_disconnect_cancels_liveness_task(monkeypatch):
|
||||
await adapter.disconnect()
|
||||
assert task.done()
|
||||
assert adapter._liveness_task is None
|
||||
|
||||
|
||||
class _DeafHealthBot(_LiveBot):
|
||||
"""Incident-2 fingerprint from #109521: ESTAB socket, keep-alive ACKs, zero frames."""
|
||||
|
||||
def __init__(self, **kwargs):
|
||||
super().__init__(**kwargs)
|
||||
self.raw_receive_events = []
|
||||
|
||||
async def deliver_raw_frame(self, payload: str = "{}"):
|
||||
"""Feed a raw gateway frame through the adapter's registered hook."""
|
||||
await self._events["on_socket_raw_receive"](payload)
|
||||
|
||||
|
||||
def _connect_deaf_bot(adapter, monkeypatch, bot_box):
|
||||
def factory(**kwargs):
|
||||
bot = _DeafHealthBot(
|
||||
intents=kwargs["intents"],
|
||||
allowed_mentions=kwargs.get("allowed_mentions"),
|
||||
)
|
||||
bot.fetch_user = AsyncMock()
|
||||
bot_box.append(bot)
|
||||
return bot
|
||||
|
||||
return _connect(adapter, monkeypatch, factory)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_connected_but_deaf_socket_is_unhealthy(monkeypatch):
|
||||
"""#109521 incident 2: a socket that stays open, ready, low-latency, and ACKing —
|
||||
but through which no gateway frame has arrived since before the silence bound —
|
||||
must read unhealthy, not healthy."""
|
||||
adapter = _make_adapter(
|
||||
monkeypatch, interval=60, threshold=1, max_ack_age=60.0, max_latency=30.0,
|
||||
max_silence=10.0,
|
||||
)
|
||||
box = []
|
||||
await _connect_deaf_bot(adapter, monkeypatch, box)
|
||||
bot = box[0]
|
||||
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
|
||||
# on_ready reset the stamp at connect; age it past the bound so the adapter is
|
||||
# "connected, all transport dimensions green, and silent" — incident 2's fingerprint.
|
||||
adapter._last_gateway_frame_at = time.perf_counter() - 60.0
|
||||
|
||||
healthy, reason = adapter._read_websocket_health(bot)
|
||||
assert healthy is False
|
||||
assert reason == "event_silence"
|
||||
|
||||
await adapter.disconnect()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_recent_raw_frame_keeps_deaf_fingerprint_socket_healthy(monkeypatch):
|
||||
"""Any raw gateway frame (heartbeat, ACK — not just messages) refreshes the dispatch
|
||||
clock: a legitimately quiet server must not be flagged by the event-age bound."""
|
||||
adapter = _make_adapter(
|
||||
monkeypatch, interval=60, threshold=1, max_ack_age=60.0, max_latency=30.0,
|
||||
)
|
||||
box = []
|
||||
await _connect_deaf_bot(adapter, monkeypatch, box)
|
||||
bot = box[0]
|
||||
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
|
||||
# on_ready resets the stamp (fresh connection); a non-message frame then advances it.
|
||||
assert adapter._last_gateway_frame_at > 0
|
||||
await bot.deliver_raw_frame('{"t":null,"op":11}')
|
||||
frame_at = adapter._last_gateway_frame_at
|
||||
assert frame_at > 0
|
||||
|
||||
healthy, reason = adapter._read_websocket_health(bot)
|
||||
assert (healthy, reason) == (True, "healthy")
|
||||
|
||||
await adapter.disconnect()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_stale_frame_over_silence_bound_reads_unhealthy(monkeypatch):
|
||||
"""A frame that arrived, then a silence longer than the bound, while transport
|
||||
dimensions still look healthy — the exact connected-but-deaf progression."""
|
||||
adapter = _make_adapter(
|
||||
monkeypatch, interval=60, threshold=1, max_ack_age=60.0, max_latency=30.0,
|
||||
max_silence=10.0,
|
||||
)
|
||||
box = []
|
||||
await _connect_deaf_bot(adapter, monkeypatch, box)
|
||||
bot = box[0]
|
||||
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
|
||||
await bot.deliver_raw_frame('{"t":null,"op":11}')
|
||||
# Simulate the silence: age the last frame past every bound, keeping the transport green.
|
||||
adapter._last_gateway_frame_at = time.perf_counter() - 60.0
|
||||
|
||||
healthy, reason = adapter._read_websocket_health(bot)
|
||||
assert (healthy, reason) == (False, "event_silence")
|
||||
|
||||
await adapter.disconnect()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_event_silence_drives_liveness_loop_to_retryable_fatal(monkeypatch, caplog):
|
||||
"""End to end: a persistent deaf socket trips the probe's failure threshold and
|
||||
surfaces as the retryable fatal the reconnect watcher already handles."""
|
||||
import logging
|
||||
|
||||
caplog.at_level(logging.WARNING, logger="plugins.platforms.discord.adapter")
|
||||
adapter = _make_adapter(
|
||||
monkeypatch, interval=0.01, threshold=2, max_ack_age=60.0, max_latency=30.0,
|
||||
max_silence=5.0,
|
||||
)
|
||||
handler = AsyncMock()
|
||||
adapter.set_fatal_error_handler(handler)
|
||||
box = []
|
||||
await _connect_deaf_bot(adapter, monkeypatch, box)
|
||||
bot = box[0]
|
||||
|
||||
# Transport dimensions stay green the whole time; frames never arrive.
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
|
||||
adapter._last_gateway_frame_at = time.perf_counter() - 60.0
|
||||
|
||||
await _wait_until(
|
||||
lambda: adapter.fatal_error_code == "discord_websocket_health_stale",
|
||||
"deaf socket never tripped the liveness probe",
|
||||
)
|
||||
assert "event_silence" in caplog.text
|
||||
# The runner-facing handler fires from the notification task (post-close); give it a beat.
|
||||
await _wait_until(
|
||||
lambda: handler.await_count >= 1,
|
||||
"fatal notification never reached the runner handler",
|
||||
)
|
||||
|
||||
await adapter.disconnect()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_silence_knob_zero_disables_event_dimension_only(monkeypatch):
|
||||
"""Setting ``websocket_event_max_silence_seconds: 0`` must opt out of the dispatch
|
||||
check alone; the transport dimensions (ack age, latency) still work."""
|
||||
adapter = DiscordAdapter(
|
||||
PlatformConfig(
|
||||
enabled=True,
|
||||
token="test-token",
|
||||
extra={
|
||||
"websocket_event_max_silence_seconds": 0,
|
||||
"websocket_heartbeat_ack_max_age_seconds": 60.0,
|
||||
"websocket_max_latency_seconds": 30.0,
|
||||
},
|
||||
)
|
||||
)
|
||||
assert adapter._event_max_silence_seconds == 0.0
|
||||
|
||||
box = []
|
||||
await _connect_deaf_bot(adapter, monkeypatch, box)
|
||||
bot = box[0]
|
||||
# Never delivered a frame and the stamp is at its connect-time value.
|
||||
adapter._last_gateway_frame_at = 0.0
|
||||
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
|
||||
healthy, reason = adapter._read_websocket_health(bot)
|
||||
assert (healthy, reason) == (True, "healthy")
|
||||
|
||||
# While the transport dimensions still catch their own failure shape:
|
||||
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=999.0)
|
||||
healthy, reason = adapter._read_websocket_health(bot)
|
||||
assert (healthy, reason) == (False, "ack_stale")
|
||||
|
||||
await adapter.disconnect()
|
||||
|
||||
@@ -84,7 +84,7 @@ This guide walks you through the full setup process — from creating your bot o
|
||||
|
||||
### Gateway WebSocket health
|
||||
|
||||
Discord REST and the Gateway WebSocket are separate transports. A successful REST response (including `fetch_user()` returning HTTP 200) does not prove that the bot can still receive Gateway events. Hermes therefore combines the ready state, client/socket closure state, socket openness, heartbeat ACK age, and finite heartbeat latency.
|
||||
Discord REST and the Gateway WebSocket are separate transports. A successful REST response (including `fetch_user()` returning HTTP 200) does not prove that the bot can still receive Gateway events. Hermes therefore combines the ready state, client/socket closure state, socket openness, heartbeat ACK age, finite heartbeat latency, and **frame recency** — the time since the last raw Gateway frame of any kind (heartbeats and ACKs count, so a legitimately quiet server is not flagged).
|
||||
|
||||
After the configured number of consecutive unhealthy samples, the adapter emits one retryable fatal event. The existing gateway reconnect watcher creates a fresh adapter; the Discord adapter does not start a second unbounded reconnect loop.
|
||||
|
||||
@@ -96,10 +96,13 @@ discord:
|
||||
websocket_liveness_failure_threshold: 2
|
||||
websocket_heartbeat_ack_max_age_seconds: 60
|
||||
websocket_max_latency_seconds: 30
|
||||
websocket_event_max_silence_seconds: 300
|
||||
```
|
||||
|
||||
The old `liveness_interval_seconds` and `liveness_failure_threshold` names remain compatibility aliases only; they no longer mean REST probing.
|
||||
|
||||
Values that fail to parse as a positive number (e.g. `15s`, `nan`, `true`) disable the corresponding knob and log one warning at adapter startup — check `gateway.log` if the probe seems inactive.
|
||||
|
||||
## Step 1: Create a Discord Application
|
||||
|
||||
1. Go to the [Discord Developer Portal](https://discord.com/developers/applications) and sign in with your Discord account.
|
||||
|
||||
Reference in New Issue
Block a user