From 535ee39053114863469f1596d11dbaf8ec6555fa Mon Sep 17 00:00:00 2001 From: imulan <18570354653@163.com> Date: Sun, 4 Oct 2026 10:10:13 +0800 Subject: [PATCH] fix(feishu): exit stalled websocket loop for supervised recovery --- frontends/fsapp.py | 34 +++++++- frontends/tests/test_fsapp_watchdog.py | 111 +++++++++++++++++++++++++ 2 files changed, 144 insertions(+), 1 deletion(-) create mode 100644 frontends/tests/test_fsapp_watchdog.py diff --git a/frontends/fsapp.py b/frontends/fsapp.py index ad959f660..f8af35c25 100644 --- a/frontends/fsapp.py +++ b/frontends/fsapp.py @@ -8,6 +8,7 @@ import traceback import lark_oapi as lark from lark_oapi.api.im.v1 import * +from lark_oapi.ws.client import loop as ws_loop def _ensure_dir(path): @@ -844,6 +845,37 @@ def handle_message(data): ).start() +def _start_ws_client(cli, loop, timeout=180, interval=5): + """Exit on a stalled SDK loop so the service manager can restart the process. + + A restart interrupts in-flight tasks; supervise this frontend (e.g. systemd). + A live loop alone does not prove that the remote connection is healthy. + """ + stopped = threading.Event() + last_beat = [time.monotonic()] + + def beat(): + if not stopped.is_set(): + last_beat[0] = time.monotonic() + loop.call_later(interval, beat) + + def watch(): + while not stopped.wait(interval): + if time.monotonic() - last_beat[0] >= timeout: + print(f"[ERROR] 飞书事件循环超过 {timeout}s 未响应,退出以交由服务管理器重启", + file=sys.stderr, flush=True) + os._exit(1) # sys.exit() would only terminate this watchdog thread. + + loop.call_soon_threadsafe(beat) + watcher = threading.Thread(target=watch, name="feishu-watchdog", daemon=True) + watcher.start() + try: + cli.start() + finally: + stopped.set() + watcher.join() + + def main(): global client, APP_ID, APP_SECRET, ALLOWED_USERS, PUBLIC_ACCESS, CONFIG_PATH APP_ID, APP_SECRET, ALLOWED_USERS, PUBLIC_ACCESS, CONFIG_PATH = _feishu_config() @@ -857,7 +889,7 @@ def main(): client = create_client() cli = lark.ws.Client(APP_ID, APP_SECRET, event_handler=handler, log_level=lark.LogLevel.INFO) print("=" * 50 + "\n飞书 Agent 已启动(长连接模式)\n" + f"App ID: {APP_ID}\n配置: {CONFIG_PATH}\n等待消息...\n" + "=" * 50, flush=True) - cli.start() + _start_ws_client(cli, ws_loop) retry_delay = 5 except KeyboardInterrupt: raise diff --git a/frontends/tests/test_fsapp_watchdog.py b/frontends/tests/test_fsapp_watchdog.py new file mode 100644 index 000000000..c66478085 --- /dev/null +++ b/frontends/tests/test_fsapp_watchdog.py @@ -0,0 +1,111 @@ +"""Watchdog checks run in disposable processes, never the pytest process.""" +import ast +from pathlib import Path +import subprocess +import sys +import textwrap +from types import SimpleNamespace + +import pytest + +SOURCE = Path(__file__).resolve().parents[1] / "fsapp.py" +TREE = ast.parse(SOURCE.read_text()) +WATCHDOG = next(node for node in TREE.body if isinstance(node, ast.FunctionDef) + and node.name == "_start_ws_client") +BOOTSTRAP = """ +import asyncio, os, sys, threading, time +from types import SimpleNamespace +loop = asyncio.new_event_loop() +asyncio.set_event_loop(loop) +""" + ast.unparse(WATCHDOG) + "\n" + + +def run_child(code): + return subprocess.run([sys.executable, "-c", BOOTSTRAP + textwrap.dedent(code)], + capture_output=True, text=True, timeout=10) + + +@pytest.mark.parametrize("phase", ["startup", "running"]) +def test_stalled_loop_exits_process(phase): + result = run_child(f""" + def start(): + if {phase!r} == 'running': + loop.call_later(0.1, time.sleep, 3) + loop.run_forever() + else: + time.sleep(3) + _start_ws_client(SimpleNamespace(start=start), loop, timeout=0.3, interval=0.02) + raise AssertionError('stalled client returned without watchdog exit') + """) + assert result.returncode == 1, result.stderr + assert "飞书事件循环超过 0.3s 未响应" in result.stderr + assert "Traceback" not in result.stderr + + +def test_idle_and_async_retry_wait_do_not_trigger_watchdog(): + result = run_child(""" + async def idle(): + # No incoming messages; longer than the watchdog deadline. + await asyncio.sleep(0.9) + for _ in range(2): + _start_ws_client(SimpleNamespace(start=lambda: loop.run_until_complete(idle())), + loop, timeout=0.3, interval=0.02) + assert not any(t.name == 'feishu-watchdog' for t in threading.enumerate()) + time.sleep(0.5) # A stopped loop must not trigger a previous watcher. + loop.close() + """) + assert result.returncode == 0, result.stderr + assert not result.stderr + + +@pytest.mark.parametrize("error", ["RuntimeError", "KeyboardInterrupt"]) +def test_start_failure_stops_watcher_and_preserves_exception(error): + result = run_child(f""" + def start(): + raise {error}('expected') + try: + _start_ws_client(SimpleNamespace(start=start), loop, timeout=0.3, interval=0.02) + except {error} as exc: + assert str(exc) == 'expected' + else: + raise AssertionError('exception swallowed') + assert not any(t.name == 'feishu-watchdog' for t in threading.enumerate()) + time.sleep(0.5) + loop.close() + """) + assert result.returncode == 0, result.stderr + assert not result.stderr + + +def test_main_wraps_sdk_start_and_keeps_existing_backoff(): + main = next(node for node in TREE.body if isinstance(node, ast.FunctionDef) and node.name == "main") + calls, delays = [], [] + client, event_loop = object(), object() + handler = object() + builder = SimpleNamespace(register_p2_im_message_receive_v1=lambda *_: SimpleNamespace(build=lambda: handler)) + + def guarded_start(actual_client, actual_loop): + calls.append((actual_client, actual_loop)) + raise RuntimeError("simulate SDK failure") + + def sleep(delay): + delays.append(delay) + if len(delays) == 7: + raise KeyboardInterrupt() + + namespace = { + "_feishu_config": lambda: ("app", "secret", set(), False, "test-config"), + "create_client": lambda: object(), "handle_message": lambda *_: None, + "_start_ws_client": guarded_start, "ws_loop": event_loop, + "time": SimpleNamespace(sleep=sleep), + "traceback": SimpleNamespace(print_exc=lambda: None), + "lark": SimpleNamespace( + ws=SimpleNamespace(Client=lambda *a, **kw: client), + EventDispatcherHandler=SimpleNamespace(builder=lambda *_: builder), + LogLevel=SimpleNamespace(INFO="info")), + } + exec(compile(ast.Module(body=[main], type_ignores=[]), str(SOURCE), "exec"), namespace) + with pytest.raises(KeyboardInterrupt): + namespace["main"]() + assert calls == [(client, event_loop)] * 7 + assert delays == [5, 10, 20, 40, 80, 120, 120]