From d694cd081c48d5306a72bc52df363005b25ab973 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 05:33:50 +0800 Subject: [PATCH] refactor: own cross-thread background tasks --- app/chain/_transfer.py | 7 +- app/runtime/tasks.py | 74 +++++++++++++++++++ .../adr/0007-background-action-reliability.md | 3 + .../backend-architecture-next-stage.md | 11 ++- scripts/architecture/task_ownership.py | 4 +- .../architecture/dependency-baseline.json | 5 +- tests/test_task_ownership_gate.py | 8 ++ tests/test_task_registry.py | 68 +++++++++++++++++ tests/test_transfer_failed_retry_buttons.py | 22 +++--- 9 files changed, 188 insertions(+), 14 deletions(-) diff --git a/app/chain/_transfer.py b/app/chain/_transfer.py index 6cdcaad61..b15ce0366 100644 --- a/app/chain/_transfer.py +++ b/app/chain/_transfer.py @@ -37,6 +37,7 @@ from app.domain.meta.metamusic import MetaMusic from app.foundation import text as text_tools from app.runtime.config import global_vars from app.runtime.log import logger +from app.runtime.tasks import get_task_registry from app.schemas.workflow import FileItem from app.schemas.message import Message from app.schemas.tmdb import TmdbEpisode @@ -1461,7 +1462,11 @@ class FailedRetryMixin: ) ) - asyncio.run_coroutine_threadsafe(_run_ai_takeover(), global_vars.loop) + get_task_registry().submit_threadsafe( + _run_ai_takeover(), + loop=global_vars.loop, + owner="chain.transfer.ai_takeover", + ) def _re_transfer( self, diff --git a/app/runtime/tasks.py b/app/runtime/tasks.py index 56b644f2c..e58d66b45 100644 --- a/app/runtime/tasks.py +++ b/app/runtime/tasks.py @@ -3,6 +3,7 @@ from __future__ import annotations import asyncio +import concurrent.futures from collections.abc import Coroutine from dataclasses import dataclass from functools import partial @@ -68,6 +69,79 @@ class TaskRegistry: cancel_on_shutdown=False, ) + def submit_threadsafe( + self, + coroutine: Coroutine[Any, Any, Any], + *, + loop: asyncio.AbstractEventLoop, + owner: str, + cancel_on_shutdown: bool = True, + ) -> concurrent.futures.Future[Any]: + """从宿主线程提交协程,并在目标循环内原子登记 owner 后执行。""" + completion: concurrent.futures.Future[Any] = concurrent.futures.Future() + task_holder: dict[str, asyncio.Task[Any]] = {} + + def mirror_completion(task: asyncio.Task[Any]) -> None: + """把登记任务的真实终态镜像给跨线程调用方。""" + if completion.done(): + return + if task.cancelled(): + completion.cancel() + return + exception = task.exception() + if exception is not None: + completion.set_exception(exception) + else: + completion.set_result(task.result()) + + def submit_on_loop() -> None: + """在目标循环内完成 accepting 检查、任务创建和 owner 登记。""" + if completion.cancelled(): + coroutine.close() + return + try: + task = self.create( + coroutine, + owner=owner, + cancel_on_shutdown=cancel_on_shutdown, + ) + except Exception as error: + if not completion.done(): + completion.set_exception(error) + loop.call_exception_handler( + { + "message": "MoviePilot 跨线程后台任务提交失败", + "exception": error, + "owner": owner, + } + ) + return + task_holder["task"] = task + task.add_done_callback(mirror_completion) + if completion.cancelled() and not task.done(): + task.cancel() + + def cancel_registered_task( + submitted: concurrent.futures.Future[Any], + ) -> None: + """调用方取消 completion 时,把取消请求转交目标循环中的真实任务。""" + if not submitted.cancelled(): + return + task = task_holder.get("task") + if task is not None and not task.done(): + try: + loop.call_soon_threadsafe(task.cancel) + except RuntimeError: + pass + + completion.add_done_callback(cancel_registered_task) + try: + loop.call_soon_threadsafe(submit_on_loop) + except RuntimeError: + coroutine.close() + raise + return completion + def register( self, task: asyncio.Task[Any], diff --git a/docs/adr/0007-background-action-reliability.md b/docs/adr/0007-background-action-reliability.md index 819fb1d17..a25056933 100644 --- a/docs/adr/0007-background-action-reliability.md +++ b/docs/adr/0007-background-action-reliability.md @@ -64,6 +64,9 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同 调度,不表示执行完成。 - Webhook E0 广播、消息入口和 Seerr 订阅入口均已迁入 lifespan TaskRegistry,具备 owner、停止接收和 有限等待语义;进程崩溃时仍允许丢失,不因此提升为 durable。 +- TaskRegistry 提供跨宿主线程的 `submit_threadsafe()`;提交回目标循环后先原子登记 owner 再执行,关停 + 竞态中要么纳入取消/等待,要么拒绝并关闭协程。整理失败按钮的 AI 接管使用 + `chain.transfer.ai_takeover` owner,不再绕过登记器直接投递主循环。 - Slack、Telegram、Discord、飞书、QQBot、企业微信与 WeChatClawBot 的渠道回环统一经 `application.messaging.ingress` 进入同一个 API/TaskRegistry 主链;需要立即返回 SDK 回调的渠道把同步 HTTP 交给宿主共享线程池,模块关闭后由线程池生命周期等待,不再创建逐消息 daemon 线程。 diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index f91fb0d01..a7e3d0a7b 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~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 的协程提交双轨。 +> 实施进度:阶段 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 接管。 ## 当前复核结论(2026-08-24) @@ -258,6 +258,15 @@ - `Scheduler.start()`、同步 `stop()`、作业定义、插件调度方法、进度 payload 和 SDK/Compat 均未修改, V2/V3 插件行为保持兼容。 +### 长期整改阶段 25:跨线程 TaskRegistry owner 收口(2026-08-24) + +- TaskRegistry 原先只有事件循环内 `create()`,同步 Chain 无法原子登记异步后台动作;新增 + `submit_threadsafe()`,在目标循环内完成 accepting 检查、任务创建和 owner 登记,再镜像真实终态。 +- shutdown 与跨线程提交竞态时,先登记的任务进入既有取消/等待;后到任务被拒绝并关闭协程,立即投递 + 失败向调用方抛出。owner 静态门禁同步覆盖新入口。 +- 整理失败按钮的 AI 接管从裸 `run_coroutine_threadsafe` 迁为 `chain.transfer.ai_takeover`,消息文案、Agent + prompt/参数、回调返回和 V2/V3 插件 ABI 保持不变;该动作仍为 E0,不宣称跨进程恢复。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: diff --git a/scripts/architecture/task_ownership.py b/scripts/architecture/task_ownership.py index 29940a0a3..796426a84 100644 --- a/scripts/architecture/task_ownership.py +++ b/scripts/architecture/task_ownership.py @@ -8,7 +8,9 @@ from pathlib import Path PROJECT_ROOT = Path(__file__).resolve().parents[2] -TASK_METHODS = frozenset({"create", "create_sync", "register"}) +TASK_METHODS = frozenset( + {"create", "create_sync", "register", "submit_threadsafe"} +) TASK_MODULE = "app.runtime.tasks" CONTEXT_MODULE = "app.api.context" TASK_FACTORIES = frozenset( diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index c72afe014..b10112719 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": 6539, - "edge_sha256": "28c7ab8544ede290e6dacf5f126ca1aad95e09627d874a0d23c6476108a88e1b", + "edge_count": 6540, + "edge_sha256": "3beb7c6b906993a8ff9d0cb9b3f7c912e1a7859785824f472547b8fecf8a4c51", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -3017,6 +3017,7 @@ "app.chain._transfer -> app.runtime", "app.chain._transfer -> app.runtime.config", "app.chain._transfer -> app.runtime.log", + "app.chain._transfer -> app.runtime.tasks", "app.chain._transfer -> app.schemas", "app.chain._transfer -> app.schemas.history", "app.chain._transfer -> app.schemas.message", diff --git a/tests/test_task_ownership_gate.py b/tests/test_task_ownership_gate.py index 0529586f1..2ca599011 100644 --- a/tests/test_task_ownership_gate.py +++ b/tests/test_task_ownership_gate.py @@ -34,6 +34,7 @@ def test_owner_gate_tracks_known_registry_without_matching_same_named_methods( task_registry.create(work()) resolve_registry(task_registry).create_sync(work, owner=dynamic_owner) get_task_registry().register(task, owner=" ") + task_registry.submit_threadsafe(work(), loop=loop) """, ) @@ -41,11 +42,13 @@ def test_owner_gate_tracks_known_registry_without_matching_same_named_methods( "create", "create_sync", "register", + "submit_threadsafe", ] assert [violation.reason for violation in violations] == [ "缺少显式 owner", "的 owner 必须是非空字符串字面量", "的 owner 必须是非空字符串字面量", + "缺少显式 owner", ] @@ -65,6 +68,11 @@ def test_owner_gate_accepts_aliases_and_stable_literal_owners(tmp_path: Path) -> task, owner="api.example.existing", ) + task_registry.submit_threadsafe( + work(), + loop=loop, + owner="api.example.threadsafe", + ) """, ) diff --git a/tests/test_task_registry.py b/tests/test_task_registry.py index 4ec9b0079..677f99de4 100644 --- a/tests/test_task_registry.py +++ b/tests/test_task_registry.py @@ -87,6 +87,74 @@ def test_task_registry_runs_sync_function_and_tracks_until_completion() -> None: asyncio.run(scenario()) +def test_task_registry_owns_threadsafe_submission_until_shutdown() -> None: + """宿主线程提交的协程应先登记 owner,并由 Registry 关停取消和等待。""" + + async def scenario() -> None: + registry = TaskRegistry() + loop = asyncio.get_running_loop() + started = asyncio.Event() + cleaned = asyncio.Event() + + async def worker() -> None: + """保持运行直到 Registry 发出取消,并记录清理已完成。""" + started.set() + try: + await asyncio.Event().wait() + finally: + cleaned.set() + + completion = await asyncio.to_thread( + registry.submit_threadsafe, + worker(), + loop=loop, + owner="test.threadsafe", + ) + await asyncio.wait_for(started.wait(), timeout=1) + assert [record.owner for record in registry.records] == [ + "test.threadsafe" + ] + + assert await registry.shutdown(timeout_seconds=1.0) is True + assert cleaned.is_set() + assert completion.cancelled() + assert registry.records == () + + asyncio.run(scenario()) + + +def test_task_registry_rejects_threadsafe_submission_after_shutdown() -> None: + """关停先赢得竞态时应关闭协程并通过 completion 报告拒绝原因。""" + + async def scenario() -> None: + registry = TaskRegistry() + loop = asyncio.get_running_loop() + reports: list[dict[str, object]] = [] + previous_handler = loop.get_exception_handler() + loop.set_exception_handler(lambda _, context: reports.append(context)) + + async def late_worker() -> None: + """模拟关停完成后从宿主线程到达的晚任务。""" + + try: + assert await registry.shutdown(timeout_seconds=1.0) is True + completion = await asyncio.to_thread( + registry.submit_threadsafe, + late_worker(), + loop=loop, + owner="test.threadsafe-late", + ) + with pytest.raises(RuntimeError, match="正在关闭"): + await asyncio.wrap_future(completion) + + assert registry.records == () + assert reports[-1]["owner"] == "test.threadsafe-late" + finally: + loop.set_exception_handler(previous_handler) + + asyncio.run(scenario()) + + def test_task_registry_keeps_timed_out_sync_owner_until_real_completion() -> None: """同步线程超过关停预算后仍应保留 owner,不能把包装任务取消成伪完成。""" diff --git a/tests/test_transfer_failed_retry_buttons.py b/tests/test_transfer_failed_retry_buttons.py index 139bf813d..ec9868f68 100644 --- a/tests/test_transfer_failed_retry_buttons.py +++ b/tests/test_transfer_failed_retry_buttons.py @@ -141,10 +141,10 @@ class TestTransferFailedRetryButtons(unittest.TestCase): ) as history_oper_cls, patch( "app.chain._transfer.build_manual_redo_prompt", return_value="retry transfer prompt", - ), patch( - "app.chain._transfer.asyncio.run_coroutine_threadsafe", - side_effect=_close_pending_coro, - ) as run_task: + ), patch("app.chain._transfer.get_task_registry") as get_registry: + get_registry.return_value.submit_threadsafe.side_effect = ( + _close_pending_coro + ) history_oper_cls.return_value.get.return_value = history with patch.object(chain, "post_message") as post_message: chain.handle_failed_transfer_callback( @@ -155,7 +155,11 @@ class TestTransferFailedRetryButtons(unittest.TestCase): username="tester", ) - run_task.assert_called_once() + get_registry.return_value.submit_threadsafe.assert_called_once() + self.assertEqual( + get_registry.return_value.submit_threadsafe.call_args.kwargs["owner"], + "chain.transfer.ai_takeover", + ) self.assertEqual(post_message.call_count, 1) self.assertEqual( post_message.call_args_list[0].args[0].title, @@ -224,10 +228,10 @@ class TestTransferFailedRetryButtons(unittest.TestCase): ), patch( "app.chain._transfer.get_running_agent_manager", return_value=manager, - ), patch( - "app.chain._transfer.asyncio.run_coroutine_threadsafe", - side_effect=_run_pending_coro, - ): + ), patch("app.chain._transfer.get_task_registry") as get_registry: + get_registry.return_value.submit_threadsafe.side_effect = ( + _run_pending_coro + ) history_oper_cls.return_value.get.return_value = history with patch.object(chain, "post_message"), patch.object( chain, "async_post_message", side_effect=fake_async_post_message