refactor: register history ai redo tasks

This commit is contained in:
jxxghp
2026-08-23 15:02:53 +08:00
parent 0459b6a662
commit e7e232d625
4 changed files with 74 additions and 9 deletions
+23 -5
View File
@@ -20,7 +20,12 @@ from app.agent.prompt.transfer_redo import (
build_manual_redo_prompt, build_manual_redo_prompt,
) )
from app.runtime.config import global_vars from app.runtime.config import global_vars
from app.api.context import get_api_runtime_config, resolve_api_runtime_config from app.api.context import (
get_api_runtime_config,
get_background_task_registry,
resolve_api_runtime_config,
resolve_background_task_registry,
)
from app.application.configuration import ApiRuntimeConfig from app.application.configuration import ApiRuntimeConfig
from app.adapters.web.security.access import verify_token from app.adapters.web.security.access import verify_token
from app.api.dependencies.auth import ( from app.api.dependencies.auth import (
@@ -39,6 +44,7 @@ from app.application.history import (
TransferHistoryMutationCommand, TransferHistoryMutationCommand,
) )
from app.runtime.log import logger from app.runtime.log import logger
from app.runtime.tasks import TaskRegistry
router = ResponseAPIRouter() router = ResponseAPIRouter()
@@ -52,7 +58,12 @@ def normalize_history_ids(history_ids: list[int]) -> list[int]:
return normalized_ids return normalized_ids
def _start_ai_redo_task(history_id: int, prompt: str, progress_key: str): def _start_ai_redo_task(
history_id: int,
prompt: str,
progress_key: str,
task_registry: TaskRegistry | None = None,
) -> None:
"""在后台任务中启动单条 AI 重新整理任务,并通过异步进度辅助类实时更新进度。""" """在后台任务中启动单条 AI 重新整理任务,并通过异步进度辅助类实时更新进度。"""
progress = AsyncProgressHelper(progress_key) progress = AsyncProgressHelper(progress_key)
@@ -99,14 +110,16 @@ def _start_ai_redo_task(history_id: int, prompt: str, progress_key: str):
finally: finally:
await progress.end() await progress.end()
asyncio.run_coroutine_threadsafe(runner(), global_vars.loop) registry = resolve_background_task_registry(task_registry)
registry.create(runner(), owner="api.history.ai_redo")
def _start_batch_ai_redo_task( def _start_batch_ai_redo_task(
history_ids: list[int], history_ids: list[int],
prompt: str, prompt: str,
progress_key: str, progress_key: str,
): task_registry: TaskRegistry | None = None,
) -> None:
"""在后台任务中启动批量 AI 重新整理任务,并通过异步进度辅助类实时更新进度。""" """在后台任务中启动批量 AI 重新整理任务,并通过异步进度辅助类实时更新进度。"""
progress = AsyncProgressHelper(progress_key) progress = AsyncProgressHelper(progress_key)
@@ -153,7 +166,8 @@ def _start_batch_ai_redo_task(
finally: finally:
await progress.end() await progress.end()
asyncio.run_coroutine_threadsafe(runner(), global_vars.loop) registry = resolve_background_task_registry(task_registry)
registry.create(runner(), owner="api.history.ai_redo_batch")
@router.get( @router.get(
@@ -248,6 +262,7 @@ async def ai_redo_transfer_history(
query: HistoryQueryService = Depends(get_history_query_service), query: HistoryQueryService = Depends(get_history_query_service),
runtime_config: ApiRuntimeConfig = Depends(get_api_runtime_config), runtime_config: ApiRuntimeConfig = Depends(get_api_runtime_config),
_: object = Depends(get_current_active_manage_user), _: object = Depends(get_current_active_manage_user),
task_registry: TaskRegistry = Depends(get_background_task_registry),
) -> Any: ) -> Any:
""" """
手动触发单条历史记录的 AI 重新整理,并返回进度键。 手动触发单条历史记录的 AI 重新整理,并返回进度键。
@@ -266,6 +281,7 @@ async def ai_redo_transfer_history(
history_id=history_id, history_id=history_id,
prompt=prompt, prompt=prompt,
progress_key=progress_key, progress_key=progress_key,
task_registry=task_registry,
) )
return _SchemaResponse(success=True, data={"progress_key": progress_key}) return _SchemaResponse(success=True, data={"progress_key": progress_key})
@@ -281,6 +297,7 @@ async def batch_ai_redo_transfer_history(
query: HistoryQueryService = Depends(get_history_query_service), query: HistoryQueryService = Depends(get_history_query_service),
runtime_config: ApiRuntimeConfig = Depends(get_api_runtime_config), runtime_config: ApiRuntimeConfig = Depends(get_api_runtime_config),
_: object = Depends(get_current_active_manage_user), _: object = Depends(get_current_active_manage_user),
task_registry: TaskRegistry = Depends(get_background_task_registry),
) -> Any: ) -> Any:
""" """
手动触发多条历史记录的 AI 批量重新整理,并返回进度键。 手动触发多条历史记录的 AI 批量重新整理,并返回进度键。
@@ -308,6 +325,7 @@ async def batch_ai_redo_transfer_history(
history_ids=history_ids, history_ids=history_ids,
prompt=prompt, prompt=prompt,
progress_key=progress_key, progress_key=progress_key,
task_registry=task_registry,
) )
return _SchemaResponse( return _SchemaResponse(
@@ -26,7 +26,7 @@
### P1:需要优先治理的真实债务 ### P1:需要优先治理的真实债务
1. **后台任务的统一所有权已覆盖 API 入口,但仍有更深层任务机制待分级。** `app/runtime/tasks.py` 已建立 lifespan 级 TaskRegistry,启动收尾、插件 Release 刷新、Webhook E0 广播、CookieCloud E1 手工调度、消息入口、Seerr 订阅和 WebAgent 断线后执行/快照保存均不再维护端点模块级任务集合或 Starlette 回调,shutdown 会停止接收、取消并有限等待,且生命周期清单明确登记其顺序。主仓 `app/` 已无裸 FastAPI `BackgroundTasks`;当前仍有约 `50` 个更底层 `create_task`/等价任务创建点,与线程池和 APScheduler 并存,后续需逐项确认 owner、取消、等待、重试、幂等和是否 durable,关键业务副作用优先接入已有 Outbox/恢复表。 1. **后台任务的统一所有权已覆盖 API 入口,但仍有更深层任务机制待分级。** `app/runtime/tasks.py` 已建立 lifespan 级 TaskRegistry,启动收尾、插件 Release 刷新、Webhook E0 广播、CookieCloud E1 手工调度、消息入口、Seerr 订阅、整理历史 AI 重做和 WebAgent 断线后执行/快照保存均不再维护端点模块级任务集合或 Starlette 回调,shutdown 会停止接收、取消并有限等待,且生命周期清单明确登记其顺序。主仓 `app/` 已无裸 FastAPI `BackgroundTasks`;当前仍有约 `50` 个更底层 `create_task`/等价任务创建点,与线程池和 APScheduler 并存,后续需逐项确认 owner、取消、等待、重试、幂等和是否 durable,关键业务副作用优先接入已有 Outbox/恢复表。
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. **查询侧数据库兼容 ABI 已完成正式装饰器清零。** 写事务装饰器和正式 `db_query/async_db_query` 均为 `0`。站点、消息、用户、订阅、下载/整理历史、工作流、MediaServer、SiteUserData、AgentChat、AgentTaskRun、TransferPending、SystemConfig、PassKey 和 SubscribeHistory 的宿主查询已迁到显式 Session 路径;对应旧插件 Model 调用由独立 `legacy_*` 外壳保留,可同时接受显式 Session 与无 Session 的位置/关键字参数。后续重点转为减少 ORM 对象跨层流转,并保持正式装饰器零回退。 3. **查询侧数据库兼容 ABI 已完成正式装饰器清零。** 写事务装饰器和正式 `db_query/async_db_query` 均为 `0`。站点、消息、用户、订阅、下载/整理历史、工作流、MediaServer、SiteUserData、AgentChat、AgentTaskRun、TransferPending、SystemConfig、PassKey 和 SubscribeHistory 的宿主查询已迁到显式 Session 路径;对应旧插件 Model 调用由独立 `legacy_*` 外壳保留,可同时接受显式 Session 与无 Session 的位置/关键字参数。后续重点转为减少 ORM 对象跨层流转,并保持正式装饰器零回退。
4. **组合根和全局状态仍形成复杂的隐式运行时图。** Singleton 实例、模块级 provider、`configure_*` 注册函数和兼容 Facade 同时存在;它们解决了旧 ABI 和启动顺序问题,但增加测试污染、重复装配、实例身份和初始化顺序风险。`app/startup/lifecycle/__init__.py:161-376` 已有声明式生命周期,`app/startup/modules_initializer.py:505-530` 也有分阶段关闭,但尚未做到所有进程级资源都只通过 typed HostRuntime 访问。后续应以“新代码禁止新增 Service Locator/Singleton 依赖、旧入口有命中观测”为 ratchet。 4. **组合根和全局状态仍形成复杂的隐式运行时图。** Singleton 实例、模块级 provider、`configure_*` 注册函数和兼容 Facade 同时存在;它们解决了旧 ABI 和启动顺序问题,但增加测试污染、重复装配、实例身份和初始化顺序风险。`app/startup/lifecycle/__init__.py:161-376` 已有声明式生命周期,`app/startup/modules_initializer.py:505-530` 也有分阶段关闭,但尚未做到所有进程级资源都只通过 typed HostRuntime 访问。后续应以“新代码禁止新增 Service Locator/Singleton 依赖、旧入口有命中观测”为 ratchet。
@@ -830,6 +830,9 @@ ADR 必须逐个映射当前 Event、BackgroundTasks、Scheduler job、Agent tas
不能因为进程内任务已统一登记而宣称崩溃可恢复。 不能因为进程内任务已统一登记而宣称崩溃可恢复。
- 订阅手工搜索、消息入口和 Seerr 订阅均已按 E0/E1 登记;其他关键业务副作用继续按等级逐项迁移, - 订阅手工搜索、消息入口和 Seerr 订阅均已按 E0/E1 登记;其他关键业务副作用继续按等级逐项迁移,
需要可靠交付的路径仍走 ARCH-251 的 Outbox/幂等切片,不扩大插件事件或 API payload。 需要可靠交付的路径仍走 ARCH-251 的 Outbox/幂等切片,不扩大插件事件或 API payload。
- 整理历史单条与批量 AI 重做分别登记为 `api.history.ai_redo`
`api.history.ai_redo_batch`;请求响应、进度键、Agent prompt、输出回调与旧直接调用入口保持不变。
两类任务随 lifespan shutdown 取消并有限等待,但仍属于进程内 E1 工作,不宣称崩溃后自动恢复。
#### ARCH-251:用现有数据库做首个 durable side-effect pilot #### ARCH-251:用现有数据库做首个 durable side-effect pilot
+3 -2
View File
@@ -13,8 +13,8 @@
"runtime_to_db": [], "runtime_to_db": [],
"workflow_to_db": [] "workflow_to_db": []
}, },
"edge_count": 6470, "edge_count": 6471,
"edge_sha256": "571b66b75b51e801a053e4164c5dbe0eb580b91f20cf84fa554e13b9469dd650", "edge_sha256": "30a14c4218ffd8e2798ab93231e7fa0d4ed3e368a5f1cdf75b09336e863e0a50",
"edges": [ "edges": [
"app -> app.runtime", "app -> app.runtime",
"app -> app.runtime.compat", "app -> app.runtime.compat",
@@ -1836,6 +1836,7 @@
"app.api.endpoints.history -> app.runtime.config", "app.api.endpoints.history -> app.runtime.config",
"app.api.endpoints.history -> app.runtime.log", "app.api.endpoints.history -> app.runtime.log",
"app.api.endpoints.history -> app.runtime.progress", "app.api.endpoints.history -> app.runtime.progress",
"app.api.endpoints.history -> app.runtime.tasks",
"app.api.endpoints.history -> app.schemas", "app.api.endpoints.history -> app.schemas",
"app.api.endpoints.history -> app.schemas.common", "app.api.endpoints.history -> app.schemas.common",
"app.api.endpoints.history -> app.schemas.history", "app.api.endpoints.history -> app.schemas.history",
+44 -1
View File
@@ -3,7 +3,7 @@
import asyncio import asyncio
from types import SimpleNamespace from types import SimpleNamespace
from app.api.endpoints import message, site, subscribe, webhook from app.api.endpoints import history, message, site, subscribe, webhook
from app.runtime.tasks import TaskRegistry from app.runtime.tasks import TaskRegistry
@@ -19,6 +19,17 @@ class _TaskRegistry(TaskRegistry):
"""保存函数、参数和 owner。""" """保存函数、参数和 owner。"""
self.calls.append((function, args, kwargs, owner)) self.calls.append((function, args, kwargs, owner))
def create(
self,
coroutine,
*,
owner: str,
cancel_on_shutdown: bool = True,
) -> None:
"""保存异步任务登记参数,并关闭未执行的 coroutine。"""
coroutine.close()
self.calls.append((None, (), {"cancel_on_shutdown": cancel_on_shutdown}, owner))
class _WebhookRequest: class _WebhookRequest:
"""提供 webhook 端点读取的最小请求接口。""" """提供 webhook 端点读取的最小请求接口。"""
@@ -139,3 +150,35 @@ def test_seerr_subscribe_uses_task_registry(monkeypatch) -> None:
"username": "tester", "username": "tester",
} }
assert owner == "api.subscribe.seerr" assert owner == "api.subscribe.seerr"
def test_history_ai_redo_uses_task_registry() -> None:
"""单条历史 AI 重做应登记宿主任务并使用稳定 owner。"""
registry = _TaskRegistry()
history._start_ai_redo_task(
history_id=7,
prompt="整理记录",
progress_key="progress-7",
task_registry=registry,
)
assert registry.calls == [
(None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo")
]
def test_history_batch_ai_redo_uses_task_registry() -> None:
"""批量历史 AI 重做应登记宿主任务并区分批量 owner。"""
registry = _TaskRegistry()
history._start_batch_ai_redo_task(
history_ids=[7, 8],
prompt="批量整理",
progress_key="progress-batch",
task_registry=registry,
)
assert registry.calls == [
(None, (), {"cancel_on_shutdown": True}, "api.history.ai_redo_batch")
]