mirror of
https://github.com/jxxghp/MoviePilot.git
synced 2026-09-04 23:17:20 +08:00
refactor: own plugin release requests
This commit is contained in:
Vendored
+11
-3
@@ -44,6 +44,7 @@ from app.adapters.system.plugin.manifest import (
|
|||||||
)
|
)
|
||||||
from app.runtime.log import logger
|
from app.runtime.log import logger
|
||||||
from app.runtime.observability import observe_compat_facade
|
from app.runtime.observability import observe_compat_facade
|
||||||
|
from app.runtime.tasks import get_task_registry
|
||||||
from app.adapters.network.http import RequestUtils, AsyncRequestUtils
|
from app.adapters.network.http import RequestUtils, AsyncRequestUtils
|
||||||
from app.foundation.singleton import WeakSingleton
|
from app.foundation.singleton import WeakSingleton
|
||||||
|
|
||||||
@@ -2336,8 +2337,12 @@ class PluginHelper(metaclass=WeakSingleton):
|
|||||||
if pending_normal_task and pending_normal_task.done():
|
if pending_normal_task and pending_normal_task.done():
|
||||||
pending_normal_task = None
|
pending_normal_task = None
|
||||||
task_key = force_task_key
|
task_key = force_task_key
|
||||||
task = loop.create_task(
|
task = get_task_registry().create(
|
||||||
self._async_refresh_plugin_repo_releases(normalized_repo_url, pending_normal_task)
|
self._async_refresh_plugin_repo_releases(
|
||||||
|
normalized_repo_url,
|
||||||
|
pending_normal_task,
|
||||||
|
),
|
||||||
|
owner="plugin.market.release_refresh",
|
||||||
)
|
)
|
||||||
self._release_tasks[task_key] = task
|
self._release_tasks[task_key] = task
|
||||||
task.add_done_callback(
|
task.add_done_callback(
|
||||||
@@ -2347,7 +2352,10 @@ class PluginHelper(metaclass=WeakSingleton):
|
|||||||
task_key = normal_task_key
|
task_key = normal_task_key
|
||||||
pending_normal_task = self._release_tasks.get(normal_task_key)
|
pending_normal_task = self._release_tasks.get(normal_task_key)
|
||||||
if pending_normal_task is None or pending_normal_task.done():
|
if pending_normal_task is None or pending_normal_task.done():
|
||||||
task = loop.create_task(self._async_get_plugin_repo_releases(normalized_repo_url))
|
task = get_task_registry().create(
|
||||||
|
self._async_get_plugin_repo_releases(normalized_repo_url),
|
||||||
|
owner="plugin.market.release_read",
|
||||||
|
)
|
||||||
self._release_tasks[task_key] = task
|
self._release_tasks[task_key] = task
|
||||||
task.add_done_callback(
|
task.add_done_callback(
|
||||||
lambda completed_task: self._remove_release_task(task_key, completed_task)
|
lambda completed_task: self._remove_release_task(task_key, completed_task)
|
||||||
|
|||||||
@@ -87,6 +87,8 @@
|
|||||||
绕过全局关停预算的第二套后台任务集合。该摘要属于可丢弃 E1 观测数据,不宣称 durable。
|
绕过全局关停预算的第二套后台任务集合。该摘要属于可丢弃 E1 观测数据,不宣称 durable。
|
||||||
模块同步清缓存的异步桥接也统一复用同一模式:Fanart 不再保留裸 `loop.create_task`,与 IMDb 一样登记
|
模块同步清缓存的异步桥接也统一复用同一模式:Fanart 不再保留裸 `loop.create_task`,与 IMDb 一样登记
|
||||||
稳定 owner,并继续保留无事件循环时 `asyncio.run` 和原同步 ABI。
|
稳定 owner,并继续保留无事件循环时 `asyncio.run` 和原同步 ABI。
|
||||||
|
插件市场 Release 合并层的普通读取和强刷子任务也已登记稳定 owner;API 外层任务被关停取消时,内部
|
||||||
|
`shield` 不再让网络请求逃逸生命周期预算,仓库级并发合并、缓存键和 V1/V2/V3 返回兼容保持不变。
|
||||||
2. **动态模块契约仍以 legacy 聚合语义为主。** 当前登记 `212` 个模块方法,其中 `194` 个仍使用 `legacy` aggregation,只有 `14` 个 `first_non_empty`、`4` 个 `ordered_list_merge`。`app/runtime/extensions/module/contracts.py:422-455` 已能登记 family、输入/结果标签和基础签名诊断,但 `193` 个方法没有 required parameters,调度器 `app/runtime/extensions/module/dispatcher.py:109-260` 仍主要依赖运行时反射、返回值形状和短路规则。未知第三方方法保留 legacy fallback 是兼容要求,不应删除;宿主高频能力则应逐族补齐可执行的输入校验、结果校验、超时和错误语义。
|
2. **动态模块契约仍以 legacy 聚合语义为主。** 当前登记 `212` 个模块方法,其中 `194` 个仍使用 `legacy` aggregation,只有 `14` 个 `first_non_empty`、`4` 个 `ordered_list_merge`。`app/runtime/extensions/module/contracts.py:422-455` 已能登记 family、输入/结果标签和基础签名诊断,但 `193` 个方法没有 required parameters,调度器 `app/runtime/extensions/module/dispatcher.py:109-260` 仍主要依赖运行时反射、返回值形状和短路规则。未知第三方方法保留 legacy fallback 是兼容要求,不应删除;宿主高频能力则应逐族补齐可执行的输入校验、结果校验、超时和错误语义。
|
||||||
3. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。
|
3. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。
|
||||||
|
|
||||||
|
|||||||
+3
-2
@@ -13,8 +13,8 @@
|
|||||||
"runtime_to_db": [],
|
"runtime_to_db": [],
|
||||||
"workflow_to_db": []
|
"workflow_to_db": []
|
||||||
},
|
},
|
||||||
"edge_count": 6501,
|
"edge_count": 6502,
|
||||||
"edge_sha256": "3c01f9aeb7703065ff447672284ddcc88d6120ac3236b9659200df9b6f921bdd",
|
"edge_sha256": "fe24e90415afe0789875f5e6d936e9b7d3536b3bfd0383d585ea35cc4be76386",
|
||||||
"edges": [
|
"edges": [
|
||||||
"app -> app.runtime",
|
"app -> app.runtime",
|
||||||
"app -> app.runtime.compat",
|
"app -> app.runtime.compat",
|
||||||
@@ -65,6 +65,7 @@
|
|||||||
"app.adapters.external.market -> app.runtime.log",
|
"app.adapters.external.market -> app.runtime.log",
|
||||||
"app.adapters.external.market -> app.runtime.observability",
|
"app.adapters.external.market -> app.runtime.observability",
|
||||||
"app.adapters.external.market -> app.runtime.settings",
|
"app.adapters.external.market -> app.runtime.settings",
|
||||||
|
"app.adapters.external.market -> app.runtime.tasks",
|
||||||
"app.adapters.external.ocr -> app.adapters",
|
"app.adapters.external.ocr -> app.adapters",
|
||||||
"app.adapters.external.ocr -> app.adapters.network",
|
"app.adapters.external.ocr -> app.adapters.network",
|
||||||
"app.adapters.external.ocr -> app.adapters.network.http",
|
"app.adapters.external.ocr -> app.adapters.network.http",
|
||||||
|
|||||||
@@ -535,6 +535,60 @@ class TestPluginHelper:
|
|||||||
assert [item["version"] for item in cached_result] == ["1.2.3"]
|
assert [item["version"] for item in cached_result] == ["1.2.3"]
|
||||||
assert request_count == 2
|
assert request_count == 2
|
||||||
|
|
||||||
|
def test_async_release_read_follows_host_task_shutdown(self, monkeypatch):
|
||||||
|
"""被 shield 的仓库级请求仍须登记 owner,并随宿主关停取消。"""
|
||||||
|
try:
|
||||||
|
from app.adapters.external import market as market_module
|
||||||
|
from app.adapters.external.market import PluginHelper
|
||||||
|
from app.runtime.tasks import TaskRegistry
|
||||||
|
except ModuleNotFoundError as exc:
|
||||||
|
pytest.skip(f"missing dependency: {exc}")
|
||||||
|
|
||||||
|
async def run_test():
|
||||||
|
"""阻塞仓库读取后关闭登记器,返回 owner 与取消收敛状态。"""
|
||||||
|
helper = PluginHelper()
|
||||||
|
registry = TaskRegistry()
|
||||||
|
started = asyncio.Event()
|
||||||
|
cancelled = asyncio.Event()
|
||||||
|
|
||||||
|
async def fake_request(*_args, **_kwargs):
|
||||||
|
"""保持网络读取运行,直到登记器发出取消。"""
|
||||||
|
started.set()
|
||||||
|
try:
|
||||||
|
await asyncio.Event().wait()
|
||||||
|
except asyncio.CancelledError:
|
||||||
|
cancelled.set()
|
||||||
|
raise
|
||||||
|
|
||||||
|
await helper.async_get_plugin_release_versions.cache_clear()
|
||||||
|
monkeypatch.setattr(
|
||||||
|
helper,
|
||||||
|
"_PluginHelper__async_request_with_fallback",
|
||||||
|
fake_request,
|
||||||
|
)
|
||||||
|
monkeypatch.setattr(
|
||||||
|
market_module,
|
||||||
|
"get_task_registry",
|
||||||
|
lambda: registry,
|
||||||
|
)
|
||||||
|
caller = asyncio.create_task(
|
||||||
|
helper.async_get_plugin_release_versions("DemoPlugin", REPO_URL)
|
||||||
|
)
|
||||||
|
await started.wait()
|
||||||
|
owners = tuple(record.owner for record in registry.records)
|
||||||
|
converged = await registry.shutdown(timeout_seconds=1.0)
|
||||||
|
results = await asyncio.gather(caller, return_exceptions=True)
|
||||||
|
await asyncio.sleep(0)
|
||||||
|
return owners, converged, cancelled.is_set(), results, helper._release_tasks
|
||||||
|
|
||||||
|
owners, converged, cancelled, results, release_tasks = asyncio.run(run_test())
|
||||||
|
|
||||||
|
assert owners == ("plugin.market.release_read",)
|
||||||
|
assert converged is True
|
||||||
|
assert cancelled is True
|
||||||
|
assert isinstance(results[0], asyncio.CancelledError)
|
||||||
|
assert release_tasks == {}
|
||||||
|
|
||||||
def test_async_normal_release_read_does_not_wait_for_pending_force_refresh(self, monkeypatch):
|
def test_async_normal_release_read_does_not_wait_for_pending_force_refresh(self, monkeypatch):
|
||||||
"""普通读取遇到后台强刷时仍优先返回已有缓存,避免页面响应被强刷阻塞。"""
|
"""普通读取遇到后台强刷时仍优先返回已有缓存,避免页面响应被强刷阻塞。"""
|
||||||
try:
|
try:
|
||||||
|
|||||||
Reference in New Issue
Block a user