From d89d29612c21a6ca734a3fae6b3a63e9215ad363 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Fri, 21 Aug 2026 19:40:20 +0800 Subject: [PATCH] refactor: enforce single-worker control plane --- app/doctor/checks.py | 45 +++++++ app/main.py | 5 + app/runtime/config.py | 2 +- app/runtime/topology.py | 40 ++++++ app/startup/lifecycle/__init__.py | 7 +- docs/architecture-overview.md | 3 +- docs/doctor.md | 5 + .../backend-architecture-next-stage.md | 2 +- docs/rules/05-architecture.md | 2 + .../architecture/dependency-baseline.json | 11 +- tests/test_process_topology.py | 116 ++++++++++++++++++ 11 files changed, 231 insertions(+), 7 deletions(-) create mode 100644 app/runtime/topology.py create mode 100644 tests/test_process_topology.py diff --git a/app/doctor/checks.py b/app/doctor/checks.py index 1c8555406..01f099e47 100644 --- a/app/doctor/checks.py +++ b/app/doctor/checks.py @@ -19,6 +19,7 @@ from urllib.request import Request, urlopen import psutil from app.runtime.config import settings +from app.runtime.topology import process_topology_issue from app.doctor.models import DoctorFinding, DoctorFindingStatus, DoctorReport, DoctorSeverity from app.adapters.system.host import SystemUtils @@ -134,6 +135,7 @@ def default_checks() -> list[CheckFunc]: return [ _check_runtime_paths, _check_config, + _check_process_topology, _check_processes_and_ports, _check_dependencies, _check_database, @@ -154,6 +156,49 @@ def _mask_text(text: str) -> str: return masked +def _check_process_topology(runner: DoctorRunnerProtocol) -> None: + """诊断 API worker 配置是否会复制全功能控制面。""" + issue = process_topology_issue( + workers=settings.API_WORKERS, + safe_mode=settings.MOVIEPILOT_SAFE_MODE, + ) + context = { + "api_workers": settings.API_WORKERS, + "safe_mode": settings.MOVIEPILOT_SAFE_MODE, + } + if issue: + runner.add( + finding_id="startup.process_topology", + severity=DoctorSeverity.Error, + status=DoctorFindingStatus.Failed, + title="进程拓扑不受支持", + detail=issue, + recommendation="将 API_WORKERS 设为 1 后重启 MoviePilot。", + context=context, + ) + return + if settings.API_WORKERS != 1: + runner.add( + finding_id="startup.process_topology", + severity=DoctorSeverity.Warn, + status=DoctorFindingStatus.Degraded, + title="安全模式临时使用多 worker", + detail="安全模式不会启动插件、调度器、监控器和工作流,因此当前不会复制完整控制面。", + recommendation="故障排除后恢复 API_WORKERS=1,再退出安全模式。", + context=context, + ) + return + runner.add( + finding_id="startup.process_topology", + severity=DoctorSeverity.Info, + status=DoctorFindingStatus.Ok, + title="进程拓扑受支持", + detail="API_WORKERS=1,插件和后台任务只会启动一份。", + recommendation="保持单 worker;扩容前需要先拆分 API 数据面和控制面。", + context=context, + ) + + def _read_json(path: Path) -> Optional[dict[str, Any]]: if not path.exists(): return None diff --git a/app/main.py b/app/main.py index 22c4aef79..a02500947 100644 --- a/app/main.py +++ b/app/main.py @@ -55,6 +55,7 @@ elif SystemUtils.is_frozen(): from app.factory import app from app.runtime.config import global_vars, settings +from app.runtime.topology import validate_process_topology from app.startup.database_initializer import prepare_database setproctitle.setproctitle(settings.PROJECT_NAME) @@ -133,6 +134,10 @@ def signal_handler(signum, frame): def run_application() -> None: """初始化进程并启动 API 服务""" + validate_process_topology( + workers=settings.API_WORKERS, + safe_mode=settings.MOVIEPILOT_SAFE_MODE, + ) signal.signal(signal.SIGTERM, signal_handler) signal.signal(signal.SIGINT, signal_handler) diff --git a/app/runtime/config.py b/app/runtime/config.py index 323e400b4..70f95d286 100644 --- a/app/runtime/config.py +++ b/app/runtime/config.py @@ -134,7 +134,7 @@ class ConfigModel(BaseModel): # 数据库连接额度校验按它换算总用量。注意:当前主程序以单进程方式启动 # (uvicorn.Config 的 workers 仅在多进程 supervisor 路径下生效), # 调大此项前需先解决调度器会在每个 worker 内重复执行的问题 - API_WORKERS: int = 1 + API_WORKERS: int = Field(default=1, ge=1) DB_TYPE: str = "sqlite" # 是否在控制台输出 SQL 语句,默认关闭 DB_ECHO: bool = False diff --git a/app/runtime/topology.py b/app/runtime/topology.py new file mode 100644 index 000000000..f85c4c31e --- /dev/null +++ b/app/runtime/topology.py @@ -0,0 +1,40 @@ +"""MoviePilot 进程拓扑约束。""" + +from typing import Optional + + +class UnsupportedProcessTopologyError(RuntimeError): + """启动配置会复制当前只能单实例运行的控制面。""" + + +def process_topology_issue(*, workers: int, safe_mode: bool) -> Optional[str]: + """ + 返回当前进程拓扑不可运行的原因。 + + :param workers: API worker 进程数 + :param safe_mode: 是否只启动安全模式数据面 + :return: 支持时返回 None,否则返回可直接展示的错误说明 + """ + if workers < 1: + return "API_WORKERS 必须大于等于 1" + if workers == 1 or safe_mode: + return None + return ( + "MoviePilot V3 全功能模式仅支持 API_WORKERS=1;" + f"当前配置为 {workers},每个 worker 都会重复启动插件、调度器、监控器和工作流。" + "请将 API_WORKERS 改为 1 后重启。故障排查可以临时启用 " + "MOVIEPILOT_SAFE_MODE=true,但安全模式不是全功能扩容方案。" + ) + + +def validate_process_topology(*, workers: int, safe_mode: bool) -> None: + """ + 在启动任何持久化或后台副作用前校验进程拓扑。 + + :param workers: API worker 进程数 + :param safe_mode: 是否只启动安全模式数据面 + :raises UnsupportedProcessTopologyError: 当前拓扑会复制全功能控制面 + """ + issue = process_topology_issue(workers=workers, safe_mode=safe_mode) + if issue: + raise UnsupportedProcessTopologyError(issue) diff --git a/app/startup/lifecycle/__init__.py b/app/startup/lifecycle/__init__.py index 82ca03522..2dfea9634 100644 --- a/app/startup/lifecycle/__init__.py +++ b/app/startup/lifecycle/__init__.py @@ -27,6 +27,7 @@ except Exception: from app.chain.system import SystemChain from app.application.plugin.runtime import get_plugin_manager from app.runtime.config import global_vars, settings +from app.runtime.topology import validate_process_topology from app.adapters.external.server import MoviePilotServerHelper from app.runtime.state import SystemHelper from app.runtime.log import logger, LoggerManager @@ -47,7 +48,7 @@ from app.startup.scheduler_initializer import ( init_scheduler, init_plugin_scheduler, ) -from app.db import check_connection_budget, get_engine, get_global_async_engine +from app.db.engine import check_connection_budget, get_engine, get_global_async_engine from app.startup.transfer_initializer import replay_pending_transfers from app.startup.workflow_initializer import init_workflow, stop_workflow from app.startup.lifecycle.components import ( @@ -277,6 +278,10 @@ async def lifespan(app: FastAPI): """ 定义应用的生命周期事件 """ + validate_process_topology( + workers=settings.API_WORKERS, + safe_mode=settings.MOVIEPILOT_SAFE_MODE, + ) print("Starting up...") # 存储当前循环 global_vars.set_loop(asyncio.get_event_loop()) diff --git a/docs/architecture-overview.md b/docs/architecture-overview.md index 488b522bc..cd302d866 100644 --- a/docs/architecture-overview.md +++ b/docs/architecture-overview.md @@ -166,7 +166,7 @@ flowchart TB |---|---|---| | `app/foundation/` | 无状态、无配置、无 I/O 的底层原语:反射/动态导入、加密、DOM、单例、文本、URL、版本比较 | `reflection.py`、`crypto.py`、`singleton.py` | | `app/domain/` | 纯 MoviePilot 业务语义:媒体上下文、识别解析、站点状态解释、磁力语义、NFO 刮削 | `context.py`、`metainfo.py`、`meta/`、`scraper.py` | -| `app/runtime/` | 进程级运行机制:配置、事件、完整日志、缓存契约与内存后端、并发、调度、限流、本地化、GC、重启状态 | `config.py`、`events.py`、`log.py`、`cache.py` | +| `app/runtime/` | 进程级运行机制:配置、进程拓扑、事件、完整日志、缓存契约与内存后端、并发、调度、限流、本地化、GC、重启状态 | `config.py`、`topology.py`、`events.py`、`log.py`、`cache.py` | | `app/runtime/extensions/` | 模块 / 插件 / 配置化服务 / 托管资源的发现、注册与生命周期适配;旧管理器文件保留稳定 ABI 门面,具体实现拆在主题子包 | `module_manager.py`、`plugin_manager.py`、`plugin/` | | `app/runtime/compat/` | 仅标准库的精确旧模块、包与符号导入路由;不是业务实现,也不是通用 re-export 层 | `manifest.py`、`imports.py` | | `app/adapters/network/` | 通用 HTTP、浏览器、DNS、Cloudflare、IP 传输机制 | `http.py`、`browser.py` | @@ -238,6 +238,7 @@ sequenceDiagram - **引擎预热 fail-fast**:同步/异步数据库引擎在单线程期完成首次创建, 避免调度器放出大量线程后再创建引擎导致连接锁竞争。 - **安全模式**:`MOVIEPILOT_SAFE_MODE` 会跳过插件、定时器、监控器、命令与工作流,用于故障自救。 +- **进程拓扑**:全功能 V3 强制 `API_WORKERS=1`,避免每个 worker 重复启动插件和后台控制面;安全模式可临时使用多 worker 诊断,但不是正式扩容方案。 - **关停隔离**:每个关停步骤由 `run_shutdown_step` 独立捕获异常,保证后续资源仍有机会释放。 --- diff --git a/docs/doctor.md b/docs/doctor.md index bfee63977..e6f44675f 100644 --- a/docs/doctor.md +++ b/docs/doctor.md @@ -37,6 +37,7 @@ Doctor 默认执行只读检查: - 运行路径:程序目录、配置目录、日志目录、Python 解释器 - 关键配置:`API_TOKEN`、`PORT`、`NGINX_PORT`、代理格式、安全模式 +- 进程拓扑:全功能模式必须使用 `API_WORKERS=1`,避免插件、调度器、监控器和工作流重复启动 - 进程与端口:后端、前端端口监听状态,runtime 文件是否过期 - 日志线索:后端日志、启动日志、前端日志和插件日志最近 24 小时内的错误 - 核心依赖:FastAPI、Pydantic、SQLAlchemy、Uvicorn、CloakBrowser 等是否可导入 @@ -71,6 +72,10 @@ Doctor 不会自动删除数据库、修改 Docker Compose、回滚迁移、禁 安全模式不修改用户配置,适合插件、调度任务或 Agent 导致后端无法启动时先恢复后台入口。修复问题后移除环境变量或使用普通 `moviepilot start` 重启即可恢复完整能力。 +MoviePilot V3 的全功能模式只支持 `API_WORKERS=1`。配置更大的值时,主入口和外部 ASGI +lifespan 都会在数据库迁移及后台任务启动前拒绝运行,Doctor 同时给出失败项。安全模式因跳过 +控制面而允许临时使用多 worker,但 Doctor 会将其标记为降级;故障排除后应恢复单 worker。 + ## Docker 诊断保活 Docker 镜像默认设置 `MOVIEPILOT_DOCKER_KEEPALIVE_ON_FAILURE=true`。当后端主进程非正常退出时,entrypoint 不会立刻退出容器,而是打印一次 doctor 报告并保持容器运行,方便执行: diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 3a1fda9fb..1ce3df85c 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(ARCH-201~203)已完成,后续阶段按任务 ID 独立提交和回滚 +> 实施进度:阶段 0(ARCH-201~203)与 ARCH-210 已完成,后续任务按 ID 独立提交和回滚 ## 1. 结论先行 diff --git a/docs/rules/05-architecture.md b/docs/rules/05-architecture.md index d88ade9b8..64e3891c6 100644 --- a/docs/rules/05-architecture.md +++ b/docs/rules/05-architecture.md @@ -82,6 +82,7 @@ create additional top-level directory categories. | Path | Ownership | |---|---| | `app/runtime/config.py` | Deployment configuration and resolved runtime settings | +| `app/runtime/topology.py` | Process topology policy shared by startup and offline diagnostics | | `app/runtime/events.py` | Event contracts, dispatch and resolver registration | | `app/runtime/log.py` | Complete console/plugin/file logging runtime and shutdown | | `app/runtime/cache.py` | Cache protocols, memory implementations, decorators and proxies | @@ -419,6 +420,7 @@ policy. `app/db` therefore has no dependency on `app/domain`. | `app/application/commands.py` | Command registry facade for Agent tools and endpoints; `Command` class registered by `app/startup/command_initializer.py` | | `app/chain/agent.py` | `AgentChain(ChainBase)`: the chain-layer entry for Agent sessions; Agent runtime stays in `app/agent/` | | `app/runtime/config.py` | `ConfigModel`, `Settings` and deployment configuration | +| `app/runtime/topology.py` | Single-worker full-runtime policy and safe-mode topology validation | | `app/runtime/events.py` | `EventManager`/`Event` compatibility facade and global `eventmanager` identity | | `app/runtime/event/registry.py` | Event subscriptions, enable/disable state and dispatch snapshots | | `app/runtime/event/binding.py` | Explicit module/plugin/host handler resolvers; unresolved classes are diagnosed and skipped, never implicitly constructed by the bus | diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index adcb6c2e2..d4aadc746 100644 --- a/tests/fixtures/architecture/dependency-baseline.json +++ b/tests/fixtures/architecture/dependency-baseline.json @@ -13,8 +13,8 @@ "runtime_to_db": [], "workflow_to_db": [] }, - "edge_count": 6076, - "edge_sha256": "98150e02051b59ca341c62a71ec050aaa3807b6f20eba40a3d1cda76682503d1", + "edge_count": 6080, + "edge_sha256": "57d091e93986ca0d8e8a93a300c3c6d0eb0c2bfe92fda747cffd829cc3c25cab", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -3507,6 +3507,7 @@ "app.doctor.checks -> app.doctor.models", "app.doctor.checks -> app.runtime", "app.doctor.checks -> app.runtime.config", + "app.doctor.checks -> app.runtime.topology", "app.doctor.formatters -> app.doctor", "app.doctor.formatters -> app.doctor.models", "app.doctor.runner -> app.adapters", @@ -3637,6 +3638,7 @@ "app.main -> app.factory", "app.main -> app.runtime", "app.main -> app.runtime.config", + "app.main -> app.runtime.topology", "app.main -> app.startup", "app.main -> app.startup.database_initializer", "app.modules -> app.runtime", @@ -5745,10 +5747,12 @@ "app.startup.lifecycle -> app.chain", "app.startup.lifecycle -> app.chain.system", "app.startup.lifecycle -> app.db", + "app.startup.lifecycle -> app.db.engine", "app.startup.lifecycle -> app.runtime", "app.startup.lifecycle -> app.runtime.config", "app.startup.lifecycle -> app.runtime.log", "app.startup.lifecycle -> app.runtime.state", + "app.startup.lifecycle -> app.runtime.topology", "app.startup.lifecycle -> app.startup", "app.startup.lifecycle -> app.startup.cache_initializer", "app.startup.lifecycle -> app.startup.command_initializer", @@ -6093,7 +6097,7 @@ "app.workflow.actions.transfer_file -> app.workflow", "app.workflow.actions.transfer_file -> app.workflow.actions" ], - "module_count": 753, + "module_count": 754, "modules": [ "app", "app.adapters", @@ -6751,6 +6755,7 @@ "app.runtime.scheduling", "app.runtime.state", "app.runtime.thread", + "app.runtime.topology", "app.scheduler", "app.schemas", "app.schemas.agent", diff --git a/tests/test_process_topology.py b/tests/test_process_topology.py new file mode 100644 index 000000000..78795ee76 --- /dev/null +++ b/tests/test_process_topology.py @@ -0,0 +1,116 @@ +"""V3 API 数据面与宿主控制面进程拓扑测试。""" + +import asyncio +from unittest.mock import MagicMock + +import pytest +from fastapi import FastAPI +from pydantic import ValidationError + +from app.doctor import checks +from app.doctor.models import DoctorFindingStatus +from app.doctor.runner import DoctorRunner +from app.runtime.config import ConfigModel, settings +from app.runtime.topology import ( + UnsupportedProcessTopologyError, + process_topology_issue, + validate_process_topology, +) +from app.startup import lifecycle + + +@pytest.mark.parametrize("safe_mode", [False, True]) +def test_single_worker_is_supported_in_every_mode(safe_mode: bool): + """单 worker 是正常模式和安全模式共同支持的默认拓扑。""" + assert process_topology_issue(workers=1, safe_mode=safe_mode) is None + validate_process_topology(workers=1, safe_mode=safe_mode) + + +def test_full_runtime_rejects_multiple_workers(): + """全功能模式不得复制插件、调度器、监控器和工作流。""" + with pytest.raises( + UnsupportedProcessTopologyError, + match="全功能模式仅支持 API_WORKERS=1", + ): + validate_process_topology(workers=2, safe_mode=False) + + +def test_safe_mode_temporarily_allows_multiple_workers(): + """安全模式跳过控制面时允许多 worker 作为故障诊断手段。""" + assert process_topology_issue(workers=2, safe_mode=True) is None + validate_process_topology(workers=2, safe_mode=True) + + +@pytest.mark.parametrize("workers", [0, -1, "invalid"]) +def test_api_workers_configuration_rejects_invalid_values(workers): + """worker 数量的边界或类型错误必须在配置解析阶段暴露。""" + with pytest.raises(ValidationError): + ConfigModel(API_WORKERS=workers) + + +def test_main_rejects_topology_before_startup_side_effects(monkeypatch): + """主入口应在注册信号、迁移数据库和启动服务器前拒绝错误拓扑。""" + from app import main + + monkeypatch.setattr(main.settings, "API_WORKERS", 2) + monkeypatch.setattr(main.settings, "MOVIEPILOT_SAFE_MODE", False) + signal_handler = MagicMock() + start_tray = MagicMock() + prepare_database = MagicMock() + server_run = MagicMock() + monkeypatch.setattr(main.signal, "signal", signal_handler) + monkeypatch.setattr(main, "start_tray", start_tray) + monkeypatch.setattr(main, "prepare_database", prepare_database) + monkeypatch.setattr(main.Server, "run", server_run) + + with pytest.raises(UnsupportedProcessTopologyError): + main.run_application() + + signal_handler.assert_not_called() + start_tray.assert_not_called() + prepare_database.assert_not_called() + server_run.assert_not_called() + + +def test_asgi_lifespan_rejects_topology_before_runtime_initialization(monkeypatch): + """外部 ASGI supervisor 也必须在生命周期副作用前执行同一校验。""" + monkeypatch.setattr(settings, "API_WORKERS", 2) + monkeypatch.setattr(settings, "MOVIEPILOT_SAFE_MODE", False) + set_loop = MagicMock() + monkeypatch.setattr(lifecycle.global_vars, "set_loop", set_loop) + + async def run_lifespan() -> None: + async with lifecycle.lifespan(FastAPI()): + pass + + with pytest.raises(UnsupportedProcessTopologyError): + asyncio.run(run_lifespan()) + + set_loop.assert_not_called() + + +def test_doctor_fails_unsupported_full_runtime_topology(monkeypatch): + """Doctor 应把正常模式多 worker 报告为影响整体状态的失败。""" + monkeypatch.setattr(settings, "API_WORKERS", 2) + monkeypatch.setattr(settings, "MOVIEPILOT_SAFE_MODE", False) + runner = DoctorRunner() + + checks._check_process_topology(runner) + + finding = runner.report.find("startup.process_topology") + assert finding is not None + assert finding.status == DoctorFindingStatus.Failed + assert finding.context == {"api_workers": 2, "safe_mode": False} + + +def test_doctor_marks_safe_mode_multi_worker_as_degraded(monkeypatch): + """安全模式多 worker 可运行,但 Doctor 必须提醒它不是正式扩容方案。""" + monkeypatch.setattr(settings, "API_WORKERS", 2) + monkeypatch.setattr(settings, "MOVIEPILOT_SAFE_MODE", True) + runner = DoctorRunner() + + checks._check_process_topology(runner) + + finding = runner.report.find("startup.process_topology") + assert finding is not None + assert finding.status == DoctorFindingStatus.Degraded