mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-08-29 03:56:43 +08:00
refactor: contain plugin catalog child tasks
This commit is contained in:
@@ -148,54 +148,62 @@ class PluginCatalogService:
|
||||
plugins = await loader(market, package_version, force)
|
||||
return task_index, result_version, plugins or []
|
||||
|
||||
tasks = []
|
||||
tasks: list[asyncio.Task[tuple[int, str, list[Any]]]] = []
|
||||
for market in markets:
|
||||
tasks.append(asyncio.create_task(
|
||||
fetch(market, None, "base_version", len(tasks))
|
||||
fetch(market, None, "base_version", len(tasks)),
|
||||
name="plugin.catalog.fetch",
|
||||
))
|
||||
for flag in compatible_flags:
|
||||
tasks.append(asyncio.create_task(
|
||||
fetch(market, flag, "higher_version", len(tasks))
|
||||
fetch(market, flag, "higher_version", len(tasks)),
|
||||
name="plugin.catalog.fetch",
|
||||
))
|
||||
|
||||
higher_plugins = []
|
||||
base_plugins = []
|
||||
if tasks:
|
||||
total_tasks = len(tasks)
|
||||
finished_tasks = 0
|
||||
task_results = {}
|
||||
if progress_callback:
|
||||
progress_callback(
|
||||
value=0,
|
||||
text=f"开始刷新插件市场,共 {total_tasks} 个请求 ...",
|
||||
data={"total": total_tasks, "finished": 0},
|
||||
)
|
||||
for completed_task in asyncio.as_completed(tasks):
|
||||
try:
|
||||
task_index, version, plugins = await completed_task
|
||||
task_results[task_index] = (version, plugins)
|
||||
except Exception as err:
|
||||
self._error(f"获取插件市场数据失败:{str(err)}")
|
||||
finished_tasks += 1
|
||||
try:
|
||||
higher_plugins = []
|
||||
base_plugins = []
|
||||
if tasks:
|
||||
total_tasks = len(tasks)
|
||||
finished_tasks = 0
|
||||
task_results = {}
|
||||
if progress_callback:
|
||||
progress_callback(
|
||||
value=finished_tasks / total_tasks * 100,
|
||||
text=(
|
||||
f"插件市场请求({finished_tasks}/{total_tasks})"
|
||||
"处理完成"
|
||||
),
|
||||
data={"total": total_tasks, "finished": finished_tasks},
|
||||
value=0,
|
||||
text=f"开始刷新插件市场,共 {total_tasks} 个请求 ...",
|
||||
data={"total": total_tasks, "finished": 0},
|
||||
)
|
||||
for task_index in sorted(task_results):
|
||||
version, plugins = task_results[task_index]
|
||||
(higher_plugins if version == "higher_version" else base_plugins).extend(
|
||||
plugins
|
||||
)
|
||||
for completed_task in asyncio.as_completed(tasks):
|
||||
try:
|
||||
task_index, version, plugins = await completed_task
|
||||
task_results[task_index] = (version, plugins)
|
||||
except Exception as err:
|
||||
self._error(f"获取插件市场数据失败:{str(err)}")
|
||||
finished_tasks += 1
|
||||
if progress_callback:
|
||||
progress_callback(
|
||||
value=finished_tasks / total_tasks * 100,
|
||||
text=(
|
||||
f"插件市场请求({finished_tasks}/{total_tasks})"
|
||||
"处理完成"
|
||||
),
|
||||
data={"total": total_tasks, "finished": finished_tasks},
|
||||
)
|
||||
for task_index in sorted(task_results):
|
||||
version, plugins = task_results[task_index]
|
||||
target = higher_plugins if version == "higher_version" else base_plugins
|
||||
target.extend(plugins)
|
||||
|
||||
result = self.merge(higher_plugins, base_plugins, markets)
|
||||
if progress_callback:
|
||||
progress_callback(value=100, text="插件市场缓存刷新完成")
|
||||
return result
|
||||
result = self.merge(higher_plugins, base_plugins, markets)
|
||||
if progress_callback:
|
||||
progress_callback(value=100, text="插件市场缓存刷新完成")
|
||||
return result
|
||||
finally:
|
||||
for task in tasks:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
if tasks:
|
||||
await asyncio.gather(*tasks, return_exceptions=True)
|
||||
|
||||
def merge(
|
||||
self,
|
||||
|
||||
@@ -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 系条目转换并清零重复代码白名单。
|
||||
> 实施进度:阶段 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 已收口插件市场请求级子任务。
|
||||
|
||||
## 当前复核结论(2026-08-24)
|
||||
|
||||
@@ -310,6 +310,16 @@
|
||||
- 依赖基线只新增 `app.application.mediaserver -> app.runtime.log` 及父包边,用于保持原有转换异常日志;模块
|
||||
保持 `806`、内部边为 `6545`,12 组禁止边与唯一隔离 TMDB SCC 均未变化。
|
||||
|
||||
### 长期整改阶段 30:插件市场请求级子任务收口(2026-08-24)
|
||||
|
||||
- 插件目录异步聚合原先用 `asyncio.create_task()` 并发读取多个市场和 V1/V2/V3 代际,正常路径会逐个
|
||||
等待,但父请求取消或进度回调抛错时会直接退出,尚未完成的 loader 继续在事件循环中运行。
|
||||
- 这些任务只服务当前 API/Agent 请求,不进入 lifespan TaskRegistry;`async_collect()` 现在通过
|
||||
`try/finally` 持有完整任务集合,异常退出时取消未完成 loader,并用 `gather(return_exceptions=True)`
|
||||
观察所有终态。单 loader 失败隔离、完成顺序进度和稳定市场/代际合并顺序均保持不变。
|
||||
- 回归测试覆盖父任务取消和进度消费者失败,均证明阻塞 loader 收到取消后才允许父调用结束;插件管理器
|
||||
方法、目录 DTO、市场参数、缓存、V1/V2/V3 索引与插件 SDK/Compat 均未修改。依赖图仍为 `806/6545`。
|
||||
|
||||
### 总体判断
|
||||
|
||||
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
|
||||
@@ -334,6 +344,8 @@
|
||||
`shield` 不再让网络请求逃逸生命周期预算,仓库级并发合并、缓存键和 V1/V2/V3 返回兼容保持不变。
|
||||
请求作用域的结构化并发不进入全局登记器:传统 WebAgent SSE 的 collection 子任务改由生成器
|
||||
`finally` 取消并等待清理,断线和 ASGI 取消均不会留下请求级 task。
|
||||
插件市场目录并发也采用同一责任边界:API/Agent 父请求取消或进度回调失败时,聚合服务会取消并等待
|
||||
全部市场/代际 loader;正常单源失败仍隔离,不把请求级读取错误登记成 lifespan 后台任务。
|
||||
订阅删除的宿主生产者也已完成 durable 分级:消息交互和远程删除不再调用裸线程统计入口,而是
|
||||
与业务删除原子暂存事件和统计 intent;旧类方法只作为插件 ABI 保留,不纳入宿主可靠性证明。
|
||||
消息渠道回环也不再各自拥有 URL、HTTP 判定和逐消息线程:七种宿主渠道共享 ingress,三种立即返回
|
||||
|
||||
@@ -111,3 +111,85 @@ async def test_async_collect_isolates_failure_and_completes_progress():
|
||||
error.assert_called_once()
|
||||
assert progress.call_args_list[0].kwargs["value"] == 0
|
||||
assert progress.call_args_list[-1].kwargs["value"] == 100
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_collect_cancels_all_loaders_when_parent_is_cancelled():
|
||||
"""请求取消时必须取消并回收全部市场 loader,不能把子任务遗留在事件循环。"""
|
||||
service = _service()
|
||||
blocker = asyncio.Event()
|
||||
all_started = asyncio.Event()
|
||||
started = 0
|
||||
cancelled = 0
|
||||
|
||||
async def loader(_market: str, _package_version: str | None, _force: bool):
|
||||
"""记录市场 loader 的启动和取消,并等待测试释放。"""
|
||||
nonlocal started, cancelled
|
||||
started += 1
|
||||
if started == 2:
|
||||
all_started.set()
|
||||
try:
|
||||
await blocker.wait()
|
||||
except asyncio.CancelledError:
|
||||
cancelled += 1
|
||||
raise
|
||||
return []
|
||||
|
||||
collect_task = asyncio.create_task(
|
||||
service.async_collect(
|
||||
markets=["https://market-a", "https://market-b"],
|
||||
compatible_flags=[],
|
||||
force=True,
|
||||
loader=loader,
|
||||
)
|
||||
)
|
||||
await asyncio.wait_for(all_started.wait(), timeout=1)
|
||||
collect_task.cancel()
|
||||
try:
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await collect_task
|
||||
await asyncio.sleep(0)
|
||||
assert cancelled == 2
|
||||
finally:
|
||||
blocker.set()
|
||||
await asyncio.sleep(0)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_async_collect_cleans_loaders_when_progress_callback_fails():
|
||||
"""进度回调异常也必须取消并回收尚未完成的市场 loader。"""
|
||||
service = _service()
|
||||
blocker = asyncio.Event()
|
||||
slow_started = asyncio.Event()
|
||||
slow_cancelled = asyncio.Event()
|
||||
|
||||
async def loader(market: str, _package_version: str | None, _force: bool):
|
||||
"""让一个市场立即完成,另一个保持阻塞以验证异常清理。"""
|
||||
if market == "https://market-a":
|
||||
await slow_started.wait()
|
||||
return []
|
||||
slow_started.set()
|
||||
try:
|
||||
await blocker.wait()
|
||||
except asyncio.CancelledError:
|
||||
slow_cancelled.set()
|
||||
raise
|
||||
return []
|
||||
|
||||
def progress(*, value: float, **_kwargs) -> None:
|
||||
"""首个市场完成时模拟进度消费者失败。"""
|
||||
if value > 0:
|
||||
raise RuntimeError("progress unavailable")
|
||||
|
||||
with pytest.raises(RuntimeError, match="progress unavailable"):
|
||||
await asyncio.wait_for(
|
||||
service.async_collect(
|
||||
markets=["https://market-a", "https://market-b"],
|
||||
compatible_flags=[],
|
||||
force=True,
|
||||
loader=loader,
|
||||
progress_callback=progress,
|
||||
),
|
||||
timeout=1,
|
||||
)
|
||||
assert slow_cancelled.is_set()
|
||||
|
||||
Reference in New Issue
Block a user