mirror of
https://github.com/NousResearch/hermes-agent.git
synced 2026-07-29 18:46:59 +00:00
fix(slack): heal wedged Socket Mode via ping/pong staleness
The Socket Mode watchdog only reconnects when is_connected() returns False or the receiver task dies. When the underlying aiohttp ClientSession is closed (e.g. after a network blip), slack_sdk gets stuck retrying "Session is closed" while is_connected() can still report healthy and the receiver task stays alive — so the watchdog never fires and the process is alive but deaf to Slack indefinitely. Add a ping/pong staleness probe: Slack sends a ping roughly every ping_interval seconds even on an idle socket, so a stale/missing last_ping_pong_time (past a first-ping grace window) is a reliable signal the transport is wedged. The watchdog now also reconnects on staleness, which rebuilds the handler with a fresh session. Guards non-numeric attributes so a mocked/partial client never triggers a spurious reconnect. 7 new tests; full test_slack.py (216) green.
This commit is contained in:
parent
7bbdabbef2
commit
caf8e2f214
2 changed files with 127 additions and 0 deletions
|
|
@ -693,6 +693,16 @@ class SlackAdapter(BasePlatformAdapter):
|
|||
self._socket_watchdog_task: Optional[asyncio.Task] = None
|
||||
self._socket_reconnect_lock = asyncio.Lock()
|
||||
self._socket_watchdog_interval_s = 15.0
|
||||
# Monotonic timestamp of the most recent Socket Mode handler (re)start,
|
||||
# used to grant a grace window for the first ping/pong after connect.
|
||||
self._socket_handler_started_monotonic: Optional[float] = None
|
||||
# Reconnect when no ping/pong has arrived for this many multiples of the
|
||||
# client's ping_interval. Slack pings roughly every ping_interval seconds
|
||||
# even on an idle socket, so prolonged silence means a wedged transport.
|
||||
self._socket_ping_stale_factor = 4
|
||||
# Allow at least this long after (re)connect before treating a missing
|
||||
# first ping/pong as evidence of a wedged transport.
|
||||
self._socket_first_ping_grace_s = 60.0
|
||||
|
||||
def _start_socket_mode_handler(self) -> None:
|
||||
"""Start the Slack Socket Mode background task."""
|
||||
|
|
@ -706,6 +716,7 @@ class SlackAdapter(BasePlatformAdapter):
|
|||
|
||||
task = asyncio.create_task(self._handler.start_async())
|
||||
self._socket_mode_task = task
|
||||
self._socket_handler_started_monotonic = time.monotonic()
|
||||
task.add_done_callback(self._on_socket_mode_task_done)
|
||||
|
||||
async def _stop_socket_mode_handler(self) -> None:
|
||||
|
|
@ -766,6 +777,37 @@ class SlackAdapter(BasePlatformAdapter):
|
|||
)
|
||||
return None
|
||||
|
||||
def _socket_ping_pong_stale(self) -> bool:
|
||||
"""True when the Socket Mode transport shows no recent ping/pong.
|
||||
|
||||
slack_sdk's Socket Mode client records ``last_ping_pong_time`` whenever
|
||||
Slack's periodic ping arrives (roughly every ``ping_interval`` seconds,
|
||||
even on an otherwise idle connection). When the underlying aiohttp
|
||||
session is closed, the client gets stuck retrying ("Session is closed")
|
||||
while ``is_connected()`` can still report healthy — so ping/pong
|
||||
staleness is the reliable signal that the socket is wedged and the
|
||||
handler must be rebuilt. Guards against non-numeric attributes so a
|
||||
mocked/partial client never triggers a spurious reconnect.
|
||||
"""
|
||||
client = getattr(self._handler, "client", None)
|
||||
if client is None:
|
||||
return False
|
||||
ping_interval = getattr(client, "ping_interval", None)
|
||||
if not isinstance(ping_interval, (int, float)) or ping_interval <= 0:
|
||||
return False
|
||||
last = getattr(client, "last_ping_pong_time", None)
|
||||
if last is None:
|
||||
# No ping yet. Healthy right after (re)connect; only suspicious once
|
||||
# the grace window elapses without ever seeing the first ping/pong.
|
||||
started = self._socket_handler_started_monotonic
|
||||
if started is None:
|
||||
return False
|
||||
grace = max(self._socket_first_ping_grace_s, ping_interval * 2)
|
||||
return (time.monotonic() - started) > grace
|
||||
if not isinstance(last, (int, float)):
|
||||
return False
|
||||
return (time.time() - last) > (ping_interval * self._socket_ping_stale_factor)
|
||||
|
||||
async def _restart_socket_mode(self, reason: str) -> None:
|
||||
"""Reconnect Socket Mode without rebuilding adapter state."""
|
||||
if not self._running:
|
||||
|
|
@ -810,6 +852,11 @@ class SlackAdapter(BasePlatformAdapter):
|
|||
connected = await self._socket_transport_connected()
|
||||
if connected is False:
|
||||
await self._restart_socket_mode("transport disconnected")
|
||||
elif self._socket_ping_pong_stale():
|
||||
# is_connected() can lie when the aiohttp session is closed
|
||||
# but the client keeps retrying; ping/pong staleness catches
|
||||
# that wedged-zombie case that the bool check above misses.
|
||||
await self._restart_socket_mode("ping/pong stale")
|
||||
except asyncio.CancelledError:
|
||||
raise
|
||||
except Exception: # pragma: no cover - defensive logging
|
||||
|
|
|
|||
|
|
@ -12,6 +12,7 @@ import asyncio
|
|||
import contextlib
|
||||
import os
|
||||
import sys
|
||||
import time
|
||||
from unittest.mock import AsyncMock, MagicMock, patch, call
|
||||
|
||||
import pytest
|
||||
|
|
@ -802,6 +803,85 @@ class TestSlackSocketWatchdog:
|
|||
finally:
|
||||
await adapter.disconnect()
|
||||
|
||||
# -- ping/pong staleness: heals the wedged transport that is_connected() misses --
|
||||
|
||||
def _adapter_with_fake_client(self, **client_attrs):
|
||||
adapter = SlackAdapter(PlatformConfig(enabled=True, token="xoxb-fake"))
|
||||
client = MagicMock()
|
||||
for key, value in client_attrs.items():
|
||||
setattr(client, key, value)
|
||||
adapter._handler = MagicMock(client=client)
|
||||
return adapter
|
||||
|
||||
def test_ping_pong_stale_when_last_ping_old(self):
|
||||
adapter = self._adapter_with_fake_client(
|
||||
ping_interval=30, last_ping_pong_time=time.time() - 1000
|
||||
)
|
||||
assert adapter._socket_ping_pong_stale() is True
|
||||
|
||||
def test_ping_pong_fresh_when_last_ping_recent(self):
|
||||
adapter = self._adapter_with_fake_client(
|
||||
ping_interval=30, last_ping_pong_time=time.time() - 5
|
||||
)
|
||||
assert adapter._socket_ping_pong_stale() is False
|
||||
|
||||
def test_ping_pong_none_within_grace_not_stale(self):
|
||||
adapter = self._adapter_with_fake_client(
|
||||
ping_interval=30, last_ping_pong_time=None
|
||||
)
|
||||
adapter._socket_handler_started_monotonic = time.monotonic()
|
||||
assert adapter._socket_ping_pong_stale() is False
|
||||
|
||||
def test_ping_pong_none_beyond_grace_is_stale(self):
|
||||
adapter = self._adapter_with_fake_client(
|
||||
ping_interval=30, last_ping_pong_time=None
|
||||
)
|
||||
adapter._socket_first_ping_grace_s = 0.0
|
||||
adapter._socket_handler_started_monotonic = time.monotonic() - 200
|
||||
assert adapter._socket_ping_pong_stale() is True
|
||||
|
||||
def test_ping_pong_no_handler_not_stale(self):
|
||||
adapter = SlackAdapter(PlatformConfig(enabled=True, token="xoxb-fake"))
|
||||
adapter._handler = None
|
||||
assert adapter._socket_ping_pong_stale() is False
|
||||
|
||||
def test_ping_pong_nonnumeric_attrs_not_stale(self):
|
||||
# A mocked/partial client (MagicMock attrs) must never trigger reconnect.
|
||||
adapter = SlackAdapter(PlatformConfig(enabled=True, token="xoxb-fake"))
|
||||
adapter._handler = MagicMock()
|
||||
assert adapter._socket_ping_pong_stale() is False
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_watchdog_reconnects_when_ping_pong_stale_despite_is_connected_true(self):
|
||||
adapter = SlackAdapter(PlatformConfig(enabled=True, token="xoxb-fake"))
|
||||
adapter._socket_watchdog_interval_s = 0.01
|
||||
factory, instances = self._make_fake_handler_factory()
|
||||
|
||||
with contextlib.ExitStack() as stack:
|
||||
for p in self._patch_stack(factory):
|
||||
stack.enter_context(p)
|
||||
|
||||
try:
|
||||
assert await adapter.connect() is True
|
||||
assert len(instances) == 1
|
||||
|
||||
# Transport lies: is_connected() stays True while ping/pong has
|
||||
# gone stale (the wedged "Session is closed" zombie).
|
||||
instances[0].client.is_connected = lambda: True
|
||||
instances[0].client.ping_interval = 30
|
||||
instances[0].client.last_ping_pong_time = time.time() - 1000
|
||||
|
||||
for _ in range(40):
|
||||
if len(instances) >= 2:
|
||||
break
|
||||
await asyncio.sleep(0.01)
|
||||
|
||||
assert len(instances) >= 2, "watchdog did not heal wedged (lying) transport"
|
||||
assert instances[0].closed is True
|
||||
assert adapter._handler is instances[-1]
|
||||
finally:
|
||||
await adapter.disconnect()
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# TestSlackProxyBehavior
|
||||
|
|
|
|||
Loading…
Add table
Add a link
Reference in a new issue