Agent 对话中断恢复用户断线后如何无缝继续会话一、用户切了个网络Agent 忘了刚才聊到哪了Agent 对话最大的不确定性不是模型幻觉不是工具调用失败是网络。用户在地铁上跟 Agent 聊产品需求分析信号断了 30 秒重连回来——对话历史全部丢失Agent 重新打招呼你好有什么可以帮你。用户心态崩溃重来一次体验灾难。HTTP 是无状态协议WebSocket 可以保持连接但在移动端切网络4G → WiFi → 4G的场景下WebSocket 照样会断开。所以对话恢复不能依赖连接不断——连接一定会断——而要依赖状态不丢。对话中断恢复的核心不是重连技术WebSocket 自动重连库很多而是对话状态的持久化和恢复。状态包括三个关键部分对话历史messages、工具调用中间态当前正在调用的工具及参数、用户上下文用户偏好、当前任务目标。只要这三个状态有 Checkpoint 持久化用户无论断多久重连都能接着聊。二、底层机制与原理剖析会话恢复的三个关键设计消息序号Sequence Number每条消息用户发送、Agent 回复都有一个严格递增的序号。重连时客户端上报自己收到的最后一条消息序号服务端从这个序号之后开始推送。这个设计解决了服务端不知道客户端收到多少的问题。Session Snapshot会话快照不只是消息列表——还包括 Agent 当前的状态。如果 Agent 正在执行工具调用如搜索恢复时不能让 Agent 重新搜一遍——它应该知道上次搜索的结果是 X。所以快照中要包含工具调用的中间结果。重连窗口会话不能永远保留——设置一个窗口。例如断开 5 分钟内的会话自动恢复5 分钟以上的视为过期Agent 任务可能已经失去时效性或者用户已经重开了一个新会话。三、生产级代码实现 Agent 对话中断恢复系统 核心数据结构SessionSnapshot 核心策略消息序号 增量同步 重连窗口 import asyncio import json import time import uuid import logging from typing import Dict, List, Optional, Any from dataclasses import dataclass, field from enum import Enum from datetime import datetime, timedelta import redis.asyncio as redis logging.basicConfig(levellogging.INFO) logger logging.getLogger(__name__) class SessionStatus(Enum): ACTIVE active INTERRUPTED interrupted # 连接断开等待恢复 EXPIRED expired # 超时未恢复已清理 COMPLETED completed class MessageRole(Enum): USER user AGENT agent SYSTEM system TOOL tool dataclass class Message: 单条消息 seq: int # 消息序号严格递增 role: MessageRole content: str timestamp: float # Unix timestamp metadata: Dict[str, Any] field(default_factorydict) # 工具调用信息等 def to_dict(self) - dict: return { seq: self.seq, role: self.role.value, content: self.content, timestamp: self.timestamp, metadata: self.metadata, } classmethod def from_dict(cls, data: dict) - Message: return cls( seqdata[seq], roleMessageRole(data[role]), contentdata[content], timestampdata[timestamp], metadatadata.get(metadata, {}), ) dataclass class SessionSnapshot: 会话快照 —— 持久化到 Redis session_id: str user_id: str status: SessionStatus SessionStatus.ACTIVE messages: List[Message] field(default_factorylist) last_seq: int 0 # Agent 中间态 pending_tool_call: Optional[Dict[str, Any]] None # 正在执行的工具调用 agent_context: Dict[str, Any] field(default_factorydict) # Agent 上下文 # 时间窗口 created_at: float 0.0 last_active_at: float 0.0 interrupted_at: Optional[float] None def __post_init__(self): now time.time() if not self.created_at: self.created_at now if not self.last_active_at: self.last_active_at now def to_dict(self) - dict: return { session_id: self.session_id, user_id: self.user_id, status: self.status.value, messages: [m.to_dict() for m in self.messages], last_seq: self.last_seq, pending_tool_call: self.pending_tool_call, agent_context: self.agent_context, created_at: self.created_at, last_active_at: self.last_active_at, interrupted_at: self.interrupted_at, } classmethod def from_dict(cls, data: dict) - SessionSnapshot: return cls( session_iddata[session_id], user_iddata[user_id], statusSessionStatus(data[status]), messages[Message.from_dict(m) for m in data.get(messages, [])], last_seqdata.get(last_seq, 0), pending_tool_calldata.get(pending_tool_call), agent_contextdata.get(agent_context, {}), created_atdata.get(created_at, 0), last_active_atdata.get(last_active_at, 0), interrupted_atdata.get(interrupted_at), ) class SessionManager: 会话管理器 设计思路 1. 所有会话状态持久化到 Redis内存不保存 2. 支持从任意断点恢复基于消息序号 3. 重连窗口5 分钟——过期会话自动清理 RECONNECT_WINDOW 300 # 重连窗口5 分钟 SESSION_TTL 3600 # Session TTL1 小时 KEY_PREFIX agent:session: def __init__(self, redis_client: redis.Redis): self.redis redis_client def _key(self, session_id: str) - str: return f{self.KEY_PREFIX}{session_id} async def create_session(self, user_id: str) - SessionSnapshot: 创建新会话 session SessionSnapshot( session_idstr(uuid.uuid4()), user_iduser_id, ) await self._save(session) logger.info(Session created: %s for user %s, session.session_id, user_id) return session async def get_session(self, session_id: str) - Optional[SessionSnapshot]: 获取会话快照 如果会话状态为 INTERRUPTED计算是否在重连窗口内 - 在窗口内状态恢复为 ACTIVE返回完整快照 - 窗口外状态标记为 EXPIRED返回 None key self._key(session_id) raw await self.redis.get(key) if not raw: return None session SessionSnapshot.from_dict(json.loads(raw)) if session.status SessionStatus.INTERRUPTED: assert session.interrupted_at is not None elapsed time.time() - session.interrupted_at if elapsed self.RECONNECT_WINDOW: # 超时——标记为过期 session.status SessionStatus.EXPIRED await self._save(session) logger.info(Session %s expired (interrupted %.0fs ago), session_id, elapsed) return None else: # 在窗口内——恢复 session.status SessionStatus.ACTIVE session.interrupted_at None session.last_active_at time.time() await self._save(session) logger.info(Session %s recovered after %.0fs, session_id, elapsed) return session async def mark_interrupted(self, session_id: str): 标记会话中断 key self._key(session_id) raw await self.redis.get(key) if not raw: return session SessionSnapshot.from_dict(json.loads(raw)) session.status SessionStatus.INTERRUPTED session.interrupted_at time.time() await self._save(session) logger.info(Session %s marked interrupted, session_id) async def add_message(self, session_id: str, role: MessageRole, content: str, metadata: Optional[Dict] None) - Optional[int]: 添加一条消息返回消息序号 key self._key(session_id) raw await self.redis.get(key) if not raw: logger.error(Session %s not found for add_message, session_id) return None session SessionSnapshot.from_dict(json.loads(raw)) session.last_seq 1 msg Message( seqsession.last_seq, rolerole, contentcontent, timestamptime.time(), metadatametadata or {}, ) session.messages.append(msg) session.last_active_at time.time() await self._save(session) return session.last_seq async def get_messages_after(self, session_id: str, after_seq: int) - List[Message]: 获取指定序号之后的消息用于重连后的增量同步 key self._key(session_id) raw await self.redis.get(key) if not raw: return [] session SessionSnapshot.from_dict(json.loads(raw)) return [m for m in session.messages if m.seq after_seq] async def update_agent_context(self, session_id: str, context: Dict[str, Any]): 更新 Agent 中间上下文 key self._key(session_id) raw await self.redis.get(key) if not raw: return session SessionSnapshot.from_dict(json.loads(raw)) session.agent_context.update(context) session.last_active_at time.time() await self._save(session) async def _save(self, session: SessionSnapshot): 持久化到 Redis 并设置 TTL key self._key(session.session_id) await self.redis.setex( key, self.SESSION_TTL, json.dumps(session.to_dict(), ensure_asciiFalse), ) # --------------------------------------------------------------------------- # 模拟 WebSocket 连接处理 # --------------------------------------------------------------------------- async def handle_websocket(websocket, session_manager: SessionManager): WebSocket 事件处理器 处理两种连接场景 1. 新会话无 session_id→ 创建会话 2. 恢复会话有 session_id→ 从断点恢复 try: # 1. 接收初始消息包含 session_id init_msg await asyncio.wait_for(websocket.receive_text(), timeout10.0) init_data json.loads(init_msg) session_id init_data.get(session_id) if session_id: # 尝试恢复会话 session await session_manager.get_session(session_id) if session: # 发送从断点之后的消息 last_received init_data.get(last_seq, 0) missed await session_manager.get_messages_after(session_id, last_received) for msg in missed: await websocket.send_text(json.dumps({ type: history, **msg.to_dict(), })) else: # 会话过期或不存在——创建新会话 session await session_manager.create_session(init_data.get(user_id, anonymous)) session_id session.session_id else: # 新会话 session await session_manager.create_session(init_data.get(user_id, anonymous)) session_id session.session_id # 2. 回执 session_id客户端保存用于下次重连 await websocket.send_text(json.dumps({ type: session_created, session_id: session_id, })) # 3. 消息循环 async for raw_message in websocket: msg_data json.loads(raw_message) content msg_data.get(content, ) role MessageRole.USER # 保存用户消息 await session_manager.add_message(session_id, role, content) # 模拟 Agent 回复生产环境调用 Agent 引擎 response fAgent 回复: {content} seq await session_manager.add_message(session_id, MessageRole.AGENT, response) await websocket.send_text(json.dumps({ type: agent_message, seq: seq, content: response, })) except asyncio.TimeoutError: logger.error(WebSocket: initial message timeout) except Exception as e: logger.error(WebSocket error: %s, e) # 连接异常断开标记会话为 INTERRUPTED if session_id: await session_manager.mark_interrupted(session_id)四、边界分析与架构权衡重连窗口的设定太短如 30 秒地铁过隧道、电梯进出这种网络短暂中断的用户无法恢复太长如 1 小时会话长期占用 Redis 内存且恢复一个 30 分钟前的会话没有意义推荐 5 分钟覆盖大部分短暂断网同时控制资源开销消息重放的性能如果对话历史有 200 条消息重连时逐条重放会导致用户端长时间刷屏优化重放时把历史消息合并为一条type: history_batch消息客户端一次性渲染对于极长对话1000 条只重放最近 50 条 摘要安全性注意事项session_id 不能是可猜测的——必须使用 UUID v4 或更高强度的随机标识用户只能恢复自己的会话——get_session时需要校验 session 的 user_id 与请求的 user_id 是否一致敏感会话如涉及支付、个人信息的对话恢复时需要二次认证五、总结Agent 对话中断恢复核心不是传输层WebSocket 重连库很多是应用层的状态持久化。消息序号 Session Snapshot 重连窗口三个机制组合序号解决增量同步快照解决 Agent 中间状态恢复窗口解决资源控制。关键是认知转变——连接一定会断要让断开无感。