From af97f2c27a634d60f7a65b1f1bc5ec12c85ad27e Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 09:57:02 +0800 Subject: [PATCH] test: drain web agent owners before loop close --- app/application/messaging/agent.py | 2 +- .../backend-architecture-next-stage.md | 13 ++++++++- pytest.ini | 1 + tests/test_web_agent_stream.py | 29 +++++++++++++++---- 4 files changed, 37 insertions(+), 8 deletions(-) diff --git a/app/application/messaging/agent.py b/app/application/messaging/agent.py index 3b816edae..5dc3a7097 100644 --- a/app/application/messaging/agent.py +++ b/app/application/messaging/agent.py @@ -212,7 +212,7 @@ async def shutdown_web_agent_background_tasks() -> None: async def wait_web_agent_background_tasks() -> None: - """等待已登记的 Web Agent 任务完成取消后的最终收尾。""" + """等待已登记的 Web Agent 任务完成最终收尾。""" tasks = tuple(_WEB_AGENT_BACKGROUND_TASKS) if tasks: await asyncio.gather(*tasks, return_exceptions=True) diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index b07a8b122..95df2d1fb 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -6,7 +6,7 @@ > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 相关文档:`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 收尾。 +> 实施进度:阶段 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 红注解。 ## 当前复核结论(2026-08-24) @@ -461,6 +461,17 @@ - `Discord` 类路径、构造参数、同步 `start_typing()`/`stop_typing()` 布尔合同、模块方法、消息格式和配置 字段均保持不变;V1/V2/V3 插件仍通过原模块能力调用。本阶段未修改插件仓、SDK 或 Compat 映射。 +### 长期整改阶段 44:WebAgent 测试后台 owner 与 CI 红注解收口(2026-08-24) + +- GitHub 单测作业虽然成功,WebAgent SSE 用例却会在临时 `asyncio.run()` 关闭事件循环时取消仍登记的快照 + owner;Linux 下 `aiosqlite` worker 随后回送结果会触发 `PytestUnhandledThreadExceptionWarning`,Actions 将 + traceback 中两处 `Event loop is closed` 各自显示为错误级 annotation,形成“绿作业带红错误”。 +- WebAgent 流测试现在复用生产 `wait_web_agent_background_tasks()` 合同,在关闭临时循环前等待正常 owner; + 断线和慢快照用例仍先证明后台执行与 `done` 发送不受落库阻塞,再释放并等待真实终态,不把生产异步语义 + 改成同步等待。`wait_web_agent_background_tasks()` 的文档也改为反映其同时服务正常与取消收尾。 +- `pytest.ini` 将未处理线程异常提升为测试失败,后续不会再以绿色结果掩盖后台线程越过事件循环生命周期。 + 本阶段不改变 WebAgent API/SSE、Agent 模块、SDK/Compat 或 V1/V2/V3 插件合同,也未修改插件仓。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/pytest.ini b/pytest.ini index 8c609116d..8ff383437 100644 --- a/pytest.ini +++ b/pytest.ini @@ -9,6 +9,7 @@ asyncio_debug = true # 仅对「无法在本仓修复根因」的已知上游/三方弃用告警做精确忽略,保持测试输出干净、 # 让本仓自身的新告警更醒目。本仓代码引发的告警一律不在此忽略,应在源码/用例处修复。 filterwarnings = + error::pytest.PytestUnhandledThreadExceptionWarning ignore:'_UnionGenericAlias' is deprecated and slated for removal in Python 3\.17:DeprecationWarning:google\.genai\.types ignore:invalid escape sequence '\\&':SyntaxWarning:oss2\.api ignore:datetime.datetime.utcfromtimestamp\(\) is deprecated:DeprecationWarning diff --git a/tests/test_web_agent_stream.py b/tests/test_web_agent_stream.py index 81d941785..c4670d056 100644 --- a/tests/test_web_agent_stream.py +++ b/tests/test_web_agent_stream.py @@ -43,6 +43,7 @@ from app.application.messaging.agent import ( detach_web_agent_message_queue, dispatch_web_agent_message_event, extract_web_agent_message_from_event_data, + wait_web_agent_background_tasks, ) from app.application.messaging.skill import skill_interaction_manager from app.chain.message import MessageChain @@ -1385,10 +1386,15 @@ def test_web_agent_stream_drops_secret_result_after_disconnect(): async def scenario(): response = await web_agent_stream(payload, request, user) - body = "".join(await _collect_streaming_response(response)) + body = "".join( + await _collect_streaming_response( + response, + wait_for_background=False, + ) + ) release_agent.set() await asyncio.wait_for(agent_completed.wait(), timeout=1) - await asyncio.sleep(0) + await wait_web_agent_background_tasks() return body try: @@ -1509,6 +1515,7 @@ def test_web_agent_stop_finishes_stream_without_error(): while '"type": "done"' not in "".join(received): received.append(await asyncio.wait_for(anext(iterator), timeout=1)) await iterator.aclose() + await wait_web_agent_background_tasks() return "".join(received) try: @@ -1623,6 +1630,7 @@ def test_web_agent_traditional_stream_keeps_alive_and_saves_after_done(): assert not snapshot_finished.is_set() snapshot_release.set() await asyncio.to_thread(snapshot_finished.wait, 1) + await wait_web_agent_background_tasks() return "".join(received) try: @@ -1759,6 +1767,7 @@ def test_web_agent_stream_sends_done_before_snapshot_persistence_finishes(): assert not snapshot_finished.is_set() snapshot_release.set() await asyncio.to_thread(snapshot_finished.wait, 1) + await wait_web_agent_background_tasks() return "".join(received) try: @@ -1795,11 +1804,19 @@ def test_web_agent_stream_sends_done_before_snapshot_persistence_finishes(): assert snapshot_finished.wait(timeout=1) -async def _collect_streaming_response(response): - """读取 StreamingResponse,便于断言 SSE 内容。""" +async def _collect_streaming_response( + response, + *, + wait_for_background: bool = True, +): + """读取 StreamingResponse,并按用例语义等待生产 owner 完成收尾。""" chunks = [] - async for chunk in response.body_iterator: - chunks.append(chunk.decode("utf-8") if isinstance(chunk, bytes) else chunk) + try: + async for chunk in response.body_iterator: + chunks.append(chunk.decode("utf-8") if isinstance(chunk, bytes) else chunk) + finally: + if wait_for_background: + await wait_web_agent_background_tasks() return chunks