From 40d4a10c99cf2f39af568c81335f9a81e6112ee6 Mon Sep 17 00:00:00 2001 From: jxxghp Date: Mon, 24 Aug 2026 05:40:35 +0800 Subject: [PATCH] refactor: own agent session cleanup tasks --- app/chain/message.py | 22 ++------- .../adr/0007-background-action-reliability.md | 2 + .../backend-architecture-next-stage.md | 13 ++++- .../architecture/dependency-baseline.json | 5 +- tests/test_agent_message_routing.py | 49 +++++++++++++++++++ 5 files changed, 71 insertions(+), 20 deletions(-) diff --git a/app/chain/message.py b/app/chain/message.py index 456ef996a..aa8cfa530 100644 --- a/app/chain/message.py +++ b/app/chain/message.py @@ -22,6 +22,7 @@ from app.chain.subscribe import SubscribeChain from app.chain.transfer import TransferChain from app.chain.interaction import MediaInteractionChain as _MediaInteractionChain from app.runtime.config import global_vars +from app.runtime.tasks import get_task_registry from app.application.messaging.agent import agent_interaction_manager, parse_agent_choice_callback from app.application.messaging.interaction import InteractionContext, InteractionDispatch from app.application.messaging.media import media_interaction_manager @@ -66,9 +67,10 @@ class MessageChain(ChainBase): clear_task = manager.clear_session( session_id=session_id, user_id=str(userid) ) - asyncio.run_coroutine_threadsafe( + get_task_registry().submit_threadsafe( clear_task, - global_vars.loop, + loop=global_vars.loop, + owner="chain.message.agent_session_clear", ) except Exception as e: if clear_task: @@ -956,21 +958,7 @@ class MessageChain(ChainBase): # 如果有会话ID,同时清除智能体的会话记忆 if session_id: - manager = get_running_agent_manager() - clear_task = None - if manager is not None: - try: - clear_task = manager.clear_session( - session_id=session_id, user_id=str(userid) - ) - asyncio.run_coroutine_threadsafe( - clear_task, - global_vars.loop, - ) - except Exception as e: - if clear_task: - clear_task.close() - logger.warning(f"清除智能体会话记忆失败: {e}") + self._schedule_agent_session_clear(session_id, userid) self.post_message( Message( diff --git a/docs/adr/0007-background-action-reliability.md b/docs/adr/0007-background-action-reliability.md index a25056933..97f3b1c4a 100644 --- a/docs/adr/0007-background-action-reliability.md +++ b/docs/adr/0007-background-action-reliability.md @@ -88,6 +88,8 @@ Event Contract Registry 是 53 个事件的逐项机器清单。下表按相同 - 已登记的周期 Agent task:E1,重启时通过任务定义重建;单次执行要有 execution 记录。 - Agent 创建/修改订阅、删除数据等工具:业务事务按 E2/E3;聊天输出不能替代业务完成证据。 - 会话 stop/cancel:E0 控制信号;被取消工具的底层阻塞 I/O 可能继续,资源所有者必须最终回收。 +- 过期会话与远程清理命令统一经 `chain.message.agent_session_clear` owner 提交 Agent 资源释放;两条入口 + 不再各自维护裸跨线程任务,宿主关停会取消并等待已登记清理。该清理仍是 E0 资源回收,不跨重启恢复。 - OpenAI/Anthropic 协议流的请求级 Agent worker 由 `api.openai.stream` / `api.anthropic.stream` 登记并在 lifespan shutdown 时取消;它们仍是 E0 请求交付,不提供跨重启恢复。 - stdio MCP 的 stderr reader 属于会话资源内部任务;会话退出时先取消并等待 reader 收口,再终止子进程,避免 diff --git a/docs/refactor/backend-architecture-next-stage.md b/docs/refactor/backend-architecture-next-stage.md index a7e3d0a7b..a4c1ffc76 100644 --- a/docs/refactor/backend-architecture-next-stage.md +++ b/docs/refactor/backend-architecture-next-stage.md @@ -267,13 +267,24 @@ - 整理失败按钮的 AI 接管从裸 `run_coroutine_threadsafe` 迁为 `chain.transfer.ai_takeover`,消息文案、Agent prompt/参数、回调返回和 V2/V3 插件 ABI 保持不变;该动作仍为 E0,不宣称跨进程恢复。 +### 长期整改阶段 26:Agent 会话清理提交路径统一(2026-08-24) + +- 过期会话回收与远程清理命令原先各自构造 `clear_session()` 协程并裸提交主循环,形成同一目标的两套 + 实现;远程命令现在复用 `_schedule_agent_session_clear()` 唯一入口。 +- 唯一入口通过 TaskRegistry 以 `chain.message.agent_session_clear` owner 跨线程提交,目标循环先登记再执行, + shutdown 可取消并等待;调度失败仍关闭协程并记录告警,不泄漏 Agent 资源清理对象。 +- `/clear_session` 命令、会话映射、成功/空会话提示、Agent manager 合同及 V2/V3 插件 ABI 均未修改;实时 + Agent 消息和停止命令仍保留各自的 Future 结果观察,不被机械改成 fire-and-forget。 +- 依赖基线仅新增 `app.chain.message -> app.runtime.tasks`,模块数保持 `806`,内部边为 `6541`;12 组禁止 + 边与唯一隔离 TMDB SCC 均未变化。 + ### 总体判断 当前架构总体合理,已经从跨层混合的遗留单体收敛为**边界清晰的模块化单体**: - 继续采用单进程控制面是正确选择,不建议现在拆成微服务;插件、调度器、工作流、事件和数据库共享进程内状态,拆分会放大部署、事务和兼容成本。 - `foundation/domain/runtime/adapters/application/chain/api/startup` 的职责方向基本成立;宿主架构基线、复杂度 ratchet、异步阻塞 ratchet 当前均通过。 -- 依赖图当前为 `806` 个 Python 模块、`6539` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 +- 依赖图当前为 `806` 个 Python 模块、`6541` 条内部导入边;唯一非平凡 SCC 位于隔离的 TMDB 第三方移植包内部,不应为了指标归零重写。 - 当前主要风险已经从“目录和依赖失控”转移到运行时协议、后台副作用的可靠性和遗留兼容面。换言之,下一阶段重点应是**语义收口和可验证性**,而不是继续搬文件或机械拆大文件。 综合评价:架构方向可持续,生产可用性较高;可演进性仍处于中等水平。现阶段没有静态审计发现必须立即推倒重来的 P0 架构问题,但存在需要按 P1/P2 计划治理的真实债务。 diff --git a/tests/fixtures/architecture/dependency-baseline.json b/tests/fixtures/architecture/dependency-baseline.json index b10112719..55b7babf0 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": 6540, - "edge_sha256": "3beb7c6b906993a8ff9d0cb9b3f7c912e1a7859785824f472547b8fecf8a4c51", + "edge_count": 6541, + "edge_sha256": "c7e2a8537efa709d30c82930fa4ee21d91fe8f24578a44b4357563fa7bfc1a0c", "edges": [ "app -> app.runtime", "app -> app.runtime.compat", @@ -3197,6 +3197,7 @@ "app.chain.message -> app.runtime", "app.chain.message -> app.runtime.config", "app.chain.message -> app.runtime.log", + "app.chain.message -> app.runtime.tasks", "app.chain.message -> app.schemas", "app.chain.message -> app.schemas.message", "app.chain.message -> app.schemas.notification", diff --git a/tests/test_agent_message_routing.py b/tests/test_agent_message_routing.py index cd040dd70..dc91c3813 100644 --- a/tests/test_agent_message_routing.py +++ b/tests/test_agent_message_routing.py @@ -35,6 +35,55 @@ def _running_loop_stub() -> Mock: ) +def test_agent_session_clear_uses_owned_threadsafe_submission(): + """Agent 会话清理应登记稳定 owner,纳入宿主关停收口。""" + loop = _running_loop_stub() + manager = Mock(clear_session=AsyncMock()) + + def submit(coroutine, **_kwargs): + """关闭测试协程,避免替身提交留下未等待警告。""" + coroutine.close() + return Future() + + with patch.object(global_vars, "CURRENT_EVENT_LOOP", loop), patch( + "app.chain.message.get_running_agent_manager", return_value=manager + ), patch("app.chain.message.get_task_registry") as get_registry: + get_registry.return_value.submit_threadsafe.side_effect = submit + MessageChain._schedule_agent_session_clear("session-1", "10001") + + manager.clear_session.assert_called_once_with( + session_id="session-1", user_id="10001" + ) + get_registry.return_value.submit_threadsafe.assert_called_once() + assert get_registry.return_value.submit_threadsafe.call_args.kwargs == { + "loop": loop, + "owner": "chain.message.agent_session_clear", + } + + +def test_remote_session_clear_reuses_owned_clear_scheduler(): + """远程清理命令应复用唯一 Agent 会话清理入口。""" + chain = MessageChain() + session_service = Mock() + session_service.clear.return_value = "session-1" + + with patch.object( + chain, "_message_session_service", return_value=session_service + ), patch.object( + chain, "_schedule_agent_session_clear" + ) as schedule_clear, patch.object(chain, "post_message") as post_message: + chain.remote_clear_session( + channel=NotificationChannel.Telegram, + userid="10001", + source="telegram-test", + ) + + schedule_clear.assert_called_once_with("session-1", "10001") + notification = post_message.call_args.args[0] + assert notification.title == "智能体会话已清除,下次将创建新的会话" + assert notification.save_history is False + + def test_explicit_ai_message_bypasses_pending_media_interaction(): """显式 /ai 消息应绕过误触发的媒体交互状态并回到 Agent 会话。""" chain = MessageChain()