Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
34 changes: 33 additions & 1 deletion frontends/fsapp.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down Expand Up @@ -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()
Expand All @@ -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
Expand Down
111 changes: 111 additions & 0 deletions frontends/tests/test_fsapp_watchdog.py
Original file line number Diff line number Diff line change
@@ -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]