fix(feishu): run the inbound dedup flush on the adapter-owned pool
A dead background event loop tears the loop's default executor down, and after that every inbound message was dropped inside the dedup gate with a RuntimeError out of asyncio.to_thread — the adapter went permanently deaf while the gateway process, websocket and service all stayed healthy. #10849 already moved the outbound SDK calls onto an adapter-owned, self-healing pool; this gives the inbound dedup-state flush the same treatment, so a default-executor teardown can no longer wedge message intake.
This commit is contained in:
@@ -3442,10 +3442,12 @@ class FeishuAdapter(BasePlatformAdapter):
|
||||
while len(self._seen_message_order) > self._dedup_cache_size:
|
||||
self._seen_message_ids.pop(self._seen_message_order.pop(0), None)
|
||||
# atomic_json_write() fsyncs; this runs on the event loop for every inbound message, so
|
||||
# offload the flush. The lock keeps flushes in mutation order (the snapshot inside the
|
||||
# worker is taken under _dedup_lock, but the write itself is not).
|
||||
# offload the flush onto the adapter-owned pool: the loop's default executor may already
|
||||
# be torn down by a dead background loop, which used to wedge every inbound message in
|
||||
# the dedup gate (#111020). The lock keeps flushes in mutation order (the snapshot
|
||||
# inside the worker is taken under _dedup_lock, but the write itself is not).
|
||||
async with self._dedup_persist_lock_or_create():
|
||||
await asyncio.to_thread(self._persist_seen_message_ids)
|
||||
await self._run_blocking(self._persist_seen_message_ids)
|
||||
return False
|
||||
|
||||
def _dedup_persist_lock_or_create(self) -> asyncio.Lock:
|
||||
|
||||
@@ -6,9 +6,11 @@ subsequent send failed permanently with "Executor shutdown has been called"
|
||||
and the gateway became a zombie. The adapter now owns its own
|
||||
ThreadPoolExecutor and recreates it on demand if it has been shut down.
|
||||
|
||||
Covers: #10849
|
||||
Covers: #10849, #111020
|
||||
"""
|
||||
import asyncio
|
||||
import concurrent.futures
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
@@ -57,3 +59,35 @@ async def test_run_blocking_executes_on_owned_pool():
|
||||
adapter._shutdown_sdk_executor()
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_is_duplicate_flush_survives_default_executor_teardown(tmp_path, monkeypatch):
|
||||
"""The inbound dedup flush must not ride the loop's default executor.
|
||||
|
||||
After the adapter's background event loop dies, its teardown shuts the
|
||||
default executor down; every subsequent inbound message was then dropped
|
||||
inside the dedup gate with a RuntimeError out of asyncio.to_thread
|
||||
(#111020). The flush now runs on the adapter-owned pool, mirroring the
|
||||
outbound SDK calls (#10849).
|
||||
"""
|
||||
monkeypatch.setenv("HERMES_HOME", str(tmp_path))
|
||||
from gateway.config import PlatformConfig
|
||||
|
||||
adapter = FeishuAdapter(PlatformConfig())
|
||||
try:
|
||||
loop = asyncio.get_running_loop()
|
||||
torn_down = concurrent.futures.ThreadPoolExecutor(max_workers=1)
|
||||
loop.set_default_executor(torn_down)
|
||||
torn_down.shutdown(wait=True)
|
||||
|
||||
# Scenario premise: with the default executor torn down, to_thread fails.
|
||||
with pytest.raises(RuntimeError):
|
||||
await asyncio.to_thread(lambda: None)
|
||||
|
||||
# The dedup gate must still admit the message and flush it to disk.
|
||||
assert await adapter._is_duplicate("om_after_executor_teardown") is False
|
||||
state = json.loads((tmp_path / "feishu_seen_message_ids.json").read_text(encoding="utf-8"))
|
||||
assert "om_after_executor_teardown" in state["message_ids"]
|
||||
finally:
|
||||
adapter._shutdown_sdk_executor()
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user