From af9a141f34496a8f4a88be3687a521699943e801 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 12:54:06 +0800 Subject: [PATCH] fix: isolate feishu websocket event loops --- app/modules/feishu/feishu.py | 64 +++++++++++++++---- .../backend-architecture-next-stage.md | 13 ++++ tests/test_feishu_ws_lifecycle.py | 63 +++++++++++++++++- 3 files changed, 128 insertions(+), 12 deletions(-) diff --git a/app/modules/feishu/feishu.py b/app/modules/feishu/feishu.py index 06f8b6431..563263a87 100644 --- a/app/modules/feishu/feishu.py +++ b/app/modules/feishu/feishu.py @@ -65,6 +65,57 @@ from app.adapters.network.http import RequestUtils from app.runtime.thread import ThreadHelper +class _ThreadLocalEventLoopProxy: + """为使用模块级 loop 的飞书 SDK 路由当前实例线程的事件循环。""" + + def __init__(self, fallback: asyncio.AbstractEventLoop) -> None: + """保存 SDK 原始循环,并初始化互不共享的线程绑定。""" + self._fallback = fallback + self._state = threading.local() + + def bind( + self, + loop: asyncio.AbstractEventLoop, + stop_event: threading.Event, + ) -> None: + """为当前飞书实例线程绑定循环和停止信号。""" + self._state.loop = loop + self._state.stop_event = stop_event + + def unbind(self) -> None: + """清除当前线程绑定,防止复用线程时误取已关闭循环。""" + self._state.__dict__.clear() + + def stop_event(self) -> Optional[threading.Event]: + """返回当前实例线程的停止信号,未绑定时返回空。""" + return getattr(self._state, "stop_event", None) + + def __getattr__(self, name: str) -> Any: + """把 SDK loop 操作转发给当前线程循环或原始兼容循环。""" + loop = getattr(self._state, "loop", self._fallback) + return getattr(loop, name) + + +_LARK_WS_ORIGINAL_SELECT = lark_ws_client_module._select +_lark_ws_loop_proxy = _ThreadLocalEventLoopProxy(lark_ws_client_module.loop) + + +async def _select_bound_ws_client() -> None: + """按当前实例的停止信号结束 SDK 阻塞选择;未绑定时保持 SDK 原行为。""" + stop_event = _lark_ws_loop_proxy.stop_event() + if stop_event is None: + await _LARK_WS_ORIGINAL_SELECT() + return + while not stop_event.is_set(): + await asyncio.sleep(1) + + +# lark_oapi 以模块全局 loop 驱动所有 Client;静态安装线程路由后,多配置实例 +# 不再在启动/退出时反复覆盖同一全局对象。 +lark_ws_client_module.loop = _lark_ws_loop_proxy +lark_ws_client_module._select = _select_bound_ws_client + + class UserOper: """兼容飞书模块存量测试的渠道用户查询门面。""" @@ -176,18 +227,11 @@ class Feishu: def _run_ws_client(self) -> None: """在后台线程中运行飞书长连接客户端。""" - original_select = lark_ws_client_module._select - original_loop = lark_ws_client_module.loop loop = asyncio.new_event_loop() original_create_task = loop.create_task self._ws_loop = loop asyncio.set_event_loop(loop) - lark_ws_client_module.loop = loop - - async def _wait_for_stop() -> None: - """等待停止信号,让 SDK 的阻塞 select 可被本地生命周期控制。""" - while not self._stop_event.is_set(): - await asyncio.sleep(1) + _lark_ws_loop_proxy.bind(loop, self._stop_event) def _create_tracked_task(coro, *args, **kwargs) -> asyncio.Task: """跟踪 SDK 后台任务,避免关闭时产生未取回的任务异常。""" @@ -202,7 +246,6 @@ class Feishu: task.add_done_callback(self._consume_ws_task_result) return task - lark_ws_client_module._select = _wait_for_stop loop.create_task = _create_tracked_task try: self._ws_client = lark.ws.Client( @@ -223,8 +266,6 @@ class Feishu: finally: if not loop.is_closed(): loop.run_until_complete(self._shutdown_ws_client()) - lark_ws_client_module._select = original_select - lark_ws_client_module.loop = original_loop loop.create_task = original_create_task pending_tasks = [ task @@ -240,6 +281,7 @@ class Feishu: loop.close() asyncio.set_event_loop(None) self._ws_loop = None + _lark_ws_loop_proxy.unbind() def _consume_ws_task_result(self, task: asyncio.Task) -> None: """取回飞书 SDK 后台任务结果,防止 asyncio 在关机时输出未消费异常。""" diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index dddb47ef0..a251d486a 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -8,6 +8,7 @@ > 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` > 实施进度:阶段 0~6 的宿主架构能力已完成收口;API/Application 公共复杂度基线已清零,启动组合根的 SystemConfigOper 构造点已由 14 降至 1;API 进程内后台任务已完成首批统一登记,插件仓适配和 Outbox 外围扩展仍按风险切片推进。Model/Base 查询与写装饰器、legacy 隐式会话外壳均已清零,插件 SDK 也不再导出宿主 Model。2026-08-23 的长期整改阶段 0 已恢复宿主、启动性能、官方插件和 SDK 契约门禁的可信基线;阶段 1a 已补齐 TaskRegistry owner 零债务门禁和诚实的关停超时语义;阶段 1b1 已收口整理 worker、pending 回放、失败通知、进程内 AI 重试、插件监控与事件投递的生命周期所有权;2026-08-24 的阶段 2 已将 212 个已观察宿主模块方法的 legacy aggregation 清零,并补齐可执行 fanout 与下载器文件 DTO 边界;阶段 3 已将消息交互和远程命令的订阅删除统一到 Application/UoW/outbox,宿主不再调用裸线程统计入口;阶段 4 已统一七种消息渠道的宿主回环与后台执行边界;阶段 5 已补齐事件窗口聚合任务的生命周期所有权;阶段 6 已统一插件文件操作的取消完成语义;阶段 7 已统一插件协程补偿的终态等待;阶段 8 已统一宿主同步函数的异步线程池入口;阶段 9 已统一工作流运行时的宿主获取路径;阶段 10 已统一模块、插件与调度运行时的显式 getter 调用;阶段 11 已清除系统配置 getter 的 Oper 形别名;阶段 12 已完成工作流域的显式 Chain 数据端口迁移;阶段 13 已收口用户、交互与消息链的数据端口;阶段 14 已收口音乐订阅数据端口;阶段 15 已收口站点数据端口;阶段 16 已收口媒体服务器数据端口;阶段 17 已收口下载数据端口;阶段 18 已收口主订阅数据端口;阶段 19 已收口整理数据端口;阶段 20 已收口 Agent 数据端口;阶段 21 已收口监控历史端口;阶段 22 已统一服务配置应用边界;阶段 23 已补齐媒体服务器 API 遗留的类形配置读取路径;阶段 24 已清除 Scheduler 内部无 owner 的协程提交双轨;阶段 25 已补齐 TaskRegistry 跨线程 owner 并迁移整理 AI 接管;阶段 26 已统一 Agent 会话清理提交;阶段 27 已统一历史 AI 进度 owner;阶段 28 已托管旧插件订阅统计线程;阶段 29 已统一 Emby 系条目转换并清零重复代码白名单;阶段 30 已收口插件市场请求级子任务;阶段 31 已托管搜索 AI 推荐任务;阶段 32 已清除事件调度器绕过生命周期 owner 的投递回退;阶段 33 已统一宿主 Agent 运行时的获取路径;阶段 34 已统一 durable-required 事件与 Outbox topic 事实源;阶段 35 已统一 LLM provider 管理 API 的运行时解析路径;阶段 36 已统一 WebAgent 音频能力访问边界;阶段 37 已统一插件输入事件发布路径;阶段 38 已统一 WebAgent 通知事件监听与队列边界;阶段 39 已补齐搜索 SSE 断线时的上游任务清理;阶段 40 已补齐异步防抖取消的终态所有权;阶段 41 已统一优雅重启兜底线程的唯一所有权;阶段 42 已补齐 Telegram typing 的多实例隔离和终态 owner;阶段 43 已统一 Discord typing 的异步 owner 和 shutdown 收尾;阶段 44 已清除 WebAgent 测试临时事件循环提前关闭产生的 CI 红注解;阶段 45 已统一影视与字幕搜索的请求级逐页任务编排;阶段 46 已收口启动性能门禁的托管 runner 假失败与诊断输出;阶段 47 已补齐 Agent 渠道流式刷新任务的重入 owner;阶段 48 已统一工件上传 action 的 Node 24 主版本;阶段 49 已统一插件安装的同步/异步代际解析事实源;阶段 50 已统一插件市场 GitHub 请求降级策略;阶段 51 已统一插件索引请求与响应三态策略;阶段 52 已统一插件 Release 分页策略;阶段 53 已统一远端插件安装模式决策;阶段 54 已补齐同步安装成功后的临时回滚备份清理;阶段 55~56 已收口官方插件观察基线与报告保留策略;阶段 57 已统一进程级运行时 Facade 门禁并补齐 ModuleManager 边界;阶段 58 已消除 AgentTask 关闭回归的跨线程零时长等待竞态。 > 当前 canonical 状态:API/Application 公共复杂度基线已清零,组合根外 `SystemConfigOper()` 构造和 Model/Oper 隐式事务均为 0;命名 Chain/Agent 数据端口、TaskRegistry owner、Module Contract V2、typed Event、Outbox durable intent、请求关联和插件运行时 getter 已形成当前路径。插件仓适配、未知第三方 fallback 和其它 E1/E3 副作用仍按风险持续治理。 +> 最新阶段:阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由。 ## 当前复核结论(2026-08-24) @@ -610,6 +611,15 @@ 生产语义。`execute_scheduled_task`、`AgentTaskExecutionService`、`AgentManager.close()` 与插件 ABI 均未修改。 +### 长期整改阶段 59:Feishu 多实例 SDK 循环路由统一(2026-08-24) + +- `FeishuModule` 支持多个配置实例,但 lark SDK 以模块级 `loop` 与 `_select` 驱动客户端;原实现由每个 + 实例线程临时覆盖并在退出时恢复这两个全局,两个实例并发启动或停止时会互相窃取循环或恢复旧全局。 +- 宿主现在只安装一份线程本地 loop 代理和停止选择器:绑定线程内继续使用实例独立事件循环,未绑定的 + SDK 调用仍委托原始 loop/select;每个实例原有 `_ws_tasks`、静默断连和有限 join 语义不变。 +- 并发回归用屏障强制两个模拟 SDK Client 同时读取模块级入口,证明二者运行在不同循环且均可独立停止。 + Feishu 配置、消息/卡片 API、模块类身份、插件 SDK/Compat 与其他 V1/V2/V3 插件行为均未修改。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: @@ -650,6 +660,9 @@ `await_task_to_terminal`;数据库 worker 的 interruptible 等待属于队列 owner,未被机械合并。 canonical 宿主原有 11 处 FastAPI `run_in_threadpool` 导入也已统一到 `runtime.execution`,AST 门禁 防止再次出现框架直连;模块局部符号及调用参数不变,不改变插件运行和 monkeypatch 接缝。 + Feishu 长连接原先由每个配置实例分别覆盖 lark SDK 的模块级 loop 与阻塞选择函数,多实例并发会把 + 先启动实例的 SDK 任务路由到后启动实例的循环;现在以单一线程本地代理分派各实例 loop 和停止信号, + 任务集合、连接清理和模块 stop 仍由各实例持有。 2. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。 Oper 内部的执行入口也已统一:最后一处 `AgentTaskOper` 直接 transaction runner 调用已迁入 diff --git a/tests/test_feishu_ws_lifecycle.py b/tests/test_feishu_ws_lifecycle.py index 7d1bb62d8..e806c5f59 100644 --- a/tests/test_feishu_ws_lifecycle.py +++ b/tests/test_feishu_ws_lifecycle.py @@ -1,4 +1,5 @@ import asyncio +import threading from types import SimpleNamespace from unittest.mock import MagicMock, patch @@ -9,7 +10,7 @@ ensure_optional_stub("psutil") ensure_optional_stub("dateparser") ensure_optional_stub("Pinyin2Hanzi", is_pinyin=lambda value: False) -from app.modules.feishu.feishu import Feishu +from app.modules.feishu.feishu import Feishu, lark_ws_client_module def _build_feishu_client() -> Feishu: @@ -30,6 +31,66 @@ async def _wait_forever() -> None: await asyncio.Future() +def test_parallel_ws_clients_keep_independent_sdk_event_loops() -> None: + """多个飞书配置并发启动时不得覆盖彼此的 SDK 事件循环。""" + clients = [_build_feishu_client(), _build_feishu_client()] + barrier = threading.Barrier(2) + constructed = threading.Event() + construction_lock = threading.Lock() + fake_clients = [] + + class _ConcurrentWsClient: + """模拟真实 SDK 通过模块级 loop 与 _select 驱动长连接。""" + + def __init__(self, *_args, **_kwargs): + """登记实例并准备关停路径需要的 SDK 私有状态。""" + self._auto_reconnect = True + self._conn = None + self._conn_url = "" + self._conn_id = "" + self._service_id = "" + self._lock = asyncio.Lock() + self.started = threading.Event() + self.observed_loop = None + with construction_lock: + fake_clients.append(self) + if len(fake_clients) == 2: + constructed.set() + + def start(self) -> None: + """强制两个线程同时解析 SDK 全局,再等待各自停止信号。""" + barrier.wait(timeout=2) + + async def run_until_stopped() -> None: + """记录真实运行循环,并调用 SDK 的模块级阻塞选择。""" + self.observed_loop = asyncio.get_running_loop() + self.started.set() + await lark_ws_client_module._select() + + lark_ws_client_module.loop.run_until_complete(run_until_stopped()) + + with patch( + "app.modules.feishu.feishu.lark.ws.Client", + _ConcurrentWsClient, + ): + try: + for client in clients: + client._start_ws_client() + assert constructed.wait(timeout=2) + assert all(fake.started.wait(timeout=2) for fake in fake_clients) + observed_loops = [fake.observed_loop for fake in fake_clients] + assert len(set(observed_loops)) == 2 + assert all(loop is not None for loop in observed_loops) + finally: + for client in clients: + client.stop() + + assert all( + client._ws_thread is None or not client._ws_thread.is_alive() + for client in clients + ) + + def test_shutdown_ws_client_cancels_sdk_tasks_before_quiet_disconnect(): """飞书关机清理应先消费后台任务,再静默关闭 WebSocket 连接。""" client = _build_feishu_client()