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_id Agent 实例标识 关联会话与实例
session_id 会话标识 恢复会话时使用
status idle/processing/error 区分状态
current_step 当前执行步骤 进度追踪
total_steps 总步骤数 进度百分比计算
cpu_percent CPU 使用率 资源监控
memory_mb 内存占用 资源监控
last_tool_call 最近一次工具调用信息 卡住定位
timestamp 心跳时间戳 判断是否过期

2.3 健康检查的判定逻辑

健康判定不是简单的"心跳在就健康"。需要根据心跳间隔和业务状态综合判断:

  • 健康:心跳间隔 < 预期间隔 × 2,status 不是 error
  • 亚健康:心跳间隔 > 预期间隔 × 2 但 < 预期间隔 × 5,status 是 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) => Promise<void>;

  constructor(
    agentId: string,
    sessionId: string,
    intervalMs: number = 5000,  // 默认 5 秒心跳间隔
    maxMissedHeartbeats: number = 3,
    sender: (payload: HeartbeatPayload) => Promise<void>
  ) {
    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: Map<string, 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 = f"checkpoints/{session_id}/step_{step_index}.json"
        self.storage.put(key, json.dumps(checkpoint))

    def recover_session(self, session_id: str) -> dict:
        """从最新检查点恢复会话"""
        # 查找该会话的所有检查点,取最新的
        pattern = f"checkpoints/{session_id}/step_*.json"
        checkpoints = self.storage.list(pattern)

        if not checkpoints:
            logger.warning(f"No checkpoints found for session {session_id}")
            return {'step_index': 0, 'step_results': {}}

        # 取最新的检查点(step_index 最大的)
        latest = max(checkpoints, key=lambda k: int(k.split('step_')[1].split('.')[0]))
        checkpoint_data = self.storage.get(latest)

        checkpoint = json.loads(checkpoint_data)
        logger.info(
            f"Recovered 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(
            f"Session {session_id} recovered: "
            f"new agent {new_agent.agent_id}, "
            f"resuming 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(f"Cleaning 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 次心跳缺失判定死亡)。检查点机制保证崩溃后断点续传:每完成一步保存结果,恢复时从最新检查点继续,不从头重跑。心跳通道与业务通道隔离,避免业务阻塞影响心跳延迟。

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐