fix(discord): dispatch-side liveness dimension detects an ACKing-but-deaf gateway socket (#109521)

Incident 2 of #109521: a Gateway socket can stay ESTABLISHED and keep
ACKing heartbeats while zero DISPATCH events are parsed, so every
transport-side liveness sample (ready/open/ack-age/latency) reads
healthy for hours. The merged #109963 deliberately dropped the
event_silence dimension: a raw-frame stamp is debug-gated
(on_socket_raw_receive needs enable_debug_events) and, since heartbeat
ACKs are frames, ack_stale always fires first by construction.

This adds the dispatch-side signal that was requested instead:

- stamp on on_socket_event_type, which discord.py 2.7.1 dispatches for
  every parsed DISPATCH frame with no debug gate (verified live against
  the real received_message path: 4/4 frames fired with
  enable_debug_events=False, on_socket_raw_receive 0/4)
- new knob websocket_event_max_silence_seconds (default 4h, the
  incident report's field-proven operator bound); 0 opts out of this
  dimension ONLY — the #109782 review failure put the knob in
  _start_liveness_probe's all-or-nothing guard, killing the whole
  watchdog; it is gated strictly inside _read_websocket_health here
- the stamp resets per connection (connect() clears it), and a None
  stamp (no event parsed yet on this connection) is not silence
- docs (en + zh-Hans) cover the new knob and the per-dimension opt-out

Fixes #109521

(cherry picked from commit b4baa97fc45794209711a45e052111d7d44d5f90)
This commit is contained in:
salch-cred
2026-09-14 21:52:31 +05:30
committed by kshitij
parent 73f808e47f
commit 101861c7f7
4 changed files with 375 additions and 1 deletions

View File

@@ -1071,10 +1071,28 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
self._max_latency_seconds = self._finite_positive_config_float(
"websocket_max_latency_seconds", 30.0,
)
# Dispatch-side liveness (#109521 incident 2): an ESTAB socket can keep ACKing
# heartbeats (op 11, no event type) while zero DISPATCH events are parsed, so every
# transport-side sample reads healthy for hours. ``socket_event_type`` fires for every
# parsed DISPATCH frame and is NOT gated behind ``enable_debug_events`` (unlike
# ``on_socket_raw_receive`` — verified against discord.py 2.7.1 ``gateway.py``:
# ``received_message`` calls ``self._dispatch('socket_event_type', event)`` before the
# op-code switch). 0 disables this dimension alone; ack-age/latency still guard.
# Default 4h mirrors the field-proven operator bound from the incident report; quiet
# guilds can go hours without a single DISPATCH event (typing/reaction/presence), so
# a short bound would force reconnect loops on healthy-but-idle installs (#109782).
self._event_max_silence_seconds = self._finite_positive_config_float(
"websocket_event_max_silence_seconds", 14400.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).
self._disconnecting = False
# Last DISPATCH frame's monotonic stamp (#109521): ticked by ``on_socket_event_type``
# (fires for every parsed DISPATCH event, not debug-gated) and read by the liveness
# probe's ``event_silence`` dimension. ``None`` means "no event yet on this connection"
# and is not treated as silence (on_ready often arrives in bursts).
self._last_dispatched_event_monotonic: Optional[float] = None
self._missed_message_backfill_task: Optional[asyncio.Task] = None
from hermes_constants import get_hermes_home
from plugins.platforms.discord.recovery import DiscordRecoveryStore
@@ -1241,6 +1259,10 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
allowed_mentions=_build_allowed_mentions(getattr(self.config, "extra", None)),
**proxy_kwargs_for_bot(proxy_url),
)
# Fresh connection, fresh dispatch-side silence window: the previous client's last
# DISPATCH stamp must not leak into this connection's liveness samples (#109521).
# READY itself is a DISPATCH event, so a healthy connection stamps almost immediately.
self._last_dispatched_event_monotonic = None
adapter_self = self # capture for closure
@self._client.event
@@ -1256,6 +1278,15 @@ 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_event_type(event_type: str):
# Dispatch-side liveness stamp (#109521): discord.py dispatches this for every
# parsed DISPATCH frame on every connection — no ``enable_debug_events`` needed
# (unlike ``on_socket_raw_receive``). Heartbeat ACKs (op 11) return before the
# dispatch, so an ACKing-but-deaf socket leaves this stamp frozen while every
# transport-side check reads healthy.
adapter_self._last_dispatched_event_monotonic = time.perf_counter()
@self._client.event
async def on_message(message: DiscordMessage):
await adapter_self._dispatch_discord_message(message)
@@ -1650,6 +1681,18 @@ class DiscordAdapter(DiscordMediaMixin, BasePlatformAdapter):
return False, "latency_non_finite"
if latency > self._max_latency_seconds:
return False, "latency_exceeded"
# Dispatch-side dimension (#109521 incident 2): transport-green + event-starved is the
# connected-but-deaf fingerprint. Gated HERE only — never in _start_liveness_probe — so an
# explicit 0 disables this dimension alone and ack-age/latency keep guarding (the #109782
# regression put the knob in the probe's all-or-nothing startup guard, killing the whole
# watchdog). ``None`` = no DISPATCH event yet on this connection: not silence (the
# not_ready check above still covers the pre-ready window).
if self._event_max_silence_seconds > 0:
last_event = self._last_dispatched_event_monotonic
if last_event is not None:
event_silence = time.perf_counter() - last_event
if not math.isfinite(event_silence) or event_silence > self._event_max_silence_seconds:
return False, "event_silence"
return True, "healthy"
async def _liveness_loop(self) -> None:
@@ -6942,6 +6985,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),
)

View File

@@ -0,0 +1,324 @@
"""Dispatch-side liveness for the Discord adapter (#109521 incident 2).
An ESTAB Gateway socket can keep ACKing heartbeats while zero DISPATCH
events are parsed — every transport-side sample reads healthy while the
adapter is deaf. The probe's ``event_silence`` dimension stamps
``on_socket_event_type`` (dispatched for every parsed DISPATCH frame,
unlike the debug-gated ``on_socket_raw_receive``) and trips after
``websocket_event_max_silence_seconds`` without one.
"""
from __future__ import annotations
import asyncio
import time
from types import SimpleNamespace
from unittest.mock import AsyncMock
import pytest
from tests.gateway.test_discord_connect import ( # noqa: E402
FakeBot,
_ensure_discord_mock,
)
_ensure_discord_mock()
import plugins.platforms.discord.adapter as discord_platform # noqa: E402
from gateway.config import PlatformConfig # noqa: E402
from plugins.platforms.discord.adapter import DiscordAdapter # noqa: E402
from tests.gateway.test_discord_liveness import ( # noqa: E402
_FakeKeepAlive,
_FakeWebSocket,
_LiveBot,
_set_websocket_health,
_wait_until,
)
class _DispatchingBot(_LiveBot):
"""A live bot that can deliver parsed DISPATCH events like the real gateway.
Real discord.py ``received_message`` parses each frame, calls
``self._dispatch('socket_event_type', event)`` for every DISPATCH op,
and returns early on heartbeat ACKs (op 11) — so a socket that only
ACKs never moves the stamp. ``deliver_dispatch`` models a parsed event
reaching ``Client.dispatch``.
"""
async def deliver_dispatch(self, event_type: str = "MESSAGE_CREATE") -> None:
handler = self._events.get("on_socket_event_type")
if handler is None:
raise AssertionError("adapter did not register on_socket_event_type")
# Client.dispatch schedules the handler as a task; awaiting it inline
# is equivalent for this trivially non-blocking handler and lets the
# caller observe the stamp immediately.
await handler(event_type)
def _make_adapter(
monkeypatch,
*,
interval: float = 0.01,
threshold: int = 1,
max_ack_age: float = 60.0,
max_latency: float = 30.0,
max_event_silence: float = 14400.0,
) -> DiscordAdapter:
monkeypatch.setenv("HERMES_DISCORD_LIVENESS_INTERVAL_SECONDS", str(interval))
monkeypatch.setenv("HERMES_DISCORD_LIVENESS_FAILURE_THRESHOLD", str(threshold))
return DiscordAdapter(
PlatformConfig(
enabled=True,
token="test-token",
extra={
"websocket_heartbeat_ack_max_age_seconds": max_ack_age,
"websocket_max_latency_seconds": max_latency,
"websocket_event_max_silence_seconds": max_event_silence,
},
)
)
async def _connect(adapter: DiscordAdapter, monkeypatch, bot_factory) -> FakeBot:
monkeypatch.setattr(
"gateway.status.acquire_scoped_lock",
lambda scope, identity, metadata=None: (True, None),
)
monkeypatch.setattr("gateway.status.release_scoped_lock", lambda scope, identity: None)
intents = SimpleNamespace(
message_content=False, dm_messages=False, guild_messages=False,
members=False, voice_states=False,
)
monkeypatch.setattr(discord_platform.Intents, "default", lambda: intents)
monkeypatch.setattr(discord_platform.commands, "Bot", bot_factory)
monkeypatch.setattr(adapter, "_resolve_allowed_usernames", AsyncMock())
assert await adapter.connect() is True
return adapter._client
def _transport_healthy(bot: _DispatchingBot) -> None:
"""Make every transport-side sample read healthy (incident 2's fingerprint)."""
_set_websocket_health(bot, ready=True, socket_open=True, latency=0.05, ack_age=0.0)
@pytest.mark.asyncio
async def test_deaf_socket_trips_event_silence_dimension(monkeypatch):
"""Incident 2 e2e: transport-green + event-starved must trip the probe.
The probe samples through ``_liveness_loop``, the real dispatch surface
(no direct ``_read_websocket_health`` call), so this also proves the
dimension is gated inside the health check and not in the startup guard.
"""
adapter = _make_adapter(monkeypatch, interval=0.01, threshold=2, max_event_silence=0.05)
handler = AsyncMock()
adapter.set_fatal_error_handler(handler)
def factory(**kwargs):
bot = _DispatchingBot(intents=kwargs["intents"], allowed_mentions=kwargs.get("allowed_mentions"))
bot.fetch_user = AsyncMock()
return bot
bot = await _connect(adapter, monkeypatch, factory)
_transport_healthy(bot)
# One early event arms the stamp; then total DISPATCH silence while every
# transport sample stays green.
await bot.deliver_dispatch("READY")
async def handler_awaited() -> None:
# _liveness_loop sets the fatal code, then hands off to
# _notify_liveness_fatal_error (a separate task that closes the client
# — up to a 1s budget — before notifying the runner), so wait for the
# handler itself, not just the code.
while True:
if handler.await_count:
return
code = getattr(adapter, "_fatal_error_code", None)
if code and adapter._liveness_notification_task is not None:
with_context = adapter._liveness_notification_task
try:
await asyncio.wait_for(asyncio.shield(with_context), timeout=3.0)
except (asyncio.TimeoutError, asyncio.CancelledError, Exception):
pass
if handler.await_count:
return
await asyncio.sleep(0.01)
await asyncio.wait_for(handler_awaited(), timeout=8.0)
assert adapter._fatal_error_code == "discord_websocket_health_stale"
assert "event_silence" in (adapter._fatal_error_message or "")
assert adapter._fatal_error_retryable is True
handler.assert_awaited_once()
@pytest.mark.asyncio
async def test_fresh_dispatch_events_keep_deaf_socket_probe_healthy(monkeypatch):
"""The converse: regular DISPATCH events must never trip event_silence."""
adapter = _make_adapter(monkeypatch, interval=0.01, threshold=1, max_event_silence=0.2)
def factory(**kwargs):
bot = _DispatchingBot(intents=kwargs["intents"], allowed_mentions=kwargs.get("allowed_mentions"))
bot.fetch_user = AsyncMock()
return bot
bot = await _connect(adapter, monkeypatch, factory)
_transport_healthy(bot)
await bot.deliver_dispatch("READY")
deadline = asyncio.get_running_loop().time() + 0.6
while asyncio.get_running_loop().time() < deadline:
await bot.deliver_dispatch("TYPING_START")
await asyncio.sleep(0.05)
assert getattr(adapter, "_fatal_error_code", None) is None
assert adapter._liveness_task is not None and not adapter._liveness_task.done()
assert adapter._running is True
await adapter.disconnect()
@pytest.mark.asyncio
async def test_zero_event_silence_knob_disables_dimension_not_probe(monkeypatch):
"""``websocket_event_max_silence_seconds: 0`` opts out of this dimension only.
Regression for the #109782 review failure: the rejected PR put the knob
in ``_start_liveness_probe``'s all-or-nothing guard, so ``0`` disabled
the whole watchdog (ack/latency stopped guarding too). Here the probe
must still run and still trip ``ack_stale``.
"""
adapter = _make_adapter(monkeypatch, interval=0.01, threshold=1, max_event_silence=0)
assert adapter._event_max_silence_seconds == 0.0
def factory(**kwargs):
bot = _DispatchingBot(intents=kwargs["intents"], allowed_mentions=kwargs.get("allowed_mentions"))
bot.fetch_user = AsyncMock()
return bot
bot = await _connect(adapter, monkeypatch, factory)
# Transport green, but the ACK clock is stale past max_ack_age (60s default
# scaled to the test's 60.0): event silence would ALSO be past bound, yet
# must not be the reason the probe trips.
bot._gateway_ready = True
bot.latency = 0.05
bot.ws = _FakeWebSocket(open=True, ack_age=10_000.0)
# No dispatch events at all — the dimension is opted out.
assert adapter._last_dispatched_event_monotonic is None
assert adapter._liveness_task is not None, "0 event-silence must not stop the probe task"
async def fatal_code() -> str | None:
while True:
code = getattr(adapter, "_fatal_error_code", None)
if code:
return code
await asyncio.sleep(0.01)
code = await asyncio.wait_for(fatal_code(), timeout=3.0)
assert code == "discord_websocket_health_stale"
assert "ack_stale" in (adapter._fatal_error_message or "")
@pytest.mark.asyncio
async def test_missing_stamp_before_first_event_is_not_silence(monkeypatch):
"""A connected client with no DISPATCH event yet must not read as deaf.
``None`` means "nothing parsed on this connection" — the pre-READY
window; ``not_ready`` owns that failure shape. Treating ``None`` as
silence would false-trip fresh reconnects on quiet guilds.
"""
adapter = _make_adapter(monkeypatch, interval=0.01, threshold=1, max_event_silence=0.05)
def factory(**kwargs):
bot = _DispatchingBot(intents=kwargs["intents"], allowed_mentions=kwargs.get("allowed_mentions"))
bot.fetch_user = AsyncMock()
return bot
bot = await _connect(adapter, monkeypatch, factory)
_transport_healthy(bot)
assert adapter._last_dispatched_event_monotonic is None
deadline = asyncio.get_running_loop().time() + 0.4
while asyncio.get_running_loop().time() < deadline:
healthy, reason = adapter._read_websocket_health(bot)
assert healthy is True, f"None-stamp window must read healthy, got {reason}"
await asyncio.sleep(0.05)
await adapter.disconnect()
@pytest.mark.asyncio
async def test_reconnect_resets_dispatch_stamp(monkeypatch):
"""A fresh client must start a fresh silence window.
Without the reset in ``connect()``, a reconnecting adapter inherits the
previous connection's last-event stamp; a new connection that has not
parsed anything yet would read as instantly past-silence.
"""
adapter = _make_adapter(monkeypatch, interval=0.01, threshold=1, max_event_silence=0.1)
handler = AsyncMock()
adapter.set_fatal_error_handler(handler)
def factory(**kwargs):
bot = _DispatchingBot(intents=kwargs["intents"], allowed_mentions=kwargs.get("allowed_mentions"))
bot.fetch_user = AsyncMock()
return bot
bot = await _connect(adapter, monkeypatch, factory)
_transport_healthy(bot)
await bot.deliver_dispatch("READY")
assert adapter._last_dispatched_event_monotonic is not None
# Second connect() builds a new client: the old stamp must be dropped.
await _connect(adapter, monkeypatch, factory)
assert adapter._last_dispatched_event_monotonic is None
await adapter.disconnect()
def test_unusable_event_silence_value_warns_and_disables_dimension(caplog):
"""``websocket_event_max_silence_seconds: "15s"`` warns and opts out of the
dimension only — consistent with the other knobs' warning contract."""
with caplog.at_level("WARNING", logger="plugins.platforms.discord.adapter"):
adapter = DiscordAdapter(
PlatformConfig(
enabled=True,
token="test-token",
extra={"websocket_event_max_silence_seconds": "15s"},
)
)
assert adapter._event_max_silence_seconds == 0.0
warned = [r.getMessage() for r in caplog.records if "liveness knob" in r.getMessage()]
assert any("websocket_event_max_silence_seconds='15s'" in w for w in warned)
def test_default_event_silence_bound(monkeypatch):
"""Default 4h mirrors the operator-proven bound; explicit config wins."""
for key in (
"HERMES_DISCORD_LIVENESS_INTERVAL_SECONDS",
"HERMES_DISCORD_LIVENESS_FAILURE_THRESHOLD",
):
monkeypatch.delenv(key, raising=False)
adapter = DiscordAdapter(PlatformConfig(enabled=True, token="test-token"))
assert adapter._event_max_silence_seconds == 14400.0
tuned = DiscordAdapter(
PlatformConfig(
enabled=True,
token="test-token",
extra={"websocket_event_max_silence_seconds": 600},
)
)
assert tuned._event_max_silence_seconds == 600.0
def test_yaml_bridge_seeds_event_silence_extra():
"""``config.yaml`` ``discord.websocket_event_max_silence_seconds`` reaches the
adapter's ``extra`` through the liveness seed loop (public key wins over
the generic ``extra`` block)."""
seeded = discord_platform._apply_yaml_config(
{"platforms": {}},
{"websocket_event_max_silence_seconds": 7200},
)
assert seeded["websocket_event_max_silence_seconds"] == 7200

View File

@@ -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 — since the dispatch-side dimension — how long it has been since the last parsed Gateway event.
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,12 +96,15 @@ discord:
websocket_liveness_failure_threshold: 2
websocket_heartbeat_ack_max_age_seconds: 60
websocket_max_latency_seconds: 30
websocket_event_max_silence_seconds: 14400
```
The old `liveness_interval_seconds` and `liveness_failure_threshold` names remain compatibility aliases only; they no longer mean REST probing.
Any knob at `0` disables the whole WebSocket liveness probe. Values that fail to parse as a positive number (e.g. `15s`, `nan`, `true`, `-1`) also disable it, and log a warning each time the adapter starts — check `gateway.log` if the probe seems inactive.
`websocket_event_max_silence_seconds` is the exception: it guards a single dimension (event dispatch), so `0` opts out of **that check only** — ready/ACK/latency keep guarding. A socket can stay ESTABLISHED and keep ACKing heartbeats while delivering zero Gateway events; heartbeat ACKs are frames without an event type, so no transport-side check can see that state. The default (4 hours) matches the outage window operators have observed in the field; a quiet guild can legitimately go hours without a single Gateway event, so keep this bound generous unless you know your traffic.
## 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.

View File

@@ -94,10 +94,13 @@ discord:
websocket_liveness_failure_threshold: 2
websocket_heartbeat_ack_max_age_seconds: 60
websocket_max_latency_seconds: 30
websocket_event_max_silence_seconds: 14400
```
旧的 `liveness_interval_seconds` / `liveness_failure_threshold` 仅作为迁移别名保留,不再表示 REST probe。
`websocket_event_max_silence_seconds` 是唯一的例外:它只守护"事件分发"这一个维度,设为 `0` 仅停用该检查——ready/ACK/延迟检测继续生效。一条连接可以在保持 ESTABLISHED 并正常应答心跳的同时,不再投递任何 Gateway 事件(心跳 ACK 是不带事件类型的帧,任何传输层检查都看不到这种状态)。默认值 4 小时对应实际故障的观察窗口;安静的服务器可能数小时没有任何事件,除非明确了解自身流量,否则请保持该阈值宽松。
## 第一步:创建 Discord 应用
1. 前往 [Discord 开发者门户](https://discord.com/developers/applications) 并使用你的 Discord 账号登录。