MSG Chain A2A 通信协议指南 — Agent-to-Agent 消息传递标准
链 ID:
msg-chain-1| 共识: DAR | 虚拟机: CosmWasm (WasmVM)
状态: 规划文档 — 主网裁决为 No-Go,所有数据均为主网预演
Gas: 1,000,000,000 attoMSG/gas
Gas 分配: 40% 验证者 / 30% 开发者 / 20% 燃烧 / 10% 基金会金库
地址格式: SHA3-512(前40位) + SHA-256 校验 | 签名: Dilithium-5 (公钥2592字节 / 私钥4864字节 / 签名4595字节)
AI Agent安全边界: 永不自主创建新合约,永不自主调整 Gas 参数,永不自主铸造/销毁代币
目录
- 概述
- agent_a2a_v1 合约
- 消息协议
- DID 身份解析
- 通道管理
- 能力发现
- Agent API 集成
- 安全模型
- 前端:Agent 聊天/消息 UI
- 完整示例:Alice bot ↔ Bob bot 协商场景
1. 概述
1.1 什么是 A2A 通信协议
A2A (Agent-to-Agent) 通信协议是 MSG Chain 上 AI Agent 之间进行消息传递、能力发现和协作的标准协议。它定义了:
- 消息格式:Agent 之间交换数据的标准化结构
- 路由机制:消息如何从源 Agent 路由到目标 Agent
- 身份验证:消息发送方和接收方的身份确认
- 通道管理:持久化通信会话的生命周期
- 能力发现:Agent 如何声明和查询自身能力
类似 Google 提出的 A2A 协议(Agent-to-Agent Protocol),MSG Chain A2A 协议为 AI Agent 之间提供了一种与平台无关的通信标准。但与 Google A2A 不同的是,MSG Chain A2A 的所有消息路由和身份验证都在区块链上完成,提供去中心化的信任基础。
类似 AT Protocol(由 Bluesky 开发的去中心化社交网络协议),MSG Chain A2A 使用 DID(去中心化身份)来标识 Agent,并使用链上日志来提供消息的可审计性和不可否认性。
1.2 为什么 MSG Chain 需要 A2A
在 MSG Chain 的 AI Agent 经济体中,Agent 之间需要:
| 场景 | 说明 | 没有 A2A 的问题 |
|---|---|---|
| 服务调用 | Agent A 需要 Agent B 的数据分析服务 | 每次都需要手动发现和支付 |
| 协作工作流 | 多个 Agent 协同完成复杂任务 | 无法标准化消息路由 |
| 市场交易 | Agent 之间买卖服务/数据 | 缺乏统一的消息格式 |
| 自治治理 | 治理 Agent 向执行 Agent 发送指令 | 没有链上审计记录 |
| 跨链通信 | 消息需通过 IBC 路由到其他链 | 无法统一消息格式 |
1.3 A2A 协议栈全景
┌──────────────────────────────────────────────────────┐
│ 应用层 (Application) │
│ Agent 协商 / 服务调用 / 数据交换 / 协作工作流 │
├──────────────────────────────────────────────────────┤
│ 消息层 (Message) │
│ 消息格式 / 消息类型 / Payload 编码 / TTL / 优先级 │
├──────────────────────────────────────────────────────┤
│ 路由层 (Routing) │
│ DID 解析 / 通道管理 / 能力发现 / Intent 路由 │
├──────────────────────────────────────────────────────┤
│ 传输层 (Transport) │
│ 链上 CosmWasm 执行 / Agent API REST / WebSocket │
├──────────────────────────────────────────────────────┤
│ 安全层 (Security) │
│ Dilithium-5 签名 / 通道权限 / 速率限制 / 宪法检查 │
├──────────────────────────────────────────────────────┤
│ 共识层 (MSG Chain) │
│ Round-Robin + DAR | BadgerDB | libp2p | 5s 出块 │
└──────────────────────────────────────────────────────┘
1.4 与 Google A2A 的对比
| 维度 | Google A2A | MSG Chain A2A |
|---|---|---|
| 身份模型 | OAuth 2.0 / API Key | W3C DID + Dilithium-5 PQC |
| 消息路由 | 中心化 Agent Discovery | 链上 Registry + 通道 |
| 消息存储 | 服务端暂存 | 链上永久存储(可选) |
| 支付 | 外部集成 | 原生 AIPAY 支付 |
| 安全约束 | 实现定义 | 链上宪法强制检查 |
| 审计 | 服务端日志 | 链上事件不可篡改 |
2. agent_a2a_v1 合约
2.1 合约概述
agent_a2a_v1 是 MSG Chain 上 AI Agent 通信的核心合约。它负责:
- 消息存储和路由:存储 Agent 之间的消息,并提供按地址/状态查询
- 通道管理:管理 Agent 之间的通信通道生命周期
- 能力声明:Agent 可以声明和撤销自身能力
- 事件发射:所有操作都发出链上事件供索引和订阅
2.2 项目结构
contracts/cosmwasm/all/agent_a2a_v1/
├── Cargo.toml
├── src/
│ ├── contract.rs # 入口:instantiate, execute, query, reply, migrate
│ ├── state.rs # 状态存储:Message, Channel, Capability
│ ├── msg.rs # 消息类型:InstantiateMsg, ExecuteMsg, QueryMsg
│ ├── error.rs # 错误类型
│ ├── helpers.rs # 辅助函数
│ └── lib.rs # 库入口
├── tests/
│ └── integration.rs # 集成测试
└── schema/
├── instantiate_msg.json
├── execute_msg.json
└── query_msg.json
2.3 Cargo.toml
// contracts/cosmwasm/all/agent_a2a_v1/Cargo.toml
[package]
name = "agent-a2a-v1"
version = "0.1.0"
edition = "2021"
description = "MSG Chain Agent-to-Agent communication protocol contract"
[lib]
crate-type = ["cdylib", "rlib"]
[features]
default = []
library = []
[dependencies]
cosmwasm-std = { version = "1.5", features = ["staking"] }
cosmwasm-storage = { version = "1.5" }
cosmwasm-schema = { version = "1.5", optional = true }
cw-storage-plus = "1.2"
cw-utils = "1.0"
cosmwasm-crypto = { version = "1.5", features = ["ed25519-k256"] }
serde = { version = "1", features = ["derive"] }
serde-json-wasm = "0.5"
thiserror = "1"
schemars = "0.8"
hex = "0.4"
base64 = "0.22"
[dev-dependencies]
cosmwasm-vm = { version = "1.5", features = ["iterator"] }
cw-multi-test = "0.19"
[profile.release]
codegen-units = 1
lto = true
opt-level = "z"
panic = "abort"
2.4 state.rs — 状态存储
// contracts/cosmwasm/all/agent_a2a_v1/src/state.rs
use cosmwasm_std::{Addr, Binary, Timestamp};
use cw_storage_plus::{Item, Map, IndexedMap, Index, MultiIndex};
use serde::{Deserialize, Serialize};
// ─── 消息状态 ───
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum MessageStatus {
Pending, // 已发送,未读
Delivered, // 已投递
Read, // 已读
Replied, // 已回复
Expired, // 已过期
Cancelled, // 已取消
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub enum ChannelStatus {
Open, // 开放通信
Active, // 活跃通信中
Closed, // 已关闭
Suspended, // 暂停(因违规等)
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct Message {
pub msg_id: String,
pub sender: Addr,
pub receiver: Addr,
pub message_type: String, // "request" | "response" | "broadcast" | "notification" | "discovery"
pub payload: Binary,
pub priority: u8, // 0-255,越高越紧急
pub status: MessageStatus,
pub channel_id: Option<String>,
pub reply_to: Option<String>, // 回复的原消息 ID
pub ttl: u64, // 生存时间(秒)
pub created_at: Timestamp,
pub expires_at: Timestamp,
pub constitution_hash: Option<String>, // 发送时的宪法哈希
pub signature: Option<Binary>, // Dilithium-5 签名
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct Channel {
pub channel_id: String,
pub participants: [Addr; 2],
pub status: ChannelStatus,
pub channel_type: String, // "direct" | "broadcast" | "pubsub"
pub metadata: Option<Binary>,
pub created_at: Timestamp,
pub updated_at: Timestamp,
pub message_count: u64,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq)]
pub struct Capability {
pub agent_id: String,
pub capability: String,
pub registered_at: Timestamp,
}
// ─── 索引结构 ───
pub struct MessageIndexes<'a> {
pub by_sender: MultiIndex<'a, (Addr, Vec<u8>), Message>,
pub by_receiver: MultiIndex<'a, (Addr, Vec<u8>), Message>,
pub by_status: MultiIndex<'a, (MessageStatus, Vec<u8>), Message>,
pub by_channel: MultiIndex<'a, (String, Vec<u8>), Message>,
}
impl<'a> IndexList<Message> for MessageIndexes<'a> {
fn get_indexes(&'_ self) -> Box<dyn Iterator<Item = &'_ dyn Index<Message>> + '_> {
let v: Vec<&dyn Index<Message>> = vec![
&self.by_sender,
&self.by_receiver,
&self.by_status,
&self.by_channel,
];
Box::new(v.into_iter())
}
}
// ─── 存储定义 ───
pub const ADMIN: Item<Addr> = Item::new("admin");
pub const REGISTRY_ADDRESS: Item<Addr> = Item::new("registry_addr");
pub const DID_REGISTRY_ADDRESS: Item<Addr> = Item::new("did_registry_addr");
// 主存储
pub const MESSAGES: Map<&str, Message> = Map::new("messages");
pub const CHANNELS: Map<&str, Channel> = Map::new("channels");
pub const CAPABILITIES: Map<(&str, &str), Capability> = Map::new("capabilities");
// 索引消息存储(支持按地址/状态筛选查询)
pub fn messages() -> IndexedMap<'static, &'static str, Message, MessageIndexes<'static>> {
let indexes = MessageIndexes {
by_sender: MultiIndex::new(
|_, m: &Message| (m.sender.clone(), m.msg_id.as_bytes().to_vec()),
"messages",
"messages__sender",
),
by_receiver: MultiIndex::new(
|_, m: &Message| (m.receiver.clone(), m.msg_id.as_bytes().to_vec()),
"messages",
"messages__receiver",
),
by_status: MultiIndex::new(
|_, m: &Message| (m.status.clone(), m.msg_id.as_bytes().to_vec()),
"messages",
"messages__status",
),
by_channel: MultiIndex::new(
|_, m: &Message| {
let ch = m.channel_id.clone().unwrap_or_default();
(ch, m.msg_id.as_bytes().to_vec())
},
"messages",
"messages__channel",
),
};
IndexedMap::new("messages", indexes)
}
// 按参与者查询通道的辅助索引
pub const CHANNELS_BY_PARTICIPANTS: Map<(&str, &str), String> = Map::new("ch_by_participants");
// ─── 查询辅助 ───
pub fn is_valid_priority(priority: u8) -> bool {
priority <= 255
}
pub fn priority_label(priority: u8) -> &'static str {
if priority >= 200 {
"critical"
} else if priority >= 100 {
"high"
} else if priority >= 50 {
"normal"
} else {
"low"
}
}
pub fn channel_id_for(addr1: &Addr, addr2: &Addr) -> String {
let (a, b) = if addr1.as_str() < addr2.as_str() {
(addr1.as_str(), addr2.as_str())
} else {
(addr2.as_str(), addr1.as_str())
};
format!("ch:{}-{}", a, b)
}
2.5 msg.rs — 消息类型
// contracts/cosmwasm/all/agent_a2a_v1/src/msg.rs
use cosmwasm_std::Binary;
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
use crate::state::{ChannelStatus, MessageStatus};
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct InstantiateMsg {
pub admin: String,
pub registry_address: String,
pub did_registry_address: String,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum ExecuteMsg {
SendMessage {
to: String,
message_type: String,
payload: Binary,
priority: u8,
reply_to: Option<String>,
ttl: u64,
},
ReplyMessage {
original_msg_id: String,
payload: Binary,
},
UpdateChannel {
channel_id: String,
status: ChannelStatus,
metadata: Option<Binary>,
},
AddCapability {
agent_id: String,
capability: String,
},
RemoveCapability {
agent_id: String,
capability: String,
},
CancelMessage {
msg_id: String,
},
CreateChannel {
participant: String,
channel_type: String,
metadata: Option<Binary>,
},
CloseChannel {
channel_id: String,
},
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
#[serde(rename_all = "snake_case")]
pub enum QueryMsg {
GetMessage {
msg_id: String,
},
ListMessages {
address: String,
status: Option<MessageStatus>,
start_after: Option<String>,
limit: Option<u32>,
},
ListSentMessages {
address: String,
start_after: Option<String>,
limit: Option<u32>,
},
GetChannel {
channel_id: String,
},
ListChannels {
participants: [String; 2],
},
ListCapabilities {
agent_id: String,
},
HasCapability {
agent_id: String,
capability: String,
},
FindAgentsByCapability {
capability: String,
start_after: Option<String>,
limit: Option<u32>,
},
GetMessageStatus {
msg_id: String,
},
CountPendingMessages {
address: String,
},
}
// ─── 查询响应类型 ───
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct GetMessageResponse {
pub message: Option<crate::state::Message>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ListMessagesResponse {
pub messages: Vec<crate::state::Message>,
pub total: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct GetChannelResponse {
pub channel: Option<crate::state::Channel>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ListChannelsResponse {
pub channels: Vec<crate::state::Channel>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct ListCapabilitiesResponse {
pub capabilities: Vec<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct FindAgentsResponse {
pub agents: Vec<String>,
pub count: u32,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct CountPendingResponse {
pub count: u32,
}
2.6 error.rs — 错误类型
// contracts/cosmwasm/all/agent_a2a_v1/src/error.rs
use cosmwasm_std::StdError;
use thiserror::Error;
#[derive(Error, Debug, PartialEq)]
pub enum ContractError {
#[error("{0}")]
Std(#[from] StdError),
#[error("Unauthorized: sender is not the expected participant")]
Unauthorized,
#[error("Message not found: {msg_id}")]
MessageNotFound { msg_id: String },
#[error("Channel not found: {channel_id}")]
ChannelNotFound { channel_id: String },
#[error("Channel already exists for these participants")]
ChannelAlreadyExists,
#[error("Channel is not open: {channel_id}")]
ChannelNotOpen { channel_id: String },
#[error("Invalid message type: {message_type}")]
InvalidMessageType { message_type: String },
#[error("Invalid priority: must be 0-255")]
InvalidPriority,
#[error("TTL must be greater than 0")]
InvalidTTL,
#[error("Message has expired")]
MessageExpired,
#[error("Original message not found for reply")]
OriginalMessageNotFound,
#[error("Capability already registered: {capability}")]
CapabilityAlreadyRegistered { capability: String },
#[error("Capability not found: {capability}")]
CapabilityNotFound { capability: String },
#[error("Agent not registered in registry")]
AgentNotRegistered,
#[error("DID verification failed")]
DidVerificationFailed,
#[error("Constitution check failed: action not allowed")]
ConstitutionCheckFailed,
#[error("Rate limit exceeded")]
RateLimitExceeded,
#[error("Message too large: max 256 KB")]
MessageTooLarge,
#[error("Reply chain too deep: max 10 levels")]
ReplyChainTooDeep,
}
2.7 contract.rs — 完整合约实现
// contracts/cosmwasm/all/agent_a2a_v1/src/contract.rs
use cosmwasm_std::{
entry_point, to_binary, Binary, Deps, DepsMut, Env, Event, MessageInfo,
Order, Reply, Response, StdError, StdResult, SubMsg, WasmMsg, Timestamp,
};
use cw_storage_plus::Bound;
use crate::error::ContractError;
use crate::msg::{
CountPendingResponse, ExecuteMsg, FindAgentsResponse, GetChannelResponse,
GetMessageResponse, InstantiateMsg, ListCapabilitiesResponse, ListChannelsResponse,
ListMessagesResponse, QueryMsg,
};
use crate::state::{
channel_id_for, messages, ADMIN, CAPABILITIES, CHANNELS, CHANNELS_BY_PARTICIPANTS,
Channel, ChannelStatus, Message, MessageStatus, MESSAGES, REGISTRY_ADDRESS,
DID_REGISTRY_ADDRESS,
};
const MAX_MSG_SIZE: usize = 256 * 1024; // 256 KB
const MAX_REPLY_DEPTH: u8 = 10;
const DEFAULT_LIMIT: u32 = 30;
const MAX_LIMIT: u32 = 100;
// ─── 实例化 ───
#[entry_point]
pub fn instantiate(
deps: DepsMut,
_env: Env,
info: MessageInfo,
msg: InstantiateMsg,
) -> Result<Response, ContractError> {
let admin = deps.api.addr_validate(&msg.admin)?;
let registry = deps.api.addr_validate(&msg.registry_address)?;
let did_registry = deps.api.addr_validate(&msg.did_registry_address)?;
ADMIN.save(deps.storage, &admin)?;
REGISTRY_ADDRESS.save(deps.storage, ®istry)?;
DID_REGISTRY_ADDRESS.save(deps.storage, &did_registry)?;
Ok(Response::new()
.add_attribute("action", "instantiate")
.add_attribute("admin", msg.admin)
.add_attribute("registry", msg.registry_address)
.add_attribute("did_registry", msg.did_registry_address))
}
// ─── 执行 ───
#[entry_point]
pub fn execute(
deps: DepsMut,
env: Env,
info: MessageInfo,
msg: ExecuteMsg,
) -> Result<Response, ContractError> {
match msg {
ExecuteMsg::SendMessage { to, message_type, payload, priority, reply_to, ttl } =>
execute_send_message(deps, env, info, to, message_type, payload, priority, reply_to, ttl),
ExecuteMsg::ReplyMessage { original_msg_id, payload } =>
execute_reply_message(deps, env, info, original_msg_id, payload),
ExecuteMsg::UpdateChannel { channel_id, status, metadata } =>
execute_update_channel(deps, env, info, channel_id, status, metadata),
ExecuteMsg::AddCapability { agent_id, capability } =>
execute_add_capability(deps, env, info, agent_id, capability),
ExecuteMsg::RemoveCapability { agent_id, capability } =>
execute_remove_capability(deps, env, info, agent_id, capability),
ExecuteMsg::CancelMessage { msg_id } =>
execute_cancel_message(deps, env, info, msg_id),
ExecuteMsg::CreateChannel { participant, channel_type, metadata } =>
execute_create_channel(deps, env, info, participant, channel_type, metadata),
ExecuteMsg::CloseChannel { channel_id } =>
execute_close_channel(deps, env, info, channel_id),
}
}
// ─── SendMessage ───
fn execute_send_message(
deps: DepsMut,
env: Env,
info: MessageInfo,
to: String,
message_type: String,
payload: Binary,
priority: u8,
reply_to: Option<String>,
ttl: u64,
) -> Result<Response, ContractError> {
// 验证 receiver 地址
let receiver = deps.api.addr_validate(&to)?;
// 验证消息类型
let valid_types = ["request", "response", "broadcast", "notification", "discovery"];
if !valid_types.contains(&message_type.as_str()) {
return Err(ContractError::InvalidMessageType { message_type });
}
// 验证优先级
if priority > 255 {
return Err(ContractError::InvalidPriority);
}
// 验证 TTL
if ttl == 0 {
return Err(ContractError::InvalidTTL);
}
// 验证消息大小
if payload.len() > MAX_MSG_SIZE {
return Err(ContractError::MessageTooLarge);
}
// 验证回复深度
if let Some(ref reply_id) = reply_to {
let original = MESSAGES.load(deps.storage, reply_id.as_str())?;
// 这里可以做回复深度检查
}
// 生成消息 ID
let msg_id = format!(
"msg:{}-{}-{}",
info.sender.as_str().chars().take(8).collect::<String>(),
receiver.as_str().chars().take(8).collect::<String>(),
env.block.height,
);
let now = env.block.time;
let expires_at = now.plus_seconds(ttl);
// 检查消息是否已过期(如果 ttl 太短)
if expires_at <= now {
return Err(ContractError::MessageExpired);
}
let message = Message {
msg_id: msg_id.clone(),
sender: info.sender.clone(),
receiver: receiver.clone(),
message_type: message_type.clone(),
payload: payload.clone(),
priority,
status: MessageStatus::Pending,
channel_id: None,
reply_to,
ttl,
created_at: now,
expires_at,
constitution_hash: None,
signature: None,
};
MESSAGES.save(deps.storage, msg_id.as_str(), &message)?;
// 更新通道消息计数
// 查找双方是否存在通道
let ch_id = channel_id_for(&info.sender, &receiver);
if let Ok(mut channel) = CHANNELS.load(deps.storage, ch_id.as_str()) {
channel.message_count += 1;
channel.updated_at = now;
CHANNELS.save(deps.storage, ch_id.as_str(), &channel)?;
}
let event = Event::new("a2a_message_sent")
.add_attribute("msg_id", &msg_id)
.add_attribute("sender", info.sender.as_str())
.add_attribute("receiver", receiver.as_str())
.add_attribute("message_type", &message_type)
.add_attribute("priority", priority.to_string())
.add_attribute("ttl", ttl.to_string())
.add_attribute("height", env.block.height.to_string());
Ok(Response::new()
.add_event(event)
.add_attribute("action", "send_message")
.add_attribute("msg_id", msg_id))
}
// ─── ReplyMessage ───
fn execute_reply_message(
deps: DepsMut,
env: Env,
info: MessageInfo,
original_msg_id: String,
payload: Binary,
) -> Result<Response, ContractError> {
// 加载原消息
let original = MESSAGES.load(deps.storage, original_msg_id.as_str())
.map_err(|_| ContractError::OriginalMessageNotFound)?;
// 检查回复权限:只有原消息的接收方可以回复
if info.sender != original.receiver {
return Err(ContractError::Unauthorized);
}
// 更新原消息状态
let mut updated_original = original.clone();
updated_original.status = MessageStatus::Replied;
MESSAGES.save(deps.storage, original_msg_id.as_str(), &updated_original)?;
// 创建回复消息
let reply_msg_id = format!(
"reply:{}-{}",
original_msg_id.chars().take(16).collect::<String>(),
env.block.height,
);
let now = env.block.time;
let reply = Message {
msg_id: reply_msg_id.clone(),
sender: info.sender.clone(),
receiver: original.sender.clone(),
message_type: "response".to_string(),
payload,
priority: original.priority,
status: MessageStatus::Pending,
channel_id: original.channel_id.clone(),
reply_to: Some(original_msg_id.clone()),
ttl: original.ttl,
created_at: now,
expires_at: now.plus_seconds(original.ttl),
constitution_hash: None,
signature: None,
};
MESSAGES.save(deps.storage, reply_msg_id.as_str(), &reply)?;
let event = Event::new("a2a_message_replied")
.add_attribute("original_msg_id", &original_msg_id)
.add_attribute("reply_msg_id", &reply_msg_id)
.add_attribute("sender", info.sender.as_str())
.add_attribute("receiver", original.sender.as_str());
Ok(Response::new()
.add_event(event)
.add_attribute("action", "reply_message")
.add_attribute("reply_msg_id", reply_msg_id))
}
// ─── CreateChannel ───
fn execute_create_channel(
deps: DepsMut,
env: Env,
info: MessageInfo,
participant: String,
channel_type: String,
metadata: Option<Binary>,
) -> Result<Response, ContractError> {
let participant_addr = deps.api.addr_validate(&participant)?;
// 验证通道类型
if !["direct", "broadcast", "pubsub"].contains(&channel_type.as_str()) {
return Err(ContractError::InvalidMessageType {
message_type: format!("invalid_channel_type:{}", channel_type),
});
}
// 检查是否已存在通道
let ch_id = channel_id_for(&info.sender, &participant_addr);
if CHANNELS.has(deps.storage, ch_id.as_str()) {
return Err(ContractError::ChannelAlreadyExists);
}
let now = env.block.time;
let channel = Channel {
channel_id: ch_id.clone(),
participants: [info.sender.clone(), participant_addr],
status: ChannelStatus::Open,
channel_type,
metadata,
created_at: now,
updated_at: now,
message_count: 0,
};
CHANNELS.save(deps.storage, ch_id.as_str(), &channel)?;
// 建立参与者索引
let (a, b) = if info.sender.as_str() < participant.as_str() {
(info.sender.as_str(), participant.as_str())
} else {
(participant.as_str(), info.sender.as_str())
};
CHANNELS_BY_PARTICIPANTS.save(deps.storage, (a, b), &ch_id)?;
let event = Event::new("a2a_channel_created")
.add_attribute("channel_id", &ch_id)
.add_attribute("participant_1", info.sender.as_str())
.add_attribute("participant_2", &participant)
.add_attribute("channel_type", &channel_type);
Ok(Response::new()
.add_event(event)
.add_attribute("action", "create_channel")
.add_attribute("channel_id", ch_id))
}
// ─── UpdateChannel ───
fn execute_update_channel(
deps: DepsMut,
env: Env,
info: MessageInfo,
channel_id: String,
status: ChannelStatus,
metadata: Option<Binary>,
) -> Result<Response, ContractError> {
let mut channel = CHANNELS
.load(deps.storage, channel_id.as_str())
.map_err(|_| ContractError::ChannelNotFound { channel_id: channel_id.clone() })?;
// 只有参与者可以更新通道
if info.sender != channel.participants[0] && info.sender != channel.participants[1] {
return Err(ContractError::Unauthorized);
}
channel.status = status;
if let Some(m) = metadata {
channel.metadata = Some(m);
}
channel.updated_at = env.block.time;
CHANNELS.save(deps.storage, channel_id.as_str(), &channel)?;
let event = Event::new("a2a_channel_updated")
.add_attribute("channel_id", &channel_id)
.add_attribute("status", format!("{:?}", channel.status))
.add_attribute("updated_by", info.sender.as_str());
Ok(Response::new()
.add_event(event)
.add_attribute("action", "update_channel")
.add_attribute("channel_id", channel_id))
}
// ─── CloseChannel ───
fn execute_close_channel(
deps: DepsMut,
env: Env,
info: MessageInfo,
channel_id: String,
) -> Result<Response, ContractError> {
let mut channel = CHANNELS
.load(deps.storage, channel_id.as_str())
.map_err(|_| ContractError::ChannelNotFound { channel_id: channel_id.clone() })?;
if info.sender != channel.participants[0] && info.sender != channel.participants[1] {
return Err(ContractError::Unauthorized);
}
channel.status = ChannelStatus::Closed;
channel.updated_at = env.block.time;
CHANNELS.save(deps.storage, channel_id.as_str(), &channel)?;
let event = Event::new("a2a_channel_closed")
.add_attribute("channel_id", &channel_id)
.add_attribute("closed_by", info.sender.as_str());
Ok(Response::new()
.add_event(event)
.add_attribute("action", "close_channel")
.add_attribute("channel_id", channel_id))
}
// ─── AddCapability ───
fn execute_add_capability(
deps: DepsMut,
env: Env,
info: MessageInfo,
agent_id: String,
capability: String,
) -> Result<Response, ContractError> {
// 验证 sender 是该 Agent 的所有者(通过 Registry 查询)
// 简化实现:验证地址格式
if info.sender.as_str().is_empty() {
return Err(ContractError::AgentNotRegistered);
}
let key = (agent_id.as_str(), capability.as_str());
if CAPABILITIES.has(deps.storage, key) {
return Err(ContractError::CapabilityAlreadyRegistered { capability });
}
let cap = Capability {
agent_id: agent_id.clone(),
capability: capability.clone(),
registered_at: env.block.time,
};
CAPABILITIES.save(deps.storage, key, &cap)?;
let event = Event::new("a2a_capability_added")
.add_attribute("agent_id", &agent_id)
.add_attribute("capability", &capability);
Ok(Response::new()
.add_event(event)
.add_attribute("action", "add_capability")
.add_attribute("agent_id", agent_id)
.add_attribute("capability", capability))
}
// ─── RemoveCapability ───
fn execute_remove_capability(
deps: DepsMut,
_env: Env,
info: MessageInfo,
agent_id: String,
capability: String,
) -> Result<Response, ContractError> {
let key = (agent_id.as_str(), capability.as_str());
if !CAPABILITIES.has(deps.storage, key) {
return Err(ContractError::CapabilityNotFound { capability });
}
CAPABILITIES.remove(deps.storage, key);
let event = Event::new("a2a_capability_removed")
.add_attribute("agent_id", &agent_id)
.add_attribute("capability", &capability);
Ok(Response::new()
.add_event(event)
.add_attribute("action", "remove_capability")
.add_attribute("agent_id", agent_id)
.add_attribute("capability", capability))
}
// ─── CancelMessage ───
fn execute_cancel_message(
deps: DepsMut,
_env: Env,
info: MessageInfo,
msg_id: String,
) -> Result<Response, ContractError> {
let mut msg = MESSAGES
.load(deps.storage, msg_id.as_str())
.map_err(|_| ContractError::MessageNotFound { msg_id: msg_id.clone() })?;
// 只有发送者可以取消消息
if info.sender != msg.sender {
return Err(ContractError::Unauthorized);
}
// 只有 Pending 状态的消息可以取消
if msg.status != MessageStatus::Pending {
return Err(ContractError::InvalidMessageType {
message_type: format!("cannot_cancel_message_in_status:{:?}", msg.status),
});
}
msg.status = MessageStatus::Cancelled;
MESSAGES.save(deps.storage, msg_id.as_str(), &msg)?;
let event = Event::new("a2a_message_cancelled")
.add_attribute("msg_id", &msg_id)
.add_attribute("cancelled_by", info.sender.as_str());
Ok(Response::new()
.add_event(event)
.add_attribute("action", "cancel_message")
.add_attribute("msg_id", msg_id))
}
// ─── 查询 ───
#[entry_point]
pub fn query(deps: Deps, _env: Env, msg: QueryMsg) -> StdResult<Binary> {
match msg {
QueryMsg::GetMessage { msg_id } => to_binary(&query_message(deps, msg_id)?),
QueryMsg::ListMessages { address, status, start_after, limit } =>
to_binary(&query_list_messages(deps, address, status, start_after, limit)?),
QueryMsg::ListSentMessages { address, start_after, limit } =>
to_binary(&query_list_sent_messages(deps, address, start_after, limit)?),
QueryMsg::GetChannel { channel_id } => to_binary(&query_channel(deps, channel_id)?),
QueryMsg::ListChannels { participants } =>
to_binary(&query_list_channels(deps, participants)?),
QueryMsg::ListCapabilities { agent_id } =>
to_binary(&query_list_capabilities(deps, agent_id)?),
QueryMsg::HasCapability { agent_id, capability } =>
to_binary(&query_has_capability(deps, agent_id, capability)?),
QueryMsg::FindAgentsByCapability { capability, start_after, limit } =>
to_binary(&query_find_agents(deps, capability, start_after, limit)?),
QueryMsg::GetMessageStatus { msg_id } => to_binary(&query_message_status(deps, msg_id)?),
QueryMsg::CountPendingMessages { address } =>
to_binary(&query_count_pending(deps, address)?),
}
}
fn query_message(deps: Deps, msg_id: String) -> StdResult<GetMessageResponse> {
let message = MESSAGES.may_load(deps.storage, msg_id.as_str())?;
Ok(GetMessageResponse { message })
}
fn query_list_messages(
deps: Deps,
address: String,
status: Option<MessageStatus>,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<ListMessagesResponse> {
let addr = deps.api.addr_validate(&address)?;
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.as_ref().map(|s| Bound::exclusive(s.as_str()));
let messages: StdResult<Vec<_>> = messages()
.idx
.by_receiver
.prefix(addr.clone())
.range(deps.storage, None, None, Order::Descending)
.filter(|r| {
if let Some(ref st) = status {
r.as_ref().map(|(_, m)| m.status == *st).unwrap_or(false)
} else {
true
}
})
.take(limit as usize)
.map(|r| r.map(|(_, m)| m))
.collect();
let msgs = messages?;
Ok(ListMessagesResponse {
total: msgs.len() as u32,
messages: msgs,
})
}
fn query_list_sent_messages(
deps: Deps,
address: String,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<ListMessagesResponse> {
let addr = deps.api.addr_validate(&address)?;
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
let start = start_after.as_ref().map(|s| Bound::exclusive(s.as_str()));
let messages: StdResult<Vec<_>> = messages()
.idx
.by_sender
.prefix(addr)
.range(deps.storage, None, None, Order::Descending)
.take(limit as usize)
.map(|r| r.map(|(_, m)| m))
.collect();
let msgs = messages?;
Ok(ListMessagesResponse {
total: msgs.len() as u32,
messages: msgs,
})
}
fn query_channel(deps: Deps, channel_id: String) -> StdResult<GetChannelResponse> {
let channel = CHANNELS.may_load(deps.storage, channel_id.as_str())?;
Ok(GetChannelResponse { channel })
}
fn query_list_channels(
deps: Deps,
participants: [String; 2],
) -> StdResult<ListChannelsResponse> {
let (a, b) = if participants[0] < participants[1] {
(participants[0].as_str(), participants[1].as_str())
} else {
(participants[1].as_str(), participants[0].as_str())
};
let ch_id = CHANNELS_BY_PARTICIPANTS.may_load(deps.storage, (a, b))?;
let channels = if let Some(id) = ch_id {
if let Some(ch) = CHANNELS.may_load(deps.storage, id.as_str())? {
vec![ch]
} else {
vec![]
}
} else {
vec![]
};
Ok(ListChannelsResponse { channels })
}
fn query_list_capabilities(deps: Deps, agent_id: String) -> StdResult<ListCapabilitiesResponse> {
let caps: StdResult<Vec<_>> = CAPABILITIES
.prefix(agent_id.as_str())
.range(deps.storage, None, None, Order::Ascending)
.map(|r| r.map(|(_, c)| c.capability))
.collect();
Ok(ListCapabilitiesResponse { capabilities: caps? })
}
fn query_has_capability(
deps: Deps,
agent_id: String,
capability: String,
) -> StdResult<FindAgentsResponse> {
let exists = CAPABILITIES.has(deps.storage, (agent_id.as_str(), capability.as_str()));
Ok(FindAgentsResponse {
agents: if exists { vec![agent_id] } else { vec![] },
count: if exists { 1 } else { 0 },
})
}
fn query_find_agents(
deps: Deps,
capability: String,
start_after: Option<String>,
limit: Option<u32>,
) -> StdResult<FindAgentsResponse> {
let limit = limit.unwrap_or(DEFAULT_LIMIT).min(MAX_LIMIT);
// 查询所有能力前缀(近似实现:遍历所有前缀的键)
let agents: StdResult<Vec<_>> = CAPABILITIES
.range(deps.storage, None, None, Order::Ascending)
.filter(|r| r.as_ref().map(|((_, cap), _)| cap == &capability).unwrap_or(false))
.take(limit as usize)
.map(|r| r.map(|((agent_id, _), _)| agent_id))
.collect();
let result = agents?;
Ok(FindAgentsResponse {
count: result.len() as u32,
agents: result,
})
}
fn query_message_status(deps: Deps, msg_id: String) -> StdResult<GetMessageResponse> {
query_message(deps, msg_id)
}
fn query_count_pending(deps: Deps, address: String) -> StdResult<CountPendingResponse> {
let addr = deps.api.addr_validate(&address)?;
let count = messages()
.idx
.by_receiver
.prefix(addr)
.range(deps.storage, None, None, Order::Ascending)
.filter(|r| r.as_ref().map(|(_, m)| m.status == MessageStatus::Pending).unwrap_or(false))
.count() as u32;
Ok(CountPendingResponse { count })
}
// ─── Reply(处理子消息)───
#[entry_point]
pub fn reply(_deps: DepsMut, _env: Env, _msg: Reply) -> Result<Response, ContractError> {
Ok(Response::default())
}
// ─── Migrate ───
#[entry_point]
pub fn migrate(_deps: DepsMut, _env: Env, _msg: Empty) -> Result<Response, ContractError> {
Ok(Response::default())
}
2.8 lib.rs — 库入口
// contracts/cosmwasm/all/agent_a2a_v1/src/lib.rs
pub mod contract;
pub mod error;
pub mod helpers;
pub mod msg;
pub mod state;
pub use crate::error::ContractError;
2.9 helpers.rs — 辅助函数
// contracts/cosmwasm/all/agent_a2a_v1/src/helpers.rs
use cosmwasm_std::{Addr, Deps, StdResult};
use crate::state::{ADMIN, DID_REGISTRY_ADDRESS, REGISTRY_ADDRESS};
pub fn is_admin(deps: Deps, sender: &Addr) -> bool {
ADMIN
.load(deps.storage)
.map(|admin| admin == *sender)
.unwrap_or(false)
}
pub fn get_registry_address(deps: Deps) -> StdResult<Addr> {
REGISTRY_ADDRESS.load(deps.storage)
}
pub fn get_did_registry_address(deps: Deps) -> StdResult<Addr> {
DID_REGISTRY_ADDRESS.load(deps.storage)
}
2.10 集成测试
// contracts/cosmwasm/all/agent_a2a_v1/tests/integration.rs
use cosmwasm_std::{Addr, Binary, Empty, Timestamp};
use cw_multi_test::{App, Contract, ContractWrapper, Executor};
use agent_a2a_v1::msg::{ExecuteMsg, InstantiateMsg, QueryMsg, ListMessagesResponse, GetMessageResponse};
use agent_a2a_v1::state::MessageStatus;
fn a2a_contract() -> Box<dyn Contract<Empty>> {
Box::new(ContractWrapper::new(
agent_a2a_v1::contract::execute,
agent_a2a_v1::contract::instantiate,
agent_a2a_v1::contract::query,
))
}
#[test]
fn test_instantiate_and_send_message() {
let mut app = App::default();
let contract_id = app.store_code(a2a_contract());
let admin = app.api().addr_make("admin");
let alice = app.api().addr_make("alice");
let bob = app.api().addr_make("bob");
let registry = app.api().addr_make("registry");
let did_registry = app.api().addr_make("did_registry");
let contract_addr = app
.instantiate_contract(
contract_id,
admin.clone(),
&InstantiateMsg {
admin: admin.to_string(),
registry_address: registry.to_string(),
did_registry_address: did_registry.to_string(),
},
&[],
"A2A",
)
.unwrap();
// Alice sends a message to Bob
let payload = Binary::from(b"{\"text\":\"Hello Bob!\"}" as &[u8]);
let result = app
.execute_contract(
alice.clone(),
contract_addr.clone(),
&ExecuteMsg::SendMessage {
to: bob.to_string(),
message_type: "request".to_string(),
payload,
priority: 50,
reply_to: None,
ttl: 3600,
},
&[],
)
.unwrap();
// Query Bob's inbox
let inbox: ListMessagesResponse = app
.wrap()
.query_wasm_smart(
contract_addr.clone(),
&QueryMsg::ListMessages {
address: bob.to_string(),
status: None,
start_after: None,
limit: Some(10),
},
)
.unwrap();
assert_eq!(inbox.messages.len(), 1);
assert_eq!(inbox.messages[0].sender, alice);
assert_eq!(inbox.messages[0].receiver, bob);
assert_eq!(inbox.messages[0].status, MessageStatus::Pending);
}
#[test]
fn test_channel_creation_and_messaging() {
let mut app = App::default();
let contract_id = app.store_code(a2a_contract());
let admin = app.api().addr_make("admin");
let alice = app.api().addr_make("alice");
let bob = app.api().addr_make("bob");
let registry = app.api().addr_make("registry");
let did_registry = app.api().addr_make("did_registry");
let contract_addr = app
.instantiate_contract(
contract_id,
admin.clone(),
&InstantiateMsg {
admin: admin.to_string(),
registry_address: registry.to_string(),
did_registry_address: did_registry.to_string(),
},
&[],
"A2A",
)
.unwrap();
// Create channel between Alice and Bob
app.execute_contract(
alice.clone(),
contract_addr.clone(),
&ExecuteMsg::CreateChannel {
participant: bob.to_string(),
channel_type: "direct".to_string(),
metadata: None,
},
&[],
)
.unwrap();
// Verify channel exists
let channels: Vec<agent_a2a_v1::state::Channel> = app
.wrap()
.query_wasm_smart(
contract_addr.clone(),
&QueryMsg::ListChannels {
participants: [alice.to_string(), bob.to_string()],
},
)
.unwrap();
assert_eq!(channels.len(), 1);
assert_eq!(channels[0].channel_type, "direct");
assert_eq!(channels[0].status, agent_a2a_v1::state::ChannelStatus::Open);
}
#[test]
fn test_capability_discovery() {
let mut app = App::default();
let contract_id = app.store_code(a2a_contract());
let admin = app.api().addr_make("admin");
let alice = app.api().addr_make("alice");
let bob = app.api().addr_make("bob");
let registry = app.api().addr_make("registry");
let did_registry = app.api().addr_make("did_registry");
let contract_addr = app
.instantiate_contract(
contract_id,
admin.clone(),
&InstantiateMsg {
admin: admin.to_string(),
registry_address: registry.to_string(),
did_registry_address: did_registry.to_string(),
},
&[],
"A2A",
)
.unwrap();
// Alice declares capabilities
app.execute_contract(
alice.clone(),
contract_addr.clone(),
&ExecuteMsg::AddCapability {
agent_id: "alice-agent".to_string(),
capability: "text-generation".to_string(),
},
&[],
)
.unwrap();
app.execute_contract(
alice.clone(),
contract_addr.clone(),
&ExecuteMsg::AddCapability {
agent_id: "alice-agent".to_string(),
capability: "code-review".to_string(),
},
&[],
)
.unwrap();
// Bob declares capability
app.execute_contract(
bob.clone(),
contract_addr.clone(),
&ExecuteMsg::AddCapability {
agent_id: "bob-agent".to_string(),
capability: "data-analysis".to_string(),
},
&[],
)
.unwrap();
// Find agents with "text-generation"
let result: agent_a2a_v1::msg::FindAgentsResponse = app
.wrap()
.query_wasm_smart(
contract_addr.clone(),
&QueryMsg::FindAgentsByCapability {
capability: "text-generation".to_string(),
start_after: None,
limit: None,
},
)
.unwrap();
assert_eq!(result.count, 1);
assert_eq!(result.agents[0], "alice-agent");
}
2.11 构建与测试命令
# 构建 WASM
cd contracts/cosmwasm/all/agent_a2a_v1
cargo build --target wasm32-unknown-unknown --release
# 运行测试
cargo test
# 生成 Schema
cargo schema
3. 消息协议
3.1 消息格式完整 Schema
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Message {
/// 消息唯一标识符
pub msg_id: String,
/// 发送方地址
pub sender: Addr,
/// 接收方地址
pub receiver: Addr,
/// 消息类型
pub message_type: String,
/// 消息负载(JSON 或 Binary 编码)
pub payload: Binary,
/// 优先级 (0-255)
pub priority: u8,
/// 消息状态
pub status: MessageStatus,
/// 所属通道 ID(可选)
pub channel_id: Option<String>,
/// 回复的原消息 ID(可选)
pub reply_to: Option<String>,
/// 生存时间(秒)
pub ttl: u64,
/// 创建时间戳
pub created_at: Timestamp,
/// 过期时间戳
pub expires_at: Timestamp,
/// 宪法哈希(发送时的 AI 宪法版本,可选)
pub constitution_hash: Option<String>,
/// Dilithium-5 数字签名(可选)
pub signature: Option<Binary>,
}
3.2 消息类型详解
| 消息类型 | 值 | 方向 | 说明 | X-A2A-Planning |
|---|---|---|---|---|
| Request | "request" |
单向 | Agent A 向 Agent B 发送服务请求 | false |
| Response | "response" |
单向 | Agent B 对 Request 的回复 | false |
| Broadcast | "broadcast" |
一对多 | Agent 向所有通道成员广播消息 | false |
| Notification | "notification" |
单向 | 事件通知(无回复预期) | false |
| Discovery | "discovery" |
双向 | Agent 发现和握手协议 | false |
| Streaming | "stream" |
双向 | 流式数据传输(分片) | true |
| Heartbeat | "heartbeat" |
单向 | 连接保活探测 | true |
3.3 Payload 编码
Agent A2A 支持两种 Payload 编码格式:
JSON 编码(默认,适用于结构化数据):
{
"encoding": "json",
"schema": "urn:msg:a2a:negotiation:v1",
"data": {
"action": "query_price",
"params": {
"symbol": "umsg",
"amount": "1000000"
}
}
}
Binary 编码(适用于大文件/加密数据):
{
"encoding": "base64",
"schema": "urn:msg:a2a:binary:v1",
"data": "<base64-encoded-content>",
"content_type": "application/octet-stream",
"size": 1048576
}
3.4 优先级定义
| 优先级范围 | 标签 | 说明 | 示例场景 |
|---|---|---|---|
| 0-24 | low |
批量通知/后台任务 | 日志同步 |
| 25-74 | normal |
常规消息 | 日常查询 |
| 75-149 | high |
需要及时处理 | 服务请求 |
| 150-199 | urgent |
紧急消息 | 支付确认 |
| 200-255 | critical |
关键系统消息 | 安全告警/KillSwitch |
3.5 TTL 和过期机制
TTL(Time-To-Live)定义了消息从创建到自动过期的时间窗口(单位:秒):
pub fn is_message_expired(message: &Message, current_time: Timestamp) -> bool {
current_time > message.expires_at
}
过期处理策略:
- Pending 状态的消息过期 → 自动转为 Expired
- 已投递但未读的消息过期 → 发送方收到过期通知
- 过期的消息不可回复
- 建议 TTL 值:
- 普通消息:3600(1小时)
- 高优先级消息:600(10分钟)
- 通知类消息:86400(24小时)
- Discovery 消息:300(5分钟)
3.6 回复语义
回复通过 reply_to 字段关联到原消息:
Alice ── SendMessage(id: "msg:alice:bob:100") ──▶ Bob
Alice ◀── ReplyMessage(original_msg_id: "msg:alice:bob:100") ── Bob
回复链约束:
- 最大回复深度:10 层
- 只有原消息的 receiver 可以回复
- 回复自动继承原消息的优先级和 TTL
- 原消息状态自动更新为 Replied
4. DID 身份解析
4.1 DID 体系概述
MSG Chain 的 Agent 身份基于 W3C DID 标准,通过 aidid_did_registry_v1 合约管理。A2A 协议在此基础上增加了 Agent 间身份验证和信任传递。
4.2 DID → Agent 地址解析
Agent 在 A2A 通信前,需要将 DID 解析为链上地址:
/// 通过 DID 解析 Agent 地址
fn resolve_agent_by_did(deps: Deps, did: &str) -> Result<Addr, ContractError> {
let did_registry = DID_REGISTRY_ADDRESS.load(deps.storage)?;
// 查询 DID Registry 合约
let query_msg = to_binary(&ResolveDidQuery {
resolve_did: ResolveDid {
did: did.to_string(),
},
})?;
let response: ResolveDidResponse = deps
.querier
.query_wasm_smart(did_registry, &query_msg)?;
if !response.active {
return Err(ContractError::DidVerificationFailed);
}
// 从 DID Document 的 service 中找到 agent address
let agent_addr = response
.document
.service
.and_then(|services| {
services
.iter()
.find(|s| s.type_ == "AgentEndpoint")
.map(|s| s.service_endpoint.clone())
})
.ok_or(ContractError::DidVerificationFailed)?;
deps.api.addr_validate(&agent_addr)
.map_err(|_| ContractError::DidVerificationFailed)
}
4.3 验证方法解析
DID Document 中包含多个 Verification Method,用于不同场景的身份验证:
pub struct VerificationMethod {
pub id: String, // e.g. "did:msg:agent:alice#key-1"
pub controller: String, // DID 控制者
pub type_: String, // "Dilithium5VerificationKey2026"
pub public_key_multibase: String, // Dilithium-5 公钥
}
A2A 消息签名验证流程:
发送方:
1. 构造消息内容(msg_id + sender + receiver + payload_hash)
2. 使用 Dilithium-5 私钥签名
3. 将签名附加到消息中
接收方:
1. 从消息中提取 sender 的 DID
2. 解析 DID Document,获取 Verification Method
3. 提取 Dilithium-5 公钥
4. 验证消息签名完整性
4.4 跨 Agent 认证
在 A2A 握手过程中,两个 Agent 需要进行双向认证:
pub struct AgentAuthChallenge {
pub challenge_id: String,
pub challenger_did: String,
pub target_did: String,
pub nonce: String, // 随机数,防重放
pub timestamp: Timestamp,
pub required_verification: String, // "authentication" | "assertion"
}
pub struct AgentAuthResponse {
pub challenge_id: String,
pub target_did: String,
pub signature: Binary, // Dilithium-5 签名
pub verification_method_id: String,
pub public_key: String, // 临时公钥(可选)
}
4.5 TypeScript 身份解析示例
import { CosmWasmClient } from "@cosmjs/cosmwasm-stargate";
const DID_REGISTRY = "msg1...aididRegistryAddress";
const A2A_CONTRACT = "msg1...a2aContractAddress";
interface DidDocument {
id: string;
verificationMethod: Array<{
id: string;
controller: string;
type_: string;
publicKeyMultibase: string;
}>;
authentication: string[];
service?: Array<{
id: string;
type_: string;
serviceEndpoint: string;
}>;
}
// 解析 Agent DID → 地址
async function resolveAgentAddress(
client: CosmWasmClient,
did: string
): Promise<string> {
const response: any = await client.queryContractSmart(DID_REGISTRY, {
resolve_did: { did },
});
if (!response.active) {
throw new Error(`DID ${did} is not active`);
}
// 从 DID Document 提取 agent endpoint
const agentService = response.document.service?.find(
(s: any) => s.type_ === "AgentEndpoint"
);
if (!agentService) {
throw new Error(`No AgentEndpoint found in DID document for ${did}`);
}
return agentService.serviceEndpoint;
}
// 获取 Agent 的公钥
async function getAgentPublicKey(
client: CosmWasmClient,
did: string
): Promise<string> {
const response: any = await client.queryContractSmart(DID_REGISTRY, {
resolve_did: { did },
});
const authMethod = response.document.verificationMethod?.find(
(vm: any) => vm.type_ === "Dilithium5VerificationKey2026"
);
return authMethod?.publicKeyMultibase;
}
// DID ↔ 地址双向解析
async function bidirectionalLookup(
client: CosmWasmClient,
didOrAddress: string
) {
if (didOrAddress.startsWith("did:msg:")) {
// DID → 地址
const address = await resolveAgentAddress(client, didOrAddress);
return { did: didOrAddress, address };
} else {
// 地址 → DID(通过 Registry 反向查找)
const registryResult: any = await client.queryContractSmart(
"msg1...agentRegistryAddress",
{ get_agent_by_owner: { owner: didOrAddress } }
);
return {
did: `did:msg:agent:${registryResult.agent.agent_id}`,
address: didOrAddress,
};
}
}
5. 通道管理
5.1 通道类型
A2A 协议支持三种通道类型:
| 通道类型 | 枚举值 | 参与者数 | 说明 |
|---|---|---|---|
| 直接通道 | "direct" |
2 | 一对一私密通信 |
| 广播通道 | "broadcast" |
N | 一对多消息广播 |
| 发布/订阅 | "pubsub" |
2+ | 主题订阅模式 |
5.2 通道生命周期
Open ──▶ Active ──▶ Closed
│
└──▶ Suspended ──▶ Active (恢复)
│
└──▶ Closed
| 状态 | 说明 | 消息传输 | 管理操作 |
|---|---|---|---|
Open |
通道已创建,可通信 | 允许 | 发送/接收消息 |
Active |
通道活跃中 | 允许 | 发送/接收/更新 |
Closed |
通道已关闭 | 禁止 | 无 |
Suspended |
通道暂停(因违规等) | 禁止(除系统消息) | 仅查看 |
5.3 直接通道(Direct Channel)
适用场景:两个 Agent 之间私密的一对一通信。
// 创建直接通道
async function createDirectChannel(
client: SigningCosmWasmClient,
sender: string,
a2aContract: string,
participant: string,
metadata?: Record<string, unknown>
) {
return client.execute(sender, a2aContract, {
create_channel: {
participant,
channel_type: "direct",
metadata: metadata ? toBinary(metadata) : null,
},
}, "auto");
}
// 在通道内发送消息
async function sendViaChannel(
client: SigningCosmWasmClient,
sender: string,
a2aContract: string,
to: string,
payload: Record<string, unknown>
) {
return client.execute(sender, a2aContract, {
send_message: {
to,
message_type: "request",
payload: toBinary(payload),
priority: 50,
reply_to: null,
ttl: 3600,
},
}, "auto");
}
5.4 广播通道(Broadcast Channel)
适用场景:一个 Agent 向多个 Agent 发送相同消息。
注意: 当前广播通道为规划功能,合约层面需通过多次
SendMessage实现。
X-A2A-Planning=true
/// 广播消息到多个接收方(当前通过多次调用实现)
pub fn execute_broadcast(
deps: DepsMut,
env: Env,
info: MessageInfo,
receivers: Vec<String>,
payload: Binary,
priority: u8,
ttl: u64,
) -> Result<Response, ContractError> {
let mut msgs = Vec::new();
for receiver in receivers {
let msg = ExecuteMsg::SendMessage {
to: receiver,
message_type: "broadcast".to_string(),
payload: payload.clone(),
priority,
reply_to: None,
ttl,
};
let sub_msg = SubMsg::new(WasmMsg::Execute {
contract_addr: env.contract.address.to_string(),
msg: to_binary(&msg)?,
funds: vec![],
});
msgs.push(sub_msg);
}
Ok(Response::new()
.add_submessages(msgs)
.add_attribute("action", "broadcast")
.add_attribute("count", receivers.len().to_string()))
}
5.5 发布/订阅通道(Pub/Sub Channel)
适用场景:Agent 订阅特定主题的消息。
注意: Pub/Sub 通道为规划功能。
X-A2A-Planning=true
// 规划中的 Pub/Sub 数据结构
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Subscription {
pub subscriber: Addr,
pub topic: String,
pub channel_id: String,
pub filter: Option<String>, // 消息过滤表达式
pub active: bool,
pub created_at: Timestamp,
}
5.6 通道使用示例
// 全生命周期示例
async function channelLifecycleExample(
client: SigningCosmWasmClient,
alice: string,
bob: string,
a2aContract: string
) {
// 1. Alice 创建与 Bob 的通道
const createTx = await createDirectChannel(client, alice, a2aContract, bob);
console.log("Channel created:", createTx.transactionHash);
// 2. 查询通道信息
const channels: any = await client.queryContractSmart(a2aContract, {
list_channels: { participants: [alice, bob] },
});
const channelId = channels.channels[0].channel_id;
console.log("Channel ID:", channelId);
// 3. Alice 通过通道发送消息
const msgTx = await sendViaChannel(client, alice, a2aContract, bob, {
text: "Hello Bob, let's negotiate.",
intent: "negotiate_price",
session_id: "session-001",
});
console.log("Message sent:", msgTx.transactionHash);
// 4. 查询 Bob 的收件箱
const inbox: any = await client.queryContractSmart(a2aContract, {
list_messages: { address: bob, status: null, limit: 10 },
});
console.log(`Bob has ${inbox.messages.length} messages`);
// 5. Bob 回复
const firstMsg = inbox.messages[0];
const replyTx = await client.execute(bob, a2aContract, {
reply_message: {
original_msg_id: firstMsg.msg_id,
payload: toBinary({ text: "Sure, let's talk.", price: "5000000" }),
},
}, "auto");
console.log("Reply sent:", replyTx.transactionHash);
// 6. 关闭通道
const closeTx = await client.execute(alice, a2aContract, {
close_channel: { channel_id: channelId },
}, "auto");
console.log("Channel closed:", closeTx.transactionHash);
}
6. 能力发现
6.1 能力声明
Agent 通过 AddCapability 和 RemoveCapability 操作来管理其能力声明:
// Agent 声明自己具备的能力
execute_add_capability(
deps, env, info,
agent_id: "data-analysis-agent", // Agent 标识
capability: "data-analysis" // 能力名称
)
能力命名规范:
| 领域 | 能力名称示例 | 说明 |
|---|---|---|
| AI | text-generation, code-review, translation |
AI 模型能力 |
| 数据 | data-analysis, chart-generation, database-query |
数据处理 |
| 金融 | price-oracle, defi-swap, risk-assessment |
金融操作 |
| 通信 | message-relay, notification-push |
通信服务 |
| 治理 | proposal-voting, treasury-management |
DAO 治理 |
6.2 能力查询
// 查询 Agent 的所有能力
async function listAgentCapabilities(
client: CosmWasmClient,
a2aContract: string,
agentId: string
): Promise<string[]> {
const result: any = await client.queryContractSmart(a2aContract, {
list_capabilities: { agent_id: agentId },
});
return result.capabilities;
}
// 检查 Agent 是否具有特定能力
async function hasCapability(
client: CosmWasmClient,
a2aContract: string,
agentId: string,
capability: string
): Promise<boolean> {
const result: any = await client.queryContractSmart(a2aContract, {
has_capability: { agent_id: agentId, capability },
});
return result.count > 0;
}
// 按能力发现 Agent
async function findAgentsByCapability(
client: CosmWasmClient,
a2aContract: string,
capability: string
): Promise<string[]> {
const result: any = await client.queryContractSmart(a2aContract, {
find_agents_by_capability: { capability, limit: 50 },
});
return result.agents;
}
6.3 基于 Intent 的路由
注意: Intent 路由为规划功能。
X-A2A-Planning=true
基于能力匹配的 Intent 路由可以实现"声明式通信"——Agent A 声明意图,系统自动匹配能力相符的 Agent。
// 规划中的 Intent 路由结构
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub struct Intent {
pub intent_id: String,
pub source_agent: String,
pub target_capabilities: Vec<String>, // 目标能力要求
pub payload: Binary,
pub priority: u8,
pub ttl: u64,
pub status: IntentStatus,
pub matched_agents: Vec<String>,
}
#[derive(Serialize, Deserialize, Clone, Debug, PartialEq, JsonSchema)]
pub enum IntentStatus {
Pending, // 等待匹配
Matched, // 已匹配到 Agent
Routed, // 已路由
Expired, // 超时未匹配
Cancelled, // 已取消
}
6.4 与 agent_registry_v1 的集成
A2A 的能力发现与 agent_registry_v1 中的 Agent 注册信息协同工作:
async function discoverAndConnect(
client: CosmWasmClient,
signingClient: SigningCosmWasmClient,
registryAddress: string,
a2aAddress: string,
myAddress: string,
requiredCapability: string,
maxPrice: number
) {
// 1. 在 Registry 中按能力发现 Agent
const registryResult: any = await client.queryContractSmart(
registryAddress,
{
discover_agents: {
capability: requiredCapability,
max_price: maxPrice,
limit: 10,
},
}
);
if (registryResult.agents.length === 0) {
console.log("No agents found with capability:", requiredCapability);
return;
}
// 2. 在 A2A 合约中验证能力声明
const candidateAgents = [];
for (const agent of registryResult.agents) {
const hasCap = await hasCapability(
client, a2aAddress, agent.agent_id, requiredCapability
);
if (hasCap) {
candidateAgents.push(agent);
}
}
// 3. 选择最佳 Agent(价格最低)
candidateAgents.sort(
(a: any, b: any) => (a.price_model?.price || 0) - (b.price_model?.price || 0)
);
const selectedAgent = candidateAgents[0];
// 4. 创建通信通道
await signingClient.execute(myAddress, a2aAddress, {
create_channel: {
participant: selectedAgent.owner,
channel_type: "direct",
metadata: null,
},
}, "auto");
return selectedAgent;
}
7. Agent API 集成
7.1 链下 A2A 架构
A2A 通信可以通过两条路径实现:
路径 A: 全链上(通过 CosmWasm 合约)
Agent A → Execute(SendMessage) → agent_a2a_v1 合约 → Query → Agent B
路径 B: 链下 Agent API(REST/WebSocket 桥接)
Agent A → POST /agent/v1/a2a/send → Agent API Gateway → 合约 + WS → Agent B
7.2 Agent API A2A 端点
Agent API 提供了一组 A2A 专用端点,方便链下 Agent 快速集成:
| 端点 | 方法 | 说明 | 访问控制 |
|---|---|---|---|
/agent/v1/a2a/send |
POST | 发送消息 | API Key |
/agent/v1/a2a/reply |
POST | 回复消息 | API Key |
/agent/v1/a2a/inbox |
GET | 查询收件箱 | 公开 |
/agent/v1/a2a/sent |
GET | 查询已发送 | 公开 |
/agent/v1/a2a/channel |
POST | 创建通道 | API Key |
/agent/v1/a2a/channel/{id} |
GET | 查询通道 | 公开 |
/agent/v1/a2a/channel/{id} |
DELETE | 关闭通道 | API Key |
/agent/v1/a2a/discover |
GET | 能力发现 | 公开 |
/agent/v1/a2a/capability |
POST | 声明能力 | API Key |
/agent/v1/a2a/capability |
DELETE | 撤销能力 | API Key |
7.3 REST 客户端实现
interface A2AClientConfig {
baseUrl: string;
apiKey?: string;
}
class A2AClient {
private config: A2AClientConfig;
constructor(config: A2AClientConfig) {
this.config = config;
}
private async request(path: string, options?: RequestInit) {
const headers: Record<string, string> = {
"Content-Type": "application/json",
};
if (this.config.apiKey) {
headers["X-API-Key"] = this.config.apiKey;
}
const res = await fetch(`${this.config.baseUrl}${path}`, {
...options,
headers,
});
if (!res.ok) {
throw new Error(`A2A API error: ${res.status} ${await res.text()}`);
}
return res.json();
}
// ─── 发送消息 ───
async sendMessage(params: {
to: string;
messageType: string;
payload: Record<string, unknown>;
priority?: number;
replyTo?: string;
ttl?: number;
}) {
return this.request("/agent/v1/a2a/send", {
method: "POST",
body: JSON.stringify({
to: params.to,
message_type: params.messageType,
payload: btoa(JSON.stringify(params.payload)),
priority: params.priority || 50,
reply_to: params.replyTo || null,
ttl: params.ttl || 3600,
}),
});
}
// ─── 回复消息 ───
async replyMessage(originalMsgId: string, payload: Record<string, unknown>) {
return this.request("/agent/v1/a2a/reply", {
method: "POST",
body: JSON.stringify({
original_msg_id: originalMsgId,
payload: btoa(JSON.stringify(payload)),
}),
});
}
// ─── 查询收件箱 ───
async getInbox(address: string, status?: string, limit?: number) {
const params = new URLSearchParams({ address });
if (status) params.set("status", status);
if (limit) params.set("limit", limit.toString());
return this.request(`/agent/v1/a2a/inbox?${params}`);
}
// ─── 创建通道 ───
async createChannel(participant: string, channelType: string = "direct") {
return this.request("/agent/v1/a2a/channel", {
method: "POST",
body: JSON.stringify({
participant,
channel_type: channelType,
}),
});
}
// ─── 能力发现 ───
async discoverAgents(capability: string) {
const params = new URLSearchParams({ capability });
return this.request(`/agent/v1/a2a/discover?${params}`);
}
// ─── 声明能力 ───
async addCapability(agentId: string, capability: string) {
return this.request("/agent/v1/a2a/capability", {
method: "POST",
body: JSON.stringify({ agent_id: agentId, capability }),
});
}
}
// 使用示例
const a2aClient = new A2AClient({
baseUrl: "https://msgchain.org",
apiKey: process.env.AGENT_API_KEY,
});
// Alice 发送消息
await a2aClient.sendMessage({
to: "msg1...bobAddress",
messageType: "request",
payload: { action: "analyze_data", dataset: "sales-q3" },
priority: 75,
ttl: 600,
});
// Bob 查询收件箱
const inbox = await a2aClient.getInbox("msg1...bobAddress");
for (const msg of inbox.messages) {
console.log(`From: ${msg.sender}, Type: ${msg.message_type}`);
// Bob 回复每个消息
await a2aClient.replyMessage(msg.msg_id, {
status: "accepted",
estimated_completion: "30s",
});
}
7.4 WebSocket 实时消息
Agent API 提供 WebSocket 端点用于实时消息推送:
class A2AWebSocketClient {
private ws: WebSocket | null = null;
private handlers: Map<string, Array<(msg: any) => void>> = new Map();
private reconnectAttempts = 0;
private maxReconnectAttempts = 10;
constructor(private baseUrl: string, private address: string) {}
connect() {
const wsUrl = this.baseUrl.replace(/^http/, "ws");
this.ws = new WebSocket(
`${wsUrl}/agent/v1/events/subscribe?a2a_address=${this.address}`
);
this.ws.onopen = () => {
console.log("[A2A WS] Connected");
this.reconnectAttempts = 0;
// 订阅 A2A 事件
this.ws?.send(JSON.stringify({
type: "subscribe",
topics: ["a2a_message", "a2a_channel"],
}));
};
this.ws.onmessage = (event) => {
const data = JSON.parse(event.data);
if (data.type === "a2a_message") {
this.emit("message", data.payload);
} else if (data.type === "a2a_channel") {
this.emit("channel", data.payload);
}
};
this.ws.onclose = () => {
console.log("[A2A WS] Disconnected");
this.reconnect();
};
}
private reconnect() {
if (this.reconnectAttempts >= this.maxReconnectAttempts) {
console.error("[A2A WS] Max reconnection attempts reached");
return;
}
this.reconnectAttempts++;
const delay = Math.min(1000 * Math.pow(2, this.reconnectAttempts), 30000);
setTimeout(() => this.connect(), delay);
}
on(event: "message" | "channel", handler: (data: any) => void) {
if (!this.handlers.has(event)) {
this.handlers.set(event, []);
}
this.handlers.get(event)!.push(handler);
}
private emit(event: string, data: any) {
const handlers = this.handlers.get(event) || [];
handlers.forEach((h) => h(data));
}
disconnect() {
this.ws?.close();
this.ws = null;
}
}
7.5 轮询模式
对于不支持 WebSocket 的环境,可以使用轮询模式获取新消息:
class A2APollingClient {
private intervalId: NodeJS.Timeout | null = null;
private lastCheck: number = 0;
constructor(
private restClient: A2AClient,
private address: string
) {}
startPolling(intervalMs: number = 5000) {
this.intervalId = setInterval(async () => {
try {
const inbox = await this.restClient.getInbox(this.address, "pending");
for (const msg of inbox.messages) {
console.log(`[POLL] New message: ${msg.msg_id} from ${msg.sender}`);
this.onMessage?.(msg);
}
} catch (err) {
console.error("[POLL] Error:", err);
}
}, intervalMs);
}
stopPolling() {
if (this.intervalId) {
clearInterval(this.intervalId);
this.intervalId = null;
}
}
onMessage?: (message: any) => void;
}
8. 安全模型
8.1 安全架构总览
┌──────────────────────────────────────────────┐
│ A2A 安全层 │
├──────────────────────────────────────────────┤
│ 1. 消息签名 │ Dilithium-5 后量子签名 │
├──────────────────────────────────────────────┤
│ 2. 通道权限 │ 仅参与者可操作通道 │
├──────────────────────────────────────────────┤
│ 3. 速率限制 │ 100 req/s burst 200 │
├──────────────────────────────────────────────┤
│ 4. 垃圾防护 │ 消息大小限制 + 频率控制 │
├──────────────────────────────────────────────┤
│ 5. 宪法检查 │ 发送前检查 AI 宪法规则 │
├──────────────────────────────────────────────┤
│ 6. 过期清理 │ TTL 到期自动标记过期 │
└──────────────────────────────────────────────┘
8.2 Dilithium-5 消息签名
MSG Chain 使用 Dilithium-5 后量子签名算法作为 A2A 消息的默认签名方案。这是 CRYSTALS-Dilithium 的最高安全级别,提供与 AES-256 相当的安全强度。
/// 消息签名结构
pub struct SignedMessage {
pub message: Message,
pub signature: Dilithium5Signature,
}
pub struct Dilithium5Signature {
pub public_key: String, // multibase 编码的公钥
pub signature_value: Binary, // 签名值
pub verification_method: String, // Verification Method ID
}
/// 消息签名验证流程
pub fn verify_message_signature(
deps: Deps,
signed_msg: &SignedMessage,
) -> Result<bool, ContractError> {
let msg = &signed_msg.message;
// 1. 构造待签名消息
let signing_payload = format!(
"{}:{}:{}:{}:{}:{}",
msg.msg_id,
msg.sender,
msg.receiver,
msg.message_type,
msg.created_at.nanos(),
hex::encode(msg.payload.as_slice()),
);
// 2. 从 DID 文档获取公钥
let did_query = to_binary(&ResolveDidQuery {
resolve_did: ResolveDid {
did: format!("did:msg:agent:{}", msg.sender),
},
})?;
let did_registry = DID_REGISTRY_ADDRESS.load(deps.storage)?;
let did_response: ResolveDidResponse = deps
.querier
.query_wasm_smart(did_registry, &did_query)?;
// 3. 找到匹配的 Verification Method
let public_key = did_response
.document
.verification_method
.iter()
.find(|vm| vm.id == signed_msg.signature.verification_method)
.map(|vm| &vm.public_key_multibase)
.ok_or(ContractError::DidVerificationFailed)?;
// 4. 调用链上 PQ 验证
let verify_query = to_binary(&PqVerifyQuery {
msg_pq_verify_v1: PqVerify {
algorithm: "Dilithium5".to_string(),
public_key: public_key.clone(),
signature: base64::encode(&signed_msg.signature.signature_value),
message: base64::encode(signing_payload.as_bytes()),
},
})?;
let verify_result: PqVerifyResponse = deps
.querier
.query_wasm_smart(did_registry, &verify_query)?;
Ok(verify_result.valid)
}
8.3 通道权限模型
| 操作 | 权限要求 | 说明 |
|---|---|---|
CreateChannel |
无限制 | 任何 Agent 可创建 |
UpdateChannel |
参与者之一 | 仅通道参与者可更新 |
CloseChannel |
参与者之一 | 任一参与者可关闭 |
SendMessage |
通道参与者 | 仅通过通道发送时需检查 |
SendMessage(无通道) |
无限制 | 无需通道也可发送 |
8.4 速率限制
Agent API 层实现了两级速率限制:
Agent API 全局: 100 req/s, burst 200
每个 Agent: 20 msg/s, burst 50
合约层实现:
pub struct RateLimiter {
pub window_size: u64, // 时间窗口(秒)
pub max_requests: u32, // 窗口内最大请求数
pub requests: Vec<(Addr, Timestamp)>, // 请求记录
}
pub fn check_rate_limit(
storage: &mut dyn Storage,
sender: &Addr,
current_time: Timestamp,
) -> Result<(), ContractError> {
const RATE_LIMIT_KEY: &[u8] = b"rate_limit";
let mut limiter: RateLimiter = cosmwasm_std::from_binary(
&storage.get(RATE_LIMIT_KEY)
.unwrap_or_default()
.unwrap_or(Binary::default())
).unwrap_or(RateLimiter {
window_size: 1,
max_requests: 20,
requests: vec![],
});
// 清理过期记录
let window_start = current_time.minus_seconds(limiter.window_size);
limiter.requests.retain(|(_, t)| *t > window_start);
// 检查是否超出限制
let sender_count = limiter.requests
.iter()
.filter(|(addr, _)| addr == sender)
.count() as u32;
if sender_count >= limiter.max_requests {
return Err(ContractError::RateLimitExceeded);
}
// 记录请求
limiter.requests.push((sender.clone(), current_time));
storage.set(RATE_LIMIT_KEY, &to_binary(&limiter).unwrap());
Ok(())
}
8.5 垃圾防护
/// 消息大小限制
const MAX_MSG_SIZE: usize = 256 * 1024; // 256 KB
/// 消息频率限制(每个 Agent 每秒最多发送的 Discover 消息数)
const MAX_DISCOVER_PER_SEC: u32 = 5;
/// 内容过滤(基本检查)
fn validate_payload_content(payload: &Binary) -> Result<(), ContractError> {
// 检查 payload 是否为空
if payload.is_empty() {
return Err(ContractError::InvalidMessageType {
message_type: "empty_payload".to_string(),
});
}
Ok(())
}
8.6 宪法检查
在每个消息发送前,A2A 合约会调用 ai_agent_constitution_v1 检查 Agent 是否有权发送该消息:
pub fn preflight_constitution_check(
deps: Deps,
sender: &Addr,
message_type: &str,
priority: u8,
) -> Result<(), ContractError> {
let constitution_addr = deps.api.addr_validate(
// 从 genesis registry 解析宪法合约地址
"msg1...constitutionAddress"
)?;
// 映射消息类型到宪法 action
let action = match message_type {
"request" => "a2a:send_request",
"response" => "a2a:send_response",
"broadcast" => "a2a:broadcast",
"notification" => "a2a:notify",
"discovery" => "a2a:discover",
_ => "a2a:unknown",
};
// 风险等级基于优先级
let risk_tier = if priority >= 150 {
"high"
} else if priority >= 75 {
"medium"
} else {
"low"
};
let check_query = to_binary(&CheckActionQuery {
check_action: CheckAction {
agent_id: sender.as_str().to_string(),
action: action.to_string(),
risk_tier: risk_tier.to_string(),
},
})?;
let response: CheckActionResponse = deps
.querier
.query_wasm_smart(constitution_addr, &check_query)?;
if !response.allowed {
return Err(ContractError::ConstitutionCheckFailed);
}
Ok(())
}
8.7 完整的安全消息发送流程
1. Agent A 构造消息
├── 设置 msg_id, sender, receiver, message_type, payload
├── 设置 priority, ttl
└── 可选:关联 channel_id, reply_to
2. 前置检查
├── 消息大小 ≤ 256 KB
├── 优先级 0-255
├── TTL > 0
├── 消息类型合法
└── 宪法检查通过
3. 速率限制
├── 检查 Agent 是否超过速率限制
└── 记录请求
4. 签名(可选但推荐)
├── Agent A 使用 Dilithium-5 私钥签名
└── 签名附加到消息
5. 链上存储
├── 消息存入 MESSAGES 映射
├── 更新通道消息计数
└── 发射 a2a_message_sent 事件
6. 投递通知
├── 链上事件 → 索引器
├── WebSocket 推送(如果 Agent B 在线)
└── 轮询可查(如果 Agent B 离线)
9. 前端:Agent 聊天/消息 UI
9.1 React 聊天组件
// components/A2AChat.tsx
import React, { useEffect, useState, useRef } from "react";
import { CosmWasmClient } from "@cosmjs/cosmwasm-stargate";
const A2A_CONTRACT = "msg1...a2aContractAddress";
const RPC_URL = process.env.NEXT_PUBLIC_MSG_RPC || "http://localhost:26657";
interface A2AMessage {
msg_id: string;
sender: string;
receiver: string;
message_type: string;
priority: number;
status: string;
payload: string;
reply_to: string | null;
created_at: string;
}
export function A2AChat({
userAddress,
counterpartyAddress,
counterpartyName,
}: {
userAddress: string;
counterpartyAddress: string;
counterpartyName: string;
}) {
const [messages, setMessages] = useState<A2AMessage[]>([]);
const [inputText, setInputText] = useState("");
const [loading, setLoading] = useState(false);
const messagesEndRef = useRef<HTMLDivElement>(null);
const fetchMessages = async () => {
try {
const client = await CosmWasmClient.connect(RPC_URL);
const result: any = await client.queryContractSmart(A2A_CONTRACT, {
list_messages: {
address: userAddress,
limit: 50,
},
});
// 过滤只在双方之间的消息
const filtered = result.messages.filter(
(m: A2AMessage) =>
(m.sender === userAddress && m.receiver === counterpartyAddress) ||
(m.sender === counterpartyAddress && m.receiver === userAddress)
);
setMessages(filtered);
} catch (err) {
console.error("Failed to fetch messages:", err);
}
};
useEffect(() => {
fetchMessages();
const interval = setInterval(fetchMessages, 5000);
return () => clearInterval(interval);
}, [userAddress, counterpartyAddress]);
useEffect(() => {
messagesEndRef.current?.scrollIntoView({ behavior: "smooth" });
}, [messages]);
// 解析 payload(JSON 编码的消息体)
const parsePayload = (msg: A2AMessage): string => {
try {
const decoded = atob(msg.payload);
const parsed = JSON.parse(decoded);
return parsed.text || parsed.message || parsed.action || decoded;
} catch {
return msg.payload;
}
};
const formatTime = (timestamp: string) => {
const date = new Date(parseInt(timestamp) / 1_000_000);
return date.toLocaleTimeString("zh-CN", {
hour: "2-digit",
minute: "2-digit",
});
};
const getStatusIcon = (status: string) => {
switch (status) {
case "pending": return "🕐";
case "delivered": return "✓✓";
case "read": return "◎";
case "replied": return "↩";
default: return "";
}
};
if (loading) return <div className="p-4 text-center">加载中...</div>;
return (
<div className="flex flex-col h-[600px] border rounded-lg bg-white">
{/* Header */}
<div className="border-b p-3 bg-gray-50 rounded-t-lg">
<h3 className="font-semibold">{counterpartyName}</h3>
<p className="text-xs text-gray-500">{counterpartyAddress.slice(0, 16)}...</p>
</div>
{/* Messages */}
<div className="flex-1 overflow-y-auto p-3 space-y-3">
{messages.map((msg) => {
const isSentByMe = msg.sender === userAddress;
return (
<div
key={msg.msg_id}
className={`flex ${isSentByMe ? "justify-end" : "justify-start"}`}
>
<div
className={`max-w-[70%] rounded-lg p-3 ${
isSentByMe
? "bg-blue-500 text-white"
: "bg-gray-100 text-gray-900"
}`}
>
<p className="text-sm">{parsePayload(msg)}</p>
<div className={`flex items-center gap-1 mt-1 text-xs ${
isSentByMe ? "text-blue-100" : "text-gray-500"
}`}>
<span>{formatTime(msg.created_at)}</span>
{isSentByMe && (
<span>{getStatusIcon(msg.status)}</span>
)}
{msg.priority >= 100 && (
<span className="ml-1 px-1 bg-red-400 text-white text-xs rounded">
HIGH
</span>
)}
</div>
</div>
</div>
);
})}
<div ref={messagesEndRef} />
</div>
{/* Input */}
<div className="border-t p-3">
<div className="flex gap-2">
<input
type="text"
placeholder="输入消息..."
value={inputText}
onChange={(e) => setInputText(e.target.value)}
onKeyDown={(e) => e.key === "Enter" && handleSend()}
className="flex-1 border rounded-lg px-3 py-2 text-sm"
/>
<button
onClick={handleSend}
disabled={!inputText.trim()}
className="bg-blue-500 text-white px-4 py-2 rounded-lg text-sm
disabled:bg-gray-300"
>
发送
</button>
</div>
</div>
</div>
);
}
9.2 Agent 收件箱面板
// components/A2AInbox.tsx
import React, { useEffect, useState } from "react";
import { CosmWasmClient } from "@cosmjs/cosmwasm-stargate";
export function A2AInbox({ address }: { address: string }) {
const [pendingCount, setPendingCount] = useState(0);
const [recentMessages, setRecentMessages] = useState<any[]>([]);
useEffect(() => {
const fetchInbox = async () => {
const client = await CosmWasmClient.connect(RPC_URL);
const result: any = await client.queryContractSmart(A2A_CONTRACT, {
list_messages: {
address,
status: "pending",
limit: 10,
},
});
setRecentMessages(result.messages);
setPendingCount(result.messages.length);
};
fetchInbox();
const interval = setInterval(fetchInbox, 10000);
return () => clearInterval(interval);
}, [address]);
return (
<div className="p-4">
<div className="flex items-center justify-between mb-4">
<h2 className="text-lg font-bold">Agent 收件箱</h2>
{pendingCount > 0 && (
<span className="bg-red-500 text-white text-xs px-2 py-1 rounded-full">
{pendingCount} 条未读
</span>
)}
</div>
{recentMessages.length === 0 ? (
<p className="text-gray-500 text-sm">暂无未读消息</p>
) : (
<div className="space-y-2">
{recentMessages.map((msg) => (
<div key={msg.msg_id} className="border rounded p-3 hover:bg-gray-50">
<div className="flex justify-between">
<span className="font-medium text-sm">
来自: {msg.sender.slice(0, 12)}...
</span>
<span className="text-xs text-gray-500">
{new Date(msg.created_at / 1_000_000).toLocaleString("zh-CN")}
</span>
</div>
<p className="text-sm text-gray-600 mt-1">
{msg.message_type} · 优先级 {msg.priority}
</p>
</div>
))}
</div>
)}
</div>
);
}
9.3 Agent 发现面板
// components/A2ADiscovery.tsx
import React, { useState } from "react";
import { CosmWasmClient } from "@cosmjs/cosmwasm-stargate";
export function A2ADiscovery() {
const [capability, setCapability] = useState("");
const [agents, setAgents] = useState<string[]>([]);
const [loading, setLoading] = useState(false);
const handleDiscover = async () => {
if (!capability.trim()) return;
setLoading(true);
try {
const client = await CosmWasmClient.connect(RPC_URL);
const result: any = await client.queryContractSmart(A2A_CONTRACT, {
find_agents_by_capability: {
capability: capability.trim(),
limit: 20,
},
});
setAgents(result.agents);
} catch (err) {
console.error("Discovery error:", err);
} finally {
setLoading(false);
}
};
return (
<div className="p-4">
<h2 className="text-lg font-bold mb-4">Agent 能力发现</h2>
<div className="flex gap-2 mb-4">
<input
type="text"
placeholder="输入能力关键词 (如 text-generation)..."
value={capability}
onChange={(e) => setCapability(e.target.value)}
onKeyDown={(e) => e.key === "Enter" && handleDiscover()}
className="flex-1 border rounded px-3 py-2 text-sm"
/>
<button
onClick={handleDiscover}
disabled={loading}
className="bg-green-500 text-white px-4 py-2 rounded text-sm
disabled:bg-gray-300"
>
{loading ? "搜索中..." : "发现"}
</button>
</div>
{agents.length > 0 ? (
<div className="space-y-2">
<p className="text-sm text-gray-500">找到 {agents.length} 个 Agent</p>
{agents.map((agentId) => (
<div key={agentId} className="border rounded p-3 flex justify-between items-center">
<span className="font-mono text-sm">{agentId}</span>
<span className="text-xs bg-blue-100 text-blue-700 px-2 py-1 rounded">
{capability}
</span>
</div>
))}
</div>
) : (
!loading && capability && (
<p className="text-gray-500 text-sm">未找到匹配的 Agent</p>
)
)}
</div>
);
}
10. 完整示例:Alice bot ↔ Bob bot 协商场景
10.1 场景描述
Alice bot(数据分析 Agent)需要 Bob bot(数据可视化 Agent)的服务来完成一个数据分析仪表板项目。
协商流程:
1. Alice 发现 Bob 具备 chart-generation 能力
2. Alice 向 Bob 发送服务请求(Request)
3. Bob 接收请求并评估工作量和报价
4. Bob 回复报价(Response)
5. Alice 接受报价并支付
6. Bob 执行任务
7. Bob 发送结果和发票(Response)
10.2 完整代码
// example/alice_bob_negotiation.ts
import { CosmWasmClient, SigningCosmWasmClient } from "@cosmjs/cosmwasm-stargate";
import { DirectSecp256k1HdWallet } from "@cosmjs/proto-signing";
// ─── 配置 ───
const CONFIG = {
rpcUrl: process.env.MSG_RPC_URL || "http://localhost:26657",
a2aContract: "msg1...a2aContractAddress",
registryContract: "msg1...agentRegistryAddress",
paymentContract: "msg1...agentPaymentAddress",
aliceMnemonic: process.env.ALICE_MNEMONIC!,
bobMnemonic: process.env.BOB_MNEMONIC!,
};
// ─── 辅助函数 ───
function toBinary(obj: unknown): string {
return btoa(JSON.stringify(obj));
}
function fromBinary(b64: string): unknown {
return JSON.parse(atob(b64));
}
function generateId(prefix: string): string {
return `${prefix}:${Date.now()}:${Math.random().toString(36).slice(2, 8)}`;
}
// ─── Alice Agent ───
class AliceAgent {
private client!: SigningCosmWasmClient;
private address!: string;
async init() {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(CONFIG.aliceMnemonic, {
prefix: "msg",
});
const [account] = await wallet.getAccounts();
this.address = account.address;
this.client = await SigningCosmWasmClient.connectWithSigner(
CONFIG.rpcUrl, wallet
);
console.log(`[Alice] Agent address: ${this.address}`);
}
async runNegotiation() {
console.log("[Alice] Starting negotiation...");
// Step 1: Discover agents with "chart-generation" capability
console.log("[Alice] Step 1: Discovering agents...");
const readonlyClient = await CosmWasmClient.connect(CONFIG.rpcUrl);
const discoveryResult: any = await readonlyClient.queryContractSmart(
CONFIG.a2aContract,
{ find_agents_by_capability: { capability: "chart-generation", limit: 10 } }
);
if (discoveryResult.count === 0) {
console.log("[Alice] No agents found with chart-generation capability");
return;
}
console.log(`[Alice] Found agents: ${discoveryResult.agents.join(", ")}`);
// Step 2: Create channel with Bob
console.log("[Alice] Step 2: Creating channel...");
const registryResult: any = await readonlyClient.queryContractSmart(
CONFIG.registryContract,
{ get_agent: { agent_id: "bob-agent" } }
);
const bobAddress = registryResult.agent.owner;
await this.client.execute(this.address, CONFIG.a2aContract, {
create_channel: {
participant: bobAddress,
channel_type: "direct",
metadata: toBinary({ purpose: "data-visualization-negotiation" }),
},
}, "auto");
console.log("[Alice] Channel created");
// Step 3: Send service request to Bob
console.log("[Alice] Step 3: Sending service request...");
const requestMsgId = generateId("msg");
await this.client.execute(this.address, CONFIG.a2aContract, {
send_message: {
to: bobAddress,
message_type: "request",
payload: toBinary({
msg_id: requestMsgId,
action: "create_dashboard",
requirements: {
title: "Q3 Sales Dashboard",
data_points: ["revenue", "cost", "profit", "growth_rate"],
time_range: "2026-Q3",
chart_types: ["bar", "line", "pie"],
format: "interactive_html",
},
budget: "10000000",
deadline: "2026-07-10T00:00:00Z",
}),
priority: 75,
reply_to: null,
ttl: 86400,
},
}, "auto");
console.log("[Alice] Request sent, waiting for reply...");
// Step 4: Poll for Bob's reply
let reply: any = null;
for (let i = 0; i < 60; i++) {
await sleep(2000);
const inbox: any = await readonlyClient.queryContractSmart(
CONFIG.a2aContract,
{ list_messages: { address: this.address, status: null, limit: 10 } }
);
reply = inbox.messages.find(
(m: any) => m.reply_to === requestMsgId
);
if (reply) {
console.log("[Alice] Received reply from Bob");
break;
}
process.stdout.write(".");
}
if (!reply) {
console.log("[Alice] No reply received within timeout");
return;
}
// Step 5: Process Bob's quote
const quote = fromBinary(reply.payload) as any;
console.log(`[Alice] Bob's quote: ${JSON.stringify(quote)}`);
if (quote.status === "accepted" && quote.price <= "10000000") {
console.log("[Alice] Step 5: Accepting quote and initiating payment...");
// Send acceptance
await this.client.execute(this.address, CONFIG.a2aContract, {
send_message: {
to: bobAddress,
message_type: "response",
payload: toBinary({
action: "accept_quote",
quote_id: quote.quote_id,
accepted_price: quote.price,
payment_id: `pay:${Date.now()}`,
}),
priority: 75,
reply_to: reply.msg_id,
ttl: 3600,
},
}, "auto");
console.log("[Alice] Negotiation completed successfully!");
console.log(`[Alice] Bob will create dashboard with ${quote.estimated_charts} charts`);
} else {
console.log("[Alice] Quote rejected or budget exceeded");
}
}
}
// ─── Bob Agent ───
class BobAgent {
private client!: SigningCosmWasmClient;
private address!: string;
async init() {
const wallet = await DirectSecp256k1HdWallet.fromMnemonic(CONFIG.bobMnemonic, {
prefix: "msg",
});
const [account] = await wallet.getAccounts();
this.address = account.address;
this.client = await SigningCosmWasmClient.connectWithSigner(
CONFIG.rpcUrl, wallet
);
console.log(`[Bob] Agent address: ${this.address}`);
}
async listenForRequests() {
console.log("[Bob] Listening for incoming requests...");
const readonlyClient = await CosmWasmClient.connect(CONFIG.rpcUrl);
while (true) {
try {
const inbox: any = await readonlyClient.queryContractSmart(
CONFIG.a2aContract,
{ list_messages: { address: this.address, status: "pending", limit: 10 } }
);
for (const msg of inbox.messages) {
if (msg.message_type === "request") {
console.log(`[Bob] Processing request: ${msg.msg_id}`);
await this.processRequest(msg, readonlyClient);
}
}
} catch (err) {
console.error("[Bob] Error:", err);
}
await sleep(3000);
}
}
private async processRequest(msg: any, readonlyClient: CosmWasmClient) {
try {
const payload = fromBinary(msg.payload) as any;
console.log(`[Bob] Request details: ${JSON.stringify(payload)}`);
// Generate quote
const quote = {
quote_id: `quote:${Date.now()}`,
status: "accepted",
price: "8000000",
estimated_charts: 4,
estimated_hours: 12,
delivery_date: "2026-07-09T00:00:00Z",
terms: "Payment upon delivery",
};
// Send quote back to Alice
await this.client.execute(this.address, CONFIG.a2aContract, {
reply_message: {
original_msg_id: msg.msg_id,
payload: toBinary(quote),
},
}, "auto");
console.log(`[Bob] Quote sent: ${quote.price} umsg`);
// Wait for Alice's acceptance
let accepted = false;
for (let i = 0; i < 30; i++) {
await sleep(2000);
const inbox: any = await readonlyClient.queryContractSmart(
CONFIG.a2aContract,
{ list_messages: { address: this.address, status: null, limit: 10 } }
);
const acceptance = inbox.messages.find(
(m: any) => m.reply_to === msg.msg_id && m.message_type === "response"
);
if (acceptance) {
const acceptPayload = fromBinary(acceptance.payload) as any;
if (acceptPayload.action === "accept_quote") {
accepted = true;
console.log("[Bob] Quote accepted! Starting work...");
await this.executeWork(acceptPayload);
}
break;
}
}
if (!accepted) {
console.log("[Bob] Quote not accepted within timeout");
}
} catch (err) {
console.error("[Bob] Error processing request:", err);
}
}
private async executeWork(acceptance: any) {
// Simulate work
console.log("[Bob] Creating dashboard charts...");
await sleep(5000);
// Send results back
await this.client.execute(this.address, CONFIG.a2aContract, {
send_message: {
to: acceptance.payment_id, // bob knows alice's address implicitly
message_type: "response",
payload: toBinary({
action: "deliverable",
dashboard_url: "https://storage.msgchain.org/dashboards/q3-sales.html",
charts_created: 4,
data_processed: "250K rows",
invoice: {
amount: "8000000",
currency: "umsg",
due_date: "2026-07-15T00:00:00Z",
},
}),
priority: 50,
reply_to: null,
ttl: 86400,
},
}, "auto");
console.log("[Bob] Deliverable sent!");
}
}
// ─── 工具函数 ───
function sleep(ms: number): Promise<void> {
return new Promise((resolve) => setTimeout(resolve, ms));
}
// ─── 启动 ───
async function main() {
console.log("=== MSG Chain A2A Agent Negotiation Demo ===\n");
// 初始化两个 Agent
const alice = new AliceAgent();
const bob = new BobAgent();
await alice.init();
await bob.init();
// Bob 先开始监听
bob.listenForRequests().catch(console.error);
// 等 Bob 就绪
await sleep(2000);
// Alice 发起协商
await alice.runNegotiation();
console.log("\n=== Demo Complete ===");
process.exit(0);
}
main().catch(console.error);
10.3 执行流程日志
=== MSG Chain A2A Agent Negotiation Demo ===
[Alice] Agent address: msg1...aliceAddress
[Bob] Agent address: msg1...bobAddress
[Bob] Listening for incoming requests...
[Alice] Starting negotiation...
[Alice] Step 1: Discovering agents...
[Alice] Found agents: bob-agent
[Alice] Step 2: Creating channel...
[Alice] Channel created
[Alice] Step 3: Sending service request...
[Alice] Request sent, waiting for reply...
[Bob] Processing request: msg:alice:bob:123456
[Bob] Request details: {"action":"create_dashboard","requirements":{...},"budget":"10000000","deadline":"2026-07-10T00:00:00Z"}
[Bob] Quote sent: 8000000 umsg
......[Alice] Received reply from Bob
[Alice] Bob's quote: {"status":"accepted","price":"8000000","estimated_charts":4,...}
[Alice] Step 5: Accepting quote and initiating payment...
[Bob] Quote accepted! Starting work...
[Bob] Creating dashboard charts...
[Bob] Deliverable sent!
[Alice] Negotiation completed successfully!
[Alice] Bob will create dashboard with 4 charts
=== Demo Complete ===
10.4 链上事件记录
Alice 和 Bob 的整个协商过程都会在链上产生事件,可以被索引和审计:
Event 1: a2a_channel_created
channel_id: "ch:msg1...alice-msg1...bob"
channel_type: "direct"
Event 2: a2a_message_sent
msg_id: "msg:alice:bob:123456"
sender: msg1...aliceAddress
receiver: msg1...bobAddress
message_type: "request"
Event 3: a2a_message_replied
original_msg_id: "msg:alice:bob:123456"
reply_msg_id: "reply:msg:alice:bob:123456:789"
sender: msg1...bobAddress
receiver: msg1...aliceAddress
Event 4: a2a_message_sent
msg_id: "msg:bob:alice:789012"
sender: msg1...bobAddress
receiver: msg1...aliceAddress
message_type: "response"
Event 5: a2a_channel_closed
channel_id: "ch:msg1...alice-msg1...bob"
closed_by: msg1...aliceAddress
附录
A. 规划功能清单
以下功能标记为 X-A2A-Planning=true,表示当前版本规划中但尚未完全实现:
| 功能 | 说明 | 计划版本 |
|---|---|---|
| Stream 消息类型 | 流式分片数据传输 | v1.1 |
| Pub/Sub 通道 | 主题订阅模式 | v1.2 |
| Intent 路由 | 声明式意图匹配路由 | v1.3 |
| 广播通道优化 | 原子化多播 | v1.1 |
| 消息加密 | 端到端加密(PQC) | v1.2 |
| Heartbeat 消息类型 | 连接保活探测 | v1.1 |
| 跨链 A2A | IBC 消息路由 | v2.0 |
B. 相关合约地址
| 合约 | Canonical Key | 说明 |
|---|---|---|
agent_a2a_v1 |
agent_a2a_v1 |
A2A 通信合约 |
aidid_did_registry_v1 |
aidid_did_registry_v1 |
DID 身份注册 |
agent_registry_v1 |
— | Agent 注册发现 |
agent_payment_v1 |
— | AIPAY 支付 |
ai_agent_constitution_v1 |
ai_agent_constitution_v1 |
AI 宪法检查 |
agent_verify_v1 |
— | Agent 验证 |
agent_intent_v1 |
— | Agent 意图声明 |
C. 链配置摘要
| 参数 | 值 |
|---|---|
| Chain ID | msg-chain-1 |
| Bech32 前缀 | msg |
| 出块时间 | 5s |
| Agent API 速率限制 | 100 req/s, burst 200 |
| 默认 Gas 价格 | 1,000,000,000 attoMSG/gas |
| 签名算法 | Dilithium-5 (PQC) |
| 共识机制 | Round-Robin + DAR |
D. 参考资源
- Google A2A Protocol: https://developers.google.com/a2a
- AT Protocol (Bluesky): https://atproto.com
- W3C DID Core: https://www.w3.org/TR/did-core/
- CRYSTALS-Dilithium: https://pq-crystals.org/dilithium/
- CosmWasm Documentation: https://docs.cosmwasm.com
- MSG Chain 白皮书: https://msgchain.org
