Agent 心跳与健康检查长连接场景下的会话状态监控Agent 连着连着就没了反应——你不知道它是真的在思考还是已经悄悄挂了。一、场景痛点你的 Agent 系统用 WebSocket 维持长连接用户发一条消息后 Agent 需要调用多个工具耗时可能 10-60 秒。问题来了这 60 秒里用户不知道 Agent 是在处理还是已经挂了。如果 Agent 进程崩溃WebSocket 连接不会自动断开——客户端一直等着等到超时才报错但此时用户已经以为系统卡死了。你加了一个进度条每 5 秒推送一条还在处理中的消息。但 Agent 真的挂了的时候进度条还在转——因为推送线程和 Agent 执行线程是独立的Agent 挂了推送线程还在跑。更棘手的是会话恢复。Agent 处理到第 3 步时崩溃了用户重连后从第 1 步重新开始——前两步的结果全部丢失。如果是付费场景每次调用消耗 token重复执行的成本直接翻倍。核心矛盾长连接场景下Agent 的存活状态和处理进度必须被持续监控不能靠连接还在就以为还活着。二、底层机制与原理剖析2.1 心跳机制的层次2.2 心跳数据模型进程心跳不只是我还活着它包含处理状态字段含义用途agent_idAgent 实例标识关联会话与实例session_id会话标识恢复会话时使用statusidle/processing/error区分状态current_step当前执行步骤进度追踪total_steps总步骤数进度百分比计算cpu_percentCPU 使用率资源监控memory_mb内存占用资源监控last_tool_call最近一次工具调用信息卡住定位timestamp心跳时间戳判断是否过期2.3 健康检查的判定逻辑健康判定不是简单的心跳在就健康。需要根据心跳间隔和业务状态综合判断健康心跳间隔 预期间隔 × 2status 不是 error亚健康心跳间隔 预期间隔 × 2 但 预期间隔 × 5status 是 processing 但 current_step 长时间不变不健康心跳间隔 预期间隔 × 5或 status 是 error或连续 3 次心跳缺失三、生产级代码实现3.1 Agent 心跳上报器// agent-heartbeat.ts —— Agent 进程心跳上报器 import { EventEmitter } from events; export enum AgentStatus { IDLE idle, // 等待用户输入 PROCESSING processing, // 正在处理用户请求 ERROR error, // 处理出错等待恢复 TERMINATING terminating, // 正在优雅关闭 } export interface HeartbeatPayload { agent_id: string; session_id: string; status: AgentStatus; current_step: number; total_steps: number; cpu_percent: number; memory_mb: number; last_tool_call: string | null; timestamp: number; // Unix 时间戳毫秒 } export class AgentHeartbeat extends EventEmitter { private agentId: string; private sessionId: string; private status: AgentStatus AgentStatus.IDLE; private currentStep: number 0; private totalSteps: number 0; private lastToolCall: string | null null; private intervalMs: number; // 心跳间隔 private maxMissedHeartbeats: number; // 允许缺失的最大心跳数 private heartbeatTimer: NodeJS.Timeout | null null; // 心跳发送通道WebSocket / HTTP / 消息队列 private sender: (payload: HeartbeatPayload) Promisevoid; constructor( agentId: string, sessionId: string, intervalMs: number 5000, // 默认 5 秒心跳间隔 maxMissedHeartbeats: number 3, sender: (payload: HeartbeatPayload) Promisevoid ) { super(); this.agentId agentId; this.sessionId sessionId; this.intervalMs intervalMs; this.maxMissedHeartbeats maxMissedHeartbeats; this.sender sender; } /** 启动心跳循环定时上报状态 */ start(): void { if (this.heartbeatTimer) return; // 已启动则不重复启动 // 定时发送心跳intervalMs 间隔 // 心跳是推模式不是拉模式——服务端不需要轮询检查 Agent 状态 this.heartbeatTimer setInterval(() { this.sendHeartbeat(); }, this.intervalMs); // 立即发送一次心跳启动时让服务端知道 Agent 已上线 this.sendHeartbeat(); } /** 停止心跳循环优雅关闭前调用 */ stop(): void { if (this.heartbeatTimer) { clearInterval(this.heartbeatTimer); this.heartbeatTimer null; } // 发送最终心跳标记为 terminating服务端知道 Agent 正在关闭 this.status AgentStatus.TERMINATING; this.sendHeartbeat(); } /** 更新处理状态Agent 每完成一步调用此方法 */ updateProgress(currentStep: number, totalSteps: number, toolCall: string): void { this.currentStep currentStep; this.totalSteps totalSteps; this.lastToolCall toolCall; this.status AgentStatus.PROCESSING; // 状态变化时立即发送一次心跳不等定时器 // 用户在等待结果状态变化应该第一时间告知服务端 this.sendHeartbeat(); } /** 标记错误状态 */ markError(): void { this.status AgentStatus.ERROR; this.sendHeartbeat(); } /** 标记空闲状态 */ markIdle(): void { this.status AgentStatus.IDLE; this.currentStep 0; this.totalSteps 0; this.lastToolCall null; this.sendHeartbeat(); } /** 发送心跳组装 payload 并调用 sender */ private async sendHeartbeat(): void { const payload: HeartbeatPayload { agent_id: this.agentId, session_id: this.sessionId, status: this.status, current_step: this.currentStep, total_steps: this.totalSteps, cpu_percent: this.getCpuUsage(), memory_mb: this.getMemoryUsage(), last_tool_call: this.lastToolCall, timestamp: Date.now(), }; try { await this.sender(payload); this.emit(heartbeat:sent, payload); } catch (err) { // 心跳发送失败不中断 Agent 处理流程 // 心跳是辅助功能核心业务不能因为心跳通道故障而停止 this.emit(heartbeat:failed, { error: err, payload }); } } /** 获取 CPU 使用率简化实现生产环境用 process.cpuUsage() */ private getCpuUsage(): number { // Node.js 的 process.cpuUsage() 返回微秒级的 CPU 时间 const usage process.cpuUsage(); const totalUsec usage.user usage.system; // 转换为百分比近似值需要采样间隔才能精确计算 return Math.min(totalUsec / 1000 / this.intervalMs, 100); } /** 获取内存使用量 */ private getMemoryUsage(): number { return process.memoryUsage().heapUsed / 1024 / 1024; // MB } }3.2 服务端健康检查监控器// health-monitor.ts —— 服务端 Agent 健康检查监控器 // 监控所有 Agent 实例的心跳判断健康状态触发告警和会话恢复 export enum HealthStatus { HEALTHY healthy, DEGRADED degraded, UNHEALTHY unhealthy, DEAD dead, } interface AgentHealthRecord { agentId: string; sessionId: string; lastHeartbeat: HeartbeatPayload; lastHeartbeatTime: number; missedHeartbeats: number; healthStatus: HealthStatus; // 停滞检测current_step 连续 N 次心跳未变化 stagnantCount: number; } export class AgentHealthMonitor extends EventEmitter { private agents: Mapstring, AgentHealthRecord new Map(); private intervalMs: number; private maxMissed: number; private stagnantThreshold: number; // 心跳停滞阈值step 不变的次数 private checkTimer: NodeJS.Timeout | null null; constructor( intervalMs: number 10000, // 每 10 秒检查一次所有 Agent maxMissed: number 3, stagnantThreshold: number 5 // 5 次心跳 step 不变判定为停滞 ) { super(); this.intervalMs intervalMs; this.maxMissed maxMissed; this.stagnantThreshold stagnantThreshold; } /** 接收 Agent 心跳更新健康记录 */ receiveHeartbeat(payload: HeartbeatPayload): void { const existing this.agents.get(payload.agent_id); if (existing) { // 检查 current_step 是否变化停滞检测 if (payload.current_step existing.lastHeartbeat.current_step payload.status AgentStatus.PROCESSING) { existing.stagnantCount; } else { existing.stagnantCount 0; } // 更新记录 existing.lastHeartbeat payload; existing.lastHeartbeatTime payload.timestamp; existing.missedHeartbeats 0; // 重新评估健康状态 this.evaluateHealth(existing); } else { // 新 Agent 上线初始化健康记录 this.agents.set(payload.agent_id, { agentId: payload.agent_id, sessionId: payload.session_id, lastHeartbeat: payload, lastHeartbeatTime: payload.timestamp, missedHeartbeats: 0, healthStatus: HealthStatus.HEALTHY, stagnantCount: 0, }); this.emit(agent:registered, { agentId: payload.agent_id }); } } /** 启动健康检查循环 */ start(): void { this.checkTimer setInterval(() { this.checkAllAgents(); }, this.intervalMs); } /** 检查所有 Agent 的健康状态 */ private checkAllAgents(): void { const now Date.now(); const expectedInterval 5000; // Agent 心跳间隔 for (const [agentId, record] of this.agents) { const elapsed now - record.lastHeartbeatTime; // 心跳缺失检测超过预期间隔则计数 1 if (elapsed expectedInterval * 2) { record.missedHeartbeats; } // 停滞检测step 不变的次数超过阈值 if (record.stagnantCount this.stagnantThreshold) { // 工具调用卡住Agent 还活着但处理停滞 this.emit(agent:stagnant, { agentId, sessionId: record.sessionId, currentStep: record.lastHeartbeat.current_step, lastToolCall: record.lastHeartbeat.last_tool_call, }); } // 重新评估健康状态 this.evaluateHealth(record); // 不健康或死亡触发告警 if (record.healthStatus HealthStatus.UNHEALTHY) { this.emit(agent:unhealthy, { agentId, sessionId: record.sessionId, missedHeartbeats: record.missedHeartbeats, }); } if (record.healthStatus HealthStatus.DEAD) { this.emit(agent:dead, { agentId, sessionId: record.sessionId, }); // 死亡 Agent 从监控列表移除不再等待心跳 // 但会话状态保留用于后续恢复 this.agents.delete(agentId); } } } /** 评估单个 Agent 的健康状态 */ private evaluateHealth(record: AgentHealthRecord): void { const previousStatus record.healthStatus; if (record.missedHeartbeats this.maxMissed * 2) { // 连续缺失超过 2 倍阈值判定死亡 record.healthStatus HealthStatus.DEAD; } else if (record.missedHeartbeats this.maxMissed) { // 连续缺失超过阈值判定不健康 record.healthStatus HealthStatus.UNHEALTHY; } else if (record.missedHeartbeats 0 || record.stagnantCount this.stagnantThreshold) { // 有缺失但未超阈值或处理停滞亚健康 record.healthStatus HealthStatus.DEGRADED; } else { // 正常心跳且处理推进中健康 record.healthStatus HealthStatus.HEALTHY; } // 状态变化时发出事件外部系统可以订阅做自动恢复 if (previousStatus ! record.healthStatus) { this.emit(health:changed, { agentId: record.agentId, from: previousStatus, to: record.healthStatus, }); } } }3.3 会话恢复与断点续传# session_recovery.py —— Agent 崩溃后的会话恢复与断点续传 import json import logging import time from datetime import datetime logger logging.getLogger(session-recovery) class SessionRecoveryManager: 会话恢复管理器Agent 崩溃后从断点继续处理 def __init__(self, storage_client, heartbeat_monitor): self.storage storage_client self.monitor heartbeat_monitor def save_checkpoint(self, session_id: str, step_index: int, step_results: dict): 保存检查点每完成一步就保存崩溃后从检查点恢复 checkpoint { session_id: session_id, step_index: step_index, step_results: step_results, timestamp: datetime.utcnow().isoformat(), } # 检查点存到对象存储比数据库更快且不影响业务表 key fcheckpoints/{session_id}/step_{step_index}.json self.storage.put(key, json.dumps(checkpoint)) def recover_session(self, session_id: str) - dict: 从最新检查点恢复会话 # 查找该会话的所有检查点取最新的 pattern fcheckpoints/{session_id}/step_*.json checkpoints self.storage.list(pattern) if not checkpoints: logger.warning(fNo checkpoints found for session {session_id}) return {step_index: 0, step_results: {}} # 取最新的检查点step_index 最大的 latest max(checkpoints, keylambda k: int(k.split(step_)[1].split(.)[0])) checkpoint_data self.storage.get(latest) checkpoint json.loads(checkpoint_data) logger.info( fRecovered session {session_id} from step {checkpoint[step_index]} ) return checkpoint def handle_dead_agent(self, agent_id: str, session_id: str): 处理死亡 Agent恢复会话并分配新 Agent # 1. 从检查点恢复会话状态 checkpoint self.recover_session(session_id) # 2. 创建新 Agent 实例传入恢复的检查点 # 新 Agent 从断点继续不从第 0 步重新开始 new_agent self.create_agent_with_checkpoint( session_id, checkpoint ) # 3. 通知用户会话恢复从第 N 步继续 logger.info( fSession {session_id} recovered: fnew agent {new_agent.agent_id}, fresuming from step {checkpoint[step_index]} ) # 4. 清理旧 Agent 的残留资源内存中的会话数据等 self.cleanup_agent_resources(agent_id) return new_agent def create_agent_with_checkpoint(self, session_id: str, checkpoint: dict): 创建新 Agent 并注入检查点数据 # 新 Agent 启动时接收检查点 # 从 checkpoint[step_index] 1 开始执行 # 前面步骤的结果从 checkpoint[step_results] 中读取 agent_config { session_id: session_id, resume_from_step: checkpoint[step_index] 1, previous_results: checkpoint[step_results], } # 调用 Agent 启动接口 return self.start_new_agent(agent_config) def cleanup_agent_resources(self, agent_id: str): 清理死亡 Agent 的残留资源 # 释放内存中的会话缓存、关闭未完成的工具调用连接等 logger.info(fCleaning up resources for dead agent {agent_id})四、边界分析与架构权衡4.1 心跳间隔的权衡心跳间隔太短1 秒网络开销大服务端处理压力大。心跳间隔太长30 秒Agent 挂了 30 秒你才知道用户已经等了 30 秒才发现系统没响应。折中基础心跳 5 秒覆盖大部分场景状态变化时立即发送一次即时心跳。这样正常情况下每 5 秒一次心跳状态变化时秒级感知。4.2 检查点的存储频率每完成一步就保存检查点意味着每步都有一次存储写入。如果步骤执行很快每步 1 秒写入频率就是每秒一次。对象存储的写入延迟约 50-100ms不影响步骤执行。但如果步骤执行很慢每步 10 秒保存频率是每 10 秒一次崩溃后最多丢失 10 秒的工作量。4.3 适用边界与禁用场景适用WebSocket/SSE 长连接的 Agent 系统、多步骤工具调用链路、需要会话恢复的付费场景禁用单次请求-响应的简单 Agent不需要心跳、短连接 HTTP API不需要长连接监控、离线批处理 Agent不需要实时状态4.4 心跳通道与业务通道的隔离心跳消息和业务消息走同一个 WebSocket 连接时如果业务消息阻塞比如大结果传输心跳也会延迟。解决方案心跳走独立连接或独立的消息类型WebSocket 的 ping 帧与数据帧是独立的。五、总结Agent 长连接场景的健康检查需要三层心跳连接层检测网络可达、进程层检测 Agent 存活与资源状态、业务层检测处理进度是否推进。单靠连接层心跳无法区分在思考和已挂掉。核心设计5 秒基础心跳 状态变化即时心跳、停滞检测step 不变的次数超阈值判定卡住、缺失检测连续 N 次心跳缺失判定死亡。检查点机制保证崩溃后断点续传每完成一步保存结果恢复时从最新检查点继续不从头重跑。心跳通道与业务通道隔离避免业务阻塞影响心跳延迟。