refactor: own legacy subscription report tasks

This commit is contained in:
jxxghp
2026-08-24 05:57:07 +08:00
parent 9abeda8a05
commit 073ff7ff81
5 changed files with 105 additions and 16 deletions
+43 -12
View File
@@ -1,12 +1,15 @@
import asyncio
import json
import platform
from collections.abc import Coroutine
from pathlib import Path
from threading import Thread
from typing import Any, Dict, List, Optional, Tuple, Union
from typing import Any, Callable, Dict, List, Optional, Tuple, Union
from urllib.parse import parse_qs, quote, urlparse, urlsplit
from app.runtime.cache import cached
from app.runtime.config import global_vars
from app.runtime.settings import RuntimeSettingsCompat
from app.runtime.tasks import get_task_registry
from app.domain.context import MediaInfo, MusicInfo
from app.domain.meta.metabase import MetaBase
from app.runtime.log import logger
@@ -955,21 +958,49 @@ class MoviePilotServerHelper:
return True
return await cls.async_sub_done(sub)
@staticmethod
def _submit_statistic_report(
coroutine: Coroutine[Any, Any, Any],
*,
submit: Callable[[Coroutine[Any, Any, Any]], object],
action: str,
) -> bool:
"""提交兼容统计上报,并在宿主不再接收任务时关闭未执行协程。"""
try:
submit(coroutine)
return True
except Exception as err:
coroutine.close()
logger.warning(f"调度{action}失败:{err}")
return False
@classmethod
def sub_reg_async(cls, sub: dict) -> bool:
"""
开线程新增订阅统计。
"""
Thread(target=cls.sub_reg, args=(sub,)).start()
return True
"""兼容旧同步入口,在宿主任务登记器中提交新增订阅统计。"""
return cls._submit_statistic_report(
asyncio.to_thread(cls.sub_reg, sub),
submit=lambda coroutine: get_task_registry().submit_threadsafe(
coroutine,
loop=global_vars.loop,
owner="compat.server.subscribe_added_report",
cancel_on_shutdown=False,
),
action="新增订阅统计上报",
)
@classmethod
def sub_done_async(cls, sub: dict) -> bool:
"""
开线程完成订阅统计。
"""
Thread(target=cls.sub_done, args=(sub,)).start()
return True
"""兼容旧同步入口,在宿主任务登记器中提交订阅完成统计。"""
return cls._submit_statistic_report(
asyncio.to_thread(cls.sub_done, sub),
submit=lambda coroutine: get_task_registry().submit_threadsafe(
coroutine,
loop=global_vars.loop,
owner="compat.server.subscribe_done_report",
cancel_on_shutdown=False,
),
action="订阅完成统计上报",
)
@classmethod
def sub_report(cls) -> bool:
@@ -74,6 +74,9 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同
timer、刷新剩余摘要并等待已启动回调,不再把 `create_task` 留给事件循环隐式回收。
- 主仓不再新增或保留裸 FastAPI `BackgroundTasks`;若任务源于已提交的用户数据且不可从数据库重建,
必须提升为 E2,进入 Outbox 或持久任务表。
- 旧插件可调用的 `MoviePilotServerHelper.sub_reg_async()` / `sub_done_async()` 保留同步 ABI,但内部不再创建
裸上报线程;任务分别登记为 `compat.server.subscribe_added_report` / `subscribe_done_report`,已开始的
同步网络工作在 shutdown 时不取消并等待完成。canonical 订阅主链继续使用 durable outbox,不回退旧入口。
### Scheduler jobs
@@ -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。
> 实施进度:阶段 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 已托管旧插件订阅统计线程
## 当前复核结论(2026-08-24
@@ -288,13 +288,24 @@
丢失,不把 UI 进度误报为 durable 完成。既有 `app.api.endpoints.history -> app.runtime.tasks` 依赖边不变,
V2/V3 插件 SDK/Compat 无改动。
### 长期整改阶段 28:插件兼容统计线程托管(2026-08-24)
- canonical 订阅新增/删除/完成统计已由事务 Outbox 驱动,但 `MoviePilotServerHelper.sub_reg_async()`
`sub_done_async()` 仍各自创建裸线程;V3 官方插件仍通过旧 `app.helper.server` 映射调用后者,不能删除 ABI。
- 两个旧同步入口保留类、方法、参数与立即布尔返回,内部改经 TaskRegistry 跨线程提交 `asyncio.to_thread`
分别登记 `compat.server.subscribe_added_report` / `compat.server.subscribe_done_report`;同步网络工作开始后
不在 shutdown 时取消,宿主会等待真实完成,生命周期不可用时拒绝并关闭未执行协程。
- 插件仓没有改动,V1/V2/V3 旧导入映射保持不变;主链也不回退兼容入口。依赖基线仅新增
`app.adapters.external.server -> app.runtime.config/tasks` 两条允许边,模块保持 `806`、内部边为 `6543`
12 组禁止边与唯一隔离 TMDB SCC 均未变化。
### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
- 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。
- `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。
- 依赖图当前为 `806` 个 Python 模块、`6541` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
- 依赖图当前为 `806` 个 Python 模块、`6543` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。
- 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。
综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。
+4 -2
View File
@@ -13,8 +13,8 @@
"runtime_to_db": [],
"workflow_to_db": []
},
"edge_count": 6541,
"edge_sha256": "c7e2a8537efa709d30c82930fa4ee21d91fe8f24578a44b4357563fa7bfc1a0c",
"edge_count": 6543,
"edge_sha256": "b9b6ba268382e7027f5557643defb09848c054cb68079345f2394d9b7faeb42c",
"edges": [
"app -> app.runtime",
"app -> app.runtime.compat",
@@ -89,9 +89,11 @@
"app.adapters.external.server -> app.domain.meta.metabase",
"app.adapters.external.server -> app.runtime",
"app.adapters.external.server -> app.runtime.cache",
"app.adapters.external.server -> app.runtime.config",
"app.adapters.external.server -> app.runtime.log",
"app.adapters.external.server -> app.runtime.observability",
"app.adapters.external.server -> app.runtime.settings",
"app.adapters.external.server -> app.runtime.tasks",
"app.adapters.external.server -> app.schemas",
"app.adapters.external.server -> app.schemas.media",
"app.adapters.external.server -> app.schemas.types",
+42
View File
@@ -0,0 +1,42 @@
"""服务端统计兼容入口的后台任务生命周期回归。"""
from unittest.mock import Mock, patch
from app.adapters.external.server import MoviePilotServerHelper
from app.runtime.config import global_vars
def test_legacy_subscription_reports_use_owned_threadsafe_tasks() -> None:
"""旧同步统计入口应保留 ABI,并让宿主等待已开始的线程工作。"""
loop = Mock(**{"is_running.return_value": True, "is_closed.return_value": False})
registry = Mock()
def submit(coroutine, **_kwargs):
"""关闭测试协程,避免替身提交留下未等待警告。"""
coroutine.close()
return Mock()
registry.submit_threadsafe.side_effect = submit
with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch(
"app.adapters.external.server.get_task_registry", return_value=registry
):
assert MoviePilotServerHelper.sub_reg_async({"media_id": "1"}) is True
assert MoviePilotServerHelper.sub_done_async({"media_id": "1"}) is True
assert registry.submit_threadsafe.call_count == 2
assert registry.submit_threadsafe.call_args_list[0].kwargs == {
"loop": loop,
"owner": "compat.server.subscribe_added_report",
"cancel_on_shutdown": False,
}
assert registry.submit_threadsafe.call_args_list[1].kwargs == {
"loop": loop,
"owner": "compat.server.subscribe_done_report",
"cancel_on_shutdown": False,
}
def test_legacy_subscription_report_rejects_without_runtime_loop() -> None:
"""宿主生命周期不可用时应拒绝提交,并保持布尔返回合同。"""
with patch.object(global_vars, "CURRENT_EVENT_LOOP", None):
assert MoviePilotServerHelper.sub_done_async({"media_id": "1"}) is False