refactor: own agent activity tasks

This commit is contained in:
jxxghp
2026-08-24 00:32:16 +08:00
parent 6b65dc5d12
commit 939e4c1cdc
4 changed files with 54 additions and 4 deletions
+8 -2
View File
@@ -40,6 +40,7 @@ from app.agent.policy.sanitizer import (
)
from app.agent.tools.tags import ToolTag
from app.runtime.log import logger
from app.runtime.tasks import TaskRegistry, get_task_registry
# 活动日志保留天数
DEFAULT_RETENTION_DAYS = 7
@@ -511,12 +512,14 @@ class ActivityLogMiddleware(AgentMiddleware[ActivityLogState, ContextT, Response
retention_days: int = DEFAULT_RETENTION_DAYS,
prompt_load_days: int = PROMPT_LOAD_DAYS,
stream_handler: Optional[Any] = None,
task_registry: Optional[TaskRegistry] = None,
) -> None:
"""初始化活动日志中间件。"""
"""初始化活动日志中间件,并绑定宿主后台任务 owner"""
self.activity_dir = activity_dir
self.retention_days = retention_days
self.prompt_load_days = prompt_load_days
self.stream_handler = stream_handler
self._task_registry = task_registry or get_task_registry()
self._background_tasks: set[asyncio.Task[None]] = set()
self._tool_provider = _ActivityLogToolProvider(activity_dir=activity_dir)
self.tools = [
@@ -626,7 +629,10 @@ class ActivityLogMiddleware(AgentMiddleware[ActivityLogState, ContextT, Response
def _schedule_activity_recording(self, messages: list) -> None:
"""提交后台活动记录任务,不阻塞当前 Agent 会话结束。"""
task = asyncio.create_task(self._record_activity(messages))
task = self._task_registry.create(
self._record_activity(messages),
owner="agent.activity_log.record",
)
self._background_tasks.add(task)
task.add_done_callback(self._on_activity_recording_done)
@@ -81,6 +81,10 @@
### P1:需要优先治理的真实债务
1. **后台任务的统一所有权已覆盖 API 入口,但仍有更深层任务机制待分级。** `app/runtime/tasks.py` 已建立 lifespan 级 TaskRegistry,启动收尾、插件 Release 刷新、Webhook E0 广播、CookieCloud E1 手工调度、消息入口、Seerr 订阅、整理历史 AI 重做、OpenAI/Anthropic 协议流和 WebAgent 断线后执行/快照保存均不再维护端点模块级任务集合或 Starlette 回调,shutdown 会停止接收、取消并有限等待,且生命周期清单明确登记其顺序。主仓 `app/` 已无裸 FastAPI `BackgroundTasks`;当前仍有约 `50` 个更底层 `create_task`/等价任务创建点,与线程池和 APScheduler 并存,后续需逐项确认 owner、取消、等待、重试、幂等和是否 durable,关键业务副作用优先接入已有 Outbox/恢复表。
Agent 活动摘要已完成一个深层 owner 切片:中间件仍保留自己的完成回调与非阻塞语义,但任务创建统一
经 lifespan `TaskRegistry` 登记为 `agent.activity_log.record`,宿主关停会取消并有限等待,不再形成
绕过全局关停预算的第二套后台任务集合。该摘要属于可丢弃 E1 观测数据,不宣称 durable。
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 隐式事务零回退。
+3 -2
View File
@@ -13,8 +13,8 @@
"runtime_to_db": [],
"workflow_to_db": []
},
"edge_count": 6499,
"edge_sha256": "54ddd2d6cd5669354a2d1303b24c95e573368e145fdd74e4240b01bb8ef1f687",
"edge_count": 6500,
"edge_sha256": "6fa263f6f10aa51b4395aa41cf32a0c24254c2131d7882341936412400a10e4a",
"edges": [
"app -> app.runtime",
"app -> app.runtime.compat",
@@ -255,6 +255,7 @@
"app.agent.middleware.activity_log -> app.agent.tools.tags",
"app.agent.middleware.activity_log -> app.runtime",
"app.agent.middleware.activity_log -> app.runtime.log",
"app.agent.middleware.activity_log -> app.runtime.tasks",
"app.agent.middleware.jobs -> app.agent",
"app.agent.middleware.jobs -> app.agent.middleware",
"app.agent.middleware.jobs -> app.agent.middleware.utils",
+39
View File
@@ -16,6 +16,7 @@ from app.agent.middleware.activity_log import (
)
from app.agent.tools.factory import MoviePilotToolFactory
from app.agent.tools.tags import ToolTag
from app.runtime.tasks import TaskRegistry
def _write_activity_log(activity_dir, date_str: str, lines: list[str]) -> None:
@@ -275,6 +276,44 @@ def test_activity_log_after_agent_does_not_wait_for_summary(tmp_path):
append_mock.assert_awaited_once_with("用户要求检查下载任务,助手调用工具完成检查。")
def test_activity_log_background_task_follows_host_shutdown(tmp_path):
"""活动摘要任务必须登记 owner,并随宿主关停取消和收敛。"""
async def _run_test():
"""启动阻塞摘要后关闭登记器,返回 owner 与最终任务状态。"""
registry = TaskRegistry()
started = asyncio.Event()
cancelled = asyncio.Event()
async def _blocked_record(_messages: list) -> None:
"""保持记录任务运行,直到宿主关停发出取消。"""
started.set()
try:
await asyncio.Event().wait()
except asyncio.CancelledError:
cancelled.set()
raise
middleware = ActivityLogMiddleware(
activity_dir=str(tmp_path),
task_registry=registry,
)
with patch.object(middleware, "_record_activity", side_effect=_blocked_record):
middleware._schedule_activity_recording([])
await started.wait()
owners = tuple(record.owner for record in registry.records)
converged = await registry.shutdown(timeout_seconds=1.0)
await asyncio.sleep(0)
return owners, converged, cancelled.is_set(), middleware._background_tasks
owners, converged, cancelled, background_tasks = asyncio.run(_run_test())
assert owners == ("agent.activity_log.record",)
assert converged is True
assert cancelled is True
assert background_tasks == set()
def test_query_activity_logs_filters_by_keyword_and_date(tmp_path):
"""活动日志查询应支持日期和关键词过滤。"""
_write_activity_log(