diff --git a/plugins/platforms/discord/adapter.py b/plugins/platforms/discord/adapter.py index 37496746f4..6437a21015 100644 --- a/plugins/platforms/discord/adapter.py +++ b/plugins/platforms/discord/adapter.py @@ -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), ) diff --git a/tests/gateway/test_discord_liveness.py b/tests/gateway/test_discord_liveness.py index 4cd87c6ddb..e91e72352a 100644 --- a/tests/gateway/test_discord_liveness.py +++ b/tests/gateway/test_discord_liveness.py @@ -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() diff --git a/website/docs/user-guide/messaging/discord.md b/website/docs/user-guide/messaging/discord.md index dfad0cc25b..240bb41193 100644 --- a/website/docs/user-guide/messaging/discord.md +++ b/website/docs/user-guide/messaging/discord.md @@ -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.