From 579226612a2cff8e34ca9cf6793081fb47933010 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 00:41:09 +0800 Subject: [PATCH] refactor: own plugin release requests --- app/adapters/external/market.py | 14 +++-- .../backend-architecture-next-stage.md | 2 + .../architecture/dependency-baseline.json | 5 +- tests/test_plugin_helper.py | 54 +++++++++++++++++++ 4 files changed, 70 insertions(+), 5 deletions(-) diff --git a/app/adapters/external/market.py b/app/adapters/external/market.py index 81c3119cc..290f561b0 100644 --- a/app/adapters/external/market.py +++ b/app/adapters/external/market.py @@ -44,6 +44,7 @@ from app.adapters.system.plugin.manifest import ( ) from app.runtime.log import logger 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.foundation.singleton import WeakSingleton @@ -2336,8 +2337,12 @@ class PluginHelper(metaclass=WeakSingleton): if pending_normal_task and pending_normal_task.done(): pending_normal_task = None task_key = force_task_key - task = loop.create_task( - self._async_refresh_plugin_repo_releases(normalized_repo_url, pending_normal_task) + task = get_task_registry().create( + self._async_refresh_plugin_repo_releases( + normalized_repo_url, + pending_normal_task, + ), + owner="plugin.market.release_refresh", ) self._release_tasks[task_key] = task task.add_done_callback( @@ -2347,7 +2352,10 @@ class PluginHelper(metaclass=WeakSingleton): task_key = normal_task_key pending_normal_task = self._release_tasks.get(normal_task_key) 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 task.add_done_callback( lambda completed_task: self._remove_release_task(task_key, completed_task) diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 7fe4b6e5d..c4cd42d47 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -87,6 +87,8 @@ 绕过全局关停预算的第二套后台任务集合。该摘要属于可丢弃 E1 观测数据,不宣称 durable。 模块同步清缓存的异步桥接也统一复用同一模式:Fanart 不再保留裸 `loop.create_task`,与 IMDb 一样登记 稳定 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 是兼容要求,不应删除;宿主高频能力则应逐族补齐可执行的输入校验、结果校验、超时和错误语义。 3. **Model/Base 的数据库装饰器和隐式会话 ABI 已全部清零。** 查询、写事务和 `legacy_*` 装饰器均为 `0`;所有 Model `db` 参数要求显式 Session,Base CRUD 仅在调用方事务内查询或 stage。可无会话构造的入口统一留在 Oper,经组合根事务执行器运行;插件 SDK 不再导出宿主 Model。后续重点转为减少 ORM 对象跨层流转,并保持 Model 隐式事务零回退。 diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index 026384b99..6bd635b87 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": 6501, - "edge_sha256": "3c01f9aeb7703065ff447672284ddcc88d6120ac3236b9659200df9b6f921bdd", + "edge_count": 6502, + "edge_sha256": "fe24e90415afe0789875f5e6d936e9b7d3536b3bfd0383d585ea35cc4be76386", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -65,6 +65,7 @@ "app.adapters.external.market -> app.runtime.log", "app.adapters.external.market -> app.runtime.observability", "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.network", "app.adapters.external.ocr -> app.adapters.network.http", diff --git a/tests/test_plugin_helper.py b/tests/test_plugin_helper.py index 61a45b046..74f3b1831 100644 --- a/tests/test_plugin_helper.py +++ b/tests/test_plugin_helper.py @@ -535,6 +535,60 @@ class TestPluginHelper: assert [item["version"] for item in cached_result] == ["1.2.3"] 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): """普通读取遇到后台强刷时仍优先返回已有缓存,避免页面响应被强刷阻塞。""" try: