GoAgentX 架构深度解析(二):AHP — 多智能体的通信基石

GoAgentX 架构深度解析(二):AHP — 多智能体的通信基石

书接上回 关于无聊自己搓一个agent 框架这件事
正好最近也有点空,就更新系列文章吧,这次的话题是 AHP,至于为啥叫这个名字,原计划是照抄HTTP协议的,但是实现的时候就不对劲了…… 后来也懒得改了就继续沿用这个AHP。

聊到多 Agent 系统,很多人第一反应是:“Agent 之间怎么说话?用 HTTP 还是 WebSocket?走消息队列?”
我的回答比较粗暴:同一个进程里跑着,聊个天还要走网络?直接用 Go channel 干不就完了。
于是就有了 AHP——一个不走网线的通信协议。

做多 Agent 系统最烦的一件事是什么?不是 Agent 不够聪明——是 Agent 之间不说话。

Leader 派了个活给 Sub,Sub 干完了想汇报,结果发现 Leader 已经超时了。Sub 想说自己干到哪一步了,发现没地方说。Leader 想知道 Sub 还活着吗,发现没有心跳机制。

我刚开始用 Python 搭的时候,用的是 Redis 队列。后来换 Go,看了两天 RabbitMQ 的文档——太TM重了。就想:同一个进程,两个 goroutine,发条消息还要经过网络?脑子有病吧。

所以我写了一个纯进程内的通信协议:不走网络、不序列化、不依赖中间件。就是 channel + 共享内存。

一、为什么自己造轮子?

GoAgentX 里有两种角色:Leader Agent(管分活的)和 Sub Agent(管干活的)。它俩之间通信需要搞定的破事:

  • 异步发消息:Leader 把任务丢出去就能干别的,不用等 Sub 搞完

  • 进度反馈:Sub 干到 50% 了,得让 Leader 知道

  • 心跳检测:Leader 得知道 Sub 是不是挂了

  • 容错处理:消息发失败了得有个兜底

我当时调研了一圈,发现要么太重(RabbitMQ),要么太慢(Redis 走网络),要么和 Go 的哲学不对付。最后决定:自己搓一个。

理由其实特简单:

  1. :channel 传东西 vs 网络 RTT——这差距大到不用比

  2. 简单:不用管序列化、网络抖动、分区容错那些分布式噩梦

  3. 好改:以后真要上微服务了,把底层从 channel 换成 gRPC 就行,上面业务代码一行不用动

AHP 的整体架构可以概括为一张图:

核心组件包括:

  • Protocol:统一门面,组合所有子组件,Facade 模式,提供一站式接口。
  • MessageQueue:每个 Agent 独立的消息队列,基于 buffered channel + backup buffer + atomic.Bool
  • HeartbeatMonitor:心跳检测 + 超时回调,共享实例,分布式系统无需额外组件
  • DLQ:死信消息存储与重试,支持自定义处理器和自动重试
  • QueueRegistry: 管理所有 Agent 的队列,懒加载 + 双检锁

二、消息模型

2.1 五种消息类型

AHP 定义了 5 种消息类型,覆盖 Agent 间通信的全部场景:

const (
    AHPMethodTask      AHPMethod = "TASK"      // 任务分配
    AHPMethodResult    AHPMethod = "RESULT"     // 任务结果
    AHPMethodProgress  AHPMethod = "PROGRESS"   // 进度反馈
    AHPMethodACK       AHPMethod = "ACK"        // 确认回复
    AHPMethodHeartbeat AHPMethod = "HEARTBEAT"  // 心跳信号
)

2.2 消息结构

type AHPMessage struct {
    MessageID   string         `json:"message_id"`
    Method      AHPMethod      `json:"method"`
    AgentID     string         `json:"agent_id"`
    TargetAgent string         `json:"target_agent"`
    TaskID      string         `json:"task_id"`
    SessionID   string         `json:"session_id"`
    Payload     map[string]any `json:"payload"`
    Timestamp   time.Time      `json:"timestamp"`
}

2.3 MessageID 生成

MessageID 的设计是一个三段式 ID:

func generateMessageID() string {
    id := atomic.AddUint64(&messageIDCounter, 1)
    randSuffix := getRandomSuffix()
    return fmt.Sprintf("%s.%d.%s",
        time.Now().Format("20060102150405.000000"), id, randSuffix)
}
  • 时间戳前缀:可读性强,方便排查问题
  • 原子计数器:同一纳秒内多个消息的序号递增
  • 随机后缀:多进程场景下避免冲突

这个方案不依赖全局协调器,在进程内就能保证唯一性。

不太理解不要担心,我只不过是借鉴了部分区块链中的设计,因为我之前参与过某条公链的设计,消息也就进行了参考,没别的意思,主要起到状态同步的作用。

三、心跳检测:HeartbeatMonitor

3.1 核心流程

  • 各 Agent 按固定间隔(默认 5s)发送心跳
  • HeartbeatMonitor 记录最近一次心跳时间
  • 超过超时时间(默认 30s)且连续错过次数达到阈值(默认 3 次),标记为离线

3.2 超时检测算法

func (m *HeartbeatMonitor) CheckTimeouts() []string {
    timedOut := m.checkAndMarkOffline()  // 写锁下检测
    for _, agentID := range timedOut {
        m.notifyCallbacks(agentID)        // 锁外执行回调
    }
    return timedOut
}

关键边界条件处理:

  1. 渐进式超时:错过 3 次心跳才判定离线,避免网络偶发延迟导致误杀
  2. 避免重复回调:已 Offline 的 Agent 不会再次触发回调
  3. 回调在锁外执行notifyCallbacks 复制回调列表后释放锁再执行,这是防止死锁的关键

3.3 两种 HeartbeatSender

  1. ahp.HeartbeatSender:发送 AHPMethodHeartbeat 消息到目标的 MessageQueue,属于带内心跳
  2. heartbeatSender(在 internal/agents/sub/):直接调用 HeartbeatMonitor.RecordHeartbeat,属于带外心跳

目前 Sub Agent 使用第二种方式,在单体部署下更高效。

四、死信队列:DLQ

Enqueue 返回错误时,Protocol.SendMessage 将失败消息路由到 DLQ:

func classifyEnqueueError(err error) string {
    switch {
    case errors.Is(err, apperrors.ErrQueueClosed):  return "queue_closed"
    case errors.Is(err, apperrors.ErrQueueFull):    return "queue_full"
    case errors.Is(err, context.Canceled):          return "context_canceled"
    case errors.Is(err, context.DeadlineExceeded):  return "context_deadline"
    default:                                        return "unknown"
    }
}

DLQProcessor 支持按错误类型注册自定义处理器,并支持自动重试:

  • MaxRetries = 0:无限重试
  • MaxRetries > 0:达到次数后标记为 exhausted
  • 当前无指数退避,这是可改进点

五、Agent 中的 AHP 集成

5.1 Messenger 接口

type Messenger interface {
    SendMessage(ctx context.Context, msg *ahp.AHPMessage) error
    ReceiveMessage(ctx context.Context) (*ahp.AHPMessage, error)
}

Leader Agent 和 Sub Agent 都实现了此接口。构造时通过依赖注入传入 MessageQueueHeartbeatMonitor

5.2 Dispatcher 的任务分发

taskDispatcher 同时支持本地执行分布式分发两种模式,核心逻辑在 executeTask 中:

if executor, ok := d.executorFuncs[task.Type]; ok {
    return executor(ctx, task, agentAddr, sessionID)  // 本地执行
}
if d.messageSender == nil { /* return error */ }
msg := ahp.NewTaskMessage(...)                        // 通过 AHP 发送
d.messageSender.Send(ctx, agentAddr, msg)
return d.waitForResult(ctx, task.TaskID)              // 阻塞等待结果

这种设计使得 Agent 通信模式可以在单体部署和分布式部署间无缝切换。

六、设计模式总结

模式 位置 说明
Facade Protocol 统一接口,组合所有组件
Registry QueueRegistry, CodecRegistry 具名实例管理,懒加载
Strategy Codec 接口 可替换的序列化策略
Observer TimeoutCallback 心跳超时回调
Dead Letter Queue DLQ + DLQProcessor 失败消息存储与重试
Double-Checked Locking GetOrCreate 兼顾性能与正确性
Panic Recovery Enqueue defer recover() 应对并发关闭
Lock-Free Read atomic.Bool 无锁读取关闭状态

七、关键设计决策

7.1 为什么非阻塞 Enqueue?

  • Agent 是多线程环境,阻塞可能导致级联等待
  • DLQ 提供了更好的容错语义,失败消息可重试
  • 调用方有更大的控制权:立即重试、稍后重试、或舍弃

7.2 TOCTOU 避免

SendMessage 有一个关键设计:不先检查 IsFull 再入队。如果先检查再入队,检查到入队之间队列可能从不满变为满(TOCTOU 竞态),导致消息丢失。直接执行操作并处理错误更健壮。

7.3 序列化预留扩展

当前 AHP 是纯进程内通信,JSON 够用。但 Codec 接口预留了两个方向的扩展:

  1. 跨进程通信:protobuf/msgpack 可提供更小的载荷
  2. 持久化:DLQ 消息若落盘,二进制格式更有优势

八、喜闻乐见的扒老底环节

说实话,AHP 不是完美的。我自己用的时候也踩过一些坑:

  1. 纯进程内:跨不了进程,真要上分布式得换 MessageQueue 实现
  2. 没有广播:要给多个 Sub 发消息?老老实实逐个发
  3. 重试策略太憨:DLQ 重试间隔是固定的,没做指数退避——持续失败的时候可能造成重试风暴
  4. 路由太死板:不支持基于内容的动态路由或者 Topic 订阅

但话说回来,这些 limitation 都是按需取舍——在单体阶段,channel 方案省掉了 90% 的分布式复杂度。将来真要上微服务,替换底层实现就行,上面的业务代码不用动。这就是抽象层的好处。单机上用的话,确实蛮丝滑的,当然如果要拿着这个项目去面试,建议你最好准备分布式情况下该如何处理,写这个的目的,就是起到抛砖引玉的作用,AI时代,除了编程基本功之外,个人建议要打磨好自己的软件工程能力,毕竟当LLM API caller 可不是咱的目标。

总结

AHP 是我给 GoAgentX 搓的通信轮子。channel 传消息、DLQ 兜底、HeartbeatMonitor 看死活——三个东西加一起,多 Agent 通信的基础设施就齐活了。

代码里留的那些接口(Codec、DLQ handler、MessageSender),说白了就是给自己留的后路:以后要上 gRPC 还是 RabbitMQ,换一层实现就完事,上面业务代码不用动。这种设计在创业项目里特别重要——你永远不知道明天架构要改成啥样

之后咱聊聊记忆蒸馏——Agent 怎么从几百条对话历史里把有用的经验提炼出来,这也是我觉得agent 面试高频会问的长上下文场景下如何确保模型的信息不失真的关键技术。

3 个赞

一核有难八核围观是吧,逮着CPU0往死里蹬啊 :joy: