refactor: own history ai progress tasks

This commit is contained in:
jxxghp
2026-08-24 05:51:17 +08:00
parent 40d4a10c99
commit 9abeda8a05
4 changed files with 124 additions and 33 deletions
+37 -20
View File
@@ -1,6 +1,6 @@
import asyncio
import time import time
from typing import List, Any, Optional from collections.abc import Coroutine
from typing import List, Any, Callable, Optional
from fastapi import Depends from fastapi import Depends
@@ -58,6 +58,21 @@ def normalize_history_ids(history_ids: list[int]) -> list[int]:
return normalized_ids return normalized_ids
def _build_progress_output_callback(
progress: AsyncProgressHelper,
data: dict[str, Any],
*,
submit: Callable[[Coroutine[Any, Any, Any]], object],
) -> Callable[[str], None]:
"""构造同步 Agent 输出回调,并把异步进度更新登记到宿主任务生命周期。"""
def update_output(text: str) -> None:
"""非阻塞提交一条进度更新,避免同步回调等待缓存 I/O。"""
submit(progress.update(text=text, data=data))
return update_output
def _start_ai_redo_task( def _start_ai_redo_task(
history_id: int, history_id: int,
prompt: str, prompt: str,
@@ -65,15 +80,17 @@ def _start_ai_redo_task(
task_registry: TaskRegistry | None = None, task_registry: TaskRegistry | None = None,
) -> None: ) -> None:
"""在后台任务中启动单条 AI 重新整理任务,并通过异步进度辅助类实时更新进度。""" """在后台任务中启动单条 AI 重新整理任务,并通过异步进度辅助类实时更新进度。"""
registry = resolve_background_task_registry(task_registry)
progress = AsyncProgressHelper(progress_key) progress = AsyncProgressHelper(progress_key)
update_output = _build_progress_output_callback(
def update_output(text: str): progress,
# 输出回调由 agent 在事件循环上同步调用,不能直接 await; {"history_id": history_id},
# 提交到全局事件循环非阻塞执行,避免同步缓存后端阻塞事件循环。 submit=lambda coroutine: registry.submit_threadsafe(
asyncio.run_coroutine_threadsafe( coroutine,
progress.update(text=text, data={"history_id": history_id}), loop=global_vars.loop,
global_vars.loop, owner="api.history.ai_redo.progress",
) ),
)
async def runner(): async def runner():
try: try:
@@ -110,7 +127,6 @@ def _start_ai_redo_task(
finally: finally:
await progress.end() await progress.end()
registry = resolve_background_task_registry(task_registry)
registry.create(runner(), owner="api.history.ai_redo") registry.create(runner(), owner="api.history.ai_redo")
@@ -121,15 +137,17 @@ def _start_batch_ai_redo_task(
task_registry: TaskRegistry | None = None, task_registry: TaskRegistry | None = None,
) -> None: ) -> None:
"""在后台任务中启动批量 AI 重新整理任务,并通过异步进度辅助类实时更新进度。""" """在后台任务中启动批量 AI 重新整理任务,并通过异步进度辅助类实时更新进度。"""
registry = resolve_background_task_registry(task_registry)
progress = AsyncProgressHelper(progress_key) progress = AsyncProgressHelper(progress_key)
update_output = _build_progress_output_callback(
def update_output(text: str): progress,
# 输出回调由 agent 在事件循环上同步调用,不能直接 await; {"history_ids": history_ids},
# 提交到全局事件循环非阻塞执行,避免同步缓存后端阻塞事件循环。 submit=lambda coroutine: registry.submit_threadsafe(
asyncio.run_coroutine_threadsafe( coroutine,
progress.update(text=text, data={"history_ids": history_ids}), loop=global_vars.loop,
global_vars.loop, owner="api.history.ai_redo_batch.progress",
) ),
)
async def runner(): async def runner():
try: try:
@@ -166,7 +184,6 @@ def _start_batch_ai_redo_task(
finally: finally:
await progress.end() await progress.end()
registry = resolve_background_task_registry(task_registry)
registry.create(runner(), owner="api.history.ai_redo_batch") registry.create(runner(), owner="api.history.ai_redo_batch")
@@ -85,6 +85,8 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同
### Agent tasks ### Agent tasks
- 流式 token、工具进度和临时展示:E0。 - 流式 token、工具进度和临时展示:E0。
- 整理历史 AI 重做的 runner 与同步输出回调共用 lifespan TaskRegistry;单条/批量进度分别登记
`api.history.ai_redo.progress` / `api.history.ai_redo_batch.progress`,不再把缓存进度更新裸投递主循环。
- 已登记的周期 Agent task:E1,重启时通过任务定义重建;单次执行要有 execution 记录。 - 已登记的周期 Agent task:E1,重启时通过任务定义重建;单次执行要有 execution 记录。
- Agent 创建/修改订阅、删除数据等工具:业务事务按 E2/E3;聊天输出不能替代业务完成证据。 - Agent 创建/修改订阅、删除数据等工具:业务事务按 E2/E3;聊天输出不能替代业务完成证据。
- 会话 stop/cancel:E0 控制信号;被取消工具的底层阻塞 I/O 可能继续,资源所有者必须最终回收。 - 会话 stop/cancel:E0 控制信号;被取消工具的底层阻塞 I/O 可能继续,资源所有者必须最终回收。
@@ -6,7 +6,7 @@
> 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本 > 审计范围:宿主后端;排除 `app/plugins/**` 运行时插件副本
> 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文 > 规范优先级:`AGENTS.md` 与 `docs/rules/` 高于本文
> 相关文档:`docs/architecture-overview.md`、`docs/refactor/backend-architecture-governance.md`、`docs/refactor/backend-module-refactor-compatibility.md` > 相关文档:`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 接管。 > 实施进度:阶段 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
## 当前复核结论(2026-08-24 ## 当前复核结论(2026-08-24
@@ -278,6 +278,16 @@
- 依赖基线仅新增 `app.chain.message -> app.runtime.tasks`,模块数保持 `806`,内部边为 `6541`12 组禁止 - 依赖基线仅新增 `app.chain.message -> app.runtime.tasks`,模块数保持 `806`,内部边为 `6541`12 组禁止
边与唯一隔离 TMDB SCC 均未变化。 边与唯一隔离 TMDB SCC 均未变化。
### 长期整改阶段 27:历史 AI 进度任务所有权统一(2026-08-24)
- 整理历史 AI 重做的外层 runner 已进入 TaskRegistry,但单条与批量输出回调仍各复制一套裸
`run_coroutine_threadsafe`,进度缓存写入绕过了同一 shutdown owner 边界。
- 两条路径现在复用唯一进度回调工厂和外层同一个 registry,分别登记
`api.history.ai_redo.progress``api.history.ai_redo_batch.progress`;停止接收、取消和有限等待语义统一。
- API 路由、权限、进度 key/payload、Agent prompt 与完成/失败文案均未修改;该进度仍按 E0 允许进程崩溃
丢失,不把 UI 进度误报为 durable 完成。既有 `app.api.endpoints.history -> app.runtime.tasks` 依赖边不变,
V2/V3 插件 SDK/Compat 无改动。
### 总体判断 ### 总体判断
当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**:
+74 -12
View File
@@ -2,10 +2,12 @@
import asyncio import asyncio
from types import SimpleNamespace from types import SimpleNamespace
from unittest.mock import AsyncMock, patch
from app.api.endpoints import anthropic, history, message, openai, site, subscribe, webhook from app.api.endpoints import anthropic, history, message, openai, site, subscribe, webhook
from app.api.dependencies import subscription as subscription_dependencies from app.api.dependencies import subscription as subscription_dependencies
from app.application.subscription.search import SubscribeSearchActor from app.application.subscription.search import SubscribeSearchActor
from app.runtime.config import global_vars
from app.runtime.tasks import TaskRegistry from app.runtime.tasks import TaskRegistry
@@ -16,6 +18,7 @@ class _TaskRegistry(TaskRegistry):
"""初始化调用记录。""" """初始化调用记录。"""
super().__init__() super().__init__()
self.calls: list[tuple] = [] self.calls: list[tuple] = []
self.threadsafe_calls: list[tuple] = []
def create_sync(self, function, *args, owner: str, **kwargs) -> None: def create_sync(self, function, *args, owner: str, **kwargs) -> None:
"""保存函数、参数和 owner。""" """保存函数、参数和 owner。"""
@@ -32,6 +35,18 @@ class _TaskRegistry(TaskRegistry):
coroutine.close() coroutine.close()
self.calls.append((None, (), {"cancel_on_shutdown": cancel_on_shutdown}, owner)) self.calls.append((None, (), {"cancel_on_shutdown": cancel_on_shutdown}, owner))
def submit_threadsafe(
self,
coroutine,
*,
loop,
owner: str,
cancel_on_shutdown: bool = True,
) -> None:
"""保存跨线程任务参数,并关闭未执行的 coroutine。"""
coroutine.close()
self.threadsafe_calls.append((loop, owner, cancel_on_shutdown))
class _RunningTaskRegistry(TaskRegistry): class _RunningTaskRegistry(TaskRegistry):
"""执行协议流任务并保留 owner,验证真实 TaskRegistry 行为。""" """执行协议流任务并保留 owner,验证真实 TaskRegistry 行为。"""
@@ -216,33 +231,80 @@ def test_manual_subscription_search_uses_task_registry() -> None:
def test_history_ai_redo_uses_task_registry() -> None: def test_history_ai_redo_uses_task_registry() -> None:
"""单条历史 AI 重做应登记宿主任务并使用稳定 owner。""" """单条历史 AI 重做应登记宿主任务并使用稳定 owner。"""
registry = _TaskRegistry() registry = _TaskRegistry()
loop = SimpleNamespace(is_running=lambda: True, is_closed=lambda: False)
history._start_ai_redo_task( with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch.object(
history_id=7, history,
prompt="整理记录", "_build_progress_output_callback",
progress_key="progress-7", return_value=lambda _text: None,
task_registry=registry, ) as build_callback:
) history._start_ai_redo_task(
history_id=7,
prompt="整理记录",
progress_key="progress-7",
task_registry=registry,
)
build_callback.call_args.kwargs["submit"](asyncio.sleep(0))
assert registry.calls == [ assert registry.calls == [
(None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo") (None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo")
] ]
assert build_callback.call_args.args[1] == {"history_id": 7}
assert registry.threadsafe_calls == [
(loop, "api.history.ai_redo.progress", True)
]
def test_history_progress_callback_uses_owned_threadsafe_submission() -> None:
"""历史 AI 输出进度应进入同一宿主登记器并保留 payload。"""
registry = _TaskRegistry()
progress = SimpleNamespace(update=AsyncMock())
loop = SimpleNamespace(is_running=lambda: True, is_closed=lambda: False)
callback = history._build_progress_output_callback(
progress,
{"history_id": 7},
submit=lambda coroutine: registry.submit_threadsafe(
coroutine,
loop=loop,
owner="api.history.ai_redo.progress",
),
)
callback("正在分析")
progress.update.assert_called_once_with(
text="正在分析", data={"history_id": 7}
)
assert registry.threadsafe_calls == [
(loop, "api.history.ai_redo.progress", True)
]
def test_history_batch_ai_redo_uses_task_registry() -> None: def test_history_batch_ai_redo_uses_task_registry() -> None:
"""批量历史 AI 重做应登记宿主任务并区分批量 owner。""" """批量历史 AI 重做应登记宿主任务并区分批量 owner。"""
registry = _TaskRegistry() registry = _TaskRegistry()
loop = SimpleNamespace(is_running=lambda: True, is_closed=lambda: False)
history._start_batch_ai_redo_task( with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch.object(
history_ids=[7, 8], history,
prompt="批量整理", "_build_progress_output_callback",
progress_key="progress-batch", return_value=lambda _text: None,
task_registry=registry, ) as build_callback:
) history._start_batch_ai_redo_task(
history_ids=[7, 8],
prompt="批量整理",
progress_key="progress-batch",
task_registry=registry,
)
build_callback.call_args.kwargs["submit"](asyncio.sleep(0))
assert registry.calls == [ assert registry.calls == [
(None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo_batch") (None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo_batch")
] ]
assert build_callback.call_args.args[1] == {"history_ids": [7, 8]}
assert registry.threadsafe_calls == [
(loop, "api.history.ai_redo_batch.progress", True)
]
def test_openai_stream_uses_task_registry(monkeypatch) -> None: def test_openai_stream_uses_task_registry(monkeypatch) -> None: