refactor: own cross-thread background tasks

This commit is contained in:
jxxghp
2026-08-24 05:33:50 +08:00
parent 4575146acd
commit d694cd081c
9 changed files with 188 additions and 14 deletions
+6 -1
View File
@@ -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,
+74
View File
@@ -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],
@@ -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 线程。
@@ -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,不宣称跨进程恢复。
### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
+3 -1
View File
@@ -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(
+3 -2
View File
@@ -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",
+8
View File
@@ -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",
)
""",
)
+68
View File
@@ -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,不能把包装任务取消成伪完成。"""
+13 -9
View File
@@ -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