From 939e4c1cdc0fa7c0f99e8484393f783eff6914bd Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 00:32:16 +0800 Subject: [PATCH] refactor: own agent activity tasks --- app/agent/middleware/activity_log.py | 10 ++++- .../backend-architecture-next-stage.md | 4 ++ .../architecture/dependency-baseline.json | 5 ++- tests/test_agent_activity_log.py | 39 +++++++++++++++++++ 4 files changed, 54 insertions(+), 4 deletions(-) diff --git a/app/agent/middleware/activity_log.py b/app/agent/middleware/activity_log.py index d3422fc46..da24312be 100644 --- a/app/agent/middleware/activity_log.py +++ b/app/agent/middleware/activity_log.py @@ -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) diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index 61197d038..1b1a0fca0 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -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 隐式事务零回退。 diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index dd3410928..78c1195b1 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": 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", diff --git a/tests/test_agent_activity_log.py b/tests/test_agent_activity_log.py index edc633ae6..c08f23088 100644 --- a/tests/test_agent_activity_log.py +++ b/tests/test_agent_activity_log.py @@ -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(