☰
Agent-Reach:轻量级智能体触达中间件,解决多Agent互相发现与通信
2026/10/6 19:42:28 网站建设 项目流程

我刚接触智能体开发时,一直没想明白一个问题:每个 Agent(智能体)都跑得好好的,但怎么让它们互相找到、互相喊话?我的环境里有一台本地工作站、一台跑批的云服务器、还有一块树莓派,上面各自挂着不同的 Agent:有做 PDF 文本抽取的、有做图像分类的、有负责定时发通知的。平时调试时我要分别记住 IP、端口、认证方式,用不同客户端去连,一旦某个机器重启或者 IP 变了,整个人就会陷入“翻配置”状态。---

后来我动手做了个周末项目Agent-Reach,把它定位成一个轻量级的智能体触达中间件。它做的事情非常简单:所有 Agent 启动后先向 Agent-Reach 注册,拿到一个唯一的 ID 和通道,之后我只需要向 Agent-Reach 发一条标准消息,它就能把消息精确投递给对应的 Agent 并拉回结果。整个过程对客户端几乎没有侵入,也不需要写一堆胶水代码。这篇文章把项目拆解、核心机制、部署过程和踩坑经验完整记录下来,适合跟我一样在折腾多 Agent 协同、家庭自动化或者想在 NAS 上统一管理智能体的朋友参考。

1. 拆解 Agent-Reach 的设计思路

1.1 分布式 Agent 协同里的“最后一公里”

现在做多 Agent 系统,大家最不缺的是单个模型的推理能力,缺的是“触达”能力。这里触达包含三层意思:一是网络层面能不能连上,二是语义层面能不能对上,三是运维层面能不能管住。很多团队直接用 WebSocket 点对点写死通信,短时间跑着没问题,但一旦 Agent 数量从两个变成十个,你就会发现有一半精力花在解决“谁掉线了”“消息没收到”“回调地址写错”这类破事上。

Agent-Reach 想解决的正是这最后一公里。它把自己藏在所有 Agent 的背后,成为一个虚拟总线:不处理业务逻辑,只负责注册、路由、投递和回执。本质上和家里的集线器类似——你不用记住每个插座在哪个房间,只要知道总开关在哪,按个按钮就能把电送到该去的地方。

1.2 为什么不是 MQTT,也不是纯 HTTP?

做这个项目之前,我认真比较过两种常见方案。

MQTT 本身非常适合低功耗设备,带有 Broker 和 QoS 机制,但它的 topic 模型比较线性,表达“给某个能力为ocr的 Agent 发任务”这类请求不够自然,而且需要额外部署一个 MQTT Broker;纯 HTTP 最直观,可以 RESTful 一把梭,但要求 Agent 必须暴露一个公网可访问的端口,在家庭网络、公司 NAT 后面基本走不通,还得自己做轮询和超时管理。

Agent-Reach 采用“注册中心 + 双向长连接”的混合架构:

  • 所有 Agent 作为客户端主动连接到 Hub,不要求 Agent 开放入站端口;
  • Hub 维护一份动态注册表,记录 Agent 的 ID、能力标签、状态;
  • 调用方(不管是人还是另一个 Agent)只和 Hub 通信,由 Hub 负责路由。
特性MQTT纯 HTTPAgent-Reach
设备入站端口不需要需要不需要
语义路由弱需自行实现原生支持
结果回执需自己叠加需自己叠加原生支持
部署复杂度需要 Broker低中
离线缓存支持 QoS无支持暂存

表格能看出取舍:Agent-Reach 不是技术上最炫的,但在“自托管多智能体触达”这个具体问题上是平衡得比较好的。这套设计思路是反复迭代后的选择,初版试过纯 HTTP 轮询,后来发现连接数量和实时性撑不住,才改成现在的长连接。

1.3 设计目标:十分钟接入一个 Agent

我给自己定的硬性指标有三个。

第一,客户端要轻。一个 Agent 只需要装一个很小的 SDK(Python 版本核心代码不到 500 行),不依赖重型框架。用 Go、Node 甚至 shell script 也能通过 HTTP API 接入,保证语言无关。

第二,要能离线缓存。Agent 断网重连后,错过的任务如果直接丢掉,那么很多自动化流程就会静默失败。所以 Hub 端需要一个简单的持久化队列,存住一段时间内无人认领的任务。

第三,状态可观测。我不能接受“消息发出去了但不知道 Agent 到底收到没有”的状态。所以每条任务必须有生命周期:已投递、已接收、已执行、已回执、已超时。

这些目标决定了后续每一个细节的实现方式。代码写起来很克制,没有刻意引入微服务,就是一个 Python 异步进程,加 SQLite 存储,一个 Docker 镜像搞定。

2. 核心机制与关键模块实现

2.1 Agent 注册中心:门牌号与电话簿

注册中心是 Agent-Reach 的心脏。每个 Agent 上线时必须上报三个东西:

  • 身份信息:全局唯一 ID,格式为agent-{namespace}-{name}-{random},例如agent-main-ocr-3f9a;
  • 元数据:能力标签、版本号、健康检查路径、回调 base URL;
  • 运行时状态:当前是否忙碌、最大并发数、最近心跳时间。

Hub 把这些信息写入一个内存 dict 的同时异步刷进 SQLite。读取路径全部走内存,保证毫秒级查询;写入路径异步落盘,避免 I/O 阻塞关键链路。

心跳机制是注册中心最容易被忽略又最关键的细节。每个 Agent 默认每 30 秒向 Hub 发送一个PING报文,Hub 返回PONG,并更新“最近心跳时间”。如果超过 90 秒(即 3 个心跳周期)没有收到任何报文,Hub 会把这个 Agent 标为quasi-offline,意思是“可能还在但我不确定”,不立即摘除;再过 60 秒仍然没心跳,就彻底标记为offline。

为什么用两段式状态?因为我踩过坑:树莓派的 Wi-Fi 偶发抖动,一个 2 秒的瞬断就会让 Agent 被误下线,然后一堆任务发给离线节点,全超时。两段式状态能避免在边缘场景下做出错误决策。

心跳报文本身是空的小 JSON,但我会在报文里附上 agent 当前的内存占用和队列深度,这样 Hub 可以做个简单的负载收集,为后续“按能力选择最优 Agent”提供依据。

2.2 消息路由:按能力调用的三种模式

注册表搞定后,路由就很简单了。Agent-Reach 支持三种投递模式:

点对点(Direct):调用方明确指定目标agent_id,Hub 直接检查该 Agent 是否在线,在线就投递。这个模式用在明确知道要调用哪个实例的场景。

按能力调度(Capability):调用方不关心具体哪个 Agent 执行,只需要声明“我要能 OCR 的”。Hub 会在在线 Agent 里过滤出带ocr标签的列表,再按负载和最近响应时间排序,选最优的一个投递。这是我最常用的模式,因为可以把 Agent 扩容、缩容完全透明化。

扇出广播(Fanout):调用方指定标签,Hub 把消息并行投递给所有带该标签的 Agent。这个用在“通知所有设备”“刷新所有缓存”这类场景。

消息体采用类 JSON-RPC 格式:

{ "id": "task-01HQZ6...", "method": "agent.reach.task", "params": { "request_id": "req_123456", "target": { "type": "capability", "value": "ocr" }, "payload": { "input": "documents/contract.pdf" }, "timeout_ms": 30000 } }

每个请求都带一个全局唯一的request_id,这是后面做去重和追踪的基础。Hub 收到消息后先落库,再路由,确保崩溃重启后还能知道“这条消息到底发出去没有”。

2.3 安全与身份:不让陌生 Agent 混进来

让 Agent 接入一个中间件,最怕的是任何人伪造身份往里灌消息。Agent-Reach 的安全设计分三层:

  • 接入层:Agent 和 Hub 之间的 WebSocket 连接强制走 TLS,Hub 侧配置自签名证书或 Let's Encrypt 证书。
  • 认证层:每个 Agent 在创建时分配一个 token,启动后用这个 token 换成短期会话票据;后续所有消息头都带这个票据,Hub 端校验签名。
  • 授权层:注册元数据里有一个acl字段,声明“这个 Agent 可以接受哪些 namespace 的任务”。比如ocrAgent 只能接收来自mainnamespace 的任务,避免被其他乱七八糟的任务源打爆。

实现上,认证用 HMAC-SHA256 对当前时间戳加 agent_id 签名,Hub 端用同一把密钥验签,防止重放。这套方案比 OAuth2.0 轻很多,对自托管工具已经够用。如果以后要暴露到外网,可以再加一层反向代理的白名单,锁来源 IP。

3. 实操过程:从零搭起 Agent-Reach

3.1 启动 Hub(注册中心)

我用 Docker Compose 把 Hub 跑起来,配置文件写得很简练。

version: "3.8" services: agent-reach-hub: image: agentreach/hub:0.9.0 container_name: agent-reach-hub restart: unless-stopped ports: - "8080:8080" # HTTP API - "8081:8081" # WebSocket 接入端口 environment: REACH_DATA_DIR: /data REACH_AUTH_SECRET: "change-me-to-a-long-random-string" REACH_DEFAULT_LEASE_TTL: "90s" REACH_TASK_TTL: "3600s" volumes: - ./data:/data healthcheck: test: ["CMD", "curl", "-f", "http://localhost:8080/healthz"] interval: 10s timeout: 3s retries: 3

启动后先检查健康接口:

curl http://localhost:8080/healthz

如果返回{"status":"ok"},说明 Hub 起来了。接着要有第一个账号才能创建 Agent。

Agent-Reach 提供一个简易的引导 API,首次启动可以通过环境变量里配置的REACH_BOOTSTRAP_TOKEN调用:

curl -X POST http://localhost:8080/v1/agents \ -H "Authorization: Bearer some-bootstrap-token" \ -H "Content-Type: application/json" \ -d '{ "name": "ocr-agent-01", "namespace": "main", "tags": ["ocr", "image"], "callback_base_url": "http://192.168.1.20:9001", "lease_ttl": 90 }'

返回值里带agent_id和agent_token,第一次显示后不再重复展示,所以建议当场存到密码管理器。

3.2 用 Python SDK 接入一个模拟 Agent

为了让读者更直观地看到接入过程,我用一个模拟 OCR 的 Python 脚本做演示。核心代码就几块:先初始化客户端,然后注册自己,之后进入消息监听循环。

import time from agent_reach_sdk import AgentClient def handle_task(message: dict) -> dict: # 模拟 OCR 过程,实际在这里调用你的推理库 payload = message["payload"] time.sleep(1) return { "status": "ok", "output": f"fake_ocr_result_for_{payload['input']}", "latency_ms": 1020, } def main(): client = AgentClient( hub_url="wss://reach.example.com:8081", agent_id="agent-main-ocr-3f9a", token="your-agent-token", tags=["ocr", "image"], healthcheck_interval=30, ) client.register() client.subscribe(handler=handle_task) client.start_heartbeat() print("Agent is online and waiting for tasks...") while True: time.sleep(1)

SDK 内部做了这些事:

  • 与 Hub 建立 WebSocket 连接;
  • 发送注册信息,等待REGISTER_ACK;
  • 启动后台心跳线程,定时发送PING;
  • 监听服务端下发的EXEC消息,调用上层业务函数;
  • 把业务函数结果封装成RESULT消息回传给 Hub。

实际跑起来后你会看到类似下面的日志:

[INFO] Connected to hub at wss://reach.example.com:8081 [INFO] Register ack received, agent_id=agent-main-ocr-3f9a [INFO] Health reporter started (interval=30s) [INFO] Waiting for tasks... [INFO] Received task req_123456, dispatching to handler [INFO] Task req_123456 completed in 1020ms, result sent.

3.3 通过 HTTP API 触达并返回结果

Agent 在线后,从任意一台能访问 Hub 的机器发任务。比如调用带ocr能力的 Agent:

curl -X POST http://localhost:8080/v1/tasks \ -H "Authorization: Bearer caller-token" \ -H "Content-Type: application/json" \ -d '{ "request_id": "req_123456", "target": {"type": "capability", "value": "ocr"}, "payload": {"input": "documents/contract.pdf"}, "timeout_ms": 30000 }'

Agent-Reach 是异步模式:这个接口会立刻返回一条任务状态记录,不会同步等待 Agent 执行完。如果你需要同步结果,有两种选择:

  • 轮询任务状态接口:GET /v1/tasks/req_123456,直到状态变成succeeded;
  • 用 WebSocket 订阅终端:Hub 会主动把结果推给调用方。

我一般测试时用轮询,生产里用订阅。轮询代码很短:

while True: resp = requests.get(f"{hub}/v1/tasks/{request_id}", headers=headers) data = resp.json() if data["status"] in ("succeeded", "failed", "timeout"): print(data) break time.sleep(0.5)

最终返回结果里有 Agent 回传的output字段:

{ "status": "succeeded", "result": { "status": "ok", "output": "fake_ocr_result_for_documents/contract.pdf", "latency_ms": 1020 }, "executed_by": "agent-main-ocr-3f9a" }

整个过程看下来,调用方根本不需要知道 Agent 的 IP、端口,甚至不需要知道它运行在哪个系统上。唯一打交道的就是 Hub。

3.4 给现有 Agent 加个健康检查 API

为了让 Hub 能主动感知业务存活状态,我在每个 Agent 内置了一个/healthzHTTP 接口。SDK 启动后会自动监听 0.0.0.0 的一个随机端口,并把这个端口上报给 Hub。Hub 在心跳报文之外,还可以每隔 60 秒主动探测一次这个接口,返回值里带{"load": 0.8, "queue_depth": 3}这样的负载数据。

这个设计让我可以在注册中心里看到每个 Agent 的真实负载,做能力路由时不只是随机挑,而是挑负载最低的。代码如下:

from agent_reach_sdk import HealthServer health = HealthServer(port=0) # port=0 表示随机端口 health.set_handler(lambda: { "load": current_load(), "queue_depth": queue.qsize() }) health.start()

不过要提醒一点:健康检查的 HTTP 端口不能和 Hub 的消息通道端口冲突。如果 Agent 运行在容器里,记得把容器端口映射出来或者让 Hub 通过 Docker 网络直连。

4. 实操中的坑与排查实录

4.1 心跳偶发丢失导致 Agent 被误下线

这个坑我在设计时已经提过解决思路,但实操中仍然会反复踩到。现象是:Agent 明明在正常运行,Hub 却隔三岔五地把状态改成quasi-offline,然后路由时跳过它,导致部分任务堆积。

排查方法:

  • 查看 Hub 日志里有没有heartbeat missed报警;
  • 用tcpdump抓包确认 Agent 是否真的把PING报文发出去了;
  • 检查 Agent 与 Hub 之间是否有 NAT 超时,比如家用路由器默认的 UDP 会话超时可能在 30 秒左右,恰好和心跳间隔接近。

我这里最终把心跳间隔从 30 秒调到 20 秒,同时把“误判阈值”从 3 次加大到 5 次。代价是误判收敛时间变长,但对稳定性要求高的场景,宁可延迟判定也不冤枉好人。

4.2 重复投递导致业务被执行两次

WebSocket 连接断开重连后,如果之前有一条消息已经发给 Agent,但 Agent 还没来得及回执,Hub 会认为投递失败,于是重试一次。这样一来,Agent 端就会把同一个request_id执行两遍。对于只读任务无所谓,但对于“发邮件”“转账”这类操作就是事故。

解决办法有两个层面:

  • Hub 端:重试前检查 Agent 是否已恢复连接,并且检查是否有未完成的回执;
  • Agent 端:SDK 内置一个按request_id去重的 set,缓存最近 10000 条已执行任务的 ID,重复消息直接返回缓存结果,不执行业务逻辑。

去重代码就几行:

_seen = set() def deduplicate(request_id): if request_id in _seen: return False _seen.add(request_id) return True

我建议无论如何都要在业务入口做一次幂等处理,再信任 Hub 的重试机制。分布式环境下“at least once”是常态,设计业务时要默认消息可能重复。

4.3 回调地址写错导致结果石沉大海

Agent 注册时有一个callback_base_url字段,SDK 在返回结果时会把结果 POST 到这个地址。我一度以为这里填的是 Agent 自己的地址,结果填成了http://localhost:9001/result。Agent 跑在树莓派上没问题,但 Hub 跑在另一台机器上,它去请求localhost只会打到 Hub 自己,结果自然找不到。

正确理解是:callback_base_url是 Hub 回调 Agent 时使用的地址,因此必须填写 Hub 能访问到的 Agent 地址。如果 Hub 和 Agent 在同一个 Docker 网络中,要填http://agent-container-name:9001,不能填 localhost。如果 Agent 在另一个局域网,需要配置反向代理或路由映射。

我后来把回调逻辑改成了走 WebSocket 通道回传,优先用通道不用 HTTP,避免这整类地址配置问题。HTTP 回调只作为降级方案。

4.4 常见问题速查表

症状可能原因处理方式
Hub 显示 Agent 在线,但任务一直 pendingAgent 的 handler 处理太慢或阻塞查看 Agent 进程 CPU 与队列深度,必要时调大timeout_ms
任务 succeeded 但 result 为空Agent 的 handler 没有返回 dict检查 handler 是否return None,SDK 会当空结果回传
多个 Agent 同时在线,但总是同一个执行负载评估逻辑过于简单在注册元数据里配置 priority 权重
重连后收不到新任务会话票据过期检查 token 有效期,SDK 应在重连时重新换取票据
时区导致心跳时间计算错误Hub 与 Agent 时区不一致全部统一使用 UTC,不读取本地时间
SQLite 锁导致任务写入超时并发写入冲突关闭 WAL 模式,或改成 Postgres 存储

这张表是我在实际跑了两周后整理出来的,每一个问题都真实遇到过。尤其是时区的那一个,看似小事,但真的会让心跳误判,排查了很久。

5. 经验补充与后续扩展

5.1 从“触达”到“编排”:让 Agent 互相调度

Agent-Reach 本身只做触达,但触达一旦稳定,多 Agent 协作就变成简单的代码逻辑。比如我做一个“合同审阅”流程:

  1. 用户上传 PDF 到主服务;
  2. 主服务通过 Agent-Reach 调用ocrAgent,提取文本;
  3. OCR 结果文本通过 Agent-Reach 发给llm-reviewerAgent,让大模型生成摘要;
  4. 摘要再广播给notify组里的所有 Agent(邮件、钉钉、Slack)。

整个过程因为每条消息都带request_id,所以串联起来非常清晰。我甚至可以查看任务拓扑图:某个摘要到底依赖哪次 OCR,在哪一步花费了多少时间。这比直接把几个 Agent 用 Python 函数链起来好维护得多,至少你不会在某个 Agent 升级后惊慌失措。

实现编排层时不需要改 Agent-Reach,只需要一个很低成本的 worker:订阅任务完成事件,然后构造新任务再发起调用。这个逻辑可以放在 Hub 之外,用任何语言写都行。

5.2 运维视角:监控与告警要趁早做

Agent-Reach 暴露了/metrics端点,基于 Prometheus 格式。我设计了五个关键指标:

  • reach_agent_online_count:在线 Agent 数量;
  • reach_task_created_total:创建任务总数;
  • reach_task_succeeded_total:成功任务总数;
  • reach_task_failed_total:失败任务总数;
  • reach_task_duration_seconds:任务执行耗时直方图。

用 Grafana 建了一个看板,实时看有没有 Agent 掉线,以及任务延迟分布。实践经验是:一旦发现某个 Agent 平均耗时翻倍,基本就是它的业务逻辑出了问题,而不是网络。这时候去翻 Agent 的日志能精准定位。

另外,我设置了最简单的告警规则:任何reach_task_failed_total在 5 分钟内增长超过 10 次,就触发钉钉通知。这让我在故障发生时不会被动。

5.3 后续扩展方向:接入大模型和 Webhook

下一步我还想把 Agent-Reach 接到 MCP 生态里,让大模型直接通过工具调用方式触达任何 Agent。目前实现了一个很薄的 MCP 服务,把大模型发出的工具请求转成Agent-Reach任务,再把结果塞回给大模型。这样模型不再是只能“聊天”,而是真的能指挥所有智能体干活。

另外计划支持 Webhook 触发:GitHub push 事件直接生成一个Agent-Reach广播,让所有监听代码变更的 Agent 自动刷新缓存或跑测试。做成这件事以后,Agent-Reach 就不只是一个人工调度的工具,而是变成一套自动化事件分发底座。

我自己实际跑下来最爽的一点是,以前要在五台机器上来回 ssh 敲命令,现在只需要在一个 Webhook 或一句 curl 里带上target=custom就能精准触达。最后再分享一个小技巧:如果你也给每台机器跑着多个 Agent,记得给 Agent 命名时带上具体位置和用途,比如agent-nas-ocr-01而不是agent-01,否则排查日志的时候你能感受到什么叫“同名混乱”。这个项目的价值不在于代码多复杂,而在于它让“Agent 之间互相找到并正确传话”这件事变得不值得再花精力去考虑。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询