EMQX A2A Registry 的 Agent Card 存活状态属性覆写机制解析(fix-17010)
【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx
导读
本篇文章围绕 EMQX 开源仓库中的变更记录 changes/ee/fix-17010.en.md 展开,深入讲解 A2A(Agent-to-Agent)Agent Card 中a2a-status与a2a-status-source两个用户属性(User Property)的覆写机制:当 Agent Card 自身携带这两个属性时,EMQX 会用自己维护的 Agent 存活状态(liveness)信息将其覆盖,以避免同一 Card 中出现重复或冲突的状态属性。读完本文,你将理解该机制的设计动机、源码级实现链路(包括 hook 注册、状态查询、属性删除与重写),以及如何通过测试用例、CLI 和 REST API 验证其行为。
变更背景:为什么需要覆写 Agent Card 中的状态属性
在 A2A 生态中,Agent Card 是描述一个 Agent 能力、接口、技能等信息的标准化文档(其 JSON Schema 见 apps/emqx_a2a_registry/priv/agent_card_schema.json)。EMQX 通过 A2A Registry 模块(apps/emqx_a2a_registry)将这些 Card 以**保留消息(retained message)**的形式发布到 discovery 主题,供其他 Agent 或客户端订阅发现。
Agent Card 的发布者(Agent 本体)在发布时可能自行携带一组 MQTT 用户属性,例如:
a2a-status:声明 Agent 当前状态;a2a-status-source:声明该状态的来源(例如agent,表示由 Agent 自己上报)。
而 EMQX 的 A2A Registry 本身也需要在每条 Card 消息上附加存活状态信息,用于向订阅者实时通告该 Agent 是否在线。这就会产生一个冲突:如果 Agent 自己上报的属性与 EMQX 计算出的存活状态同时存在,订阅方会看到两份语义可能不一致的a2a-status,造成歧义。
正是为了解决这一问题,fix-17010 引入了覆写逻辑。变更记录原文如下:
Now,
a2a-statusanda2a-status-sourceuser properties present in A2A Agent Cards are overridden with EMQX's liveness information to avoid duplicate properties.
该变更同样被收录进版本更新说明 changes/6.2.1.en.md(第 226 行,对应 PR #17010),属于 6.2.1 版本的行为变更。
覆写规则:删除旧值、写入 EMQX 权威状态
覆写逻辑的核心是先删除、再重写,而不是简单拼接。从源码实现来看,整个处理发生在message.delivered钩子回调中,位于 apps/emqx_a2a_registry/src/emqx_a2a_registry_hookcb.erl 的augment_message_metadata/1函数(第 135–148 行):
augment_message_metadata(Msg0) -> ClientId = emqx_message:from(Msg0), Props0 = emqx_message:get_header(properties, Msg0, #{}), UserProperties0 = maps:get('User-Property', Props0, []), UserProperties1 = proplists:delete(?A2A_PROP_STATUS_KEY, UserProperties0), UserProperties2 = proplists:delete(?A2A_PROP_STATUS_SOURCE_KEY, UserProperties1), Status = emqx_a2a_registry:lookup_agent_status(ClientId), UserProperties = [ {?A2A_PROP_STATUS_KEY, Status}, {?A2A_PROP_STATUS_SOURCE_KEY, <<"emqx">>} | UserProperties2 ], Props = maps:put('User-Property', UserProperties, Props0), emqx_message:set_header(properties, Props, Msg0).处理流程可拆解为四步:
- 读取原始用户属性:从消息头的
User-Property中取出属性列表; - 删除冲突键:分别用
proplists:delete移除原消息中携带的a2a-status与a2a-status-source,确保不会出现重复属性; - 查询权威状态:调用
emqx_a2a_registry:lookup_agent_status(ClientId)获取该 Agent 的实时存活状态; - 前置写入新值:将
{a2a-status, Status}与{a2a-status-source, <<"emqx">>}放在属性列表头部(作为第一个出现值,MQTT 协议中同名 User Property 可重复,订阅方一般取首个),完成覆写。
属性名与取值的常量定义
上述代码中的?A2A_PROP_STATUS_KEY等宏定义在 apps/emqx_a2a_registry/src/emqx_a2a_registry_internal.hrl:
| 宏 | 值 | 语义 |
|---|---|---|
?A2A_PROP_STATUS_KEY | <<"a2a-status">> | 状态属性键名 |
?A2A_PROP_STATUS_SOURCE_KEY | <<"a2a-status-source">> | 状态来源属性键名 |
?A2A_PROP_ONLINE_VAL | <<"online">> | 在线状态值 |
?A2A_PROP_OFFLINE_VAL | <<"offline">> | 离线状态值 |
可以看到,覆写后a2a-status-source被固定为emqx,明确告知订阅方:该状态由 EMQX Broker 权威计算得出,而非 Agent 自报。
存活状态如何计算
lookup_agent_status/1定义在 apps/emqx_a2a_registry/src/emqx_a2a_registry.erl(第 109–116 行):
lookup_agent_status(ClientId) -> FormatFn = undefined, case emqx_cm:lookup_client_and_format({clientid, ClientId}, FormatFn) of [{_Chan, #{conn_state := connected} = _Info, _Stats}] -> ?A2A_PROP_ONLINE_VAL; _ -> ?A2A_PROP_OFFLINE_VAL end.它通过emqx_cm:lookup_client_and_format查询连接管理表中的客户端会话:只要该 ClientId 存在且conn_state为connected,即判定为online,否则一律返回offline。也就是说,状态是 Broker 在消息投递时刻实时计算的,而非依赖 Agent 在 Card 中自报的值——这正是"避免重复属性、以 EMQX 为准"的意义所在。
触发时机:message.delivered 钩子与作用范围
覆写逻辑并非对任意消息生效,而是有严格的作用范围与触发条件,由 apps/emqx_a2a_registry/src/emqx_a2a_registry_hookcb.erl 的钩子机制控制:
register_hooks() -> ok = emqx_hooks:add('message.publish', ?MESSAGE_PUBLISH_HOOK, ?HP_RETAINER + 1), ok = emqx_hooks:add('message.delivered', ?MESSAGE_DELIVERED_HOOK, ?HP_RETAINER - 1), ok.message.publish钩子:负责在消息发布入站时校验Card 消息——只有保留消息(retain标志为真)且主题匹配 A2A discovery 主题才会放行;非保留消息会被拒绝(设置allow_publish => false),非 discovery 主题的消息则原样放行。同时还会校验发布者 ClientId 必须与主题中的org_id/unit_id/agent_id三段一致,并校验 Card JSON 符合 Schema(validate_card_schema/1)。message.delivered钩子:在消息投递给订阅者时触发,负责给 discovery 主题的消息附加/覆写存活状态属性(即本次变更的核心逻辑),其他主题消息不做任何处理。
两个钩子的执行优先级分别锚定在 Retainer 前后(?HP_RETAINER ± 1),确保与保留消息的存储、读取时序正确配合。也就是说,覆写发生在"从保留消息存储中取出、投递给订阅者"的路径上,因此订阅者在任何时刻收到的最新 Card 消息都携带 Broker 计算的实时状态。
触发前提:功能开关与保留消息
覆写逻辑受开关控制,两个钩子的回调入口都会先检查emqx_a2a_registry_config:is_enabled(),未启用时直接放行原消息。开关定义在 apps/emqx_a2a_registry/src/emqx_a2a_registry_schema.erl:
a2a_registry { enable = false # 总开关,默认关闭 validate_schema = true # 是否校验 Agent Card 符合 JSON Schema,默认开启 }另外,由于 Agent Card 依赖保留消息存储(emqx_retainer),该功能还要求 Retainer 已启用(retainer.enable = true),否则 CLI 与 API 会返回retainer_disabled错误提示。
作用范围:仅限 A2A discovery 主题
do_on_message_delivered会先调用maybe_augment_message_metadata解析消息主题:只有能成功解析为 A2A discovery 主题(格式为$a2a/v1/discovery/{org_id}/{unit_id}/{agent_id},支持命名空间前缀)的消息才会被覆写属性;其他主题的消息原样透传,不受影响。discovery 主题的拼接逻辑见 apps/emqx_a2a_registry/src/emqx_a2a_registry.erl 的discovery_topic/4(第 45–63 行),主题命名空间常量定义于 apps/emqx_a2a_registry/src/emqx_a2a_registry_internal.hrl($a2a/v1/discovery)。
测试验证:属性被覆写的完整行为样例
仓库的集成测试对本次覆写行为有直接且完整的覆盖,见 apps/emqx_a2a_registry/test/emqx_a2a_registry_SUITE.erl(第 268–289 行)。测试构造了一个"故意携带冲突属性"的发布:
%% If `a2a-status` and/or `a2a-source` keys is present in user properties, they are %% overridden. {ok, _} = emqtt:publish( Agent2, discovery_topic(?ORG_ID, ?UNIT_ID, ?AGENT_ID), #{ 'User-Property' => [ {?A2A_PROP_STATUS_KEY, <<"whatever">>}, {?A2A_PROP_STATUS_SOURCE_KEY, <<"agent">>} ] }, sample_card_bin(), [{qos, 1}, {retain, true}] ), {publish, #{properties := #{'User-Property' := Props3}}} = ?assertReceive({publish, _}), ?assertMatch( #{ ?A2A_PROP_STATUS_KEY := [?A2A_PROP_ONLINE_VAL], ?A2A_PROP_STATUS_SOURCE_KEY := [<<"emqx">>] }, get_a2a_props(Props3) ),测试要点:
- Agent 携带
a2a-status = "whatever"、a2a-status-source = "agent"发布 Card(QoS 1 + retain); - 订阅方收到的消息中,这两个属性已被覆写为
a2a-status = "online"(该测试中 Agent 处于在线状态)与a2a-status-source = "emqx"; - 同一测试用例(第 233–247 行)还验证了 Agent 断开后重新订阅会收到
a2a-status = "offline",证明状态随 Broker 侧连接状态实时变化。
此外,测试还覆盖了非保留 Card 消息被拒绝、非 discovery 主题消息正常透传、功能关闭后客户端可任意发布等边界行为,可作为理解整套 A2A Registry 消息处理语义的参考。
如何观察与使用覆写后的状态
覆写发生在message.delivered路径,因此订阅 discovery 主题的任何客户端都能直接观察到最终属性。此外,A2A Registry 提供的查询入口也统一使用同一套状态语义:
- CLI(apps/emqx_a2a_registry/src/emqx_a2a_registry_cli.erl):
a2a_registry list [--status online|offline]、a2a_registry get ORG_ID UNIT_ID AGENT_ID、a2a_registry delete ...、a2a_registry register ORG_ID UNIT_ID AGENT_ID <agent-card.json>、a2a_registry stats。其中list的--status参数只接受online/offline两个取值(第 331–338 行校验),与属性值常量保持一致; - REST API(apps/emqx_a2a_registry/src/emqx_a2a_registry_api.erl):Card 输出结构
card_out中的status字段同样是online/offline枚举(第 157 行),该值由 apps/emqx_a2a_registry/src/emqx_a2a_registry.erl 的format_card/4(第 179–207 行)在读取保留消息时通过lookup_agent_status实时计算,与投递路径上的覆写逻辑保持一致。
也就是说,无论你是通过 MQTT 订阅 discovery 主题、通过emqx ctl a2a_registry命令,还是通过 REST API 查询,看到的a2a-status/status都是同一份"EMQX 视角的存活状态",不会再出现 Agent 自报状态与 Broker 实际状态不一致导致的重复或矛盾。
总结
fix-17010 通过message.delivered钩子对 A2A discovery 主题消息执行"删除旧值 + 前置写入新值"的覆写策略,将a2a-status与a2a-status-source统一为 EMQX 基于连接管理表实时计算的存活状态(online/offline,来源固定为emqx)。这一机制既保证了订阅者始终拿到权威、一致的 Agent 状态,又避免了属性重复带来的歧义,是整个 A2A Registry 发现链路上关键的一环。核心实现位于 apps/emqx_a2a_registry/src/emqx_a2a_registry_hookcb.erl 与 apps/emqx_a2a_registry/src/emqx_a2a_registry.erl,行为契约由 apps/emqx_a2a_registry/test/emqx_a2a_registry_SUITE.erl 中的集成测试锁定,可作为后续二次开发或排障时的直接参考。
【免费下载链接】emqxThe most scalable and reliable MQTT broker for AI, IoT, IIoT and connected vehicles项目地址: https://gitcode.com/gh_mirrors/em/emqx
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考