From 01d29d952ffaab19308bae0b6cfcc2b703cc1f97 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 08:29:07 +0800 Subject: [PATCH] fix: deduplicate graceful restart fallback --- app/runtime/state.py | 45 +++- .../backend-architecture-next-stage.md | 14 +- tests/test_system_utils.py | 233 +++++++++++------- 3 files changed, 191 insertions(+), 101 deletions(-) diff --git a/app/runtime/state.py b/app/runtime/state.py index 338933c31..ac5cf149c 100644 --- a/app/runtime/state.py +++ b/app/runtime/state.py @@ -38,6 +38,8 @@ class SystemHelper(ConfigReloadMixin): __local_restart_log_file = settings.LOG_PATH / "moviepilot.restart.stdout.log" __one_shot_update_flag_file = settings.TEMP_PATH / "moviepilot.pending_update" __docker_restart_intent_file = settings.TEMP_PATH / "moviepilot.intentional_restart" + __graceful_shutdown_monitor_lock = threading.Lock() + __graceful_shutdown_monitor: Optional[threading.Thread] = None def on_config_changed(self): """配置变化后重新应用日志设置。""" @@ -360,21 +362,44 @@ class SystemHelper(ConfigReloadMixin): @staticmethod def _start_graceful_shutdown_monitor(): """ - 启动优雅退出超时监控 - 如果30秒内进程没有退出,则使用Docker API强制重启 + 启动唯一的优雅退出超时监控。 + + 如果 180 秒内进程没有退出,则使用 Docker API 强制重启;重复重启请求 + 复用当前 monitor,避免并行触发多次容器重启。 """ def monitor_thread(): - time.sleep(180) # 等待180秒 - logger.warning("优雅退出超时180秒,使用Docker API强制重启...") try: - SystemHelper._docker_api_restart() - except Exception as e: - logger.error(f"强制重启失败: {str(e)}") + time.sleep(180) + logger.warning("优雅退出超时180秒,使用Docker API强制重启...") + try: + SystemHelper._docker_api_restart() + except Exception as e: + logger.error(f"强制重启失败: {str(e)}") + finally: + with SystemHelper.__graceful_shutdown_monitor_lock: + if ( + SystemHelper.__graceful_shutdown_monitor + is threading.current_thread() + ): + SystemHelper.__graceful_shutdown_monitor = None - # 在后台线程中启动监控 - thread = threading.Thread(target=monitor_thread, daemon=True) - thread.start() + with SystemHelper.__graceful_shutdown_monitor_lock: + running = SystemHelper.__graceful_shutdown_monitor + if running is not None and running.is_alive(): + logger.debug("优雅退出超时监控已在运行,跳过重复启动") + return + thread = threading.Thread( + target=monitor_thread, + name="MoviePilot-GracefulRestartFallback", + daemon=True, + ) + SystemHelper.__graceful_shutdown_monitor = thread + try: + thread.start() + except BaseException: + SystemHelper.__graceful_shutdown_monitor = None + raise @staticmethod def _docker_api_restart() -> Tuple[bool, str]: diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 5a1f60b89..ace6dcc7f 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 已补齐异步防抖取消的终态所有权。 +> 实施进度:阶段 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 已统一优雅重启兜底线程的唯一所有权。 ## 当前复核结论(2026-08-24) @@ -426,6 +426,18 @@ - 回归测试覆盖旧任务在取消清理中阻塞、后续任务同时存在的场景。`debounce()` 装饰器、leading/trailing 触发规则、同步 `Debouncer`、`app.utils.debounce` 兼容映射和 V1/V2/V3 插件导入路径均未修改。 +### 长期整改阶段 41:优雅重启兜底线程所有权收口(2026-08-24) + +- API、命令链、升级和启动期资源更新都可能调用 `SystemHelper.restart()`;Docker 优雅退出路径原先每次 + 请求都会创建一个独立 daemon monitor,连续请求会在 180 秒后并行触发多次强制容器重启。现在 monitor + 由 `SystemHelper` 以锁和唯一句柄持有,存活期间的重复请求复用同一 owner,线程启动失败或执行结束都会 + 释放句柄,后续真实重启仍可创建新的兜底监控。 +- 回归测试用事件屏障验证首线程存活期间不会重复创建,并验证强制重启只调用一次、终态 owner 被释放; + 同时将触发说明从过期的 30 秒更正为实际的 180 秒。SIGTERM 优雅退出、Docker API 回退、意图标记、 + 本地 CLI 与一次性升级语义均未改变。 +- 该 daemon monitor 故意跨越正常 shutdown drain,以便进程卡死时仍能请求 Docker 重启,因此不纳入 + `TaskRegistry` 或普通线程池等待。`SystemHelper` 公开方法、SDK/Compat 映射和 V1/V2/V3 插件 ABI 均未修改。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/tests/test_system_utils.py b/tests/test_system_utils.py index df4bd6107..4e8093c45 100644 --- a/tests/test_system_utils.py +++ b/tests/test_system_utils.py @@ -6,9 +6,9 @@ import struct import subprocess import sys import tempfile +import threading import time from pathlib import Path -from unittest import TestCase from unittest.mock import MagicMock, call, patch import psutil @@ -19,108 +19,161 @@ from app.runtime.config import ConfigModel, settings from app.adapters.system.host import SystemUtils -class SystemUtilsTest(TestCase): +def test_get_config_path_uses_repository_config_for_source_runtime(): + """源码运行时应从项目根目录读取配置,不能随适配器目录层级偏移。""" + expected = Path(__file__).resolve().parents[1] / "config" - def test_get_config_path_uses_repository_config_for_source_runtime(self): - """源码运行时应从项目根目录读取配置,不能随适配器目录层级偏移。""" - expected = Path(__file__).resolve().parents[1] / "config" + with patch.dict(os.environ, {}, clear=True), \ + patch.object(SystemUtils, "is_docker", return_value=False), \ + patch.object(SystemUtils, "is_frozen", return_value=False): + assert SystemUtils.get_config_path() == expected + assert SystemUtils.get_env_path() == expected / "app.env" - with patch.dict(os.environ, {}, clear=True), \ - patch.object(SystemUtils, "is_docker", return_value=False), \ - patch.object(SystemUtils, "is_frozen", return_value=False): - self.assertEqual(SystemUtils.get_config_path(), expected) - self.assertEqual(SystemUtils.get_env_path(), expected / "app.env") - def test_get_config_path_preserves_explicit_config_dir(self): - """显式配置目录始终优先于运行环境推导。""" - explicit_path = Path("/custom/moviepilot-config") +def test_get_config_path_preserves_explicit_config_dir(): + """显式配置目录始终优先于运行环境推导。""" + explicit_path = Path("/custom/moviepilot-config") - with patch.dict(os.environ, {"CONFIG_DIR": "/ignored"}, clear=True), \ - patch.object(SystemUtils, "is_docker", return_value=True): - self.assertEqual( - SystemUtils.get_config_path(str(explicit_path)), - explicit_path, - ) + with patch.dict(os.environ, {"CONFIG_DIR": "/ignored"}, clear=True), \ + patch.object(SystemUtils, "is_docker", return_value=True): + assert SystemUtils.get_config_path(str(explicit_path)) == explicit_path - def test_get_config_path_preserves_runtime_specific_defaults(self): - """容器和冻结程序继续使用各自稳定的配置目录。""" - with patch.dict(os.environ, {}, clear=True), \ - patch.object(SystemUtils, "is_docker", return_value=True): - self.assertEqual(SystemUtils.get_config_path(), Path("/config")) - with patch.dict(os.environ, {}, clear=True), \ - patch.object(SystemUtils, "is_docker", return_value=False), \ - patch.object(SystemUtils, "is_frozen", return_value=True), \ - patch("app.adapters.system.host.sys.executable", "/opt/moviepilot/moviepilot"): - self.assertEqual( - SystemUtils.get_config_path(), - Path("/opt/moviepilot/config"), - ) +def test_get_config_path_preserves_runtime_specific_defaults(): + """容器和冻结程序继续使用各自稳定的配置目录。""" + with patch.dict(os.environ, {}, clear=True), \ + patch.object(SystemUtils, "is_docker", return_value=True): + assert SystemUtils.get_config_path() == Path("/config") - def test_execute_with_subprocess_keeps_stdout_when_command_fails(self): - """ - 命令失败时如果原因只写入 stdout,也需要回传给调用方用于错误提示。 - """ - error = subprocess.CalledProcessError( - returncode=1, - cmd=["pip", "check"], - output="demo requires pkg>=2, but you have pkg 1\n", - stderr="", + with patch.dict(os.environ, {}, clear=True), \ + patch.object(SystemUtils, "is_docker", return_value=False), \ + patch.object(SystemUtils, "is_frozen", return_value=True), \ + patch("app.adapters.system.host.sys.executable", "/opt/moviepilot/moviepilot"): + assert SystemUtils.get_config_path() == Path("/opt/moviepilot/config") + + +def test_execute_with_subprocess_keeps_stdout_when_command_fails(): + """命令失败时如果原因只写入 stdout,也需要回传给调用方用于错误提示。""" + error = subprocess.CalledProcessError( + returncode=1, + cmd=["pip", "check"], + output="demo requires pkg>=2, but you have pkg 1\n", + stderr="", + ) + + with patch("app.adapters.system.host.subprocess.run", side_effect=error): + success, message = SystemUtils.execute_with_subprocess(["pip", "check"]) + + assert not success + assert "返回码:1" in message + assert "标准输出:demo requires pkg>=2, but you have pkg 1" in message + + +def test_execute_with_subprocess_reports_empty_failure_output(): + """命令失败且没有输出时应给出明确占位信息,避免错误原因看起来被截断。""" + error = subprocess.CalledProcessError( + returncode=2, + cmd=["pip", "check"], + output="", + stderr="", + ) + + with patch("app.adapters.system.host.subprocess.run", side_effect=error): + success, message = SystemUtils.execute_with_subprocess(["pip", "check"]) + + assert not success + assert "返回码:2" in message + assert "无标准输出或错误输出" in message + + +def test_docker_restart_policy_marks_intent_before_sigterm(): + """Docker 优雅重启前应写入意图标记,避免 entrypoint 误进入 doctor 保活。""" + with tempfile.TemporaryDirectory() as temp_dir: + original_config_dir = settings.CONFIG_DIR + original_intent_file = SystemHelper._SystemHelper__docker_restart_intent_file + settings.CONFIG_DIR = temp_dir + SystemHelper._SystemHelper__docker_restart_intent_file = ( + settings.TEMP_PATH / "moviepilot.intentional_restart" ) + try: + with patch("app.runtime.state.is_docker", return_value=True), \ + patch.object(SystemHelper, "_check_restart_policy", return_value=True), \ + patch.object(SystemHelper, "_start_graceful_shutdown_monitor"), \ + patch("app.runtime.state.os.kill") as kill_mock: + ret, msg = SystemHelper.restart() - with patch("app.adapters.system.host.subprocess.run", side_effect=error): - success, message = SystemUtils.execute_with_subprocess(["pip", "check"]) - - self.assertFalse(success) - self.assertIn("返回码:1", message) - self.assertIn("标准输出:demo requires pkg>=2, but you have pkg 1", message) - - def test_execute_with_subprocess_reports_empty_failure_output(self): - """ - 命令失败且没有任何输出时,给出明确占位信息,避免错误原因看起来被截断。 - """ - error = subprocess.CalledProcessError( - returncode=2, - cmd=["pip", "check"], - output="", - stderr="", - ) - - with patch("app.adapters.system.host.subprocess.run", side_effect=error): - success, message = SystemUtils.execute_with_subprocess(["pip", "check"]) - - self.assertFalse(success) - self.assertIn("返回码:2", message) - self.assertIn("无标准输出或错误输出", message) + assert ret + assert msg == "" + assert (settings.TEMP_PATH / "moviepilot.intentional_restart").exists() + kill_mock.assert_called_once() + finally: + SystemHelper._SystemHelper__docker_restart_intent_file = original_intent_file + settings.CONFIG_DIR = original_config_dir -class SystemHelperRestartTest(TestCase): +def test_graceful_shutdown_monitor_has_single_owner_and_releases_it(monkeypatch): + """重复重启请求应共享唯一兜底线程,线程结束后必须释放 owner。""" + sleep_started = threading.Event() + release_sleep = threading.Event() + restart = MagicMock(return_value=(True, "")) + monitor_attr = "_SystemHelper__graceful_shutdown_monitor" + original_monitor = getattr(SystemHelper, monitor_attr) + thread = None - def test_docker_restart_policy_marks_intent_before_sigterm(self): - """ - Docker 内置重启走优雅退出时,应写入意图标记,避免 entrypoint 误进入 doctor 保活。 - """ - with tempfile.TemporaryDirectory() as temp_dir: - original_config_dir = settings.CONFIG_DIR - original_intent_file = SystemHelper._SystemHelper__docker_restart_intent_file - settings.CONFIG_DIR = temp_dir - SystemHelper._SystemHelper__docker_restart_intent_file = ( - settings.TEMP_PATH / "moviepilot.intentional_restart" - ) - try: - with patch("app.runtime.state.is_docker", return_value=True), \ - patch.object(SystemHelper, "_check_restart_policy", return_value=True), \ - patch.object(SystemHelper, "_start_graceful_shutdown_monitor"), \ - patch("app.runtime.state.os.kill") as kill_mock: - ret, msg = SystemHelper.restart() + def wait_for_shutdown(_seconds: float) -> None: + """用事件屏障模拟 180 秒等待,确保第二次启动发生在首线程存活期间。""" + sleep_started.set() + release_sleep.wait(timeout=1) - self.assertTrue(ret) - self.assertEqual(msg, "") - self.assertTrue((settings.TEMP_PATH / "moviepilot.intentional_restart").exists()) - kill_mock.assert_called_once() - finally: - SystemHelper._SystemHelper__docker_restart_intent_file = original_intent_file - settings.CONFIG_DIR = original_config_dir + setattr(SystemHelper, monitor_attr, None) + try: + monkeypatch.setattr("app.runtime.state.time.sleep", wait_for_shutdown) + monkeypatch.setattr(SystemHelper, "_docker_api_restart", restart) + + SystemHelper._start_graceful_shutdown_monitor() + assert sleep_started.wait(timeout=1) + thread = getattr(SystemHelper, monitor_attr) + SystemHelper._start_graceful_shutdown_monitor() + + assert getattr(SystemHelper, monitor_attr) is thread + release_sleep.set() + thread.join(timeout=1) + + assert thread.is_alive() is False + assert getattr(SystemHelper, monitor_attr) is None + restart.assert_called_once_with() + finally: + release_sleep.set() + if thread is not None: + thread.join(timeout=1) + setattr(SystemHelper, monitor_attr, original_monitor) + + +def test_graceful_shutdown_monitor_releases_owner_when_thread_start_fails(monkeypatch): + """兜底线程启动失败时必须释放 owner,允许后续请求重试。""" + monitor_attr = "_SystemHelper__graceful_shutdown_monitor" + original_monitor = getattr(SystemHelper, monitor_attr) + + class FailingThread: + """模拟在登记 owner 后启动失败的线程对象。""" + + def __init__(self, **_kwargs): + """接收真实 Thread 构造参数,但不创建系统线程。""" + + def start(self): + """模拟底层线程资源不足导致的启动失败。""" + raise RuntimeError("thread start failed") + + setattr(SystemHelper, monitor_attr, None) + try: + monkeypatch.setattr("app.runtime.state.threading.Thread", FailingThread) + + with pytest.raises(RuntimeError, match="thread start failed"): + SystemHelper._start_graceful_shutdown_monitor() + + assert getattr(SystemHelper, monitor_attr) is None + finally: + setattr(SystemHelper, monitor_attr, original_monitor) def test_execute_with_subprocess_passes_env_to_subprocess():