refactor: own agent session cleanup tasks

This commit is contained in:
jxxghp
2026-08-24 05:40:35 +08:00
parent d694cd081c
commit 40d4a10c99
5 changed files with 71 additions and 20 deletions
+5 -17
View File
@@ -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(
@@ -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 收口,再终止子进程,避免
@@ -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 计划治理的真实债务。
+3 -2
View File
@@ -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",
+49
View File
@@ -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()