From ad45bfac5d11e757f84058d7babfd8ffb4cabd91 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 14:25:59 +0800 Subject: [PATCH] refactor: unify message queue shutdown convergence --- app/application/messaging/message.py | 46 +++++++++++++--- .../backend-architecture-next-stage.md | 21 ++++++- docs/rules/05-architecture.md | 2 + tests/test_lifecycle_shutdown.py | 12 ++++ tests/test_message_queue_shutdown.py | 55 ++++++++++++++++++- 5 files changed, 123 insertions(+), 13 deletions(-) diff --git a/app/application/messaging/message.py b/app/application/messaging/message.py index 69ded86ff..a0e03f0bc 100644 --- a/app/application/messaging/message.py +++ b/app/application/messaging/message.py @@ -36,6 +36,7 @@ from app.foundation.crypto import HashUtils _ALBUM_TRAILING_YEAR_RE = re.compile( r"(?:[\s\u3000]*[\(\[(【]\s*(?:19|20)\d{2}\s*[\)\])】])+$" ) +_MESSAGE_QUEUE_STOP_TIMEOUT_SECONDS = 10.0 class AsyncMessageQueryRepository(Protocol): @@ -996,7 +997,7 @@ class MessageQueueManager(metaclass=SingletonClass): while self._running: current_time = datetime.now() if self._is_in_scheduled_time(current_time): - while not self.queue.empty(): + while self._running and not self.queue.empty(): if global_vars.is_system_stopped: break if not self._is_in_scheduled_time(datetime.now()): @@ -1010,15 +1011,28 @@ class MessageQueueManager(metaclass=SingletonClass): if self._stop_event.wait(self.check_interval): break - def stop(self) -> None: + def stop( + self, + timeout: float = _MESSAGE_QUEUE_STOP_TIMEOUT_SECONDS, + ) -> bool: """ - 停止队列管理器 + 在有限时间内停止队列管理器。 + + :param timeout: 等待监控线程退出的最长秒数 + :return: 监控线程已经终止时返回 True,超时或线程自停时返回 False """ self._running = False self._stop_event.set() logger.info("正在停止消息队列...") - self.thread.join() + if self.thread is threading.current_thread(): + logger.error("消息队列不能在自身监控线程内等待退出") + return False + self.thread.join(timeout=max(0.0, timeout)) + if self.thread.is_alive(): + logger.error(f"消息队列在 {timeout:g} 秒内未停止") + return False logger.info("消息队列已停止") + return True class MessageHelper(metaclass=Singleton): @@ -1095,12 +1109,28 @@ class MessageHelper(metaclass=Singleton): return None -def stop_message(): +def stop_message( + timeout: float = _MESSAGE_QUEUE_STOP_TIMEOUT_SECONDS, +) -> bool: """ - 停止消息服务 + 停止已启动的消息服务并返回全部资源是否收敛。 + + :param timeout: 等待消息队列监控线程退出的最长秒数 + :return: 已启动资源均完成关闭时返回 True,否则返回 False """ + all_converged = True # 只关闭已启动的服务,避免清理路径反向创建后台线程和缓存 if queue_manager := MessageQueueManager.get_existing_instance(): - queue_manager.stop() + try: + if queue_manager.stop(timeout=timeout) is False: + all_converged = False + except Exception as err: + logger.error(f"停止消息队列失败:{str(err)}") + all_converged = False if template_helper := TemplateHelper.get_existing_instance(): - template_helper.close() + try: + template_helper.close() + except Exception as err: + logger.error(f"关闭消息模板缓存失败:{str(err)}") + all_converged = False + return all_converged diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 95597e05e..5479014d8 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -2,13 +2,13 @@ > 文档性质:当前架构复核、优秀 Python 后端实践对标、AI 可执行任务手册 > 适用仓库:`MoviePilot`,分支 `v3` -> 审计基线:`009631b8`(2026-08-24) +> 审计基线:`415335b2`(2026-08-24) > 审计范围:宿主后端;排除 `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 收尾;阶段 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 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同。 +> 实施进度:阶段 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 关闭回归的跨线程零时长等待竞态;阶段 59 已统一 Feishu 多实例长连接的 SDK 循环路由;阶段 60 已清除命令服务虚假的关停 owner 声明;阶段 61 已统一 Capability Runtime 同步/异步关闭的诚实收敛结果;阶段 62 已统一消息渠道长连接的多实例关闭收敛合同;阶段 63 已补齐应用消息队列线程的关闭收敛合同。 > 当前 canonical 状态:API/Application 公共复杂度基线已清零,组合根外 `SystemConfigOper()` 构造和 Model/Oper 隐式事务均为 0;命名 Chain/Agent 数据端口、TaskRegistry owner、Module Contract V2、typed Event、Outbox durable intent、请求关联和插件运行时 getter 已形成当前路径。插件仓适配、未知第三方 fallback 和其它 E1/E3 副作用仍按风险持续治理。 -> 最新阶段:阶段 62 已统一消息渠道长连接的关闭收敛合同。 +> 最新阶段:阶段 63 已补齐应用消息队列线程的关闭收敛合同。 ## 当前复核结论(2026-08-24) @@ -657,6 +657,21 @@ 只扩展为可选布尔结果,既有返回 `None` 的宿主模块和 V1/V2/V3 插件仍按成功兼容处理。未修改插件仓、 SDK/Compat 映射或事件 payload。 +### 长期整改阶段 63:应用消息队列线程关闭收敛合同补齐(2026-08-24) + +- `MessageQueueManager` 的监控线程仍使用无界 `join()`;渠道同步发送回调一旦阻塞,`stop_modules()` 所在 + 事件循环也会同步卡住,外层生命周期的 300 秒协程预算无法中断它。`stop_message()` 同时丢弃关闭结果, + 形成阶段 61 布尔收敛事实源之外的遗漏。 +- 当前停止入口使用 10 秒默认预算并返回真实线程终态;超时时保留原线程 owner,后续可在回调返回后重试。 + 监控循环收到停止请求后不再继续取新的排队消息,模板缓存仍会独立尽力关闭,任一资源失败统一向 + `stop_modules()` 返回 `False`,且不跳过其余模块资源收口。 +- `MessageQueueManager()`、无参数 `stop()`、`stop_message()` 与 `app.helper.message` 精确映射全部保留;新增 + 返回值对忽略结果的旧调用兼容,未修改消息内容、调度时段、发送回调、SDK/Compat 映射或插件仓, + V1/V2/V3 插件边界不变。 +- 阻塞回调、有限返回、保留 owner、重试终止、模板缓存继续清理和 startup 失败传播均有故障注入; + 消息/生命周期专项 58 项、架构与兼容专项 122 项、Pylint 10.00/10、strict mypy 39 文件、宿主与质量 + ratchet 以及四分片全量 `5941 passed, 3 skipped` 均通过。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index 5727653a0..390727fb5 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -197,6 +197,8 @@ ModuleManager 与 startup 组合根继续关闭其余资源但必须向上返回 等领域关闭入口必须直接传播 Runtime 的整体结果,不得以单个能力快照或无返回包装器覆盖失败。 消息渠道模块必须通过 `_MessageChannelModuleBase._stop_service_instances()` 聚合多实例关闭结果; 长连接、轮询或 Socket 服务只有在真实终止后才能返回成功,超时 owner 不得清空句柄。 +应用消息队列的监控线程遵守同一收敛语义:停止必须有限等待,回调阻塞导致线程仍存活时保留 owner +并向 startup 返回 `False`,不得用无界 `join()` 阻塞生命周期或把日志当作成功。 API 中允许丢失或可重建的进程内任务必须登记到 `app/runtime/tasks.py`;登记器先于其他 运行资源启动,并在资源释放前停止接收、取消和有限等待。需要崩溃恢复的 E2/E3 副作用仍应 进入 Outbox 或持久任务表,不能把 TaskRegistry 当成 durable queue。 diff --git a/tests/test_lifecycle_shutdown.py b/tests/test_lifecycle_shutdown.py index 007b1d9e9..d0e20b1af 100644 --- a/tests/test_lifecycle_shutdown.py +++ b/tests/test_lifecycle_shutdown.py @@ -957,6 +957,18 @@ def test_stop_modules_propagates_false_without_skipping_later_cleanup(monkeypatc _assert_completed_once(dependency) +def test_stop_modules_propagates_message_queue_nonconvergence(monkeypatch): + """消息队列线程未终止时必须由模块服务关闭结果向上暴露。""" + dependencies = _patch_module_shutdown_dependencies(monkeypatch) + dependencies["stop_message"].return_value = False + + converged = asyncio.run(modules_initializer.stop_modules()) + + assert converged is False + for dependency in dependencies.values(): + _assert_completed_once(dependency) + + def test_stop_modules_drains_web_agent_tasks_before_persistence(monkeypatch): """关闭时先收口 Web Agent,再关闭持久化准入和数据库任务。""" order = [] diff --git a/tests/test_message_queue_shutdown.py b/tests/test_message_queue_shutdown.py index 199fc4656..6c5b746b5 100644 --- a/tests/test_message_queue_shutdown.py +++ b/tests/test_message_queue_shutdown.py @@ -1,4 +1,6 @@ +import threading import time +from unittest.mock import MagicMock from app.application.messaging.message import MessageQueueManager, TemplateHelper, stop_message from app.foundation.singleton import SingletonClass @@ -11,20 +13,69 @@ def test_message_queue_stop_wakes_idle_monitor(monkeypatch): manager.__init__(check_interval=10) started_at = time.monotonic() - manager.stop() + converged = manager.stop() elapsed = time.monotonic() - started_at + assert converged is True assert elapsed < 1 assert not manager.thread.is_alive() +def test_message_queue_stop_reports_blocked_callback_and_supports_retry(monkeypatch): + """发送回调仍阻塞时应有限返回 False,并保留线程供后续重试。""" + entered = threading.Event() + release = threading.Event() + + def blocked_send_callback(*_args, **_kwargs) -> None: + """模拟无法由消息队列强制取消的同步渠道调用。""" + entered.set() + release.wait() + + monkeypatch.setattr(MessageQueueManager, "init_config", lambda self: None) + manager = object.__new__(MessageQueueManager) + manager.__init__(send_callback=blocked_send_callback, check_interval=0) + manager.queue.put({"args": ("payload",), "kwargs": {}}) + + assert entered.wait(timeout=1) + try: + started_at = time.monotonic() + assert manager.stop(timeout=0.01) is False + assert time.monotonic() - started_at < 1 + assert manager.thread.is_alive() + finally: + release.set() + + assert manager.stop(timeout=1) is True + assert not manager.thread.is_alive() + + def test_stop_message_does_not_initialize_absent_services(monkeypatch): """消息服务未初始化时,关闭入口不应为了清理而创建后台资源""" monkeypatch.setattr(SingletonClass, "_instances", {}) assert MessageQueueManager.get_existing_instance() is None assert TemplateHelper.get_existing_instance() is None - stop_message() + assert stop_message() is True assert MessageQueueManager not in SingletonClass._instances assert TemplateHelper not in SingletonClass._instances + + +def test_stop_message_aggregates_queue_failure_and_closes_template(monkeypatch): + """消息队列未收敛时仍应关闭模板缓存并向生命周期返回 False。""" + queue_manager = MagicMock() + queue_manager.stop.return_value = False + template_helper = MagicMock() + monkeypatch.setattr( + SingletonClass, + "_instances", + { + MessageQueueManager: queue_manager, + TemplateHelper: template_helper, + }, + ) + + assert stop_message(timeout=0.25) is False + + queue_manager.stop.assert_called_once_with(timeout=0.25) + template_helper.close.assert_called_once_with()