
1. 为什么我需要一个叫 Agent-Reach 的东西先聊个真实场景。我手上现在维护着十多个自动化 Agent有的是爬数据的有的是做内容审核的有的是定时触发报表的。它们各管一摊平时倒也相安无事。直到有一天我需要让其中一个 Agent 在跑完任务之后去通知另外三个 Agent 接力干活还要把结果推给一个外部服务的 Webhook。问题一下就来了这些 Agent 彼此之间怎么发现对方消息走什么协议状态怎么同步失败重试怎么办以前的做法很笨——在每个 Agent 的代码里写死对其他 Agent 的 HTTP 地址用轮询去拉状态失败了就 throw 出去让上层人工处理。跑起来能用但维护成本高得吓人。每加一个新 Agent就得改一堆配置文件还得祈祷端口没冲突、地址没写错。我一度怀疑自己写的不是业务代码而是通讯录维护工具。后来我决定做一套专门解决“Agent 之间如何可靠地触达彼此”的基础设施代号就叫Agent-Reach。它不复杂也不是什么分布式中间件大作它只解决一个问题让 Agent 像打电话一样找到对方、传递消息、确认结果。这篇文章我会把整个思路、核心设计、部署配置、踩坑过程都写出来希望能给同样被“Agent 互联”折磨过的人一个参考。先说清楚Agent-Reach 不是什么新协议它只是把现有成熟组件和模式组合起来做成一个轻量级的服务网格。核心由三部分组成注册发现Registry Health Check、消息通道异步队列、结果回执Callback 状态机。听起来平平无奇但把这三点做成一套顺手的东西实际节省的时间远超预期。2. Agent-Reach 的定位与核心设计思路2.1 它不是消息队列也不是 RPC 框架做技术选型之前我先把 Agent-Reach 的边界划清楚。市面上的消息中间件比如 Kafka、RabbitMQ、RocketMQ能力很强能处理海量吞吐、复杂路由、持久化这些事但对“Agent 互连”这个场景来说太重了。RPC 框架gRPC、Dubbo讲究的是接口定义和调用链治理也不太适合“一个 Agent 自主决定要不要响应”这种松散耦合的互动模式。Agent 之间协作的特点是什么是自主性和异步性。Agent A 发出一条请求Agent B 不一定立刻处理它可能忙别的、可能离线、可能要先跑一个长时间任务。这种情况下用同步 RPC 会非常痛苦一个 Agent 超时能把整条链路堵死。所以 Agent-Reach 的定位是一个轻量的、面向 Agent 协作的触达协议与运行时。它不像 Kafka 那样追求吞吐也不像 gRPC 那样追求强一致它追求的是——在 Agent 数量不多几十个这个量级、但协作关系复杂多对多、动态变更的场景下用最少的机器成本和维护成本实现可靠的“触达 回执 重试”。核心消息模型就两条Command一个 Agent 请求另一个 Agent 做某件事比如request_translate、request_summary。它表示“我希望你做什么”。Event一个 Agent 通知其他 Agent 某件事已经发生了比如task_finished、data_updated。它表示“有件事你们要知道”。Command 需要回执Event 不需要。这个区分是 Agent-Reach 整个设计的基石后面所有功能都是为了服务这两类消息的差异化处理。2.2 注册中心Agent 也要有通讯录既然要“触达”第一件事就是解决“怎么找到对方”。我给每个 Agent 在启动时自动注册到 Agent-Reach 的 Registry 组件注册内容包含Agent 名称全局唯一比如agent-translator能力标签capabilities比如translate:zh-en基础地址Base URL比如http://10.1.2.3:8080负载权重和健康状态这个设计参考了 Consul 和 Eureka 的思路但做得更薄——不需要复杂的数据中心、ACL、多租户只需要一个内存 Map TTL 过期就够用了。Agent 每隔 15 秒上报一次心跳超过 45 秒没有心跳就标记为 unhealthy超过 90 秒就自动摘除。为什么这些数字选这个值因为 15 秒心跳相比 5 秒的激进策略能明显减少无意义的网络包45 秒失联容忍则能在 Agent 短暂卡顿比如 GC 暂停时不至于误杀。这些都是实际调试中调出来的不是拍脑袋定的。还有一个细节Agent 注册时不是直接注册到某个中心节点而是注册到一个虚拟分组namespace。比如default组放通用 Agentcrawler组放爬虫相关 Agent。跨组通信默认是禁止的需要显式申请。这样做的原因是防止 Agent 数量多了之后消息在一个池子里互相污染——翻译 Agent 不应该收到爬虫 Agent 的爬取指令。2.3 触达协议轻量级 JSON over HTTP 可选异步通道Registry 负责“找到人”接下来是“把话递到”。Agent-Reach 的默认触达协议是 JSON over HTTP命令以POST /reach/command发送响应统一返回一个requestId。这个选择很简单HTTP 生态最成熟任何语言都能轻松实现调试工具也多肉眼就能看协议内容。但纯同步 HTTP 在“Agent B 处理时间很长”的场景下撑不住。所以我加了异步模式如果 Agent B 注册时声明mode: asyncAgent-Reach 会把请求投递到内置的轻量队列中Agent B 自己来拉取pull任务完成后通过 Callback 接口回执。这样避免了 HTTP 长连接、超时、线程阻塞等一堆问题。实现上Agent-Reach 的队列组件没有用外部的消息中间件而是自己实现了一个基于SQLite HTTP 轮询的简易消息存储。为什么不用 Redis Stream 或者 RabbitMQ因为我希望 Agent-Reach 是一个非常容易部署的单元不引入额外依赖SQLite 文件即用即走特别适合中小规模的自建部署。数据量上来之后理论上会有性能瓶颈但前面说了这个量级几十个 Agent、每秒几百条消息SQLite 完全可以扛得住。我实测过在普通 SATA SSD 上每秒两千条消息的写入和消费毫无压力。异步模式下消息生命周期如下Agent A 构造指令并 POST 给 Agent-ReachAgent-Reach 校验 Agent A 的 token将消息写入队列返回requestIdAgent B 以长轮询long-poll方式从队列领取消息Agent B 处理完毕POST 结果给 Agent-Reach 的回执接口Agent-Reach 将结果通知 Agent A如果 Agent A 关注回调这个模式在代码里落地其实就两百来行核心逻辑但边界条件多后面我会专门讲踩坑过程。3. 部署与接入实战3.1 单机快速部署Docker 一把梭Agent-Reach 的部署设计目标是“一条命令拉起”。我在根目录放了docker-compose.yml核心服务就一个镜像version: 3.8 services: agent-reach: image: reachhub/agent-reach:0.4.2 container_name: agent-reach-core restart: unless-stopped ports: - 8086:8086 - 9090:9090 environment: - REACH_REGISTRY_TTL45 - REACH_REGISTRY_HEARTBEAT15 - REACH_QUEUE_ENGINEsqlite - REACH_QUEUE_DB_PATH/data/agent-reach.db - REACH_TOKEN_ENABLEDtrue - REACH_TOKEN_STORE/data/tokens.json volumes: - ./data:/data几个环境变量值得解释一下REACH_QUEUE_ENGINEsqlite队列引擎选择 SQLite默认即此值当前没有别的可选项。后续想扩展 Redis 引擎也方便接口都抽象好了。REACH_TOKEN_ENABLEDtrue开启 Agent token 鉴权。这个强烈建议开启。刚开始我为了省事没开结果局域网里任何主机都能向 Agent-Reach 提交指令某次被一个测试脚本灌了一堆垃圾任务排查了半天才发现是没鉴权导致的。REACH_QUEUE_DB_PATHSQLite 文件路径放在 volume 中持久化重启不丢任务。启动之后用 curl 验证一下服务状态curl -s http://localhost:8086/healthz # 期望输出: {status:ok,time:2025-01-15T10:00:00Z}3.2 一个 Python Agent 的接入示例接入 Agent-Reach 不需要 SDK只要在项目里加一个线程做三件事心跳注册、长轮询拉取、回执上报。下面是我常用的一套最小实现可以直接抄import json import requests import threading import time import logging REACH_BASE http://localhost:8086 AGENT_NAME agent-translator AGENT_TOKEN tok_xxx AGENT_CAPABILITIES [translate:zh-en, translate:en-zh] def register_and_heartbeat(): 启动时注册之后每 15 秒心跳一次。 payload { name: AGENT_NAME, capabilities: AGENT_CAPABILITIES, base_url: http://10.0.0.5:8081, mode: async, metadata: {owner: ops-team, lang: python3.11} } headers {Authorization: fBearer {AGENT_TOKEN}} # 首次注册若已存在则更新幂等 r requests.post(f{REACH_BASE}/registry/register, jsonpayload, headersheaders, timeout5) r.raise_for_status() while True: time.sleep(15) try: hb requests.post( f{REACH_BASE}/registry/heartbeat, json{name: AGENT_NAME}, headersheaders, timeout5 ) # 如果心跳返回 404说明注册信息已过期被清理重新注册 if hb.status_code 404: requests.post(f{REACH_BASE}/registry/register, jsonpayload, headersheaders, timeout5) except requests.RequestException as e: logging.warning(fHeartbeat failed: {e}) def long_poll_loop(): 长轮询领取任务处理完成后上报结果。 headers {Authorization: fBearer {AGENT_TOKEN}} while True: try: poll requests.post( f{REACH_BASE}/queue/poll, json{name: AGENT_NAME, timeout: 20}, headersheaders, timeout30 ) if poll.status_code ! 200: continue task poll.json() if not task: continue logging.info(fReceived task: {task[requestId]} {task[action]}) result handle_task(task) # 上报回执成功失败都要上报 ack { requestId: task[requestId], agentName: AGENT_NAME, status: success if result[ok] else failed, output: result[data], durationMs: result.get(duration, 0) } requests.post(f{REACH_BASE}/callback/complete, jsonack, headersheaders, timeout5) except Exception as e: logging.error(fPoll loop error: {e}) time.sleep(2) def handle_task(task): 实际业务逻辑此处仅示意。 action task.get(action) payload task.get(payload, {}) if action translate: text payload.get(text, ) # 调用真实的翻译模型/服务 return {ok: True, data: {translated_text: f[{text}] translated}, duration: 120} return {ok: False, data: {error: unsupported action}}这个线程跑起来之后这个 Agent 就具备了被 Agent-Reach 触达的能力。接入一个 Agent 的实际耗时大约 30 分钟主要是适配handle_task里的业务逻辑和配置 token。3.3 触达调用方的写法Fire-and-Forget 与 Wait-Result作为调用方Orchestrator发指令也极其简单。Fire-and-forget 场景import requests import uuid REACH_BASE http://localhost:8086 ORCH_TOKEN tok_orch def send_command(target_agent, action, payload, timeout_ms5000): cmd { target: target_agent, action: action, payload: payload, requestId: str(uuid.uuid4()), timeoutMs: timeout_ms } headers {Authorization: fBearer {ORCH_TOKEN}} r requests.post(f{REACH_BASE}/command/send, jsoncmd, headersheaders, timeout10) return r.json() # {requestId: ..., status: accepted} # 发送后不等待结果Agent 完成后自行回调 send_command(agent-translator, translate, {text: hello})如果需要等待结果需要配合 Agent-Reach 的WaitResult机制——发指令时带上waitResult: trueAgent-Reach 会阻塞等待回执内部实现是 CompletableFuture 挂起 超时控制直到 Agent 上报完成或超时返回。这个机制在编排链路中非常有用比如先翻译再摘要再保存每一步都需要上一步的产物。4. 核心功能深入指令路由、回执状态机与重试策略4.1 路由决策按能力标签匹配还是按 Agent 名精确匹配Agent-Reach 的路由支持两种方式按优先级从高到低精确目标target直接指定 Agent 名称比如target: agent-translator。这种模式下Agent-Reach 直接把消息发给指定的 Agent失败则报错。能力标签匹配target不指定名称而是指定capability: translate:zh-en。Agent-Reach 会到 Registry 里查找所有拥有该能力标签且为 healthy 的 Agent然后按负载权重做轮询或随机选择。第二种方式对于“某种活谁都可能干”的场景非常合适。比如我有三个爬虫 Agent都可以爬微博热搜那么提交任务时只声明capability: crawler:weibo三个 Agent 轮着干活哪个忙就少发点。路由逻辑里有几个参数实际调优过参数默认值说明健康阈值healthy只有 healthy 状态才参与路由负载策略round-robin按注册顺序轮询也可切 weighted-random死信次数3消息投递失败超过 3 次后进入 dead-letter 队列回执超时5 分钟超过视为失败触发重试或告警4.2 回执状态机的关键状态在设计 Agent-Reach 时我最小心的是避免把消息状态模型搞得太复杂。最终只保留了五个核心状态足够覆盖绝大多数场景PENDING消息已投递到队列Agent 尚未领取。CLAIMEDAgent 已领取正在处理中。COMPLETEDAgent 处理成功回执已确认。FAILEDAgent 处理失败回执为失败状态。DEAD多次投递/处理失败进入死信队列人工介入。状态迁移图很简单PENDING - CLAIMED - COMPLETED/FAILED若 CLAIMED 超时Agent 领取后未在规定时间内回执状态回到 PENDING 并重新投递。重新投递次数超过REACH_MAX_DELIVERIES默认 3则状态置为 DEAD。这个状态机的实现让我避免了一个常见坑漏掉“已领取但处理途中进程崩溃”的情况。进程崩溃后消息既没有 COMPLETED 也没有 FAILED如果不做超时回收这条任务就永远卡在 CLAIMED 状态。Agent-Reach 在实现时每次 CLAIMED 都会记录一个claimedAt时间戳调度线程每隔 30 秒扫描一次发现claimedAt超过REACH_CLAIM_TIMEOUT默认 3 分钟的消息强制重置为 PENDING 重新投递。这个机制我用“延迟回收 等幂令牌”来保证不重复执行——重新投递时带同一个requestIdAgent 侧执行前先检查该 requestId 是否已经处理过。4.3 重试策略不是每次都无脑重试最开始我想简单点失败就重试 3 次。但跑了一段时间发现有些任务失败是确定性失败比如数据格式非法、权限不足重试多少次都会失败反而浪费资源、产生大量死信。后来我设计了三档重试策略在发送指令时可选NO_RETRY不重试失败即返回错误。适用于格式校验类任务失败没有意义。RETRY_ON_TIMEOUT仅在超时包括网络超时和 Agent 处理超时情况下重试。适用于外部服务依赖类任务可能只是对方暂时不可用。RETRY_ALWAYS任何失败都重试直到成功或达到最大次数。适用于数据可重放、对一致性要求高的任务。默认值是 RETRY_ON_TIMEOUT。这个策略调整之后整个系统的死信数量下降了 70% 左右Agent 之间的噪音也小了很多。5. 我在真实部署中踩过的一串坑5.1 心跳误删导致的 Agent 幽灵注册这是最恶心的问题。现象是Agent A 已经停了但 Registry 里它的注册信息迟迟不消失甚至 Agent A 重启之后新的注册信息反而被拒绝提示“agent already registered”。排查链路是这样的通过curl http://localhost:8086/registry/agents查看注册信息发现 Agent A 有两个不同 URL 的记录。再查心跳日志发现 Agent 新实例注册时Agent-Reach 按“agent name base_url”作为唯一索引所以旧实例的 base_url 没变时应该直接覆盖但如果 base_url 变了比如容器重启后 IP 变了Agent-Reach 就认为是新 Agent旧记录又没被删除于是出现了“同名双活”。从根上解决我把唯一索引改成仅按name作为唯一键。注册时若存在同名 Agent强制先删除旧记录再插入新记录同时通知消息通道中所有排队消息进入“目标不存在”状态因为旧实例已失效让调用方可以重新发起。这在业务上不一定合理但至少不会留下幽灵注册。5.2 SQLite 并发写导致的 SQLITE_BUSYAgent-Reach 的队列用 SQLite 存储多线程提交消息时频繁出现SQLITE_BUSY: database is locked。开始我以为是自己代码 bug后来用PRAGMA busy_timeout5000加了个写锁等待时间解决了 90% 的问题。剩余 10% 是某些长时间事务比如死信批量归档把写锁占用太久我额外加了一个连接池最大连接数限制并保证队列消息的插入操作使用PRAGMA journal_modeWAL模式。WAL 模式下读不阻塞写写不阻塞读并发能力大幅提升。如果你要自己实现类似系统记住这个参数比我全部代码都重要PRAGMA journal_modeWAL; PRAGMA synchronousNORMAL; PRAGMA busy_timeout5000;5.3 Agent 端长轮询线程被打满我的 Agent 会开 3 个长轮询线程去拉取消息但某个时刻大量任务涌入时这 3 个线程全部卡在业务处理上没有空闲线程继续拉取新任务导致队列积压。这个问题的本质是消费者线程数与业务处理线程数的耦合。我后来把长轮询线程和业务处理线程分离长轮询线程只负责拿到消息就丢进一个有界队列BlockingQueue业务线程池负责消费。这样长轮询线程永远空闲用于拉取不会因为业务卡顿而死锁。同时在 Agent-Reach 的队列领取接口里我加了一个capacity字段——Agent 上报自己的当前任务容量Agent-Reach 按容量决定要不要继续投递。这在业务高峰期很有用避免任务发给一个已经忙不过来的 Agent。5.4 回执丢数据的“隐性问题”某次排查发现任务明明 COMPLETED 了调用方却收到了失败结果。查看日志发现callback/complete接口收到 Agent 回执后Agent-Reach 把结果直接转发给调用方的回调地址但转发用的是请求线程池的同一个连接——当任务量暴涨时回调连接池被占满回调 HTTP 请求排队超时调用方误以为失败。修复方式是回调转发改为异步重试队列Agent-Reach 收到回执后先把结果持久化写库这一步很关键保证不丢然后异步执行 HTTP 回调失败时最多重试 3 次重试之间指数退避。这样回执数据不丢调用方最终一定能收到结果代价是回调延迟偶尔会增加几百毫秒但这个场景下可以接受。6. 进阶玩法反向触达、动态编排与多环境隔离6.1 反向触达前面讲的是 Agent-Reach 向 Agent 发送指令但实际协作中Agent 也需要主动向 Agent-Reach 汇报状态、申请资源、请求调用其他 Agent 的服务。比如一个内容审核 Agent 审核完一段文本后希望自动通知翻译 Agent 把审核结果翻译成多语言版本。实现方式非常简单Agent 直接往 Agent-Reach 的/command/send发消息带自己的 token 和 target。因为 Agent-Reach 对调用方和被执行方采用同样的鉴权体系所以 Agent 天然可以成为调用方。我用这种方式做了不少“Agent 链式协作”的场景——A 做完主动通知 BB 做完通知 C整条链路完全由事件驱动中间不需要一个中心调度器。6.2 动态编排用状态机管理跨 Agent 的流程当任务涉及多个 Agent 协作时Agent-Reach 只负责单次触达跨 Agent 的编排逻辑需要上层编排器。我实现了一个轻量级编排器它本身就是注册在 Agent-Reach 里的一个特殊 Agent接收pipeline类型指令内部用有限状态机推进每一步pipeline_start - step_1 (call agent-translator) - step_2 (call agent-summarizer) - pipeline_end编排器的状态存储也放在 SQLite 里这样即使编排器重启也能从上次步骤继续执行不至于整个流程归零。这里有两个坑提醒一下编排器拉取“下一步该交给谁”时不应该硬编码 Agent 名称否则 Agent 换了地址你就得改代码。用 capability 匹配让 Agent-Reach 动态路由会省心很多。编排器的状态推进和 Agent 的回执之间一定是异步的。不要试图把回执和状态推进做在同一个事务里否则会出现“Agent 已回执成功但编排器没来得及推进状态”的尴尬局面。我这里的实现是回执先落地为“事件”编排器轮询“可推进步骤”时再消费事件。6.3 多环境隔离用 namespace 分隔开发、测试、生产我一开始把开发环境和测试环境的 Agent 都注册在同一个 Agent-Reach 实例里结果开发 Agent 的调试任务偶尔会发给测试环境的 Agent两边数据格式对不上排查起来一头雾水。后来启用了 namespace 机制Agent 注册时带namespace字段如dev/test/prod路由时严格限定只能在同一 namespace 内。这个隔离做得很轻只是在注册表的 key 中加了 namespace 前缀消息队列的 topic 也带上 namespace 前缀。代价是每个环境需要独立部署一套 Agent-Reach但换来的是极其清晰的边界。我现在本地开发时就开一个容器测试环境开一个生产环境独立一个数据互不干扰。namespace 还有一个用处共享 Agent 模式。有些常用 Agent比如通用翻译只在生产环境部署开发环境要借用的话可以把生产环境的某些 Agent 标记为shared允许其他 namespace 访问。实现上就是在鉴权 token 里加一个访问范围字段Agent-Reach 判断允许跨 namespace 时才放行。7. Agent-Reach 的监控、日志与运维建议7.1 几个必看的监控指标Agent-Reach 内置了/metrics端点暴露 Prometheus 格式指标。我在 Grafana 里配了几张关键面板强烈建议你重点盯这几个指标消息队列积压数某个 Agent 的未消费消息数持续上涨说明 Agent 处理速度跟不上生产速度。重新投递次数分布如果发现大量消息的重新投递次数为 2-3说明 Agent 不稳定经常领取后超时。死信数量死信骤增时优先看 FAILED 占比和 RETRY 策略。如果都是确定性失败应该改 NO_RETRY否则死信堆积很快。P95 回执延迟从指令发出到收到回执的延迟这个指标影响整体编排链路的耗时。7.2 日志规范requestId 贯穿全链路Agent-Reach 所有日志都强制带requestId。在这个系统里requestId就像一次协作的“身份证号”——调用方发指令时会生成一个Agent 领取消息、处理完成、回执上报都会带上同一个requestId。排查问题时只要拿这个 id 一搜日志链路就能看到消息在哪个环节卡住了。这一点一开始我没做好。那时日志按 Agent 名称各自记录跨 Agent 排查问题时要把多个日志文件拿来对时间戳痛苦得不行。后来统一了 requestId 传递规范任何跨 Agent 调用都要求透传排查效率提升了一个量级。7.3 运维备份与快速恢复Agent-Reach 的数据是一个 SQLite 文件运维上非常简单。我的备份做法每天凌晨sqlite3 agent-reach.db .backup agent-reach-$(date %F).db保留 30 天同时把tokens.json也一并备份。恢复时只需停掉容器替换 db 文件再起容器10 秒内恢复服务。要特别注意备份时要用 SQLite 的 backup API 或.backup命令不要直接 cp 文件。直接复制正在写入的 SQLite 文件会得到损坏的备份。我自己踩过这个坑恢复过一次备份后发现队列数据不完整后来才知道 SQLite 在线热备要用专用命令。8. 一些经验感悟与后续扩展方向Agent-Reach 从第一行代码到稳定运行大概花了两周业余时间。这个过程中我最大的感受是Agent 协作的基础设施不在于功能多炫而在于把“找得到、递得到、答得回”这三件事做可靠。找得到靠一个带健康检查的注册中心递得到靠一套带路由和重试的消息通道答得回靠一个不回丢的回执状态推进机制。这三件事每一项都不难做难的是把它们结合成一个对用户友好的、可以直接上生产的闭环。如果你也想做一个类似的东西我给三点建议第一千万别一上来就引入微服务全家桶。Consul、gRPC、Kafka 一套组合拳下来你的 Agent 还没互联先把自己绊倒了。从最简单的 SQLite HTTP JSON 开始跑通链路后再一点点增强。第二先想清楚 Agent 之间是真的需要强一致还是最终一致。大多数 Agent 协作都是最终一致就够了。为了强一致引入分布式事务在这个量级是纯负担。第三测试环境和生产环境从第一天就分开。我因为省事一开始全共用一套后来清理数据、改 namespace 反而花了更多时间。后续我想做的是把 Agent-Reach 的多集群联邦支持加上让不同机器上的 Agent-Reach 能互通消息这样 Agent 部署在多地也能互相触达。架构上可以在每个实例上加一个peer列表转发跨实例消息时通过 HTTP 隧道不需要额外组件。不过这个版本目前还在设计阶段等做完稳定了我再写一篇分享。有点说多了如果你也正在被一堆 Agent 的互联问题折磨希望这套设计和这些踩坑记录能帮你少走点弯路。实现一个够用的 Agent 触达层真的不难难的是提前想到那些边界情况。把这篇文章提到的几个坑都避开你的 Agent 协作链路至少能稳一半。